Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions Directory.Packages.props
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,8 @@
<PackageVersion Include="Temporalio.Extensions.Aws.Lambda" Version="1.20.0" />
<PackageVersion Include="Temporalio.Extensions.Aws.Lambda.OpenTelemetry" Version="1.20.0" />
<PackageVersion Include="Temporalio.Extensions.DiagnosticSource" Version="1.20.0" />
<PackageVersion Include="Temporalio.Extensions.Gcp.CloudRun.Id" Version="1.20.0" />
<PackageVersion Include="Temporalio.Extensions.Gcp.CloudRun.OpenTelemetry" Version="1.20.0" />
<PackageVersion Include="Temporalio.Extensions.Hosting" Version="1.20.0" />
<PackageVersion Include="Temporalio.Extensions.OpenTelemetry" Version="1.20.0" />
<PackageVersion Include="TemporalCommunity.Aspire.Hosting" Version="0.1.0" />
Expand Down
1 change: 1 addition & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@ Prerequisites:
* [EagerWorkflowStart](src/EagerWorkflowStart) - Demonstrates usage of Eager Workflow Start to reduce latency for workflows that start with a local activity.
* [Encryption](src/Encryption) - End-to-end encryption with Temporal payload codecs.
* [EnvConfig](src/EnvConfig) - Load client configuration from TOML files with programmatic overrides
* [Gcp/CloudRun](src/Gcp/CloudRun) - Run a Temporal Worker on a Google Cloud Run worker pool using the Cloud Run Id and OpenTelemetry plugins.
* [LambdaWorker](src/LambdaWorker) - Run a Temporal Worker inside an AWS Lambda function.
* [Mutex](src/Mutex) - How to implement a mutex as a workflow. Demonstrates how to avoid race conditions or parallel mutually exclusive operations on the same resource.
* [NexusCancellation](src/NexusCancellation) - Demonstrates how to cancel a running Nexus operation from a caller workflow.
Expand Down
18 changes: 18 additions & 0 deletions TemporalioSamples.sln
Original file line number Diff line number Diff line change
Expand Up @@ -145,6 +145,10 @@ Project("{2150E333-8FDC-42A3-9474-1A3956D46DE8}") = "NexusStandaloneActivity", "
EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "TemporalioSamples.NexusStandaloneActivity", "src\NexusStandaloneActivity\TemporalioSamples.NexusStandaloneActivity.csproj", "{4D8C9F9B-F8E3-4160-9286-32C3966A0125}"
EndProject
Project("{2150E333-8FDC-42A3-9474-1A3956D46DE8}") = "Gcp", "Gcp", "{30714183-DB89-4930-B9B4-8E35068EF641}"
EndProject
Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "TemporalioSamples.Gcp.CloudRun", "src\Gcp\CloudRun\TemporalioSamples.Gcp.CloudRun.csproj", "{5A82F4CB-2D7A-4483-BED6-9A5C5D5C6E3C}"
EndProject
Global
GlobalSection(SolutionConfigurationPlatforms) = preSolution
Debug|Any CPU = Debug|Any CPU
Expand Down Expand Up @@ -815,6 +819,18 @@ Global
{4D8C9F9B-F8E3-4160-9286-32C3966A0125}.Release|x64.Build.0 = Release|Any CPU
{4D8C9F9B-F8E3-4160-9286-32C3966A0125}.Release|x86.ActiveCfg = Release|Any CPU
{4D8C9F9B-F8E3-4160-9286-32C3966A0125}.Release|x86.Build.0 = Release|Any CPU
{5A82F4CB-2D7A-4483-BED6-9A5C5D5C6E3C}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
{5A82F4CB-2D7A-4483-BED6-9A5C5D5C6E3C}.Debug|Any CPU.Build.0 = Debug|Any CPU
{5A82F4CB-2D7A-4483-BED6-9A5C5D5C6E3C}.Debug|x64.ActiveCfg = Debug|Any CPU
{5A82F4CB-2D7A-4483-BED6-9A5C5D5C6E3C}.Debug|x64.Build.0 = Debug|Any CPU
{5A82F4CB-2D7A-4483-BED6-9A5C5D5C6E3C}.Debug|x86.ActiveCfg = Debug|Any CPU
{5A82F4CB-2D7A-4483-BED6-9A5C5D5C6E3C}.Debug|x86.Build.0 = Debug|Any CPU
{5A82F4CB-2D7A-4483-BED6-9A5C5D5C6E3C}.Release|Any CPU.ActiveCfg = Release|Any CPU
{5A82F4CB-2D7A-4483-BED6-9A5C5D5C6E3C}.Release|Any CPU.Build.0 = Release|Any CPU
{5A82F4CB-2D7A-4483-BED6-9A5C5D5C6E3C}.Release|x64.ActiveCfg = Release|Any CPU
{5A82F4CB-2D7A-4483-BED6-9A5C5D5C6E3C}.Release|x64.Build.0 = Release|Any CPU
{5A82F4CB-2D7A-4483-BED6-9A5C5D5C6E3C}.Release|x86.ActiveCfg = Release|Any CPU
{5A82F4CB-2D7A-4483-BED6-9A5C5D5C6E3C}.Release|x86.Build.0 = Release|Any CPU
EndGlobalSection
GlobalSection(SolutionProperties) = preSolution
HideSolutionNode = FALSE
Expand Down Expand Up @@ -887,5 +903,7 @@ Global
{97376F57-BA10-464B-AFBD-583E187DE947} = {34ADAC00-5559-6993-F779-D71CE5A0F93F}
{5D5AFBCD-B4B7-4A4E-91C8-FD80E4C1E27B} = {1A647B41-53D0-4638-AE5A-6630BAAE45FC}
{4D8C9F9B-F8E3-4160-9286-32C3966A0125} = {5D5AFBCD-B4B7-4A4E-91C8-FD80E4C1E27B}
{30714183-DB89-4930-B9B4-8E35068EF641} = {1A647B41-53D0-4638-AE5A-6630BAAE45FC}
{5A82F4CB-2D7A-4483-BED6-9A5C5D5C6E3C} = {30714183-DB89-4930-B9B4-8E35068EF641}
EndGlobalSection
EndGlobal
14 changes: 14 additions & 0 deletions src/Gcp/CloudRun/Dockerfile
Original file line number Diff line number Diff line change
@@ -0,0 +1,14 @@
# syntax=docker/dockerfile:1
# Build from the repo root: docker build -f src/Gcp/CloudRun/Dockerfile -t <image> .
FROM mcr.microsoft.com/dotnet/sdk:8.0 AS build

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I would suggest bumping to .NET 10. .NET 8 and 9 are EOL in Novemeber. And the global.json file is pinned to 10.0.400, which will require .NET 10 SDK.

