diff --git a/contrib/temporal-gcp-cloud-run-worker-id/README.md b/contrib/temporal-gcp-cloud-run-worker-id/README.md new file mode 100644 index 0000000000..256fb6d2da --- /dev/null +++ b/contrib/temporal-gcp-cloud-run-worker-id/README.md @@ -0,0 +1,88 @@ +# Temporal Google Cloud Run worker identity support + +This module configures a Temporal worker for Google Cloud Run from instance metadata, for both Cloud Run **worker pools** and Cloud Run **services**. It derives the worker's Temporal identity and its `WorkerDeploymentVersion` from Cloud Run instance metadata, so every Cloud Run revision registers as a distinct, `PINNED` Worker Deployment Version. + +The primary API is `WorkerIdPlugin`. Register it once on your workflow client and it propagates to every worker created from that client, setting the client identity and the worker deployment version automatically. This mirrors the `CloudRunOpenTelemetryPlugin` in the companion `temporal-gcp-cloud-run` module. + +> Experimental: Google Cloud Run support is experimental and may change without notice. + +## Quick start + +Add `temporal-gcp-cloud-run-worker-id` next to your Temporal SDK dependency, then register the plugin on the workflow client options: + +```java +import io.temporal.client.WorkflowClient; +import io.temporal.client.WorkflowClientOptions; +import io.temporal.gcp.cloudrun.workerid.WorkerIdPlugin; +import io.temporal.serviceclient.WorkflowServiceStubs; +import io.temporal.serviceclient.WorkflowServiceStubsOptions; +import io.temporal.worker.Worker; +import io.temporal.worker.WorkerFactory; + +public final class Main { + public static void main(String[] args) { + WorkflowServiceStubs service = + WorkflowServiceStubs.newServiceStubs( + WorkflowServiceStubsOptions.newBuilder() + .setTarget("my-namespace.tmprl.cloud:7233") + .build()); + + // Registering the plugin on the client: + // - reads Cloud Run instance metadata once while the client is configured, and + // - sets the client identity to the derived worker identity (unless you set one yourself). + WorkflowClient client = + WorkflowClient.newInstance( + service, + WorkflowClientOptions.newBuilder() + .setNamespace("my-namespace") + .setPlugins(new WorkerIdPlugin()) + .build()); + + WorkerFactory factory = WorkerFactory.newInstance(client); + + // The plugin propagates from the client to workers and sets each worker's deployment version + // (with worker versioning enabled and a PINNED default behavior). No per-worker wiring needed. + Worker worker = factory.newWorker("orders"); + worker.registerWorkflowImplementationTypes(OrderWorkflowImpl.class); + worker.registerActivitiesImplementations(new OrderActivitiesImpl()); + + factory.start(); + } +} +``` + +You can also register the plugin on `WorkflowServiceStubsOptions.Builder.setPlugins(...)`; from there it propagates to the client and workers as well. + +## How it works + +`WorkerIdPlugin` reads Cloud Run instance metadata through `GoogleCloudRunMetadata`, which 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. + +The plugin then applies the metadata through the SDK's plugin hooks: + +- **Client** (`configureWorkflowClient`): sets the client identity to `@` (falling back to `@` and then the bare ``), but only when you have not already set an identity, so a user-provided identity always wins. The metadata is fetched here, once, and cached. +- **Worker** (`configureWorker`): sets the worker deployment version — the name becomes the deployment name and the revision becomes the build id — with worker versioning enabled and `VersioningBehavior.PINNED` as the default, so in-flight workflows stay on the Cloud Run revision that started them (a per-workflow `@WorkflowVersioningBehavior` takes precedence). + +Because the metadata server is only reachable from a Cloud Run instance, the plugin **fails fast**: the fetch in `configureWorkflowClient` throws `IllegalStateException` when the metadata server cannot be reached (which usually means the process is not running on Google Cloud Run), and `configureWorker` throws `IllegalStateException` when the name or revision is not set (which usually means the process is not running on a Cloud Run worker pool or service). The plugin does not silently no-op off-platform. + +## Reading the metadata directly + +If you prefer to read the values yourself, or to fetch the metadata once and pass it in, use `GoogleCloudRunMetadata` directly: + +```java +GoogleCloudRunMetadata metadata = GoogleCloudRunMetadata.fetch(); +String identity = metadata.workerIdentity(); +WorkerDeploymentVersion version = metadata.workerDeploymentVersion(); + +// Or hand the already-fetched metadata to the plugin to skip its own fetch: +WorkerIdPlugin plugin = new WorkerIdPlugin(metadata); +``` + +`GoogleCloudRunMetadata.fetch(String metadataUrl, Duration timeout)` overrides 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-worker-id/build.gradle b/contrib/temporal-gcp-cloud-run-worker-id/build.gradle new file mode 100644 index 0000000000..f93f338ea2 --- /dev/null +++ b/contrib/temporal-gcp-cloud-run-worker-id/build.gradle @@ -0,0 +1,12 @@ +description = '''Temporal Java SDK Google Cloud Run Worker Identity Support Module''' + +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-worker-id/src/main/java/io/temporal/gcp/cloudrun/workerid/GoogleCloudRunMetadata.java b/contrib/temporal-gcp-cloud-run-worker-id/src/main/java/io/temporal/gcp/cloudrun/workerid/GoogleCloudRunMetadata.java new file mode 100644 index 0000000000..5375a960bd --- /dev/null +++ b/contrib/temporal-gcp-cloud-run-worker-id/src/main/java/io/temporal/gcp/cloudrun/workerid/GoogleCloudRunMetadata.java @@ -0,0 +1,234 @@ +package io.temporal.gcp.cloudrun.workerid; + +import io.temporal.common.Experimental; +import io.temporal.common.WorkerDeploymentVersion; +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. Most applications register {@link WorkerIdPlugin} + * on their workflow client instead of using this class directly; the plugin fetches this metadata + * and applies the derived identity and deployment version to the client and workers. Use this class + * directly to read the {@linkplain #workerIdentity() worker identity} or {@linkplain + * #workerDeploymentVersion() worker deployment version} yourself. + * + *

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); + } + + 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-worker-id/src/main/java/io/temporal/gcp/cloudrun/workerid/WorkerIdPlugin.java b/contrib/temporal-gcp-cloud-run-worker-id/src/main/java/io/temporal/gcp/cloudrun/workerid/WorkerIdPlugin.java new file mode 100644 index 0000000000..a216c04403 --- /dev/null +++ b/contrib/temporal-gcp-cloud-run-worker-id/src/main/java/io/temporal/gcp/cloudrun/workerid/WorkerIdPlugin.java @@ -0,0 +1,157 @@ +package io.temporal.gcp.cloudrun.workerid; + +import io.temporal.client.WorkflowClientOptions; +import io.temporal.common.Experimental; +import io.temporal.common.SimplePlugin; +import io.temporal.common.VersioningBehavior; +import io.temporal.worker.WorkerDeploymentOptions; +import io.temporal.worker.WorkerOptions; +import java.util.Objects; +import java.util.function.Supplier; + +/** + * Plugin that configures a Temporal worker for Google Cloud Run from instance metadata, for both + * Cloud Run worker pools and Cloud Run services. + * + *

Register the plugin once on the workflow client and it propagates to every worker created from + * that client. It reads {@link GoogleCloudRunMetadata Cloud Run instance metadata} once while the + * client is configured, caches it, and then: + * + *

    + *
  • sets the workflow client identity to the {@linkplain + * GoogleCloudRunMetadata#workerIdentity() derived worker identity}, but only when the caller + * has not already set an identity (a user-provided identity always wins); + *
  • sets each worker's {@link WorkerDeploymentOptions} to the {@linkplain + * GoogleCloudRunMetadata#workerDeploymentVersion() derived deployment version} with worker + * versioning enabled and a {@link VersioningBehavior#PINNED PINNED} default behavior, so + * in-flight workflows stay on the Cloud Run revision that started them. + *
+ * + *

The metadata is fetched lazily at client-configure time rather than in the constructor, + * because the fetch performs a network request to the Cloud Run metadata server that belongs at + * connect time. The metadata server is only reachable from a Cloud Run instance, so the fetch fails + * fast with an {@link IllegalStateException} when this process is not running on Cloud Run. + * + *

Register the plugin with {@link WorkflowClientOptions.Builder#setPlugins}: + * + *

{@code
+ * WorkflowClient client =
+ *     WorkflowClient.newInstance(
+ *         service,
+ *         WorkflowClientOptions.newBuilder()
+ *             .setNamespace(namespace)
+ *             .setPlugins(new WorkerIdPlugin())
+ *             .build());
+ *
+ * WorkerFactory factory = WorkerFactory.newInstance(client);
+ * Worker worker = factory.newWorker("my-task-queue");
+ * }
+ * + *

Advanced / testing: {@link #WorkerIdPlugin(GoogleCloudRunMetadata)} accepts an + * already-resolved {@link GoogleCloudRunMetadata} instance, which skips the lazy fetch entirely. + * This is useful when the application fetches the metadata itself (for example to log it) or when a + * test injects fixed metadata. + * + *

Experimental: Google Cloud Run support is experimental and may change without notice. + */ +@Experimental +public final class WorkerIdPlugin extends SimplePlugin { + /** Unique plugin name, used for logging and duplicate detection. */ + public static final String NAME = "io.temporal.gcp.cloudrun.workerid.WorkerIdPlugin"; + + private final Supplier metadataSupplier; + private volatile GoogleCloudRunMetadata metadata; + + /** + * Creates a plugin that fetches Cloud Run instance metadata from the {@linkplain + * GoogleCloudRunMetadata#DEFAULT_METADATA_URL default metadata server} while the workflow client + * is configured. + */ + public WorkerIdPlugin() { + this(GoogleCloudRunMetadata::fetch); + } + + /** + * Creates a plugin that uses an already-resolved {@link GoogleCloudRunMetadata} instance instead + * of fetching it. No request is made to the Cloud Run metadata server. + * + * @param metadata previously fetched Cloud Run instance metadata. + */ + public WorkerIdPlugin(GoogleCloudRunMetadata metadata) { + this(pinnedSupplier(metadata)); + } + + /** + * Package-private test seam that supplies the {@link GoogleCloudRunMetadata} lazily. It lets unit + * tests point the fetch at an in-process metadata server and injected environment through the + * {@link GoogleCloudRunMetadata#fetch(String, java.time.Duration, java.util.function.Function)} + * seam, and to exercise the off-platform fail-fast path. It is not part of the public API; use + * {@link #WorkerIdPlugin()} or {@link #WorkerIdPlugin(GoogleCloudRunMetadata)} instead. + * + * @param metadataSupplier supplier invoked once, at client-configure time, to resolve the + * metadata. + */ + WorkerIdPlugin(Supplier metadataSupplier) { + super(NAME); + this.metadataSupplier = Objects.requireNonNull(metadataSupplier, "metadataSupplier"); + } + + /** + * Fetches (once) and caches the Cloud Run instance metadata, then sets the derived worker + * identity on the client options when the caller has not already set an identity. + * + * @param builder the workflow client options builder to configure. + * @throws IllegalStateException if the Cloud Run metadata server cannot be reached, which usually + * means this process is not running on Google Cloud Run. + */ + @Override + public void configureWorkflowClient(WorkflowClientOptions.Builder builder) { + GoogleCloudRunMetadata resolved = metadata(); + if (isBlank(builder.build().getIdentity())) { + builder.setIdentity(resolved.workerIdentity()); + } + } + + /** + * Sets the worker's {@link WorkerDeploymentOptions} from the cached Cloud Run metadata, enabling + * worker versioning with a {@link VersioningBehavior#PINNED PINNED} default behavior. + * + * @param taskQueue the task queue name for the worker being created. + * @param builder the worker options builder to configure. + * @throws IllegalStateException if the Cloud Run name or revision is not set, which usually means + * this process is not running on a Cloud Run worker pool or service. + */ + @Override + public void configureWorker(String taskQueue, WorkerOptions.Builder builder) { + GoogleCloudRunMetadata resolved = metadata(); + builder.setDeploymentOptions( + WorkerDeploymentOptions.newBuilder() + .setUseVersioning(true) + .setVersion(resolved.workerDeploymentVersion()) + .setDefaultVersioningBehavior(VersioningBehavior.PINNED) + .build()); + } + + private GoogleCloudRunMetadata metadata() { + GoogleCloudRunMetadata local = metadata; + if (local == null) { + synchronized (this) { + local = metadata; + if (local == null) { + local = Objects.requireNonNull(metadataSupplier.get(), "Cloud Run metadata"); + metadata = local; + } + } + } + return local; + } + + private static Supplier pinnedSupplier(GoogleCloudRunMetadata metadata) { + Objects.requireNonNull(metadata, "metadata"); + return () -> metadata; + } + + private static boolean isBlank(String value) { + return value == null || value.trim().isEmpty(); + } +} diff --git a/contrib/temporal-gcp-cloud-run-worker-id/src/test/java/io/temporal/gcp/cloudrun/workerid/GoogleCloudRunMetadataTest.java b/contrib/temporal-gcp-cloud-run-worker-id/src/test/java/io/temporal/gcp/cloudrun/workerid/GoogleCloudRunMetadataTest.java new file mode 100644 index 0000000000..3f7e8acd7d --- /dev/null +++ b/contrib/temporal-gcp-cloud-run-worker-id/src/test/java/io/temporal/gcp/cloudrun/workerid/GoogleCloudRunMetadataTest.java @@ -0,0 +1,233 @@ +package io.temporal.gcp.cloudrun.workerid; + +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.common.WorkerDeploymentVersion; +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)); + } + + 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/contrib/temporal-gcp-cloud-run-worker-id/src/test/java/io/temporal/gcp/cloudrun/workerid/WorkerIdPluginTest.java b/contrib/temporal-gcp-cloud-run-worker-id/src/test/java/io/temporal/gcp/cloudrun/workerid/WorkerIdPluginTest.java new file mode 100644 index 0000000000..a1c31e8f93 --- /dev/null +++ b/contrib/temporal-gcp-cloud-run-worker-id/src/test/java/io/temporal/gcp/cloudrun/workerid/WorkerIdPluginTest.java @@ -0,0 +1,192 @@ +package io.temporal.gcp.cloudrun.workerid; + +import static org.junit.Assert.assertEquals; +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 java.util.function.Supplier; +import org.junit.After; +import org.junit.Before; +import org.junit.Test; + +/** + * Unit tests for {@link WorkerIdPlugin}. + * + *

The metadata request is served by an in-process {@link HttpServer} and the environment lookup + * is injected through the {@link GoogleCloudRunMetadata#fetch(String, Duration, + * java.util.function.Function)} test seam, so these tests touch neither the network nor the real + * process environment. The plugin's package-private {@link WorkerIdPlugin#WorkerIdPlugin(Supplier)} + * seam lets each test point the plugin at that in-process server (or at an unreachable address, to + * exercise the off-platform fail-fast path). + */ +public class WorkerIdPluginTest { + private static final Duration TIMEOUT = Duration.ofSeconds(2); + + private HttpServer server; + private final AtomicReference responseBody = 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 -> { + byte[] body = responseBody.get().getBytes(StandardCharsets.UTF_8); + exchange.sendResponseHeaders(200, body.length == 0 ? -1 : body.length); + try (OutputStream out = exchange.getResponseBody()) { + out.write(body); + } + }); + server.start(); + } + + @After + public void stopServer() { + server.stop(0); + } + + @Test + public void configureWorkflowClientSetsDerivedIdentityWhenUnset() { + 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"); + + WorkflowClientOptions.Builder builder = WorkflowClientOptions.newBuilder(); + pluginFor(env).configureWorkflowClient(builder); + + assertEquals("instance-1@revision-1", builder.build().getIdentity()); + } + + @Test + public void configureWorkflowClientPreservesUserProvidedIdentity() { + 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"); + + WorkflowClientOptions.Builder builder = + WorkflowClientOptions.newBuilder().setIdentity("user-set"); + pluginFor(env).configureWorkflowClient(builder); + + assertEquals("user-set", builder.build().getIdentity()); + } + + @Test + public void configureWorkerEnablesPinnedVersioning() { + 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.Builder builder = WorkerOptions.newBuilder(); + pluginFor(env).configureWorker("orders", builder); + + WorkerDeploymentOptions deploymentOptions = builder.build().getDeploymentOptions(); + assertTrue(deploymentOptions.isUsingVersioning()); + assertEquals( + new WorkerDeploymentVersion("worker-pool", "revision-1"), deploymentOptions.getVersion()); + assertEquals(VersioningBehavior.PINNED, deploymentOptions.getDefaultVersioningBehavior()); + } + + @Test + public void configureWorkflowClientFailsFastOffCloudRun() { + String unreachableUrl = + "http://127.0.0.1:" + reserveUnusedPort() + "/computeMetadata/v1/instance/id"; + WorkerIdPlugin plugin = + new WorkerIdPlugin( + () -> GoogleCloudRunMetadata.fetch(unreachableUrl, TIMEOUT, name -> null)); + + IllegalStateException e = + assertThrows( + IllegalStateException.class, + () -> plugin.configureWorkflowClient(WorkflowClientOptions.newBuilder())); + assertTrue(e.getMessage().contains("metadata server")); + } + + @Test + public void configureWorkerFailsFastWhenNotWorkerPoolOrService() { + responseBody.set("instance-1"); + + // Metadata server is reachable (instance id is present) but no name/revision env is set, so the + // deployment version cannot be built. This is the "on some other platform" case. + WorkerIdPlugin plugin = new WorkerIdPlugin(metadata(new HashMap<>())); + + assertThrows( + IllegalStateException.class, + () -> plugin.configureWorker("orders", WorkerOptions.newBuilder())); + } + + @Test + public void metadataIsFetchedOnceAndSharedByBothHooks() { + 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"); + + GoogleCloudRunMetadata resolved = metadata(env); + AtomicInteger supplierCalls = new AtomicInteger(); + Supplier countingSupplier = + () -> { + supplierCalls.incrementAndGet(); + return resolved; + }; + WorkerIdPlugin plugin = new WorkerIdPlugin(countingSupplier); + + plugin.configureWorkflowClient(WorkflowClientOptions.newBuilder()); + plugin.configureWorker("orders", WorkerOptions.newBuilder()); + + assertEquals(1, supplierCalls.get()); + } + + @Test + public void injectedMetadataIsUsedWithoutFetching() { + 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"); + + WorkerIdPlugin plugin = new WorkerIdPlugin(metadata(env)); + + WorkflowClientOptions.Builder builder = WorkflowClientOptions.newBuilder(); + plugin.configureWorkflowClient(builder); + + assertEquals("instance-1@revision-1", builder.build().getIdentity()); + } + + private WorkerIdPlugin pluginFor(Map env) { + return new WorkerIdPlugin(() -> metadata(env)); + } + + private GoogleCloudRunMetadata metadata(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 6cbf879490..276f024743 100644 --- a/settings.gradle +++ b/settings.gradle @@ -17,6 +17,8 @@ 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-gcp-cloud-run-worker-id' +project(':temporal-gcp-cloud-run-worker-id').projectDir = file('contrib/temporal-gcp-cloud-run-worker-id') include 'temporal-spring-boot-autoconfigure' include 'temporal-spring-boot-starter' include 'temporal-remote-data-encoder' diff --git a/temporal-bom/build.gradle b/temporal-bom/build.gradle index c79c4df44a..6ed9f82ebc 100644 --- a/temporal-bom/build.gradle +++ b/temporal-bom/build.gradle @@ -11,6 +11,7 @@ dependencies { api project(':temporal-opentracing') api project(':temporal-aws-lambda') api project(':temporal-gcp-cloud-run') + api project(':temporal-gcp-cloud-run-worker-id') api project(':temporal-remote-data-encoder') api project(':temporal-sdk') api project(':temporal-serviceclient')