From 91d786713b0bc00be23011aaa3d0934373e3f49e Mon Sep 17 00:00:00 2001 From: seanbollin Date: Fri, 21 Aug 2026 12:22:49 -0700 Subject: [PATCH 1/4] Add Google Cloud Run worker identity/deployment helper Adds an experimental Google Cloud Run helper, mirroring the existing AWS Lambda module's worker-ID behavior. Because Cloud Run runs a long-lived container (unlike Lambda's per-invocation model), this is a metadata helper rather than a worker wrapper: it reads the Cloud Run instance metadata -- the instance id from the metadata server, plus the worker pool/service name and revision from CLOUD_RUN_WORKER_POOL / CLOUD_RUN_REVISION (worker pools) or K_SERVICE / K_REVISION (services) -- and derives a worker identity and a WorkerDeploymentVersion to apply to a normal long-lived worker. Covers both Cloud Run worker pools and services. Co-Authored-By: Claude Opus 4.8 --- contrib/temporal-gcp-cloud-run/README.md | 71 +++++ contrib/temporal-gcp-cloud-run/build.gradle | 7 + .../gcp/cloudrun/GoogleCloudRunMetadata.java | 249 ++++++++++++++++++ settings.gradle | 2 + 4 files changed, 329 insertions(+) create mode 100644 contrib/temporal-gcp-cloud-run/README.md create mode 100644 contrib/temporal-gcp-cloud-run/build.gradle create mode 100644 contrib/temporal-gcp-cloud-run/src/main/java/io/temporal/gcp/cloudrun/GoogleCloudRunMetadata.java diff --git a/contrib/temporal-gcp-cloud-run/README.md b/contrib/temporal-gcp-cloud-run/README.md new file mode 100644 index 0000000000..a0e3da1763 --- /dev/null +++ b/contrib/temporal-gcp-cloud-run/README.md @@ -0,0 +1,71 @@ +# Temporal Google Cloud Run support + +This module derives a Temporal worker identity and a `WorkerDeploymentVersion` from Google Cloud Run instance metadata, for both Cloud Run **worker pools** and Cloud Run **services**. + +Cloud Run runs a long-lived container, so there is no per-request handler to wrap. This module is a small metadata helper rather than a worker wrapper: fetch the metadata once during startup and apply it to your client and worker option builders. + +> Experimental: Google Cloud Run support is experimental and may change without notice. + +## Quick start + +Add `temporal-gcp-cloud-run` next to your Temporal SDK dependency, then fetch the metadata and apply it while the worker starts up: + +```java +import io.temporal.client.WorkflowClient; +import io.temporal.client.WorkflowClientOptions; +import io.temporal.gcp.cloudrun.GoogleCloudRunMetadata; +import io.temporal.serviceclient.WorkflowServiceStubs; +import io.temporal.serviceclient.WorkflowServiceStubsOptions; +import io.temporal.worker.Worker; +import io.temporal.worker.WorkerFactory; +import io.temporal.worker.WorkerOptions; + +public final class Main { + public static void main(String[] args) { + // Read Cloud Run instance metadata once during startup. + GoogleCloudRunMetadata metadata = GoogleCloudRunMetadata.fetch(); + + WorkflowServiceStubs service = + WorkflowServiceStubs.newServiceStubs( + WorkflowServiceStubsOptions.newBuilder() + .setTarget("my-namespace.tmprl.cloud:7233") + .build()); + + // applyTo(...) sets the derived worker identity on the client options. + WorkflowClient client = + WorkflowClient.newInstance( + service, metadata.applyTo(WorkflowClientOptions.newBuilder()).build()); + + WorkerFactory factory = WorkerFactory.newInstance(client); + + // applyTo(...) sets the deployment version and enables worker versioning on the worker options. + WorkerOptions workerOptions = metadata.applyTo(WorkerOptions.newBuilder()).build(); + + Worker worker = factory.newWorker("orders", workerOptions); + worker.registerWorkflowImplementationTypes(OrderWorkflowImpl.class); + worker.registerActivitiesImplementations(new OrderActivitiesImpl()); + + factory.start(); + } +} +``` + +Both `applyTo(...)` methods return the builder they were given, so they compose with the rest of your builder configuration. + +## How it works + +`GoogleCloudRunMetadata.fetch()` resolves three values: + +- **name** (the Temporal deployment name): the first non-empty of `CLOUD_RUN_WORKER_POOL` (set on Cloud Run worker pools) then `K_SERVICE` (set on Cloud Run services). +- **revision**: the first non-empty of `CLOUD_RUN_REVISION` (worker pools) then `K_REVISION` (services). +- **instanceId**: read from the Cloud Run metadata server with a single HTTP `GET` to `http://metadata.google.internal/computeMetadata/v1/instance/id` with the required `Metadata-Flavor: Google` header. The metadata server is available on both worker pools and services. + +Worker pools receive `CLOUD_RUN_WORKER_POOL` and `CLOUD_RUN_REVISION` and no `K_*` variables, while services receive `K_SERVICE` and `K_REVISION`, so resolving each value from the worker-pool variable first and the service variable second supports both. + +`workerIdentity()` returns `@`, falling back to `@` and then to the bare `` when those values are blank. `workerDeploymentVersion()` maps the name to the deployment name and the revision to the build id, so each Cloud Run revision becomes a distinct `WorkerDeploymentVersion`. + +The two `applyTo(...)` overloads mirror the SDK's "apply defaults to your options" idiom: `applyTo(WorkflowClientOptions.Builder)` sets the worker identity on the client side, and `applyTo(WorkerOptions.Builder)` sets the deployment version (with versioning enabled) on the worker side. Each returns the builder for chaining. If you prefer to read the values yourself, call `workerIdentity()` and `workerDeploymentVersion()` directly. + +Because the metadata server is only reachable from a Cloud Run instance, `fetch()` throws `IllegalStateException` when it cannot be reached, and `workerDeploymentVersion()` (and therefore `applyTo(WorkerOptions.Builder)`) throws `IllegalStateException` when the name or revision is not set. Use `GoogleCloudRunMetadata.fetch(String metadataUrl, Duration timeout)` to override the metadata URL or the request timeout. + +This module depends only on the Temporal SDK at compile time and uses the JDK's `HttpURLConnection` for the metadata request, so it adds no additional runtime dependencies. diff --git a/contrib/temporal-gcp-cloud-run/build.gradle b/contrib/temporal-gcp-cloud-run/build.gradle new file mode 100644 index 0000000000..d08b6289f1 --- /dev/null +++ b/contrib/temporal-gcp-cloud-run/build.gradle @@ -0,0 +1,7 @@ +description = '''Temporal Google Cloud Run support''' + +dependencies { + // This module shouldn't carry temporal-sdk with it, especially for situations when users may + // be using a shaded artifact. + compileOnly project(':temporal-sdk') +} diff --git a/contrib/temporal-gcp-cloud-run/src/main/java/io/temporal/gcp/cloudrun/GoogleCloudRunMetadata.java b/contrib/temporal-gcp-cloud-run/src/main/java/io/temporal/gcp/cloudrun/GoogleCloudRunMetadata.java new file mode 100644 index 0000000000..40c69adecf --- /dev/null +++ b/contrib/temporal-gcp-cloud-run/src/main/java/io/temporal/gcp/cloudrun/GoogleCloudRunMetadata.java @@ -0,0 +1,249 @@ +package io.temporal.gcp.cloudrun; + +import io.temporal.client.WorkflowClientOptions; +import io.temporal.common.Experimental; +import io.temporal.common.WorkerDeploymentVersion; +import io.temporal.worker.WorkerDeploymentOptions; +import io.temporal.worker.WorkerOptions; +import java.io.ByteArrayOutputStream; +import java.io.IOException; +import java.io.InputStream; +import java.net.HttpURLConnection; +import java.net.URL; +import java.nio.charset.StandardCharsets; +import java.time.Duration; +import java.util.Objects; + +/** + * Reads Google Cloud Run instance metadata and derives a Temporal worker identity and a {@link + * WorkerDeploymentVersion} from it. + * + *