WORKDIR /src
COPY global.json Directory.Build.props Directory.Packages.props .editorconfig nuget.config ./

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

There is no nuget.config file in this repo. This would likely fail the docker run.

COPY src/Gcp/CloudRun/ ./src/Gcp/CloudRun/
RUN dotnet publish src/Gcp/CloudRun/TemporalioSamples.Gcp.CloudRun.csproj -c Release -o /app

FROM mcr.microsoft.com/dotnet/runtime:8.0
RUN useradd --create-home --uid 10001 worker
WORKDIR /app
COPY --from=build --chown=worker:worker /app ./
USER 10001
ENTRYPOINT ["dotnet", "TemporalioSamples.Gcp.CloudRun.dll"]
14 changes: 14 additions & 0 deletions src/Gcp/CloudRun/GreetingActivities.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,14 @@
namespace TemporalioSamples.Gcp.CloudRun;

using Microsoft.Extensions.Logging;
using Temporalio.Activities;

public static class GreetingActivities
{
[Activity]
public static string SayHello(string name)
{
ActivityExecutionContext.Current.Logger.LogInformation("SayHello activity: {Name}", name);
return $"Hello, {name}!";
}
}
19 changes: 19 additions & 0 deletions src/Gcp/CloudRun/GreetingWorkflow.workflow.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,19 @@
namespace TemporalioSamples.Gcp.CloudRun;

using Microsoft.Extensions.Logging;
using Temporalio.Workflows;

[Workflow]
public class GreetingWorkflow

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I understand that testing CloudRun in the CI would be difficult. But it would be good to have some aspect tested, such as the workflow and activity do what is expected.

{
[WorkflowRun]
public async Task<string> RunAsync(string name)
{
Workflow.Logger.LogInformation("GreetingWorkflow started: {Name}", name);
var result = await Workflow.ExecuteActivityAsync(
() => GreetingActivities.SayHello(name),
new() { StartToCloseTimeout = TimeSpan.FromSeconds(10) });
Workflow.Logger.LogInformation("GreetingWorkflow completed: {Result}", result);
return result;
}
}
76 changes: 76 additions & 0 deletions src/Gcp/CloudRun/Program.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,76 @@
using System.Runtime.InteropServices;
using Microsoft.Extensions.Logging;
using Temporalio.Client;
using Temporalio.Common.EnvConfig;
using Temporalio.Extensions.Gcp.CloudRun.Id;
using Temporalio.Extensions.Gcp.CloudRun.OpenTelemetry;
using Temporalio.Worker;
using TemporalioSamples.Gcp.CloudRun;

