Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
71 changes: 71 additions & 0 deletions contrib/temporal-gcp-cloud-run/README.md
Original file line number Diff line number Diff line change
@@ -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 `<instanceId>@<revision>`, falling back to `<instanceId>@<name>` and then to the bare `<instanceId>` 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.
7 changes: 7 additions & 0 deletions contrib/temporal-gcp-cloud-run/build.gradle
Original file line number Diff line number Diff line change
@@ -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')
}
Original file line number Diff line number Diff line change
@@ -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.
*
* <p>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)}.
*
* <p>The deployment name and revision are resolved from environment variables Cloud Run injects
* into every instance. Cloud Run <b>worker pools</b> set {@code CLOUD_RUN_WORKER_POOL} and {@code
* CLOUD_RUN_REVISION}; Cloud Run <b>services</b> 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.
*
* <p><b>Experimental:</b> 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.
*
* <p>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.
*
* <p>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.
*
* <p>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();
}
}
2 changes: 2 additions & 0 deletions settings.gradle
Original file line number Diff line number Diff line change
Expand Up @@ -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'
Expand Down
Loading