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'