// Connect from environment config (TEMPORAL_ADDRESS, TEMPORAL_NAMESPACE, TEMPORAL_API_KEY, ...).
var connectOptions = ClientEnvConfig.LoadClientConnectOptions();
connectOptions.TargetHost ??= "localhost:7233";
connectOptions.LoggerFactory = LoggerFactory.Create(builder =>
builder.
AddSimpleConsole(options => options.TimestampFormat = "[HH:mm:ss] ").
SetMinimumLevel(LogLevel.Information));

var taskQueue = Environment.GetEnvironmentVariable("TEMPORAL_TASK_QUEUE") ?? "cloud-run-worker";

// @@@SNIPSTART dotnet-cloud-run
// The Cloud Run Id plugin sets the client identity to "{instanceId}@{revision}" from Cloud Run
// metadata at connect time; every Worker created from the client inherits it.
connectOptions.Plugins = new ITemporalClientPlugin[] { new CloudRunIdPlugin() };

// ApplyGoogleCloudRunOpenTelemetryDefaults adds the tracing interceptor and a runtime exporting Core
// metrics + traces over OTLP to the collector sidecar; the returned handle owns the tracer provider.
// Both plugins configure the same connect options.
using var telemetry = connectOptions.ApplyGoogleCloudRunOpenTelemetryDefaults();

var client = await TemporalClient.ConnectAsync(connectOptions);
// @@@SNIPEND

// --starter runs a single workflow, e.g. to kick off work against the same server the worker polls.
if (args.Contains("--starter"))

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This fails if not running on CloudRun because of the CloudRunIdPlugin requiring CloudRun environment. But this seems like it is supposed to run on the local machine? Also, I don't see this being referenced in the rest of the example.

{
var greeting = await client.ExecuteWorkflowAsync(
(GreetingWorkflow wf) => wf.RunAsync("Temporal"),
new($"cloud-run-worker-{Guid.NewGuid():N}", taskQueue));
Console.WriteLine("Workflow result: {0}", greeting);
await telemetry.FlushAsync(TimeSpan.FromSeconds(2));
return;
}

using var cts = new CancellationTokenSource();
Console.CancelKeyPress += (_, eventArgs) =>
{
eventArgs.Cancel = true;
cts.Cancel();
};

// Cloud Run sends SIGTERM ~10s before SIGKILL.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The default for GracefulShutdownTimeout is 0s. When the cancellation token is signaled, the worker is going to immediately shutdown, abandoning in-flight work. Probably want to set that so something close-ish to the Cloud Run difference, maybe 5 seconds.

using var sigterm = PosixSignalRegistration.Create(PosixSignal.SIGTERM, _ => cts.Cancel());

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This needs to set ctx.Cancel = true in the callback to allow graceful shutdown.


using var worker = new TemporalWorker(
client,
new TemporalWorkerOptions(taskQueue).
AddWorkflow<GreetingWorkflow>().
AddActivity(GreetingActivities.SayHello));

Console.WriteLine(
"Worker running: taskQueue={0} address={1} namespace={2}",
taskQueue,
connectOptions.TargetHost,
connectOptions.Namespace ?? "default");
try
{
await worker.ExecuteAsync(cts.Token);
}
catch (OperationCanceledException)
{
Console.WriteLine("Worker shutting down");
}

// Flush traces within the shutdown grace window (Core metrics export periodically, no explicit flush).
await telemetry.FlushAsync(TimeSpan.FromSeconds(2));
Console.WriteLine("Worker stopped");
55 changes: 55 additions & 0 deletions src/Gcp/CloudRun/README.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,55 @@
# Cloud Run Worker

