From e9b6faa1e3ab3d216079f3b7633ba29bef570f60 Mon Sep 17 00:00:00 2001 From: seanbollin Date: Wed, 19 Aug 2026 15:01:04 -0700 Subject: [PATCH] Add GCP Cloud Run OpenTelemetry support Adds a Cloud Run serverless-worker OpenTelemetry integration, mirroring the approved .NET SDK common-core refactor (temporalio/sdk-dotnet#844) and building on #2955. - New module temporal-gcp-cloud-run (io.temporal.gcp.cloudrun): a CloudRunOpenTelemetryPlugin that exports Core metrics and tracing spans over OTLP/gRPC to a local OpenTelemetry Collector sidecar. Cloud Run specifics only (service name from CLOUD_RUN_WORKER_POOL/K_SERVICE, 60s report interval, deferred shutdown flush); depends only on temporal-opentelemetry, no Google client libraries. - Extract a shared OpenTelemetryWorker.resolveServiceName(env, default, fallbackEnvVars...) into the provider-neutral core and add OpenTelemetryPlugin.Builder.getMetricsReportInterval() (both additive). - Refactor the AWS Lambda OtelLambdaWorkerConfigurationHelper onto the shared resolver with no public API or behavior change (existing tests unmodified). Co-authored-by: Edward Amsden Co-Authored-By: Claude Opus 4.8 --- .../OtelLambdaWorkerConfigurationHelper.java | 16 +- contrib/temporal-gcp-cloud-run/README.md | 99 +++++++++ contrib/temporal-gcp-cloud-run/build.gradle | 17 ++ .../cloudrun/CloudRunOpenTelemetryPlugin.java | 196 ++++++++++++++++++ .../CloudRunOpenTelemetryPluginTest.java | 123 +++++++++++ .../opentelemetry/OpenTelemetryPlugin.java | 4 + .../opentelemetry/OpenTelemetryWorker.java | 31 ++- settings.gradle | 2 + temporal-bom/build.gradle | 1 + 9 files changed, 474 insertions(+), 15 deletions(-) create mode 100644 contrib/temporal-gcp-cloud-run/README.md create mode 100644 contrib/temporal-gcp-cloud-run/build.gradle create mode 100644 contrib/temporal-gcp-cloud-run/src/main/java/io/temporal/gcp/cloudrun/CloudRunOpenTelemetryPlugin.java create mode 100644 contrib/temporal-gcp-cloud-run/src/test/java/io/temporal/gcp/cloudrun/CloudRunOpenTelemetryPluginTest.java diff --git a/contrib/temporal-aws-lambda/src/main/java/io/temporal/aws/lambda/OtelLambdaWorkerConfigurationHelper.java b/contrib/temporal-aws-lambda/src/main/java/io/temporal/aws/lambda/OtelLambdaWorkerConfigurationHelper.java index 235352b7ac..d087b1db49 100644 --- a/contrib/temporal-aws-lambda/src/main/java/io/temporal/aws/lambda/OtelLambdaWorkerConfigurationHelper.java +++ b/contrib/temporal-aws-lambda/src/main/java/io/temporal/aws/lambda/OtelLambdaWorkerConfigurationHelper.java @@ -125,12 +125,8 @@ public static void configureFlushHook( } static String resolveServiceName(Map env) { - String serviceName = nonEmptyEnv(env, OTEL_SERVICE_NAME); - if (serviceName != null) { - return serviceName; - } - serviceName = nonEmptyEnv(env, AWS_LAMBDA_FUNCTION_NAME); - return serviceName == null ? DEFAULT_SERVICE_NAME : serviceName; + return OpenTelemetryWorker.resolveServiceName( + env, DEFAULT_SERVICE_NAME, AWS_LAMBDA_FUNCTION_NAME); } public static final class Builder { @@ -258,14 +254,6 @@ OpenTelemetry create( IdGenerator idGenerator); } - private static String nonEmptyEnv(Map env, String name) { - if (env == null) { - return null; - } - String value = env.get(name); - return value == null || value.trim().isEmpty() ? null : value; - } - private static void appendServiceStubsPlugin( WorkflowServiceStubsOptions.Builder options, WorkflowServiceStubsPlugin plugin) { WorkflowServiceStubsPlugin[] existing = options.build().getPlugins(); diff --git a/contrib/temporal-gcp-cloud-run/README.md b/contrib/temporal-gcp-cloud-run/README.md new file mode 100644 index 0000000000..01a6c15184 --- /dev/null +++ b/contrib/temporal-gcp-cloud-run/README.md @@ -0,0 +1,99 @@ +# Temporal Google Cloud Run module + +This module provides an OpenTelemetry plugin with defaults for Temporal Java SDK workers running on Google Cloud Run. Cloud Run worker pools are the recommended deployment because Temporal workers are continuous, pull-based background workloads. + +> **Collector required by default:** The plugin exports metrics and traces to an OTLP collector at `http://localhost:4317`. It does not export directly to Google Cloud. Deploy the Google-Built OpenTelemetry Collector as a sidecar, configure another collector endpoint, or provide an application-owned `OpenTelemetry` instance. Without a collector at the configured endpoint, telemetry is not delivered to Google Cloud. + +This integration is for container-based Cloud Run workloads. It does not implement a Cloud Run functions invocation lifecycle. + +A Cloud Run service can also host a Temporal worker, but it must use instance-based billing so CPU is available outside request handling, keep at least one instance active through minimum instances or manual scaling, and run an ingress container that listens on `PORT`. These are deployment requirements; the plugin cannot configure them from inside the worker process. + +## Usage + +Add `temporal-gcp-cloud-run` next to your Temporal SDK dependency, then install the plugin on service stubs options before creating clients and workers: + +```java +CloudRunOpenTelemetryPlugin plugin = CloudRunOpenTelemetryPlugin.newBuilder().build(); + +WorkflowServiceStubs service = + WorkflowServiceStubs.newServiceStubs( + WorkflowServiceStubsOptions.newBuilder() + .setPlugins(plugin) + .build()); +WorkflowClient client = WorkflowClient.newInstance(service); +WorkerFactory factory = WorkerFactory.newInstance(client); +``` + +The plugin configures the SDK metrics scope, tracing interceptors, and OTLP metric and trace exporters through `temporal-opentelemetry`. Do not install both `CloudRunOpenTelemetryPlugin` and `OpenTelemetryPlugin` on the same service stubs. + +## Shutdown lifecycle + +`WorkerFactory.shutdown()` initiates asynchronous shutdown. The Cloud Run plugin therefore does not flush from the worker-factory shutdown callback by default: doing so could miss telemetry emitted while activities and workflows finish. Wait for the factory to terminate, then flush with the time remaining before Cloud Run sends `SIGKILL`. For example, the following JVM shutdown hook reserves six seconds for graceful shutdown, one second for forced shutdown, and two seconds for telemetry flushing within Cloud Run's ten-second termination window: + +```java +Runtime.getRuntime() + .addShutdownHook( + new Thread( + () -> { + factory.shutdown(); + factory.awaitTermination(6, TimeUnit.SECONDS); + if (!factory.isTerminated()) { + factory.shutdownNow(); + factory.awaitTermination(1, TimeUnit.SECONDS); + } + plugin.newFlushHook().run(Duration.ofSeconds(2)); + })); +``` + +Applications with an existing lifecycle manager should perform the same sequence there instead of registering another JVM hook. `Builder.setFlushOnWorkerFactoryShutdown(true)` restores the underlying plugin's immediate, best-effort flush, but it should only be used when no work can emit telemetry after the shutdown request. + +The OTLP endpoint is resolved in this order: + +1. `Builder.setEndpoint(...)`. +2. `OTEL_EXPORTER_OTLP_ENDPOINT`. +3. `http://localhost:4317`. + +When the plugin creates the OpenTelemetry SDK, metrics are reported and exported every sixty seconds by default. This matches the coordinated GCP plugin default across Temporal SDKs and exceeds Google Cloud's five-second minimum export interval. If you use `Builder.setMetricsReportInterval(...)`, keep the interval above that minimum. The collector also needs the unbatched metrics pipeline described below to make a forced shutdown flush safe regardless of its timing relative to the last periodic export. + +With an application-owned `OpenTelemetry` instance, the setting only controls how often the Temporal metrics scope reports into that instance. Configure the instance's metric reader to export at an interval above the Google Cloud minimum as well. + +The OpenTelemetry service name is resolved in this order: + +1. `Builder.setServiceName(...)`. +2. `OTEL_SERVICE_NAME`. +3. `CLOUD_RUN_WORKER_POOL` for a Cloud Run worker pool. +4. `K_SERVICE` for a Cloud Run service. +5. `temporal-worker`. + +The collector should use its GCP resource detector to add the Google Cloud attributes it recognizes. Do not rely on the detector to infer Cloud Run worker-pool-specific location or revision attributes; configure those explicitly with a collector resource processor if they are required. This module does not call the Google Cloud metadata server and adds no Google Cloud client libraries or exporters to the worker process. + +## Collector sidecar + +Google publishes the Google-Built OpenTelemetry Collector as a container image. Configure it as a second Cloud Run container, listen for OTLP gRPC on `localhost:4317`, and use its GCP exporters for metrics and traces. For the image, recommended collector configuration, IAM roles, health check, and Secret Manager mount, see [Deploy Google-Built OpenTelemetry Collector on Cloud Run](https://cloud.google.com/stackdriver/docs/instrumentation/opentelemetry-collector-cloud-run). That guide demonstrates a Cloud Run service; adapt its collector container and configuration when deploying a worker pool. + +Do not put a batch processor in the `googlemanagedprometheus` metrics pipeline. A periodic cumulative metric export followed closely by a forced shutdown flush can otherwise put two points for the same time series in one request, which Managed Service for Prometheus rejects. Pass metrics through the memory limiter, GCP resource detection, and any collision transforms directly to `googlemanagedprometheus`. Keep a dedicated five-second batch processor on the traces pipeline: + +```yaml +processors: + batch/traces: + send_batch_max_size: 200 + send_batch_size: 200 + timeout: 5s + +service: + pipelines: + metrics: + receivers: [otlp] + processors: [memory_limiter, resourcedetection, transform/collision] + exporters: [googlemanagedprometheus] + traces: + receivers: [otlp] + processors: [memory_limiter, resourcedetection, transform/set_project_id, batch/traces] + exporters: [otlp] +``` + +Cloud Run worker pools support sidecar containers over localhost and are intended for continuous background work. The deployment should start the collector before the Temporal worker and use the collector health extension as its startup probe. + +To use an external collector instead, set `OTEL_EXPORTER_OTLP_ENDPOINT` or call `Builder.setEndpoint(...)`. + +To use an application-owned provider, call `Builder.setOpenTelemetry(...)`. In that path, no exporters are created; the plugin installs the Temporal metrics scope, tracing interceptors, and shutdown flush hook around the supplied provider. diff --git a/contrib/temporal-gcp-cloud-run/build.gradle b/contrib/temporal-gcp-cloud-run/build.gradle new file mode 100644 index 0000000000..35f2e781d4 --- /dev/null +++ b/contrib/temporal-gcp-cloud-run/build.gradle @@ -0,0 +1,17 @@ +description = '''Temporal Java SDK Google Cloud Run 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-serviceclient') + compileOnly project(':temporal-sdk') + compileOnly "javax.annotation:javax.annotation-api:$annotationApiVersion" + + api project(':temporal-opentelemetry') + + testImplementation project(':temporal-sdk') + testImplementation project(':temporal-serviceclient') + 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/CloudRunOpenTelemetryPlugin.java b/contrib/temporal-gcp-cloud-run/src/main/java/io/temporal/gcp/cloudrun/CloudRunOpenTelemetryPlugin.java new file mode 100644 index 0000000000..a985e66326 --- /dev/null +++ b/contrib/temporal-gcp-cloud-run/src/main/java/io/temporal/gcp/cloudrun/CloudRunOpenTelemetryPlugin.java @@ -0,0 +1,196 @@ +package io.temporal.gcp.cloudrun; + +import io.opentelemetry.api.OpenTelemetry; +import io.temporal.client.WorkflowClientOptions; +import io.temporal.common.Experimental; +import io.temporal.common.SimplePlugin; +import io.temporal.opentelemetry.OpenTelemetryPlugin; +import io.temporal.opentelemetry.OpenTelemetryWorker; +import io.temporal.opentelemetry.TimedShutdownHook; +import io.temporal.serviceclient.WorkflowServiceStubsOptions; +import io.temporal.worker.WorkerFactory; +import io.temporal.worker.WorkerFactoryOptions; +import java.time.Duration; +import java.util.Map; +import java.util.Objects; +import java.util.function.Consumer; +import javax.annotation.Nonnull; + +/** + * OpenTelemetry plugin with defaults for Temporal workers running on Google Cloud Run, primarily in + * worker pools. + */ +@Experimental +public final class CloudRunOpenTelemetryPlugin extends SimplePlugin { + public static final String NAME = OpenTelemetryPlugin.NAME; + public static final String OTEL_EXPORTER_OTLP_ENDPOINT = + OpenTelemetryWorker.OTEL_EXPORTER_OTLP_ENDPOINT; + public static final String OTEL_SERVICE_NAME = OpenTelemetryWorker.OTEL_SERVICE_NAME; + public static final String CLOUD_RUN_WORKER_POOL = "CLOUD_RUN_WORKER_POOL"; + public static final String K_SERVICE = "K_SERVICE"; + public static final String DEFAULT_OTLP_ENDPOINT = OpenTelemetryWorker.DEFAULT_OTLP_ENDPOINT; + public static final String DEFAULT_SERVICE_NAME = OpenTelemetryWorker.DEFAULT_SERVICE_NAME; + public static final Duration DEFAULT_METRICS_REPORT_INTERVAL = Duration.ofSeconds(60); + + private final OpenTelemetryPlugin delegate; + + private CloudRunOpenTelemetryPlugin(Builder builder) { + super(NAME); + this.delegate = builder.buildDelegate(); + } + + public static Builder newBuilder() { + return new Builder(System.getenv()); + } + + public static Builder newBuilder(@Nonnull Map env) { + return new Builder(env); + } + + public String getEndpoint() { + return delegate.getEndpoint(); + } + + public String getServiceName() { + return delegate.getServiceName(); + } + + public OpenTelemetry getOpenTelemetry() { + return delegate.getOpenTelemetry(); + } + + /** + * Creates a flush hook that reports buffered Temporal metrics before force-flushing OpenTelemetry + * providers. + */ + public TimedShutdownHook newFlushHook() { + return delegate.newFlushHook(); + } + + @Override + public void configureServiceStubs(@Nonnull WorkflowServiceStubsOptions.Builder builder) { + delegate.configureServiceStubs(builder); + } + + @Override + public void configureWorkflowClient(@Nonnull WorkflowClientOptions.Builder builder) { + delegate.configureWorkflowClient(builder); + } + + @Override + public void configureWorkerFactory(@Nonnull WorkerFactoryOptions.Builder builder) { + delegate.configureWorkerFactory(builder); + } + + @Override + public void shutdownWorkerFactory( + @Nonnull WorkerFactory factory, @Nonnull Consumer next) { + delegate.shutdownWorkerFactory(factory, next); + } + + /** Builder for {@link CloudRunOpenTelemetryPlugin}. */ + public static final class Builder { + private final Map env; + private final OpenTelemetryPlugin.Builder delegate; + private String serviceName; + + private Builder(Map env) { + this.env = Objects.requireNonNull(env, "env"); + this.delegate = + OpenTelemetryPlugin.newBuilder(env) + .setMetricsReportInterval(DEFAULT_METRICS_REPORT_INTERVAL) + .setFlushOnWorkerFactoryShutdown(false); + } + + /** + * Uses an application-owned OpenTelemetry instance instead of creating an SDK and exporters. + */ + public Builder setOpenTelemetry(@Nonnull OpenTelemetry openTelemetry) { + delegate.setOpenTelemetry(openTelemetry); + return this; + } + + /** Sets the OTLP metric and trace exporter endpoint used by the default SDK setup. */ + public Builder setEndpoint(@Nonnull String endpoint) { + delegate.setEndpoint(endpoint); + return this; + } + + /** Sets the service name used by the default SDK resource and Temporal metrics reporter. */ + public Builder setServiceName(@Nonnull String serviceName) { + this.serviceName = Objects.requireNonNull(serviceName, "serviceName"); + return this; + } + + /** + * Sets the interval used by the Temporal metrics scope and, when the plugin creates the + * OpenTelemetry instance, its periodic metric reader. + * + *

An application-owned OpenTelemetry instance must configure its metric reader separately. + */ + public Builder setMetricsReportInterval(@Nonnull Duration metricsReportInterval) { + delegate.setMetricsReportInterval(metricsReportInterval); + return this; + } + + /** Sets how long the OpenTelemetry flush hook waits for provider flushing. */ + public Builder setFlushTimeout(@Nonnull Duration flushTimeout) { + delegate.setFlushTimeout(flushTimeout); + return this; + } + + /** Overrides the OpenTelemetry provider flush hook. */ + public Builder setFlushHook(@Nonnull Runnable flushHook) { + delegate.setFlushHook(flushHook); + return this; + } + + /** + * Controls whether the plugin flushes immediately after {@link WorkerFactory#shutdown()} or + * {@link WorkerFactory#shutdownNow()} initiates shutdown. + * + *

This is disabled by default because worker-factory shutdown is asynchronous. Cloud Run + * applications should wait for worker termination and then run {@link + * CloudRunOpenTelemetryPlugin#newFlushHook()} with the remaining shutdown time. + */ + public Builder setFlushOnWorkerFactoryShutdown(boolean flushOnWorkerFactoryShutdown) { + delegate.setFlushOnWorkerFactoryShutdown(flushOnWorkerFactoryShutdown); + return this; + } + + public String getEndpoint() { + return delegate.getEndpoint(); + } + + public String getServiceName() { + return serviceName == null ? resolveServiceName(env) : serviceName; + } + + public Duration getMetricsReportInterval() { + return delegate.getMetricsReportInterval(); + } + + public OpenTelemetry createOpenTelemetry() { + applyServiceNameDefault(); + return delegate.createOpenTelemetry(); + } + + public CloudRunOpenTelemetryPlugin build() { + return new CloudRunOpenTelemetryPlugin(this); + } + + private OpenTelemetryPlugin buildDelegate() { + applyServiceNameDefault(); + return delegate.build(); + } + + private void applyServiceNameDefault() { + delegate.setServiceName(getServiceName()); + } + } + + static String resolveServiceName(Map env) { + return OpenTelemetryWorker.resolveServiceName( + env, DEFAULT_SERVICE_NAME, CLOUD_RUN_WORKER_POOL, K_SERVICE); + } +} diff --git a/contrib/temporal-gcp-cloud-run/src/test/java/io/temporal/gcp/cloudrun/CloudRunOpenTelemetryPluginTest.java b/contrib/temporal-gcp-cloud-run/src/test/java/io/temporal/gcp/cloudrun/CloudRunOpenTelemetryPluginTest.java new file mode 100644 index 0000000000..3be4529168 --- /dev/null +++ b/contrib/temporal-gcp-cloud-run/src/test/java/io/temporal/gcp/cloudrun/CloudRunOpenTelemetryPluginTest.java @@ -0,0 +1,123 @@ +package io.temporal.gcp.cloudrun; + +import static org.junit.Assert.*; + +import io.opentelemetry.api.OpenTelemetry; +import io.temporal.client.WorkflowClientOptions; +import io.temporal.opentelemetry.OpenTelemetryPlugin; +import io.temporal.serviceclient.WorkflowServiceStubsOptions; +import io.temporal.worker.WorkerFactoryOptions; +import java.time.Duration; +import java.util.HashMap; +import java.util.Map; +import java.util.concurrent.atomic.AtomicInteger; +import org.junit.Test; + +public class CloudRunOpenTelemetryPluginTest { + @Test + public void defaultsToLocalCollectorAndGenericServiceName() { + CloudRunOpenTelemetryPlugin.Builder builder = + CloudRunOpenTelemetryPlugin.newBuilder(new HashMap<>()); + + assertEquals("http://localhost:4317", builder.getEndpoint()); + assertEquals("temporal-worker", builder.getServiceName()); + assertEquals(Duration.ofSeconds(60), builder.getMetricsReportInterval()); + } + + @Test + public void resolvesCloudRunServiceNameWithExpectedPrecedence() { + Map env = new HashMap<>(); + env.put(CloudRunOpenTelemetryPlugin.K_SERVICE, "cloud-run-service"); + assertEquals("cloud-run-service", CloudRunOpenTelemetryPlugin.newBuilder(env).getServiceName()); + + env.put(CloudRunOpenTelemetryPlugin.CLOUD_RUN_WORKER_POOL, "worker-pool"); + assertEquals("worker-pool", CloudRunOpenTelemetryPlugin.newBuilder(env).getServiceName()); + + env.put(CloudRunOpenTelemetryPlugin.OTEL_SERVICE_NAME, "otel-service"); + assertEquals("otel-service", CloudRunOpenTelemetryPlugin.newBuilder(env).getServiceName()); + + assertEquals( + "builder-service", + CloudRunOpenTelemetryPlugin.newBuilder(env) + .setServiceName("builder-service") + .getServiceName()); + } + + @Test + public void ignoresEmptyEnvironmentValues() { + Map env = new HashMap<>(); + env.put(CloudRunOpenTelemetryPlugin.OTEL_SERVICE_NAME, " "); + env.put(CloudRunOpenTelemetryPlugin.CLOUD_RUN_WORKER_POOL, ""); + env.put(CloudRunOpenTelemetryPlugin.K_SERVICE, "cloud-run-service"); + + assertEquals("cloud-run-service", CloudRunOpenTelemetryPlugin.newBuilder(env).getServiceName()); + } + + @Test + public void buildAppliesResolvedEndpointAndServiceName() { + Map env = new HashMap<>(); + env.put(CloudRunOpenTelemetryPlugin.OTEL_EXPORTER_OTLP_ENDPOINT, "http://collector:4317"); + env.put(CloudRunOpenTelemetryPlugin.CLOUD_RUN_WORKER_POOL, "worker-pool"); + + CloudRunOpenTelemetryPlugin plugin = + CloudRunOpenTelemetryPlugin.newBuilder(env).setOpenTelemetry(OpenTelemetry.noop()).build(); + + assertEquals("http://collector:4317", plugin.getEndpoint()); + assertEquals("worker-pool", plugin.getServiceName()); + } + + @Test + public void installsMetricsScopeAndTracingInterceptors() { + CloudRunOpenTelemetryPlugin plugin = + CloudRunOpenTelemetryPlugin.newBuilder(new HashMap<>()) + .setOpenTelemetry(OpenTelemetry.noop()) + .build(); + WorkflowServiceStubsOptions.Builder serviceOptions = WorkflowServiceStubsOptions.newBuilder(); + WorkflowClientOptions.Builder clientOptions = WorkflowClientOptions.newBuilder(); + WorkerFactoryOptions.Builder factoryOptions = WorkerFactoryOptions.newBuilder(); + + plugin.configureServiceStubs(serviceOptions); + plugin.configureWorkflowClient(clientOptions); + plugin.configureWorkerFactory(factoryOptions); + + assertEquals(OpenTelemetryPlugin.NAME, plugin.getName()); + assertNotNull(serviceOptions.build().getMetricsScope()); + assertEquals(1, clientOptions.build().getInterceptors().length); + assertEquals(1, factoryOptions.build().getWorkerInterceptors().length); + } + + @Test + public void workerFactoryShutdownDefersFlushByDefault() { + AtomicInteger flushes = new AtomicInteger(); + AtomicInteger shutdowns = new AtomicInteger(); + CloudRunOpenTelemetryPlugin plugin = + CloudRunOpenTelemetryPlugin.newBuilder(new HashMap<>()) + .setOpenTelemetry(OpenTelemetry.noop()) + .setFlushHook(flushes::incrementAndGet) + .build(); + + plugin.shutdownWorkerFactory(null, factory -> shutdowns.incrementAndGet()); + + assertEquals(1, shutdowns.get()); + assertEquals(0, flushes.get()); + + plugin.newFlushHook().run(); + + assertEquals(1, flushes.get()); + } + + @Test + public void workerFactoryShutdownFlushCanBeEnabled() { + AtomicInteger flushes = new AtomicInteger(); + CloudRunOpenTelemetryPlugin plugin = + CloudRunOpenTelemetryPlugin.newBuilder(new HashMap<>()) + .setOpenTelemetry(OpenTelemetry.noop()) + .setFlushHook(flushes::incrementAndGet) + .setFlushOnWorkerFactoryShutdown(true) + .build(); + + plugin.shutdownWorkerFactory(null, factory -> {}); + + assertEquals(1, flushes.get()); + } +} diff --git a/contrib/temporal-opentelemetry/src/main/java/io/temporal/opentelemetry/OpenTelemetryPlugin.java b/contrib/temporal-opentelemetry/src/main/java/io/temporal/opentelemetry/OpenTelemetryPlugin.java index 2604e205f1..3b1e2f1604 100644 --- a/contrib/temporal-opentelemetry/src/main/java/io/temporal/opentelemetry/OpenTelemetryPlugin.java +++ b/contrib/temporal-opentelemetry/src/main/java/io/temporal/opentelemetry/OpenTelemetryPlugin.java @@ -235,6 +235,10 @@ public String getServiceName() { return serviceName == null ? OpenTelemetryWorker.resolveServiceName(env) : serviceName; } + public Duration getMetricsReportInterval() { + return metricsReportInterval; + } + Builder setTelemetryFactory(OpenTelemetryWorker.TelemetryFactory telemetryFactory) { this.telemetryFactory = Objects.requireNonNull(telemetryFactory, "telemetryFactory"); return this; diff --git a/contrib/temporal-opentelemetry/src/main/java/io/temporal/opentelemetry/OpenTelemetryWorker.java b/contrib/temporal-opentelemetry/src/main/java/io/temporal/opentelemetry/OpenTelemetryWorker.java index 90e7697aac..8d1eabfd15 100644 --- a/contrib/temporal-opentelemetry/src/main/java/io/temporal/opentelemetry/OpenTelemetryWorker.java +++ b/contrib/temporal-opentelemetry/src/main/java/io/temporal/opentelemetry/OpenTelemetryWorker.java @@ -171,8 +171,37 @@ static String resolveEndpoint(Map env) { } static String resolveServiceName(Map env) { + return resolveServiceName(env, DEFAULT_SERVICE_NAME); + } + + /** + * Resolves an OpenTelemetry service name from the environment. + * + *

Resolution prefers {@code OTEL_SERVICE_NAME}, then each of {@code fallbackEnvVars} in order, + * and finally {@code defaultServiceName}. Environment values that are unset, empty, or blank are + * treated as absent. + * + * @param env environment variables to resolve from + * @param defaultServiceName value returned when no environment variable provides a service name + * @param fallbackEnvVars environment variable names consulted, in order, after {@code + * OTEL_SERVICE_NAME} + * @return the resolved service name + */ + public static String resolveServiceName( + Map env, String defaultServiceName, String... fallbackEnvVars) { String serviceName = nonEmptyEnv(env, OTEL_SERVICE_NAME); - return serviceName == null ? DEFAULT_SERVICE_NAME : serviceName; + if (serviceName != null) { + return serviceName; + } + if (fallbackEnvVars != null) { + for (String fallbackEnvVar : fallbackEnvVars) { + serviceName = nonEmptyEnv(env, fallbackEnvVar); + if (serviceName != null) { + return serviceName; + } + } + } + return defaultServiceName; } public static final class Builder { 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' diff --git a/temporal-bom/build.gradle b/temporal-bom/build.gradle index 031633473e..c79c4df44a 100644 --- a/temporal-bom/build.gradle +++ b/temporal-bom/build.gradle @@ -10,6 +10,7 @@ dependencies { api project(':temporal-opentelemetry') api project(':temporal-opentracing') api project(':temporal-aws-lambda') + api project(':temporal-gcp-cloud-run') api project(':temporal-remote-data-encoder') api project(':temporal-sdk') api project(':temporal-serviceclient')