diff --git a/Directory.Packages.props b/Directory.Packages.props index 04aaa6f..7c7fa3f 100644 --- a/Directory.Packages.props +++ b/Directory.Packages.props @@ -21,6 +21,8 @@ + + diff --git a/README.md b/README.md index 7002fc7..7cc60b1 100644 --- a/README.md +++ b/README.md @@ -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. diff --git a/TemporalioSamples.sln b/TemporalioSamples.sln index c945182..02fc71e 100644 --- a/TemporalioSamples.sln +++ b/TemporalioSamples.sln @@ -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 @@ -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 @@ -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 diff --git a/src/Gcp/CloudRun/Dockerfile b/src/Gcp/CloudRun/Dockerfile new file mode 100644 index 0000000..2c01989 --- /dev/null +++ b/src/Gcp/CloudRun/Dockerfile @@ -0,0 +1,14 @@ +# syntax=docker/dockerfile:1 +# Build from the repo root: docker build -f src/Gcp/CloudRun/Dockerfile -t . +FROM mcr.microsoft.com/dotnet/sdk:8.0 AS build +WORKDIR /src +COPY global.json Directory.Build.props Directory.Packages.props .editorconfig nuget.config ./ +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"] diff --git a/src/Gcp/CloudRun/GreetingActivities.cs b/src/Gcp/CloudRun/GreetingActivities.cs new file mode 100644 index 0000000..7c66e06 --- /dev/null +++ b/src/Gcp/CloudRun/GreetingActivities.cs @@ -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}!"; + } +} diff --git a/src/Gcp/CloudRun/GreetingWorkflow.workflow.cs b/src/Gcp/CloudRun/GreetingWorkflow.workflow.cs new file mode 100644 index 0000000..e3e8297 --- /dev/null +++ b/src/Gcp/CloudRun/GreetingWorkflow.workflow.cs @@ -0,0 +1,19 @@ +namespace TemporalioSamples.Gcp.CloudRun; + +using Microsoft.Extensions.Logging; +using Temporalio.Workflows; + +[Workflow] +public class GreetingWorkflow +{ + [WorkflowRun] + public async Task 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; + } +} diff --git a/src/Gcp/CloudRun/Program.cs b/src/Gcp/CloudRun/Program.cs new file mode 100644 index 0000000..18c353c --- /dev/null +++ b/src/Gcp/CloudRun/Program.cs @@ -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")) +{ + 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. +using var sigterm = PosixSignalRegistration.Create(PosixSignal.SIGTERM, _ => cts.Cancel()); + +using var worker = new TemporalWorker( + client, + new TemporalWorkerOptions(taskQueue). + AddWorkflow(). + 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"); diff --git a/src/Gcp/CloudRun/README.md b/src/Gcp/CloudRun/README.md new file mode 100644 index 0000000..14cbb59 --- /dev/null +++ b/src/Gcp/CloudRun/README.md @@ -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 + +## 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=@.iam.gserviceaccount.com +export WORKER_IMAGE=$REGION-docker.pkg.dev/$(gcloud config get-value project)/samples/cloud-run +export TEMPORAL_ADDRESS= TEMPORAL_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 + +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"`. diff --git a/src/Gcp/CloudRun/TemporalioSamples.Gcp.CloudRun.csproj b/src/Gcp/CloudRun/TemporalioSamples.Gcp.CloudRun.csproj new file mode 100644 index 0000000..3e9c317 --- /dev/null +++ b/src/Gcp/CloudRun/TemporalioSamples.Gcp.CloudRun.csproj @@ -0,0 +1,12 @@ + + + + Exe + + + + + + + + diff --git a/src/Gcp/CloudRun/collector-config.yaml b/src/Gcp/CloudRun/collector-config.yaml new file mode 100644 index 0000000..dc7f7fe --- /dev/null +++ b/src/Gcp/CloudRun/collector-config.yaml @@ -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 diff --git a/src/Gcp/CloudRun/worker-pool.yaml b/src/Gcp/CloudRun/worker-pool.yaml new file mode 100644 index 0000000..8688311 --- /dev/null +++ b/src/Gcp/CloudRun/worker-pool.yaml @@ -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