Run a Temporal Worker on a
[Google Cloud Run worker pool](https://cloud.google.com/run/docs/deploy-worker-pools) using both GCP
Cloud Run extensions together:

- `Temporalio.Extensions.Gcp.CloudRun.Id` derives the client identity `{instanceId}@{revision}` from
Cloud Run instance metadata, so each instance is identifiable in Temporal.
- `Temporalio.Extensions.Gcp.CloudRun.OpenTelemetry` exports Core SDK metrics and traces to a
[Google-Built OpenTelemetry Collector](https://cloud.google.com/stackdriver/docs/instrumentation/opentelemetry-collector-cloud-run)
sidecar (metrics to Managed Service for Prometheus, traces to Cloud Trace).

`Program.cs` registers a `CloudRunIdPlugin` on `TemporalClientConnectOptions.Plugins` and calls
`ApplyGoogleCloudRunOpenTelemetryDefaults()`, which adds the tracing interceptor and an OTLP exporter
aimed at the sidecar, then runs a greeting workflow and activity until SIGTERM.

## Prerequisites

- A Temporal server the worker pool can reach (`TEMPORAL_ADDRESS` / `TEMPORAL_NAMESPACE`)
- A Google Cloud project with the Cloud Run and Artifact Registry APIs enabled
- [`gcloud`](https://cloud.google.com/sdk/docs/install), authenticated with the project set
- The [Temporal CLI](https://docs.temporal.io/cli) and .NET 8

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Should be .NET 10 since that is what the repo requires.


## Deploy

Set the placeholders used by `worker-pool.yaml`, then build the image, store the collector config as
a secret, and deploy. Run from the repo root:

```bash
export REGION=us-central1 WORKER_POOL=temporal-cloud-run-worker INSTANCE_COUNT=1
export SERVICE_ACCOUNT_EMAIL=<sa>@<project>.iam.gserviceaccount.com
export WORKER_IMAGE=$REGION-docker.pkg.dev/$(gcloud config get-value project)/samples/cloud-run
export TEMPORAL_ADDRESS=<host:7233> TEMPORAL_NAMESPACE=<namespace> TEMPORAL_TASK_QUEUE=cloud-run-worker
export COLLECTOR_CONFIG_SECRET=otel-collector-config COLLECTOR_CONFIG_SECRET_VERSION=latest

gcloud artifacts repositories create samples --repository-format=docker --location "$REGION"
docker build -f src/Gcp/CloudRun/Dockerfile -t "$WORKER_IMAGE" . && docker push "$WORKER_IMAGE"
gcloud secrets create "$COLLECTOR_CONFIG_SECRET" --data-file=src/Gcp/CloudRun/collector-config.yaml

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This likely fails if run a second time.


envsubst < src/Gcp/CloudRun/worker-pool.yaml > /tmp/worker-pool.yaml
gcloud run worker-pools replace /tmp/worker-pool.yaml --region "$REGION"
```

The service account needs the `monitoring.metricWriter`, `cloudtrace.agent`, and
`secretmanager.secretAccessor` roles. For Temporal Cloud, add `TEMPORAL_API_KEY` from a Secret
Manager `secretKeyRef` (see `worker-pool.yaml`).

## Run a workflow

```bash
temporal workflow execute --task-queue cloud-run-worker --type GreetingWorkflow --input '"Temporal"'
```

Metrics appear in Metrics Explorer and traces in Trace Explorer, tagged with the Cloud Run identity.
Delete the pool with `gcloud run worker-pools delete "$WORKER_POOL" --region "$REGION"`.
12 changes: 12 additions & 0 deletions src/Gcp/CloudRun/TemporalioSamples.Gcp.CloudRun.csproj
Original file line number Diff line number Diff line change
@@ -0,0 +1,12 @@
<Project Sdk="Microsoft.NET.Sdk">

<PropertyGroup>
<OutputType>Exe</OutputType>
</PropertyGroup>

<ItemGroup>
<PackageReference Include="Temporalio.Extensions.Gcp.CloudRun.Id" />
<PackageReference Include="Temporalio.Extensions.Gcp.CloudRun.OpenTelemetry" />
</ItemGroup>

</Project>
82 changes: 82 additions & 0 deletions src/Gcp/CloudRun/collector-config.yaml
Original file line number Diff line number Diff line change
@@ -0,0 +1,82 @@
# Google-Built OpenTelemetry Collector config for the sidecar: metrics -> Managed Service for Prometheus, traces -> Cloud Trace. Auth via the worker-pool service account (ADC).
receivers:
otlp:
protocols:
grpc:
endpoint: localhost:4317

processors:
# Batch traces for throughput. Do NOT add a batch processor to the cumulative-metrics pipeline: a
# shutdown flush could be batched with a recent periodic export of the same series and rejected as
# a duplicate time series.
batch/traces:
send_batch_max_size: 200
send_batch_size: 200
timeout: 5s
memory_limiter:
check_interval: 1s
limit_percentage: 65
spike_limit_percentage: 20
resourcedetection:
detectors: [gcp]
timeout: 10s
# Rename Temporal datapoint labels that collide with the target labels Managed Service for Prometheus injects (e.g. `namespace`).
transform/collision:
metric_statements:
- context: datapoint
statements:
- set(attributes["exported_location"], attributes["location"])
- delete_key(attributes, "location")
- set(attributes["exported_cluster"], attributes["cluster"])
- delete_key(attributes, "cluster")
- set(attributes["exported_namespace"], attributes["namespace"])
- delete_key(attributes, "namespace")
- set(attributes["exported_job"], attributes["job"])
- delete_key(attributes, "job")
- set(attributes["exported_instance"], attributes["instance"])
- delete_key(attributes, "instance")
- set(attributes["exported_project_id"], attributes["project_id"])
- delete_key(attributes, "project_id")
# The Telemetry API expects the Google Cloud project in gcp.project_id.
transform/set_project_id:
error_mode: ignore
trace_statements:
- set(resource.attributes["gcp.project_id"], resource.attributes["gcp.project.id"]) where resource.attributes["gcp.project.id"] != nil
- set(resource.attributes["gcp.project_id"], resource.attributes["cloud.account.id"]) where resource.attributes["gcp.project_id"] == nil and resource.attributes["cloud.account.id"] != nil

exporters:
googlemanagedprometheus:
otlp:
endpoint: telemetry.googleapis.com:443
compression: none
balancer_name: pick_first
auth:
authenticator: googleclientauth

extensions:
health_check:
endpoint: 0.0.0.0:13133
googleclientauth:

service:
extensions:
- health_check
- googleclientauth
pipelines:
metrics:
receivers: [otlp]
processors: [memory_limiter, resourcedetection, transform/collision]
exporters: [googlemanagedprometheus]
traces:
receivers: [otlp]
processors: [memory_limiter, resourcedetection, transform/set_project_id, batch/traces]
exporters: [otlp]
telemetry:
metrics:
readers:
- periodic:
exporter:
otlp:
protocol: grpc
endpoint: http://localhost:4317
insecure: true
56 changes: 56 additions & 0 deletions src/Gcp/CloudRun/worker-pool.yaml
Original file line number Diff line number Diff line change
@@ -0,0 +1,56 @@
# Cloud Run WorkerPool: a Temporal Worker + a Collector sidecar. Render placeholders with `envsubst` (see README) then `gcloud run worker-pools replace`. For Temporal Cloud, add TEMPORAL_API_KEY from a Secret Manager secretKeyRef.
apiVersion: run.googleapis.com/v1
kind: WorkerPool
metadata:
name: "${WORKER_POOL}"
labels:
cloud.googleapis.com/location: "${REGION}"
annotations:
run.googleapis.com/scalingMode: manual
run.googleapis.com/manualInstanceCount: "${INSTANCE_COUNT}"
spec:
template:
metadata:
annotations:
run.googleapis.com/container-dependencies: '{"worker":["collector"]}'
run.googleapis.com/execution-environment: gen2
spec:
containerConcurrency: 0
serviceAccountName: "${SERVICE_ACCOUNT_EMAIL}"
containers:
- name: worker
image: "${WORKER_IMAGE}"
env:
- name: TEMPORAL_ADDRESS
value: "${TEMPORAL_ADDRESS}"
- name: TEMPORAL_NAMESPACE
value: "${TEMPORAL_NAMESPACE}"
- name: TEMPORAL_TASK_QUEUE
value: "${TEMPORAL_TASK_QUEUE}"
- name: OTEL_EXPORTER_OTLP_ENDPOINT
value: http://localhost:4317
resources:
limits:
cpu: "1"
memory: 512Mi
- name: collector
image: us-docker.pkg.dev/cloud-ops-agents-artifacts/google-cloud-opentelemetry-collector/otelcol-google:0.156.0
args:
- --config=env:OTELCOL_CONFIG
env:
- name: OTELCOL_CONFIG
valueFrom:
secretKeyRef:
key: "${COLLECTOR_CONFIG_SECRET_VERSION}"
name: "${COLLECTOR_CONFIG_SECRET}"
startupProbe:
httpGet:
path: /
port: 13133
timeoutSeconds: 1
periodSeconds: 2
failureThreshold: 30
resources:
limits:
cpu: "1"
memory: 512Mi
Loading