diff --git a/gcp/cloud-run/Dockerfile b/gcp/cloud-run/Dockerfile new file mode 100644 index 00000000..b91aa346 --- /dev/null +++ b/gcp/cloud-run/Dockerfile @@ -0,0 +1,16 @@ +FROM eclipse-temurin:17-jdk-jammy AS build + +WORKDIR /workspace +COPY . . + +RUN ./gradlew --no-daemon :gcp:cloud-run:installDist + +FROM eclipse-temurin:17-jre-jammy + +RUN useradd --create-home --uid 10001 temporal +WORKDIR /app +COPY --from=build --chown=temporal:temporal \ + /workspace/gcp/cloud-run/build/install/cloud-run-worker/ /app/ + +USER 10001 +ENTRYPOINT ["/app/bin/cloud-run-worker"] diff --git a/gcp/cloud-run/README.md b/gcp/cloud-run/README.md new file mode 100644 index 00000000..f9aea074 --- /dev/null +++ b/gcp/cloud-run/README.md @@ -0,0 +1,77 @@ +# Temporal Cloud Run worker + +A Temporal Worker running in a Google Cloud Run **worker pool** that uses both GCP Cloud Run +plugins together. Worker pools keep CPU allocated, unlike request-driven Cloud Run services. Both +plugins are registered on the service stubs, so they propagate to the client and to workers created +from it. + +- `CloudRunIdPlugin` (`io.temporal:temporal-gcp-cloud-run-id`) sets the Temporal client identity + from Cloud Run instance metadata as `{instanceId}@{revision}`, unless an identity is already set. +- `CloudRunOpenTelemetryPlugin` (`io.temporal:temporal-gcp-cloud-run-opentelemetry`) configures the + SDK metrics scope, tracing interceptors, OTLP exporters (default `http://localhost:4317`), and + shutdown flushing, exporting to a + [Google-Built OpenTelemetry Collector](https://cloud.google.com/stackdriver/docs/instrumentation/opentelemetry-collector-cloud-run) + sidecar. It derives `service.name` from `CLOUD_RUN_WORKER_POOL` and reports metrics every 60 + seconds, matching the OpenTelemetry SDK default across the Temporal SDKs. + +> Google Cloud Run support is experimental and may change without notice. + +## Build + +`temporal-gcp-cloud-run-id` and `temporal-gcp-cloud-run-opentelemetry` are published on Maven Central. + +```bash +./gradlew :gcp:cloud-run:build +``` + +## Files + +- `.../cloudrun/CloudRunWorker.java` — both plugins, client, Temporal Worker, bounded `SIGTERM` shutdown. +- `collector-config.yaml` — collector for cumulative Prometheus metrics and batched traces. +- `worker-pool.yaml` — worker and collector containers sharing localhost, config from Secret Manager. +- `Dockerfile` — packages the Gradle application as the worker container. + +## Deploy + +Run from the repository root, with a Temporal Cloud namespace and API key. + +1. Store the API key and collector config in Secret Manager (one-time setup): + + ```bash + printf '%s' "$TEMPORAL_API_KEY" | \ + gcloud secrets create temporal-api-key --data-file=- --project="$PROJECT_ID" + gcloud secrets create temporal-collector-config \ + --data-file=gcp/cloud-run/collector-config.yaml --project="$PROJECT_ID" + ``` + +2. Create the Artifact Registry repo (one-time setup), then build and push the image: + + ```bash + gcloud artifacts repositories create temporal-samples --repository-format=docker \ + --location="$REGION" --project="$PROJECT_ID" + IMAGE="$REGION-docker.pkg.dev/$PROJECT_ID/temporal-samples/cloud-run-worker:latest" + docker build --platform linux/amd64 -f gcp/cloud-run/Dockerfile -t "$IMAGE" . && docker push "$IMAGE" + ``` + +3. Replace the placeholders in `worker-pool.yaml`, then deploy: + + ```bash + gcloud run worker-pools replace gcp/cloud-run/worker-pool.yaml --project="$PROJECT_ID" + ``` + +The service account needs `roles/monitoring.metricWriter`, `roles/telemetry.tracesWriter`, and +`roles/secretmanager.secretAccessor`. Set `TEMPORAL_ADDRESS`, `TEMPORAL_NAMESPACE`, and +`TEMPORAL_API_KEY` to your Temporal Cloud values, then start a workflow on task queue +`cloud-run-worker`: + +```bash +temporal workflow execute --type GreetingWorkflow --task-queue cloud-run-worker \ + --workflow-id cloud-run-greeting --input '"Google Cloud"' \ + --address "$TEMPORAL_ADDRESS" --namespace "$TEMPORAL_NAMESPACE" \ + --api-key "$TEMPORAL_API_KEY" --tls +``` + +It prints `Hello Google Cloud!`, confirming the deployed worker ran the task. + +The collector does **not** batch cumulative metrics: a shutdown flush batched with a recent periodic +export would collide on the same Prometheus series and be rejected as `Duplicate TimeSeries`. diff --git a/gcp/cloud-run/build.gradle b/gcp/cloud-run/build.gradle new file mode 100644 index 00000000..b5cd037f --- /dev/null +++ b/gcp/cloud-run/build.gradle @@ -0,0 +1,28 @@ +apply plugin: 'application' + +dependencies { + // The GCP Cloud Run modules first shipped in 1.40.0; pin the io.temporal deps to that + // released version here rather than bumping the repo-wide javaSDKVersion. The explicit + // temporal-sdk dep is required because temporal-gcp-cloud-run-id declares it as compileOnly. + implementation "io.temporal:temporal-sdk:1.40.0" + implementation "io.temporal:temporal-envconfig:1.40.0" + implementation "io.temporal:temporal-gcp-cloud-run-opentelemetry:1.40.0" + implementation "io.temporal:temporal-gcp-cloud-run-id:1.40.0" + runtimeOnly group: 'ch.qos.logback', name: 'logback-classic', version: '1.5.6' + + testImplementation "io.temporal:temporal-testing:$javaSDKVersion" + testImplementation "junit:junit:4.13.2" + testImplementation(platform("org.junit:junit-bom:5.10.3")) + testRuntimeOnly "org.junit.vintage:junit-vintage-engine" + + dependencies { + errorproneJavac('com.google.errorprone:javac:9+181-r4173-1') + errorprone('com.google.errorprone:error_prone_core:2.28.0') + } +} + +application { + mainClass = 'io.temporal.samples.gcp.cloudrun.CloudRunWorker' + // Stable launcher name the Dockerfile relies on, independent of the Gradle project path. + applicationName = 'cloud-run-worker' +} diff --git a/gcp/cloud-run/collector-config.yaml b/gcp/cloud-run/collector-config.yaml new file mode 100644 index 00000000..b4973f86 --- /dev/null +++ b/gcp/cloud-run/collector-config.yaml @@ -0,0 +1,87 @@ +# @@@SNIPSTART java-cloud-run-collector-config +receivers: + otlp: + protocols: + grpc: + endpoint: localhost:4317 + +processors: + # Batch traces for throughput. Do not add this processor to the cumulative metrics pipeline: + # a shutdown flush can otherwise be batched with a recent periodic export of the same series. + batch/traces: + send_batch_max_size: 200 + send_batch_size: 200 + timeout: 5s + memory_limiter: + # This is the collector's memory polling cadence, not the SDK metric export interval. + check_interval: 1s + limit_percentage: 65 + spike_limit_percentage: 20 + resource_detection: + detectors: [gcp] + timeout: 10s + # Avoid collisions with labels that Google Managed Service for Prometheus adds. + 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: + # Google Cloud's supported OTLP path for traces is the Telemetry API. + otlp_grpc: + endpoint: telemetry.googleapis.com:443 + compression: none + balancer_name: pick_first + auth: + authenticator: googleclientauth + +extensions: + # Cloud Run container dependencies require a startup probe. This endpoint is also used for the + # collector liveness probe in worker-pool.yaml. + health_check: + endpoint: 0.0.0.0:13133 + googleclientauth: + +service: + extensions: + - health_check + - googleclientauth + pipelines: + metrics/otlp: + receivers: [otlp] + processors: [memory_limiter, resource_detection, transform/collision] + exporters: [googlemanagedprometheus] + traces: + receivers: [otlp] + processors: [memory_limiter, resource_detection, transform/set_project_id, batch/traces] + exporters: [otlp_grpc] + # Feed collector self-metrics back through the metrics pipeline. + telemetry: + metrics: + readers: + - periodic: + exporter: + otlp: + protocol: grpc + endpoint: http://localhost:4317 + insecure: true +# @@@SNIPEND diff --git a/gcp/cloud-run/src/main/java/io/temporal/samples/gcp/cloudrun/CloudRunWorker.java b/gcp/cloud-run/src/main/java/io/temporal/samples/gcp/cloudrun/CloudRunWorker.java new file mode 100644 index 00000000..bd1845fb --- /dev/null +++ b/gcp/cloud-run/src/main/java/io/temporal/samples/gcp/cloudrun/CloudRunWorker.java @@ -0,0 +1,71 @@ +package io.temporal.samples.gcp.cloudrun; + +import io.temporal.client.WorkflowClient; +import io.temporal.envconfig.ClientConfigProfile; +import io.temporal.gcp.cloudrun.id.CloudRunIdPlugin; +import io.temporal.gcp.cloudrun.opentelemetry.CloudRunOpenTelemetryPlugin; +import io.temporal.serviceclient.WorkflowServiceStubs; +import io.temporal.serviceclient.WorkflowServiceStubsOptions; +import io.temporal.worker.Worker; +import io.temporal.worker.WorkerFactory; +import java.io.IOException; +import java.time.Duration; +import java.util.concurrent.TimeUnit; + +/** A Temporal Worker for a Cloud Run worker pool. */ +public final class CloudRunWorker { + public static final String DEFAULT_TASK_QUEUE = "cloud-run-worker"; + + private CloudRunWorker() {} + + public static void main(String[] args) throws IOException { + // @@@SNIPSTART java-cloud-run-worker + ClientConfigProfile profile = ClientConfigProfile.load(); + CloudRunOpenTelemetryPlugin otelPlugin = CloudRunOpenTelemetryPlugin.newBuilder().build(); + CloudRunIdPlugin idPlugin = new CloudRunIdPlugin(); + + WorkflowServiceStubsOptions serviceOptions = + WorkflowServiceStubsOptions.newBuilder(profile.toWorkflowServiceStubsOptions()) + .setPlugins(otelPlugin, idPlugin) + .build(); + WorkflowServiceStubs service = WorkflowServiceStubs.newServiceStubs(serviceOptions); + // @@@SNIPEND + WorkflowClient client = WorkflowClient.newInstance(service, profile.toWorkflowClientOptions()); + WorkerFactory factory = WorkerFactory.newInstance(client); + + String taskQueue = taskQueue(); + Worker worker = factory.newWorker(taskQueue); + worker.registerWorkflowImplementationTypes(GreetingWorkflowImpl.class); + worker.registerActivitiesImplementations(new GreetingActivitiesImpl()); + + Runtime.getRuntime() + .addShutdownHook( + new Thread(() -> shutdown(factory, service, otelPlugin), "temporal-worker-shutdown")); + + factory.start(); + System.out.printf( + "Temporal worker started: taskQueue=%s, otelEndpoint=%s, serviceName=%s%n", + taskQueue, otelPlugin.getEndpoint(), otelPlugin.getServiceName()); + + // Keep the process alive until Cloud Run sends SIGTERM. + factory.awaitTermination(Long.MAX_VALUE, TimeUnit.DAYS); + } + + private static String taskQueue() { + String configured = System.getenv("TEMPORAL_TASK_QUEUE"); + return configured == null || configured.trim().isEmpty() ? DEFAULT_TASK_QUEUE : configured; + } + + private static void shutdown( + WorkerFactory factory, WorkflowServiceStubs service, CloudRunOpenTelemetryPlugin otelPlugin) { + // Flush after worker shutdown to capture finishing tasks; Cloud Run allows 10s. + factory.shutdown(); + factory.awaitTermination(6, TimeUnit.SECONDS); + if (!factory.isTerminated()) { + factory.shutdownNow(); + factory.awaitTermination(1, TimeUnit.SECONDS); + } + otelPlugin.newFlushHook().run(Duration.ofSeconds(2)); + service.shutdown(); + } +} diff --git a/gcp/cloud-run/src/main/java/io/temporal/samples/gcp/cloudrun/GreetingActivities.java b/gcp/cloud-run/src/main/java/io/temporal/samples/gcp/cloudrun/GreetingActivities.java new file mode 100644 index 00000000..64d9fa24 --- /dev/null +++ b/gcp/cloud-run/src/main/java/io/temporal/samples/gcp/cloudrun/GreetingActivities.java @@ -0,0 +1,8 @@ +package io.temporal.samples.gcp.cloudrun; + +import io.temporal.activity.ActivityInterface; + +@ActivityInterface +public interface GreetingActivities { + String composeGreeting(String name); +} diff --git a/gcp/cloud-run/src/main/java/io/temporal/samples/gcp/cloudrun/GreetingActivitiesImpl.java b/gcp/cloud-run/src/main/java/io/temporal/samples/gcp/cloudrun/GreetingActivitiesImpl.java new file mode 100644 index 00000000..9d473a0d --- /dev/null +++ b/gcp/cloud-run/src/main/java/io/temporal/samples/gcp/cloudrun/GreetingActivitiesImpl.java @@ -0,0 +1,8 @@ +package io.temporal.samples.gcp.cloudrun; + +public final class GreetingActivitiesImpl implements GreetingActivities { + @Override + public String composeGreeting(String name) { + return "Hello " + name + "!"; + } +} diff --git a/gcp/cloud-run/src/main/java/io/temporal/samples/gcp/cloudrun/GreetingWorkflow.java b/gcp/cloud-run/src/main/java/io/temporal/samples/gcp/cloudrun/GreetingWorkflow.java new file mode 100644 index 00000000..262e9977 --- /dev/null +++ b/gcp/cloud-run/src/main/java/io/temporal/samples/gcp/cloudrun/GreetingWorkflow.java @@ -0,0 +1,10 @@ +package io.temporal.samples.gcp.cloudrun; + +import io.temporal.workflow.WorkflowInterface; +import io.temporal.workflow.WorkflowMethod; + +@WorkflowInterface +public interface GreetingWorkflow { + @WorkflowMethod + String getGreeting(String name); +} diff --git a/gcp/cloud-run/src/main/java/io/temporal/samples/gcp/cloudrun/GreetingWorkflowImpl.java b/gcp/cloud-run/src/main/java/io/temporal/samples/gcp/cloudrun/GreetingWorkflowImpl.java new file mode 100644 index 00000000..60c71cfa --- /dev/null +++ b/gcp/cloud-run/src/main/java/io/temporal/samples/gcp/cloudrun/GreetingWorkflowImpl.java @@ -0,0 +1,17 @@ +package io.temporal.samples.gcp.cloudrun; + +import io.temporal.activity.ActivityOptions; +import io.temporal.workflow.Workflow; +import java.time.Duration; + +public final class GreetingWorkflowImpl implements GreetingWorkflow { + private final GreetingActivities activities = + Workflow.newActivityStub( + GreetingActivities.class, + ActivityOptions.newBuilder().setStartToCloseTimeout(Duration.ofSeconds(10)).build()); + + @Override + public String getGreeting(String name) { + return activities.composeGreeting(name); + } +} diff --git a/gcp/cloud-run/src/main/resources/logback.xml b/gcp/cloud-run/src/main/resources/logback.xml new file mode 100644 index 00000000..28eb2cba --- /dev/null +++ b/gcp/cloud-run/src/main/resources/logback.xml @@ -0,0 +1,14 @@ + + + + %d{HH:mm:ss.SSS} %-5level [%thread] %logger{36} - %msg%n + + + + + + + + + + diff --git a/gcp/cloud-run/src/test/java/io/temporal/samples/gcp/cloudrun/GreetingWorkflowTest.java b/gcp/cloud-run/src/test/java/io/temporal/samples/gcp/cloudrun/GreetingWorkflowTest.java new file mode 100644 index 00000000..3035e975 --- /dev/null +++ b/gcp/cloud-run/src/test/java/io/temporal/samples/gcp/cloudrun/GreetingWorkflowTest.java @@ -0,0 +1,32 @@ +package io.temporal.samples.gcp.cloudrun; + +import static org.junit.Assert.assertEquals; + +import io.temporal.client.WorkflowOptions; +import io.temporal.testing.TestWorkflowRule; +import org.junit.Rule; +import org.junit.Test; + +public class GreetingWorkflowTest { + @Rule + public TestWorkflowRule testWorkflowRule = + TestWorkflowRule.newBuilder() + .setWorkflowTypes(GreetingWorkflowImpl.class) + .setDoNotStart(true) + .build(); + + @Test + public void completesGreeting() { + testWorkflowRule.getWorker().registerActivitiesImplementations(new GreetingActivitiesImpl()); + testWorkflowRule.getTestEnvironment().start(); + + GreetingWorkflow workflow = + testWorkflowRule + .getWorkflowClient() + .newWorkflowStub( + GreetingWorkflow.class, + WorkflowOptions.newBuilder().setTaskQueue(testWorkflowRule.getTaskQueue()).build()); + + assertEquals("Hello Google Cloud!", workflow.getGreeting("Google Cloud")); + } +} diff --git a/gcp/cloud-run/worker-pool.yaml b/gcp/cloud-run/worker-pool.yaml new file mode 100644 index 00000000..0136d974 --- /dev/null +++ b/gcp/cloud-run/worker-pool.yaml @@ -0,0 +1,68 @@ +apiVersion: run.googleapis.com/v1 +kind: WorkerPool +metadata: + name: temporal-cloud-run-worker + labels: + cloud.googleapis.com/location: REGION + annotations: + run.googleapis.com/scalingMode: manual + run.googleapis.com/manualInstanceCount: "1" +spec: + template: + metadata: + annotations: + # Cloud Run starts the worker only after the collector startup probe succeeds. + run.googleapis.com/container-dependencies: '{"worker":["collector"]}' + run.googleapis.com/execution-environment: gen2 + spec: + containerConcurrency: 0 + serviceAccountName: temporal-cloud-run-worker@PROJECT_ID.iam.gserviceaccount.com + containers: + - name: worker + image: REGION-docker.pkg.dev/PROJECT_ID/temporal-samples/cloud-run-worker:latest + env: + - name: TEMPORAL_ADDRESS + value: NAMESPACE.tmprl.cloud:7233 + - name: TEMPORAL_NAMESPACE + value: NAMESPACE + - name: TEMPORAL_API_KEY + valueFrom: + secretKeyRef: + key: "1" + name: temporal-api-key + - name: TEMPORAL_TASK_QUEUE + value: cloud-run-worker + - 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: "1" + name: temporal-collector-config + startupProbe: + httpGet: + path: / + port: 13133 + timeoutSeconds: 5 + periodSeconds: 10 + failureThreshold: 12 + livenessProbe: + httpGet: + path: / + port: 13133 + timeoutSeconds: 5 + periodSeconds: 30 + failureThreshold: 3 + resources: + limits: + cpu: "1" + memory: 512Mi diff --git a/settings.gradle b/settings.gradle index b19f2fd7..e66f00a2 100644 --- a/settings.gradle +++ b/settings.gradle @@ -9,3 +9,4 @@ include 'springboot' include 'springboot-basic' include 'lambda-worker:starter' include 'lambda-worker:worker' +include 'gcp:cloud-run'