Cloud Run runs a long-lived container rather than a per-request handler, so this class is a + * metadata helper rather than a worker wrapper. Fetch the metadata once while a worker starts up, + * then apply it to your client and worker option builders with {@link + * #applyTo(WorkflowClientOptions.Builder)} and {@link #applyTo(WorkerOptions.Builder)}. + * + *

The deployment name and revision are resolved from environment variables Cloud Run injects + * into every instance. Cloud Run worker pools set {@code CLOUD_RUN_WORKER_POOL} and {@code + * CLOUD_RUN_REVISION}; Cloud Run services set {@code K_SERVICE} and {@code K_REVISION}. The + * name is the first non-empty of {@code CLOUD_RUN_WORKER_POOL} then {@code K_SERVICE}, and the + * revision is the first non-empty of {@code CLOUD_RUN_REVISION} then {@code K_REVISION}. The unique + * instance id is only available from the Cloud Run metadata server, so {@link #fetch()} performs a + * single HTTP request against it. + * + *

Experimental: Google Cloud Run support is experimental and may change without notice. + */ +@Experimental +public final class GoogleCloudRunMetadata { + /** Name of the environment variable Cloud Run worker pools set to the worker pool name. */ + public static final String CLOUD_RUN_WORKER_POOL = "CLOUD_RUN_WORKER_POOL"; + + /** Name of the environment variable Cloud Run worker pools set to the revision name. */ + public static final String CLOUD_RUN_REVISION = "CLOUD_RUN_REVISION"; + + /** Name of the environment variable Cloud Run services set to the deployed service name. */ + public static final String K_SERVICE = "K_SERVICE"; + + /** Name of the environment variable Cloud Run services set to the deployed revision name. */ + public static final String K_REVISION = "K_REVISION"; + + /** Default Cloud Run metadata server URL that returns the unique instance id. */ + public static final String DEFAULT_METADATA_URL = + "http://metadata.google.internal/computeMetadata/v1/instance/id"; + + /** Default connect and read timeout used when contacting the metadata server. */ + public static final Duration DEFAULT_TIMEOUT = Duration.ofSeconds(2); + + private static final String METADATA_FLAVOR_HEADER = "Metadata-Flavor"; + private static final String METADATA_FLAVOR_VALUE = "Google"; + + private final String instanceId; + private final String name; + private final String revision; + + private GoogleCloudRunMetadata(String instanceId, String name, String revision) { + this.instanceId = instanceId; + this.name = name; + this.revision = revision; + } + + /** + * Fetches Cloud Run instance metadata using the {@linkplain #DEFAULT_METADATA_URL default + * metadata URL} and the {@linkplain #DEFAULT_TIMEOUT default timeout}. + * + * @return metadata describing the current Cloud Run instance. + * @throws IllegalStateException if the metadata server cannot be reached, which usually means the + * process is not running on Google Cloud Run. + */ + public static GoogleCloudRunMetadata fetch() { + return fetch(DEFAULT_METADATA_URL, DEFAULT_TIMEOUT); + } + + /** + * Fetches Cloud Run instance metadata from the supplied metadata server URL. + * + *

The deployment name is read from {@code CLOUD_RUN_WORKER_POOL} then {@code K_SERVICE}, and + * the revision from {@code CLOUD_RUN_REVISION} then {@code K_REVISION}. The unique instance id is + * read from {@code metadataUrl} with the required {@code Metadata-Flavor: Google} request header. + * + * @param metadataUrl URL of the Cloud Run metadata endpoint that returns the instance id. + * @param timeout connect and read timeout applied to the metadata request. + * @return metadata describing the current Cloud Run instance. + * @throws IllegalStateException if the metadata server cannot be reached, which usually means the + * process is not running on Google Cloud Run. + */ + public static GoogleCloudRunMetadata fetch(String metadataUrl, Duration timeout) { + Objects.requireNonNull(metadataUrl, "metadataUrl"); + Objects.requireNonNull(timeout, "timeout"); + + String name = firstNonBlank(System.getenv(CLOUD_RUN_WORKER_POOL), System.getenv(K_SERVICE)); + String revision = firstNonBlank(System.getenv(CLOUD_RUN_REVISION), System.getenv(K_REVISION)); + + HttpURLConnection connection = null; + try { + connection = (HttpURLConnection) new URL(metadataUrl).openConnection(); + connection.setRequestMethod("GET"); + connection.setRequestProperty(METADATA_FLAVOR_HEADER, METADATA_FLAVOR_VALUE); + int timeoutMillis = timeoutMillis(timeout); + connection.setConnectTimeout(timeoutMillis); + connection.setReadTimeout(timeoutMillis); + + String instanceId = readBody(connection).trim(); + return new GoogleCloudRunMetadata(instanceId, name, revision); + } catch (IOException e) { + throw new IllegalStateException( + "Unable to read the Cloud Run instance id from the metadata server at " + + metadataUrl + + "; this process may not be running on Google Cloud Run", + e); + } finally { + if (connection != null) { + connection.disconnect(); + } + } + } + + /** + * @return the unique Cloud Run instance id read from the metadata server. + */ + public String getInstanceId() { + return instanceId; + } + + /** + * @return the Cloud Run deployment name, resolved from {@code CLOUD_RUN_WORKER_POOL} then {@code + * K_SERVICE}, or {@code null} when neither was set. + */ + public String getName() { + return name; + } + + /** + * @return the Cloud Run revision name, resolved from {@code CLOUD_RUN_REVISION} then {@code + * K_REVISION}, or {@code null} when neither was set. + */ + public String getRevision() { + return revision; + } + + /** + * Builds a Temporal worker identity for this Cloud Run instance. + * + *

The identity is {@code instanceId@revision}. When the revision is blank the name is used + * instead, and when both are blank the bare instance id is returned. + * + * @return a worker identity string suitable for {@code WorkflowClientOptions} and {@code + * WorkerOptions}. + */ + public String workerIdentity() { + if (!isBlank(revision)) { + return instanceId + "@" + revision; + } + if (!isBlank(name)) { + return instanceId + "@" + name; + } + return instanceId; + } + + /** + * Builds a {@link WorkerDeploymentVersion} from the Cloud Run name and revision. + * + *

The name becomes the deployment name and the revision becomes the build id, so each Cloud + * Run revision maps to a distinct worker deployment version. + * + * @return a worker deployment version derived from the resolved name and revision. + * @throws IllegalStateException if the name or revision is blank, which usually means the process + * is not running on a Cloud Run worker pool or service. + */ + public WorkerDeploymentVersion workerDeploymentVersion() { + if (isBlank(name) || isBlank(revision)) { + throw new IllegalStateException( + "A Cloud Run name and revision are required to build a WorkerDeploymentVersion; " + + "this process may not be running on a Cloud Run worker pool or service"); + } + return new WorkerDeploymentVersion(name, revision); + } + + /** + * Applies the derived {@linkplain #workerIdentity() worker identity} to a workflow client options + * builder. + * + * @param builder the workflow client options builder to configure. + * @return the same builder, for chaining. + */ + public WorkflowClientOptions.Builder applyTo(WorkflowClientOptions.Builder builder) { + Objects.requireNonNull(builder, "builder"); + builder.setIdentity(workerIdentity()); + return builder; + } + + /** + * Applies the derived {@linkplain #workerDeploymentVersion() worker deployment version} to a + * worker options builder, enabling worker versioning. + * + * @param builder the worker options builder to configure. + * @return the same builder, for chaining. + * @throws IllegalStateException if the name or revision is blank, which usually means the process + * is not running on a Cloud Run worker pool or service. + */ + public WorkerOptions.Builder applyTo(WorkerOptions.Builder builder) { + Objects.requireNonNull(builder, "builder"); + builder.setDeploymentOptions( + WorkerDeploymentOptions.newBuilder() + .setUseVersioning(true) + .setVersion(workerDeploymentVersion()) + .build()); + return builder; + } + + private static String readBody(HttpURLConnection connection) throws IOException { + try (InputStream in = connection.getInputStream()) { + ByteArrayOutputStream out = new ByteArrayOutputStream(); + byte[] chunk = new byte[512]; + int read; + while ((read = in.read(chunk)) != -1) { + out.write(chunk, 0, read); + } + return new String(out.toByteArray(), StandardCharsets.UTF_8); + } + } + + private static int timeoutMillis(Duration timeout) { + long millis = timeout.toMillis(); + if (millis < 0) { + throw new IllegalArgumentException("timeout must not be negative"); + } + return (int) Math.min(millis, Integer.MAX_VALUE); + } + + private static String firstNonBlank(String first, String second) { + if (!isBlank(first)) { + return first; + } + return isBlank(second) ? null : second; + } + + private static boolean isBlank(String value) { + return value == null || value.trim().isEmpty(); + } +} diff --git a/settings.gradle b/settings.gradle index 3699ff1508..6cbf879490 100644 --- a/settings.gradle +++ b/settings.gradle @@ -15,6 +15,8 @@ include 'temporal-workflowstreams' project(':temporal-workflowstreams').projectDir = file('contrib/temporal-workflowstreams') include 'temporal-aws-lambda' project(':temporal-aws-lambda').projectDir = file('contrib/temporal-aws-lambda') +include 'temporal-gcp-cloud-run' +project(':temporal-gcp-cloud-run').projectDir = file('contrib/temporal-gcp-cloud-run') include 'temporal-spring-boot-autoconfigure' include 'temporal-spring-boot-starter' include 'temporal-remote-data-encoder' From 923be731bc0d799751ee17a6e209a49b006f0515 Mon Sep 17 00:00:00 2001 From: seanbollin Date: Tue, 25 Aug 2026 13:17:23 -0700 Subject: [PATCH 2/4] Set worker versioning behavior to PINNED in the Cloud Run worker apply The worker-side apply helper enabled versioning and set the deployment version but left the default versioning behavior unset, so a versioned worker with a plain (un-annotated) workflow failed to register. Default it to PINNED; a per-workflow versioning behavior still takes precedence. Co-Authored-By: Claude Opus 4.8 --- contrib/temporal-gcp-cloud-run/README.md | 2 +- .../java/io/temporal/gcp/cloudrun/GoogleCloudRunMetadata.java | 4 +++- 2 files changed, 4 insertions(+), 2 deletions(-) diff --git a/contrib/temporal-gcp-cloud-run/README.md b/contrib/temporal-gcp-cloud-run/README.md index a0e3da1763..c078d4fc42 100644 --- a/contrib/temporal-gcp-cloud-run/README.md +++ b/contrib/temporal-gcp-cloud-run/README.md @@ -64,7 +64,7 @@ Worker pools receive `CLOUD_RUN_WORKER_POOL` and `CLOUD_RUN_REVISION` and no `K_ `workerIdentity()` returns `@`, falling back to `@` and then to the bare `` when those values are blank. `workerDeploymentVersion()` maps the name to the deployment name and the revision to the build id, so each Cloud Run revision becomes a distinct `WorkerDeploymentVersion`. -The two `applyTo(...)` overloads mirror the SDK's "apply defaults to your options" idiom: `applyTo(WorkflowClientOptions.Builder)` sets the worker identity on the client side, and `applyTo(WorkerOptions.Builder)` sets the deployment version (with versioning enabled) on the worker side. Each returns the builder for chaining. If you prefer to read the values yourself, call `workerIdentity()` and `workerDeploymentVersion()` directly. +The two `applyTo(...)` overloads mirror the SDK's "apply defaults to your options" idiom: `applyTo(WorkflowClientOptions.Builder)` sets the worker identity on the client side, and `applyTo(WorkerOptions.Builder)` sets the deployment version (with versioning enabled, pinning workflows to this version by default via `VersioningBehavior.PINNED`; a per-workflow behavior takes precedence) on the worker side. Each returns the builder for chaining. If you prefer to read the values yourself, call `workerIdentity()` and `workerDeploymentVersion()` directly. Because the metadata server is only reachable from a Cloud Run instance, `fetch()` throws `IllegalStateException` when it cannot be reached, and `workerDeploymentVersion()` (and therefore `applyTo(WorkerOptions.Builder)`) throws `IllegalStateException` when the name or revision is not set. Use `GoogleCloudRunMetadata.fetch(String metadataUrl, Duration timeout)` to override the metadata URL or the request timeout. diff --git a/contrib/temporal-gcp-cloud-run/src/main/java/io/temporal/gcp/cloudrun/GoogleCloudRunMetadata.java b/contrib/temporal-gcp-cloud-run/src/main/java/io/temporal/gcp/cloudrun/GoogleCloudRunMetadata.java index 40c69adecf..0edd847afe 100644 --- a/contrib/temporal-gcp-cloud-run/src/main/java/io/temporal/gcp/cloudrun/GoogleCloudRunMetadata.java +++ b/contrib/temporal-gcp-cloud-run/src/main/java/io/temporal/gcp/cloudrun/GoogleCloudRunMetadata.java @@ -2,6 +2,7 @@ import io.temporal.client.WorkflowClientOptions; import io.temporal.common.Experimental; +import io.temporal.common.VersioningBehavior; import io.temporal.common.WorkerDeploymentVersion; import io.temporal.worker.WorkerDeploymentOptions; import io.temporal.worker.WorkerOptions; @@ -199,7 +200,7 @@ public WorkflowClientOptions.Builder applyTo(WorkflowClientOptions.Builder build /** * Applies the derived {@linkplain #workerDeploymentVersion() worker deployment version} to a - * worker options builder, enabling worker versioning. + * worker options builder, enabling worker versioning with a PINNED default behavior. * * @param builder the worker options builder to configure. * @return the same builder, for chaining. @@ -212,6 +213,7 @@ public WorkerOptions.Builder applyTo(WorkerOptions.Builder builder) { WorkerDeploymentOptions.newBuilder() .setUseVersioning(true) .setVersion(workerDeploymentVersion()) + .setDefaultVersioningBehavior(VersioningBehavior.PINNED) .build()); return builder; } From 31a6d7efcbd300082eb248874a46ba8b67f8fbd4 Mon Sep 17 00:00:00 2001 From: seanbollin Date: Tue, 25 Aug 2026 13:42:58 -0700 Subject: [PATCH 3/4] Add unit tests for GoogleCloudRunMetadata Cover the Cloud Run metadata helper: environment-variable precedence (CLOUD_RUN_WORKER_POOL over K_SERVICE, CLOUD_RUN_REVISION over K_REVISION), worker identity fallbacks, WorkerDeploymentVersion mapping and its empty name/revision error, the metadata HTTP request (Metadata-Flavor: Google header, body trimming, non-200 and unreachable errors), and the applyTo(...) methods (client identity, and PINNED worker deployment versioning). The metadata request is served by an in-process com.sun.net.httpserver HttpServer and the environment lookup is injected through a new package-private fetch(String, Duration, Function) test seam, so the tests touch neither the network nor the real process environment. The seam is not part of the public API and does not change public behavior. Add the matching testImplementation dependencies (temporal-sdk, junit) to the module, mirroring the temporal-aws-lambda module. Co-Authored-By: Claude Opus 4.8 --- contrib/temporal-gcp-cloud-run/build.gradle | 5 + .../gcp/cloudrun/GoogleCloudRunMetadata.java | 22 +- .../cloudrun/GoogleCloudRunMetadataTest.java | 281 ++++++++++++++++++ 3 files changed, 306 insertions(+), 2 deletions(-) create mode 100644 contrib/temporal-gcp-cloud-run/src/test/java/io/temporal/gcp/cloudrun/GoogleCloudRunMetadataTest.java diff --git a/contrib/temporal-gcp-cloud-run/build.gradle b/contrib/temporal-gcp-cloud-run/build.gradle index d08b6289f1..831e558db5 100644 --- a/contrib/temporal-gcp-cloud-run/build.gradle +++ b/contrib/temporal-gcp-cloud-run/build.gradle @@ -4,4 +4,9 @@ dependencies { // This module shouldn't carry temporal-sdk with it, especially for situations when users may // be using a shaded artifact. compileOnly project(':temporal-sdk') + + testImplementation project(':temporal-sdk') + testImplementation "junit:junit:${junitVersion}" + + testRuntimeOnly group: 'ch.qos.logback', name: 'logback-classic', version: "${logbackVersion}" } diff --git a/contrib/temporal-gcp-cloud-run/src/main/java/io/temporal/gcp/cloudrun/GoogleCloudRunMetadata.java b/contrib/temporal-gcp-cloud-run/src/main/java/io/temporal/gcp/cloudrun/GoogleCloudRunMetadata.java index 0edd847afe..c0cd07091f 100644 --- a/contrib/temporal-gcp-cloud-run/src/main/java/io/temporal/gcp/cloudrun/GoogleCloudRunMetadata.java +++ b/contrib/temporal-gcp-cloud-run/src/main/java/io/temporal/gcp/cloudrun/GoogleCloudRunMetadata.java @@ -14,6 +14,7 @@ import java.nio.charset.StandardCharsets; import java.time.Duration; import java.util.Objects; +import java.util.function.Function; /** * Reads Google Cloud Run instance metadata and derives a Temporal worker identity and a {@link @@ -94,11 +95,28 @@ public static GoogleCloudRunMetadata fetch() { * process is not running on Google Cloud Run. */ public static GoogleCloudRunMetadata fetch(String metadataUrl, Duration timeout) { + return fetch(metadataUrl, timeout, System::getenv); + } + + /** + * Package-private test seam that injects the environment-variable lookup used to resolve the + * deployment name and revision. This lets unit tests exercise the environment-variable precedence + * and the metadata HTTP request deterministically, without depending on the real process + * environment. It is not part of the public API and must not be relied on outside of tests; use + * {@link #fetch(String, Duration)} instead. + * + * @param metadataUrl URL of the Cloud Run metadata endpoint that returns the instance id. + * @param timeout connect and read timeout applied to the metadata request. + * @param getenv environment-variable lookup, normally {@code System::getenv}. + */ + static GoogleCloudRunMetadata fetch( + String metadataUrl, Duration timeout, Function getenv) { Objects.requireNonNull(metadataUrl, "metadataUrl"); Objects.requireNonNull(timeout, "timeout"); + Objects.requireNonNull(getenv, "getenv"); - String name = firstNonBlank(System.getenv(CLOUD_RUN_WORKER_POOL), System.getenv(K_SERVICE)); - String revision = firstNonBlank(System.getenv(CLOUD_RUN_REVISION), System.getenv(K_REVISION)); + String name = firstNonBlank(getenv.apply(CLOUD_RUN_WORKER_POOL), getenv.apply(K_SERVICE)); + String revision = firstNonBlank(getenv.apply(CLOUD_RUN_REVISION), getenv.apply(K_REVISION)); HttpURLConnection connection = null; try { diff --git a/contrib/temporal-gcp-cloud-run/src/test/java/io/temporal/gcp/cloudrun/GoogleCloudRunMetadataTest.java b/contrib/temporal-gcp-cloud-run/src/test/java/io/temporal/gcp/cloudrun/GoogleCloudRunMetadataTest.java new file mode 100644 index 0000000000..fa0f3af2f6 --- /dev/null +++ b/contrib/temporal-gcp-cloud-run/src/test/java/io/temporal/gcp/cloudrun/GoogleCloudRunMetadataTest.java @@ -0,0 +1,281 @@ +package io.temporal.gcp.cloudrun; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertNull; +import static org.junit.Assert.assertThrows; +import static org.junit.Assert.assertTrue; + +import com.sun.net.httpserver.HttpServer; +import io.temporal.client.WorkflowClientOptions; +import io.temporal.common.VersioningBehavior; +import io.temporal.common.WorkerDeploymentVersion; +import io.temporal.worker.WorkerDeploymentOptions; +import io.temporal.worker.WorkerOptions; +import java.io.IOException; +import java.io.OutputStream; +import java.net.InetSocketAddress; +import java.net.ServerSocket; +import java.nio.charset.StandardCharsets; +import java.time.Duration; +import java.util.HashMap; +import java.util.Map; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicReference; +import org.junit.After; +import org.junit.Before; +import org.junit.Test; + +/** + * Unit tests for {@link GoogleCloudRunMetadata}. + * + *

The metadata request is served by an in-process {@link HttpServer} and the environment lookup + * is injected through the package-private {@link GoogleCloudRunMetadata#fetch(String, Duration, + * java.util.function.Function)} test seam, so these tests touch neither the network nor the real + * process environment. + */ +public class GoogleCloudRunMetadataTest { + private static final Duration TIMEOUT = Duration.ofSeconds(2); + + private HttpServer server; + private final AtomicReference responseBody = new AtomicReference<>(""); + private final AtomicInteger responseStatus = new AtomicInteger(200); + private final AtomicReference capturedMetadataFlavor = new AtomicReference<>(); + private final AtomicReference capturedMethod = new AtomicReference<>(); + + @Before + public void startServer() throws IOException { + server = HttpServer.create(new InetSocketAddress("127.0.0.1", 0), 0); + server.createContext( + "/computeMetadata/v1/instance/id", + exchange -> { + capturedMetadataFlavor.set(exchange.getRequestHeaders().getFirst("Metadata-Flavor")); + capturedMethod.set(exchange.getRequestMethod()); + byte[] body = responseBody.get().getBytes(StandardCharsets.UTF_8); + exchange.sendResponseHeaders(responseStatus.get(), body.length == 0 ? -1 : body.length); + try (OutputStream out = exchange.getResponseBody()) { + out.write(body); + } + }); + server.start(); + } + + @After + public void stopServer() { + server.stop(0); + } + + // --- Environment-variable precedence --- + + @Test + public void cloudRunWorkerPoolWinsOverKService() { + responseBody.set("instance-1"); + Map env = new HashMap<>(); + env.put(GoogleCloudRunMetadata.CLOUD_RUN_WORKER_POOL, "worker-pool"); + env.put(GoogleCloudRunMetadata.K_SERVICE, "service"); + env.put(GoogleCloudRunMetadata.CLOUD_RUN_REVISION, "worker-pool-revision"); + env.put(GoogleCloudRunMetadata.K_REVISION, "service-revision"); + + GoogleCloudRunMetadata metadata = fetch(env); + + assertEquals("worker-pool", metadata.getName()); + assertEquals("worker-pool-revision", metadata.getRevision()); + } + + @Test + public void kServiceUsedWhenWorkerPoolAbsent() { + responseBody.set("instance-1"); + Map env = new HashMap<>(); + env.put(GoogleCloudRunMetadata.K_SERVICE, "service"); + env.put(GoogleCloudRunMetadata.K_REVISION, "service-revision"); + + GoogleCloudRunMetadata metadata = fetch(env); + + assertEquals("service", metadata.getName()); + assertEquals("service-revision", metadata.getRevision()); + } + + @Test + public void blankWorkerPoolVariablesFallThroughToKService() { + responseBody.set("instance-1"); + Map env = new HashMap<>(); + env.put(GoogleCloudRunMetadata.CLOUD_RUN_WORKER_POOL, " "); + env.put(GoogleCloudRunMetadata.CLOUD_RUN_REVISION, ""); + env.put(GoogleCloudRunMetadata.K_SERVICE, "service"); + env.put(GoogleCloudRunMetadata.K_REVISION, "service-revision"); + + GoogleCloudRunMetadata metadata = fetch(env); + + assertEquals("service", metadata.getName()); + assertEquals("service-revision", metadata.getRevision()); + } + + @Test + public void nameAndRevisionAreNullWhenNoEnvSet() { + responseBody.set("instance-1"); + + GoogleCloudRunMetadata metadata = fetch(new HashMap<>()); + + assertNull(metadata.getName()); + assertNull(metadata.getRevision()); + } + + // --- Worker identity --- + + @Test + public void workerIdentityCombinesInstanceIdAndRevision() { + responseBody.set("instance-1"); + Map env = new HashMap<>(); + env.put(GoogleCloudRunMetadata.CLOUD_RUN_WORKER_POOL, "worker-pool"); + env.put(GoogleCloudRunMetadata.CLOUD_RUN_REVISION, "revision-1"); + + assertEquals("instance-1@revision-1", fetch(env).workerIdentity()); + } + + @Test + public void workerIdentityFallsBackToNameWhenRevisionBlank() { + responseBody.set("instance-1"); + Map env = new HashMap<>(); + env.put(GoogleCloudRunMetadata.CLOUD_RUN_WORKER_POOL, "worker-pool"); + + assertEquals("instance-1@worker-pool", fetch(env).workerIdentity()); + } + + @Test + public void workerIdentityFallsBackToInstanceIdWhenNameAndRevisionBlank() { + responseBody.set("instance-1"); + + assertEquals("instance-1", fetch(new HashMap<>()).workerIdentity()); + } + + // --- Worker deployment version --- + + @Test + public void workerDeploymentVersionMapsNameToDeploymentAndRevisionToBuildId() { + responseBody.set("instance-1"); + Map env = new HashMap<>(); + env.put(GoogleCloudRunMetadata.CLOUD_RUN_WORKER_POOL, "worker-pool"); + env.put(GoogleCloudRunMetadata.CLOUD_RUN_REVISION, "revision-1"); + + WorkerDeploymentVersion version = fetch(env).workerDeploymentVersion(); + + assertEquals("worker-pool", version.getDeploymentName()); + assertEquals("revision-1", version.getBuildId()); + } + + @Test + public void workerDeploymentVersionRequiresName() { + responseBody.set("instance-1"); + Map env = new HashMap<>(); + env.put(GoogleCloudRunMetadata.CLOUD_RUN_REVISION, "revision-1"); + + IllegalStateException e = + assertThrows(IllegalStateException.class, () -> fetch(env).workerDeploymentVersion()); + assertTrue(e.getMessage().contains("name and revision")); + } + + @Test + public void workerDeploymentVersionRequiresRevision() { + responseBody.set("instance-1"); + Map env = new HashMap<>(); + env.put(GoogleCloudRunMetadata.CLOUD_RUN_WORKER_POOL, "worker-pool"); + + assertThrows(IllegalStateException.class, () -> fetch(env).workerDeploymentVersion()); + } + + // --- Metadata HTTP request --- + + @Test + public void fetchSendsMetadataFlavorHeaderAndTrimsBody() { + responseBody.set(" instance-42\n"); + Map env = new HashMap<>(); + env.put(GoogleCloudRunMetadata.CLOUD_RUN_WORKER_POOL, "worker-pool"); + env.put(GoogleCloudRunMetadata.CLOUD_RUN_REVISION, "revision-1"); + + GoogleCloudRunMetadata metadata = fetch(env); + + assertEquals("instance-42", metadata.getInstanceId()); + assertEquals("Google", capturedMetadataFlavor.get()); + assertEquals("GET", capturedMethod.get()); + } + + @Test + public void fetchThrowsOnNonSuccessStatus() { + responseStatus.set(500); + responseBody.set("boom"); + + IllegalStateException e = + assertThrows(IllegalStateException.class, () -> fetch(new HashMap<>())); + assertTrue(e.getMessage().contains("metadata server")); + } + + @Test + public void fetchThrowsWhenServerUnreachable() { + String unreachableUrl = + "http://127.0.0.1:" + reserveUnusedPort() + "/computeMetadata/v1/instance/id"; + Map env = new HashMap<>(); + + assertThrows( + IllegalStateException.class, + () -> GoogleCloudRunMetadata.fetch(unreachableUrl, TIMEOUT, env::get)); + } + + // --- Apply methods --- + + @Test + public void applyToWorkflowClientOptionsSetsDerivedIdentity() { + responseBody.set("instance-1"); + Map env = new HashMap<>(); + env.put(GoogleCloudRunMetadata.CLOUD_RUN_WORKER_POOL, "worker-pool"); + env.put(GoogleCloudRunMetadata.CLOUD_RUN_REVISION, "revision-1"); + + // NOTE: The shared cross-SDK design and the Go and .NET helpers set the client identity only + // when it is unset, so a user-provided identity is preserved. The Java + // applyTo(WorkflowClientOptions.Builder) currently sets it unconditionally. This test asserts + // only the "identity unset" case, which holds under either behavior; the user-provided-identity + // case is intentionally not asserted here while that divergence is resolved. + WorkflowClientOptions options = fetch(env).applyTo(WorkflowClientOptions.newBuilder()).build(); + + assertEquals("instance-1@revision-1", options.getIdentity()); + } + + @Test + public void applyToWorkerOptionsEnablesPinnedVersioning() { + responseBody.set("instance-1"); + Map env = new HashMap<>(); + env.put(GoogleCloudRunMetadata.CLOUD_RUN_WORKER_POOL, "worker-pool"); + env.put(GoogleCloudRunMetadata.CLOUD_RUN_REVISION, "revision-1"); + + WorkerOptions options = fetch(env).applyTo(WorkerOptions.newBuilder()).build(); + + WorkerDeploymentOptions deploymentOptions = options.getDeploymentOptions(); + assertTrue(deploymentOptions.isUsingVersioning()); + assertEquals( + new WorkerDeploymentVersion("worker-pool", "revision-1"), deploymentOptions.getVersion()); + assertEquals(VersioningBehavior.PINNED, deploymentOptions.getDefaultVersioningBehavior()); + } + + @Test + public void applyToWorkerOptionsThrowsWhenDeploymentVersionCannotBeBuilt() { + responseBody.set("instance-1"); + + assertThrows( + IllegalStateException.class, + () -> fetch(new HashMap<>()).applyTo(WorkerOptions.newBuilder())); + } + + private GoogleCloudRunMetadata fetch(Map env) { + return GoogleCloudRunMetadata.fetch(metadataUrl(), TIMEOUT, env::get); + } + + private String metadataUrl() { + return "http://127.0.0.1:" + server.getAddress().getPort() + "/computeMetadata/v1/instance/id"; + } + + private static int reserveUnusedPort() { + try (ServerSocket socket = new ServerSocket(0)) { + return socket.getLocalPort(); + } catch (IOException e) { + throw new RuntimeException(e); + } + } +} From b6552db6bb7a962703988a4d87f49b579639274e Mon Sep 17 00:00:00 2001 From: seanbollin Date: Tue, 25 Aug 2026 15:01:38 -0700 Subject: [PATCH 4/4] Avoid deprecated URL(String) constructor in GoogleCloudRunMetadata new URL(String) is deprecated since Java 20 and fails the SDK's -Werror build on newer JDKs (the Java 23 "Edge" CI job). Use URI.create(...).toURL() instead, the recommended non-deprecated replacement (MalformedURLException is still an IOException and stays caught). Co-Authored-By: Claude Opus 4.8 --- .../java/io/temporal/gcp/cloudrun/GoogleCloudRunMetadata.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/contrib/temporal-gcp-cloud-run/src/main/java/io/temporal/gcp/cloudrun/GoogleCloudRunMetadata.java b/contrib/temporal-gcp-cloud-run/src/main/java/io/temporal/gcp/cloudrun/GoogleCloudRunMetadata.java index c0cd07091f..9d34bfc0e0 100644 --- a/contrib/temporal-gcp-cloud-run/src/main/java/io/temporal/gcp/cloudrun/GoogleCloudRunMetadata.java +++ b/contrib/temporal-gcp-cloud-run/src/main/java/io/temporal/gcp/cloudrun/GoogleCloudRunMetadata.java @@ -10,7 +10,7 @@ import java.io.IOException; import java.io.InputStream; import java.net.HttpURLConnection; -import java.net.URL; +import java.net.URI; import java.nio.charset.StandardCharsets; import java.time.Duration; import java.util.Objects; @@ -120,7 +120,7 @@ static GoogleCloudRunMetadata fetch( HttpURLConnection connection = null; try { - connection = (HttpURLConnection) new URL(metadataUrl).openConnection(); + connection = (HttpURLConnection) URI.create(metadataUrl).toURL().openConnection(); connection.setRequestMethod("GET"); connection.setRequestProperty(METADATA_FLAVOR_HEADER, METADATA_FLAVOR_VALUE); int timeoutMillis = timeoutMillis(timeout);