diff --git a/contrib/temporal-gcp-cloud-run/README.md b/contrib/temporal-gcp-cloud-run/README.md new file mode 100644 index 0000000000..c078d4fc42 --- /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, 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. + +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..831e558db5 --- /dev/null +++ b/contrib/temporal-gcp-cloud-run/build.gradle @@ -0,0 +1,12 @@ +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') + + 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 new file mode 100644 index 0000000000..9d34bfc0e0 --- /dev/null +++ b/contrib/temporal-gcp-cloud-run/src/main/java/io/temporal/gcp/cloudrun/GoogleCloudRunMetadata.java @@ -0,0 +1,269 @@ +package io.temporal.gcp.cloudrun; + +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; +import java.io.ByteArrayOutputStream; +import java.io.IOException; +import java.io.InputStream; +import java.net.HttpURLConnection; +import java.net.URI; +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 + * 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) { + 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(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 { + connection = (HttpURLConnection) URI.create(metadataUrl).toURL().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 with a PINNED default behavior. + * + * @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()) + .setDefaultVersioningBehavior(VersioningBehavior.PINNED) + .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/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); + } + } +} 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'