-
Notifications
You must be signed in to change notification settings - Fork 42
Add Cloud Run worker sample (CloudRunId + OpenTelemetry plugins) #236
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: main
Are you sure you want to change the base?
Changes from all commits
a31cce4
cf46e7e
fb54d30
01c6db8
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| 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 | ||
| WORKDIR /src | ||
| COPY global.json Directory.Build.props Directory.Packages.props .editorconfig nuget.config ./ | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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"] | ||
| 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}!"; | ||
| } | ||
| } |
| 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 | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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; | ||
| } | ||
| } | ||
| 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")) | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. This fails if not running on CloudRun because of the |
||
| { | ||
| 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. | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. The default for |
||
| using var sigterm = PosixSignalRegistration.Create(PosixSignal.SIGTERM, _ => cts.Cancel()); | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. This needs to set |
||
|
|
||
| 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"); | ||
| 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 | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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 | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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"`. | ||
| 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> |
| 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 |
| 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 |
There was a problem hiding this comment.
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.