diff --git a/pyproject.toml b/pyproject.toml index d6397a7b6..b0cc90008 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -44,6 +44,16 @@ lambda-worker-otel = [ "opentelemetry-semantic-conventions>=0.40b0,<1", "opentelemetry-sdk-extension-aws>=2.0.0,<3", ] +cloud-run-worker-otel = [ + "opentelemetry-api>=1.26,<2", + "opentelemetry-sdk>=1.26,<2", + "opentelemetry-exporter-otlp-proto-grpc>=1.11.1,<2", + # The OTLP exporter's generated protobuf code is not compatible with + # protobuf 7.x yet; cap it so `uv lock --upgrade` (test-latest-deps) does not + # downgrade the exporter to an ancient version to pull protobuf 7. A low + # exporter floor is kept so the protobuf<4 check-protos job can still resolve. + "protobuf<7", +] aioboto3 = ["aioboto3>=10.4.0", "types-aioboto3[s3]>=10.4.0"] google-genai = ["google-genai>=2.10.0,<3.0.0"] strands-agents = ["strands-agents>=1.39.0"] diff --git a/temporalio/contrib/aws/lambda_worker/otel.py b/temporalio/contrib/aws/lambda_worker/otel.py index 216e80f48..5b5dc91e1 100644 --- a/temporalio/contrib/aws/lambda_worker/otel.py +++ b/temporalio/contrib/aws/lambda_worker/otel.py @@ -11,18 +11,20 @@ from __future__ import annotations import logging -import os from dataclasses import dataclass, field from datetime import timedelta from opentelemetry.sdk.resources import Resource -from opentelemetry.sdk.trace.export import BatchSpanProcessor from opentelemetry.semconv.attributes.service_attributes import SERVICE_NAME from opentelemetry.trace import get_tracer_provider, set_tracer_provider from temporalio.contrib.aws.lambda_worker._configure import LambdaWorkerConfig -from temporalio.contrib.opentelemetry import OpenTelemetryPlugin, create_tracer_provider -from temporalio.runtime import OpenTelemetryConfig, Runtime, TelemetryConfig +from temporalio.contrib.opentelemetry import ( + OpenTelemetryPlugin, + _serverless, + create_tracer_provider, +) +from temporalio.runtime import Runtime, TelemetryConfig logger = logging.getLogger(__name__) @@ -51,23 +53,15 @@ class OtelOptions: def _resolve_service_name(options: OtelOptions) -> str: - service_name = options.service_name - if not service_name: - service_name = os.environ.get("OTEL_SERVICE_NAME", "") - if not service_name: - service_name = os.environ.get("AWS_LAMBDA_FUNCTION_NAME", "") - if not service_name: - service_name = "temporal-lambda-worker" - return service_name + return _serverless.resolve_service_name( + options.service_name, + ["AWS_LAMBDA_FUNCTION_NAME"], + "temporal-lambda-worker", + ) def _resolve_endpoint(options: OtelOptions) -> str: - endpoint = options.collector_endpoint - if not endpoint: - endpoint = os.environ.get("OTEL_EXPORTER_OTLP_ENDPOINT", "") - if not endpoint: - endpoint = "http://localhost:4317" - return endpoint + return _serverless.resolve_endpoint(options.collector_endpoint) def apply_defaults( @@ -135,12 +129,8 @@ def apply_defaults( # Use OTLP gRPC exporter if available, otherwise skip trace export. try: - from opentelemetry.exporter.otlp.proto.grpc.trace_exporter import ( - OTLPSpanExporter, - ) - tracer_provider.add_span_processor( - BatchSpanProcessor(OTLPSpanExporter(endpoint=endpoint, insecure=True)) + _serverless.build_otlp_span_processor(endpoint, insecure=True) ) except ImportError: logger.warning( @@ -195,23 +185,12 @@ def build_metrics_telemetry_config( A ``TelemetryConfig`` ready to pass to :py:class:`temporalio.runtime.Runtime`. """ - if not endpoint: - endpoint = "http://localhost:4317" - - otel_config = OpenTelemetryConfig( - url=endpoint, + return _serverless.build_metrics_telemetry_config( + endpoint=endpoint, + service_name=service_name, metric_periodicity=metric_periodicity, ) - global_tags: dict[str, str] = {} - if service_name: - global_tags["service_name"] = service_name - - return TelemetryConfig( - metrics=otel_config, - global_tags=global_tags, - ) - def apply_tracing(config: LambdaWorkerConfig) -> None: """Configure only OTel tracing (no metrics) on the Lambda worker config. diff --git a/temporalio/contrib/gcp/__init__.py b/temporalio/contrib/gcp/__init__.py new file mode 100644 index 000000000..f87ca7c64 --- /dev/null +++ b/temporalio/contrib/gcp/__init__.py @@ -0,0 +1 @@ +"""Google Cloud integrations for Temporal SDK.""" diff --git a/temporalio/contrib/gcp/cloud_run/README.md b/temporalio/contrib/gcp/cloud_run/README.md new file mode 100644 index 000000000..35dab08a6 --- /dev/null +++ b/temporalio/contrib/gcp/cloud_run/README.md @@ -0,0 +1,150 @@ +# Temporal Google Cloud Run integration + +`temporalio.contrib.gcp.cloud_run` provides an OpenTelemetry plugin with +defaults for Temporal Python 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 runtime +> and tracer provider. 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. + +## Installation + +Install the SDK with the Cloud Run OpenTelemetry dependencies: + +```bash +python -m pip install 'temporalio[cloud-run-worker-otel]' +``` + +## Usage + +Create one plugin and install it when connecting the Temporal client. Client +plugins automatically propagate to workers created from that client. + +```python +import asyncio +from datetime import timedelta + +from temporalio.client import Client +from temporalio.contrib.gcp.cloud_run import OpenTelemetryPlugin +from temporalio.worker import Worker + +otel_plugin = OpenTelemetryPlugin() +client = await Client.connect( + "localhost:7233", + plugins=[otel_plugin], +) +worker = Worker( + client, + task_queue="my-task-queue", + workflows=[MyWorkflow], + activities=[my_activity], +) + +try: + await worker.run() +finally: + # Run this only after every worker using the plugin has stopped. + await asyncio.to_thread(otel_plugin.shutdown, timedelta(seconds=2)) +``` + +The plugin configures Temporal Core metrics through a +`temporalio.runtime.Runtime` and delegates tracing to +`temporalio.contrib.opentelemetry.OpenTelemetryPlugin`. Do not install both the +Cloud Run and generic OpenTelemetry plugins on the same client. + +## Shutdown lifecycle + +`Worker.shutdown()` waits for worker shutdown to complete. Stop every worker +using the plugin, then call `OpenTelemetryPlugin.shutdown()` with the time +remaining before Cloud Run sends `SIGKILL`. The call force-flushes Python traces +and shuts down a tracer provider created by the plugin. An application-owned +provider is force-flushed but remains the application's responsibility. + +The Python SDK Core runtime exports metrics periodically and currently has no +explicit metrics-flush API. The default metric periodicity is 60 seconds, +matching the upstream OpenTelemetry SDK default. This avoids sending multiple +cumulative snapshots of the same Temporal series too frequently. Tune it with +`metric_periodicity` when the metrics backend semantics support a different +interval. + +Do not use a collector batch processor on the Google Managed Service for +Prometheus metrics pipeline. A runtime shutdown can force a metrics export +immediately after a periodic export, regardless of the periodicity. If the +collector combines both cumulative snapshots into one write, Google Managed +Service for Prometheus can reject it with `Duplicate TimeSeries`. Send each +OTLP metrics export directly through the memory limiter, GCP resource detection, +and any required resource/metric transforms to the Google Managed Prometheus +exporter. Keep a separate batch processor on the traces pipeline; trace batching +does not have the cumulative time-series collision. + +`flush_on_worker_stop=True` enables a best-effort trace flush after an +individual worker stops. It is disabled by default because one plugin can be +used by multiple workers and the remaining workers can emit telemetry after the +first worker stops. + +The OTLP endpoint is resolved in this order: + +1. `OpenTelemetryPlugin(endpoint=...)`. +2. `OTEL_EXPORTER_OTLP_ENDPOINT`. +3. `http://localhost:4317`. + +The OpenTelemetry service name is resolved in this order: + +1. `OpenTelemetryPlugin(service_name=...)`. +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 plugin owns the runtime and replay-safe tracer provider it creates. To use +an application-owned runtime or provider, pass `runtime=` or `tracer_provider=`. +An application-owned provider must be created with +`temporalio.contrib.opentelemetry.create_tracer_provider`. Use +`build_metrics_telemetry_config()` to compose GCP metrics defaults with custom +runtime logging or other telemetry settings. Do not also pass a different +`runtime` to `Client.connect()`. + +OpenTelemetry has one global tracer provider per process. Create the Cloud Run +plugin before installing another provider, or pass the already-installed +replay-safe provider through `tracer_provider=`. The plugin raises instead of +silently using a different global provider. + +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. Configure +separate pipelines: metrics without a batch processor, and traces with a +dedicated batch processor. For the image, 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. + +Cloud Run worker pools support sidecar containers over localhost and are +intended for continuous background work. 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 pass +`endpoint=` to the plugin. diff --git a/temporalio/contrib/gcp/cloud_run/__init__.py b/temporalio/contrib/gcp/cloud_run/__init__.py new file mode 100644 index 000000000..2771dd7fd --- /dev/null +++ b/temporalio/contrib/gcp/cloud_run/__init__.py @@ -0,0 +1,36 @@ +"""OpenTelemetry integration for Temporal workers running on Google Cloud Run. + +This package is designed for container-based Cloud Run workers, especially +worker pools, and exports telemetry to an OpenTelemetry collector sidecar by +default. + +.. warning:: + This package is experimental and may change in future versions. + Use with caution in production environments. +""" + +from temporalio.contrib.gcp.cloud_run._opentelemetry import ( + CLOUD_RUN_SERVICE_ENV_VAR, + CLOUD_RUN_WORKER_POOL_ENV_VAR, + DEFAULT_FLUSH_TIMEOUT, + DEFAULT_METRIC_PERIODICITY, + DEFAULT_OTLP_ENDPOINT, + DEFAULT_SERVICE_NAME, + OTEL_EXPORTER_OTLP_ENDPOINT_ENV_VAR, + OTEL_SERVICE_NAME_ENV_VAR, + OpenTelemetryPlugin, + build_metrics_telemetry_config, +) + +__all__ = [ + "CLOUD_RUN_SERVICE_ENV_VAR", + "CLOUD_RUN_WORKER_POOL_ENV_VAR", + "DEFAULT_FLUSH_TIMEOUT", + "DEFAULT_METRIC_PERIODICITY", + "DEFAULT_OTLP_ENDPOINT", + "DEFAULT_SERVICE_NAME", + "OTEL_EXPORTER_OTLP_ENDPOINT_ENV_VAR", + "OTEL_SERVICE_NAME_ENV_VAR", + "OpenTelemetryPlugin", + "build_metrics_telemetry_config", +] diff --git a/temporalio/contrib/gcp/cloud_run/_opentelemetry.py b/temporalio/contrib/gcp/cloud_run/_opentelemetry.py new file mode 100644 index 000000000..e8fff526f --- /dev/null +++ b/temporalio/contrib/gcp/cloud_run/_opentelemetry.py @@ -0,0 +1,351 @@ +"""OpenTelemetry plugin with Google Cloud Run defaults.""" + +from __future__ import annotations + +import asyncio +import logging +import threading +from collections.abc import Awaitable, Callable +from datetime import timedelta + +from opentelemetry.sdk.resources import Resource +from opentelemetry.trace import ( + ProxyTracerProvider, + TracerProvider, + get_tracer_provider, + set_tracer_provider, +) + +from temporalio.contrib.opentelemetry import ( + OpenTelemetryPlugin as TemporalOpenTelemetryPlugin, +) +from temporalio.contrib.opentelemetry import _serverless, create_tracer_provider +from temporalio.contrib.opentelemetry._tracer_provider import ( + ReplaySafeTracerProvider, +) +from temporalio.runtime import Runtime, TelemetryConfig +from temporalio.service import ConnectConfig, ServiceClient +from temporalio.worker import Worker + +logger = logging.getLogger(__name__) + +OTEL_EXPORTER_OTLP_ENDPOINT_ENV_VAR = _serverless.OTEL_EXPORTER_OTLP_ENDPOINT_ENV_VAR +"""Standard OpenTelemetry environment variable for the common OTLP endpoint.""" + +OTEL_SERVICE_NAME_ENV_VAR = _serverless.OTEL_SERVICE_NAME_ENV_VAR +"""Standard OpenTelemetry environment variable for ``service.name``.""" + +CLOUD_RUN_WORKER_POOL_ENV_VAR = "CLOUD_RUN_WORKER_POOL" +"""Cloud Run environment variable containing the worker-pool name.""" + +CLOUD_RUN_SERVICE_ENV_VAR = "K_SERVICE" +"""Cloud Run environment variable containing the service name.""" + +DEFAULT_OTLP_ENDPOINT = _serverless.DEFAULT_OTLP_ENDPOINT +"""Default local OTLP gRPC collector endpoint.""" + +DEFAULT_SERVICE_NAME = "temporal-worker" +"""Service name used outside a recognized Cloud Run environment.""" + +DEFAULT_METRIC_PERIODICITY = timedelta(seconds=60) +"""Default interval between Temporal Core metric exports.""" + +DEFAULT_FLUSH_TIMEOUT = timedelta(seconds=10) +"""Default timeout for tracing force-flush operations.""" + + +class OpenTelemetryPlugin(TemporalOpenTelemetryPlugin): + """OpenTelemetry plugin for Temporal workers running on Google Cloud Run. + + The default configuration sends Temporal Core metrics and Python traces over + OTLP gRPC to a collector on ``localhost:4317``. The collector is responsible + for Google Cloud resource detection, authentication, and export. + + The plugin creates and installs a replay-safe global tracer provider unless + ``tracer_provider`` is supplied. It also creates a Temporal runtime for Core + metrics unless ``runtime`` is supplied. Call :py:meth:`shutdown` after every + worker using the plugin has stopped. + + .. warning:: + This class is experimental and may change in future versions. + Use with caution in production environments. + """ + + def __init__( + self, + *, + endpoint: str | None = None, + service_name: str | None = None, + metric_periodicity: timedelta | None = None, + flush_timeout: timedelta = DEFAULT_FLUSH_TIMEOUT, + flush_on_worker_stop: bool = False, + tracer_provider: TracerProvider | None = None, + runtime: Runtime | None = None, + add_temporal_spans: bool = False, + ) -> None: + """Initialize the Google Cloud OpenTelemetry plugin. + + Args: + endpoint: OTLP gRPC endpoint. Falls back to + ``OTEL_EXPORTER_OTLP_ENDPOINT``, then ``http://localhost:4317``. + service_name: OpenTelemetry service name. Falls back to + ``OTEL_SERVICE_NAME``, ``CLOUD_RUN_WORKER_POOL``, ``K_SERVICE``, + then ``temporal-worker``. + metric_periodicity: How often Temporal Core metrics are exported. + Defaults to 60 seconds. Cannot be used with ``runtime``. + flush_timeout: Default tracing force-flush timeout. + flush_on_worker_stop: Whether to force-flush traces after each + worker stops. Disabled by default because one plugin can be used + by multiple workers. + tracer_provider: Application-owned replay-safe tracer provider. It + must have been created with + :py:func:`temporalio.contrib.opentelemetry.create_tracer_provider`. + When supplied, the plugin does not create an exporter or shut + down the provider. + runtime: Application-owned Temporal runtime. When supplied, the + plugin does not create or modify Core metrics configuration. Use + :py:func:`build_metrics_telemetry_config` to build a composable + telemetry configuration for a custom runtime. + add_temporal_spans: Whether the underlying Temporal OpenTelemetry + plugin should add Temporal-specific operation spans. + """ + _validate_positive_duration("flush_timeout", flush_timeout) + if metric_periodicity is not None: + _validate_positive_duration("metric_periodicity", metric_periodicity) + if runtime is not None and metric_periodicity is not None: + raise ValueError("metric_periodicity cannot be set with runtime") + + resolved_endpoint = _resolve_endpoint(endpoint) + resolved_service_name = _resolve_service_name(service_name) + resolved_metric_periodicity = ( + metric_periodicity + if metric_periodicity is not None + else DEFAULT_METRIC_PERIODICITY + ) + + owns_tracer_provider = tracer_provider is None + if tracer_provider is None: + resolved_tracer_provider = _create_tracer_provider( + resolved_endpoint, resolved_service_name + ) + elif isinstance(tracer_provider, ReplaySafeTracerProvider): + resolved_tracer_provider = tracer_provider + else: + raise TypeError( + "tracer_provider must be created with " + "temporalio.contrib.opentelemetry.create_tracer_provider" + ) + + try: + if runtime is None: + runtime = Runtime( + telemetry=_serverless.build_metrics_telemetry_config( + endpoint=resolved_endpoint, + service_name=resolved_service_name, + metric_periodicity=resolved_metric_periodicity, + ) + ) + _install_global_tracer_provider(resolved_tracer_provider) + except BaseException: + if owns_tracer_provider: + resolved_tracer_provider.shutdown() + raise + + self._endpoint = resolved_endpoint + self._service_name = resolved_service_name + self._flush_timeout = flush_timeout + self._flush_on_worker_stop = flush_on_worker_stop + self._tracer_provider = resolved_tracer_provider + self._owns_tracer_provider = owns_tracer_provider + self._runtime = runtime + self._shutdown_lock = threading.Lock() + self._shutdown = False + self._shutdown_succeeded = True + + super().__init__(add_temporal_spans=add_temporal_spans) + + @property + def endpoint(self) -> str: + """Resolved OTLP collector endpoint.""" + return self._endpoint + + @property + def service_name(self) -> str: + """Resolved OpenTelemetry service name.""" + return self._service_name + + @property + def tracer_provider(self) -> ReplaySafeTracerProvider: + """Replay-safe tracer provider used by the plugin.""" + return self._tracer_provider + + @property + def runtime(self) -> Runtime: + """Temporal runtime used for Core metrics.""" + return self._runtime + + async def connect_service_client( + self, + config: ConnectConfig, + next: Callable[[ConnectConfig], Awaitable[ServiceClient]], + ) -> ServiceClient: + """Install the metrics runtime before connecting the service client.""" + if config.runtime is not None and config.runtime is not self._runtime: + raise ValueError( + "OpenTelemetryPlugin runtime conflicts with the runtime passed " + "to Client.connect; pass that runtime to the plugin instead" + ) + config.runtime = self._runtime + return await super().connect_service_client(config, next) + + async def run_worker( + self, worker: Worker, next: Callable[[Worker], Awaitable[None]] + ) -> None: + """Run a worker and optionally force-flush traces after it stops.""" + try: + await super().run_worker(worker, next) + finally: + if self._flush_on_worker_stop: + try: + succeeded = await asyncio.to_thread( + self.force_flush, self._flush_timeout + ) + if not succeeded: + logger.warning( + "OpenTelemetry trace flush timed out after worker stop" + ) + except Exception: + logger.exception( + "OpenTelemetry trace flush failed after worker stop" + ) + + def force_flush(self, timeout: timedelta | None = None) -> bool: + """Export buffered Python traces without shutting down the provider. + + Temporal Core metrics are exported periodically and the Python runtime + currently has no explicit metrics-flush API. + + Args: + timeout: Maximum time to wait. Defaults to ``flush_timeout`` from + the constructor. + + Returns: + ``True`` when the tracer provider reports a successful flush. + """ + resolved_timeout = timeout if timeout is not None else self._flush_timeout + _validate_positive_duration("timeout", resolved_timeout) + return self._tracer_provider.force_flush( + _duration_to_milliseconds(resolved_timeout) + ) + + def shutdown(self, timeout: timedelta | None = None) -> bool: + """Flush traces and shut down a provider created by the plugin. + + Application-owned tracer providers are force-flushed but are not shut + down. This method is idempotent. Stop every worker using the plugin + before calling it. + + Args: + timeout: Maximum time allowed for the tracing force-flush. Provider + shutdown happens after that flush. + + Returns: + ``True`` when the tracing force-flush succeeded. + """ + with self._shutdown_lock: + if self._shutdown: + return self._shutdown_succeeded + + succeeded = self.force_flush(timeout) + if self._owns_tracer_provider: + self._tracer_provider.shutdown() + self._shutdown = True + self._shutdown_succeeded = succeeded + return succeeded + + +def build_metrics_telemetry_config( + *, + endpoint: str | None = None, + service_name: str | None = None, + metric_periodicity: timedelta | None = None, +) -> TelemetryConfig: + """Build Temporal Core telemetry with Google Cloud Run defaults. + + This helper is useful when the application needs to customize runtime + logging or other telemetry settings before constructing a + :py:class:`temporalio.runtime.Runtime`. + + Args: + endpoint: OTLP gRPC endpoint, with the same resolution order as + :py:class:`OpenTelemetryPlugin`. + service_name: Service name, with the same resolution order as + :py:class:`OpenTelemetryPlugin`. + metric_periodicity: Metric export interval. Defaults to 60 seconds. + + Returns: + Telemetry configuration ready for a Temporal runtime. + """ + resolved_periodicity = ( + metric_periodicity + if metric_periodicity is not None + else DEFAULT_METRIC_PERIODICITY + ) + _validate_positive_duration("metric_periodicity", resolved_periodicity) + return _serverless.build_metrics_telemetry_config( + endpoint=_resolve_endpoint(endpoint), + service_name=_resolve_service_name(service_name), + metric_periodicity=resolved_periodicity, + ) + + +def _create_tracer_provider( + endpoint: str, service_name: str +) -> ReplaySafeTracerProvider: + provider = create_tracer_provider( + resource=Resource.create({"service.name": service_name}) + ) + try: + provider.add_span_processor(_serverless.build_otlp_span_processor(endpoint)) + except ImportError as err: + raise RuntimeError( + "The OTLP gRPC exporter is required. Install the " + "'temporalio[cloud-run-worker-otel]' extra." + ) from err + return provider + + +def _install_global_tracer_provider(provider: ReplaySafeTracerProvider) -> None: + current = get_tracer_provider() + if current is provider: + return + if not isinstance(current, ProxyTracerProvider): + raise RuntimeError( + "The global OpenTelemetry tracer provider is already configured. " + "Pass that provider as tracer_provider if it was created with " + "temporalio.contrib.opentelemetry.create_tracer_provider." + ) + set_tracer_provider(provider) + if get_tracer_provider() is not provider: + raise RuntimeError("Failed to install the OpenTelemetry tracer provider") + + +def _resolve_endpoint(explicit: str | None) -> str: + return _serverless.resolve_endpoint(explicit) + + +def _resolve_service_name(explicit: str | None) -> str: + return _serverless.resolve_service_name( + explicit, + [CLOUD_RUN_WORKER_POOL_ENV_VAR, CLOUD_RUN_SERVICE_ENV_VAR], + DEFAULT_SERVICE_NAME, + ) + + +def _validate_positive_duration(name: str, value: timedelta) -> None: + if value <= timedelta(0): + raise ValueError(f"{name} must be positive") + + +def _duration_to_milliseconds(value: timedelta) -> int: + return max(1, round(value.total_seconds() * 1000)) diff --git a/temporalio/contrib/opentelemetry/_serverless.py b/temporalio/contrib/opentelemetry/_serverless.py new file mode 100644 index 000000000..37fdd5e86 --- /dev/null +++ b/temporalio/contrib/opentelemetry/_serverless.py @@ -0,0 +1,147 @@ +"""Provider-neutral OpenTelemetry helpers shared by serverless integrations. + +This is a private module. It is intentionally **not** re-exported from +``temporalio.contrib.opentelemetry`` so that importing that package gains no new +imports. The OTLP gRPC span exporter is imported lazily inside +:py:func:`build_otlp_span_processor`, so importing this module does not require +the ``opentelemetry-exporter-otlp-proto-grpc`` dependency. + +The AWS Lambda and GCP Cloud Run integrations layer their provider-specific +policy (ID generators, default service names, environment fallbacks) on top of +these helpers. Empty environment/argument values are skipped with a plain +truthiness check and are not stripped of whitespace, matching the pre-existing +AWS Lambda resolution behavior. +""" + +from __future__ import annotations + +import os +from collections.abc import Mapping, Sequence +from datetime import timedelta + +from opentelemetry.sdk.trace.export import BatchSpanProcessor + +from temporalio.runtime import OpenTelemetryConfig, TelemetryConfig + +OTEL_EXPORTER_OTLP_ENDPOINT_ENV_VAR = "OTEL_EXPORTER_OTLP_ENDPOINT" +"""Standard OpenTelemetry environment variable for the common OTLP endpoint.""" + +OTEL_SERVICE_NAME_ENV_VAR = "OTEL_SERVICE_NAME" +"""Standard OpenTelemetry environment variable for ``service.name``.""" + +DEFAULT_OTLP_ENDPOINT = "http://localhost:4317" +"""Default local OTLP gRPC collector endpoint.""" + + +def resolve_endpoint( + explicit: str | None, + *, + env: Mapping[str, str] = os.environ, + default: str = DEFAULT_OTLP_ENDPOINT, +) -> str: + """Resolve the OTLP collector endpoint. + + Resolution order: ``explicit`` -> ``OTEL_EXPORTER_OTLP_ENDPOINT`` -> + ``default``. Empty values are skipped with a truthiness check, without + stripping whitespace. + + Args: + explicit: Endpoint supplied directly by the caller, if any. + env: Environment mapping. Defaults to the live process environment. + default: Endpoint used when nothing else is provided. + + Returns: + The resolved endpoint. + """ + return explicit or env.get(OTEL_EXPORTER_OTLP_ENDPOINT_ENV_VAR) or default + + +def resolve_service_name( + explicit: str | None, + fallback_env_vars: Sequence[str], + default: str, + *, + env: Mapping[str, str] = os.environ, +) -> str: + """Resolve the OpenTelemetry service name. + + Resolution order: ``explicit`` -> ``OTEL_SERVICE_NAME`` -> each name in + ``fallback_env_vars`` in order -> ``default``. Empty values are skipped with + a truthiness check, without stripping whitespace. + + Args: + explicit: Service name supplied directly by the caller, if any. + fallback_env_vars: Provider-specific environment variable names checked, + in order, after ``OTEL_SERVICE_NAME``. + default: Service name used when nothing else is provided. + env: Environment mapping. Defaults to the live process environment. + + Returns: + The resolved service name. + """ + if explicit: + return explicit + otel_service_name = env.get(OTEL_SERVICE_NAME_ENV_VAR) + if otel_service_name: + return otel_service_name + for name in fallback_env_vars: + value = env.get(name) + if value: + return value + return default + + +def build_metrics_telemetry_config( + *, + endpoint: str, + service_name: str, + metric_periodicity: timedelta | None, +) -> TelemetryConfig: + """Build Core telemetry configuration for OTLP metrics export. + + Args: + endpoint: OTLP collector endpoint. Falls back to + :py:data:`DEFAULT_OTLP_ENDPOINT` when empty. + service_name: Service name added as the ``service_name`` global tag. + When empty, no global tag is added. + metric_periodicity: Metric export interval, passed through unchanged. + + Returns: + A :py:class:`temporalio.runtime.TelemetryConfig` with metrics pointed at + the collector. + """ + return TelemetryConfig( + metrics=OpenTelemetryConfig( + url=endpoint or DEFAULT_OTLP_ENDPOINT, + metric_periodicity=metric_periodicity, + ), + global_tags={"service_name": service_name} if service_name else {}, + ) + + +def build_otlp_span_processor( + endpoint: str, + *, + insecure: bool = True, +) -> BatchSpanProcessor: + """Build a batch span processor backed by the OTLP gRPC exporter. + + The exporter is imported lazily so that importing this module does not + require ``opentelemetry-exporter-otlp-proto-grpc``. + + Args: + endpoint: OTLP collector endpoint. + insecure: Whether to use an insecure (non-TLS) gRPC channel. + + Returns: + A batch span processor that exports to the OTLP collector. + + Raises: + ImportError: If the OTLP gRPC exporter is not installed. The caller + decides whether to warn and continue or re-raise. + """ + from opentelemetry.exporter.otlp.proto.grpc.trace_exporter import ( # type: ignore[reportMissingTypeStubs] + OTLPSpanExporter, + ) + + return BatchSpanProcessor(OTLPSpanExporter(endpoint=endpoint, insecure=insecure)) diff --git a/tests/contrib/gcp/__init__.py b/tests/contrib/gcp/__init__.py new file mode 100644 index 000000000..e69de29bb diff --git a/tests/contrib/gcp/cloud_run/__init__.py b/tests/contrib/gcp/cloud_run/__init__.py new file mode 100644 index 000000000..e69de29bb diff --git a/tests/contrib/gcp/cloud_run/test_opentelemetry.py b/tests/contrib/gcp/cloud_run/test_opentelemetry.py new file mode 100644 index 000000000..b69eeabca --- /dev/null +++ b/tests/contrib/gcp/cloud_run/test_opentelemetry.py @@ -0,0 +1,314 @@ +"""Tests for the Google Cloud Run OpenTelemetry plugin.""" + +from __future__ import annotations + +from datetime import timedelta +from typing import Any, cast +from unittest.mock import Mock, call + +import pytest +from opentelemetry.trace import NoOpTracerProvider, set_tracer_provider + +from temporalio.client import ClientConfig +from temporalio.contrib.gcp.cloud_run import ( + CLOUD_RUN_SERVICE_ENV_VAR, + CLOUD_RUN_WORKER_POOL_ENV_VAR, + DEFAULT_METRIC_PERIODICITY, + DEFAULT_OTLP_ENDPOINT, + DEFAULT_SERVICE_NAME, + OTEL_EXPORTER_OTLP_ENDPOINT_ENV_VAR, + OTEL_SERVICE_NAME_ENV_VAR, + OpenTelemetryPlugin, + build_metrics_telemetry_config, +) +from temporalio.contrib.opentelemetry import ( + OpenTelemetryInterceptor, + create_tracer_provider, +) +from temporalio.contrib.opentelemetry._tracer_provider import ( + ReplaySafeTracerProvider, +) +from temporalio.runtime import OpenTelemetryConfig, Runtime +from temporalio.service import ConnectConfig, ServiceClient +from temporalio.worker import Worker + + +@pytest.fixture(autouse=True) +def _reset_global_provider( # pyright: ignore[reportUnusedFunction] + reset_otel_tracer_provider: None, # pyright: ignore[reportUnusedParameter] +) -> None: + pass + + +@pytest.fixture +def tracer_provider() -> ReplaySafeTracerProvider: + return create_tracer_provider(shutdown_on_exit=False) + + +def _application_owned_plugin( + tracer_provider: ReplaySafeTracerProvider, + **kwargs: Any, +) -> OpenTelemetryPlugin: + return OpenTelemetryPlugin( + tracer_provider=tracer_provider, + runtime=cast(Runtime, Mock(spec=Runtime)), + **kwargs, + ) + + +def _clear_environment(monkeypatch: pytest.MonkeyPatch) -> None: + for name in ( + OTEL_EXPORTER_OTLP_ENDPOINT_ENV_VAR, + OTEL_SERVICE_NAME_ENV_VAR, + CLOUD_RUN_WORKER_POOL_ENV_VAR, + CLOUD_RUN_SERVICE_ENV_VAR, + ): + monkeypatch.delenv(name, raising=False) + + +def test_defaults_and_plugin_integration( + tracer_provider: ReplaySafeTracerProvider, monkeypatch: pytest.MonkeyPatch +) -> None: + _clear_environment(monkeypatch) + + plugin = _application_owned_plugin(tracer_provider) + + assert plugin.endpoint == DEFAULT_OTLP_ENDPOINT + assert plugin.service_name == DEFAULT_SERVICE_NAME + assert plugin.tracer_provider is tracer_provider + assert plugin.name() == "OpenTelemetryPlugin" + + config = plugin.configure_client( + cast(ClientConfig, cast(object, {"interceptors": []})) + ) + assert len(config["interceptors"]) == 1 + assert isinstance(config["interceptors"][0], OpenTelemetryInterceptor) + + +def test_resolution_precedence( + tracer_provider: ReplaySafeTracerProvider, monkeypatch: pytest.MonkeyPatch +) -> None: + _clear_environment(monkeypatch) + monkeypatch.setenv(CLOUD_RUN_SERVICE_ENV_VAR, "cloud-run-service") + plugin = _application_owned_plugin(tracer_provider) + assert plugin.service_name == "cloud-run-service" + + monkeypatch.setenv(CLOUD_RUN_WORKER_POOL_ENV_VAR, "worker-pool") + plugin = _application_owned_plugin(tracer_provider) + assert plugin.service_name == "worker-pool" + + monkeypatch.setenv(OTEL_SERVICE_NAME_ENV_VAR, "otel-service") + monkeypatch.setenv(OTEL_EXPORTER_OTLP_ENDPOINT_ENV_VAR, "http://collector:4317") + plugin = _application_owned_plugin(tracer_provider) + assert plugin.service_name == "otel-service" + assert plugin.endpoint == "http://collector:4317" + + plugin = _application_owned_plugin( + tracer_provider, + endpoint="https://explicit-collector:4317", + service_name="explicit-service", + ) + assert plugin.service_name == "explicit-service" + assert plugin.endpoint == "https://explicit-collector:4317" + + +def test_ignores_empty_environment_values( + tracer_provider: ReplaySafeTracerProvider, monkeypatch: pytest.MonkeyPatch +) -> None: + # Empty environment values are falsy and are skipped, falling through to + # the next source in the chain. Whitespace is not stripped, so only truly + # empty values are ignored (matching the shared serverless behavior). + monkeypatch.setenv(OTEL_SERVICE_NAME_ENV_VAR, "") + monkeypatch.setenv(CLOUD_RUN_WORKER_POOL_ENV_VAR, "") + monkeypatch.setenv(CLOUD_RUN_SERVICE_ENV_VAR, "cloud-run-service") + + plugin = _application_owned_plugin(tracer_provider) + + assert plugin.service_name == "cloud-run-service" + + +def test_build_metrics_telemetry_config(monkeypatch: pytest.MonkeyPatch) -> None: + _clear_environment(monkeypatch) + monkeypatch.setenv(CLOUD_RUN_WORKER_POOL_ENV_VAR, "worker-pool") + config = build_metrics_telemetry_config( + endpoint="http://collector:4317", + metric_periodicity=timedelta(seconds=30), + ) + + assert isinstance(config.metrics, OpenTelemetryConfig) + assert config.metrics.url == "http://collector:4317" + assert config.metrics.metric_periodicity == timedelta(seconds=30) + assert config.global_tags == {"service_name": "worker-pool"} + + +def test_default_metric_periodicity_is_sixty_seconds( + monkeypatch: pytest.MonkeyPatch, +) -> None: + _clear_environment(monkeypatch) + + config = build_metrics_telemetry_config() + + assert DEFAULT_METRIC_PERIODICITY == timedelta(seconds=60) + assert isinstance(config.metrics, OpenTelemetryConfig) + assert config.metrics.metric_periodicity == timedelta(seconds=60) + + +@pytest.mark.asyncio +async def test_connect_service_client_installs_runtime( + tracer_provider: ReplaySafeTracerProvider, +) -> None: + runtime = cast(Runtime, Mock(spec=Runtime)) + plugin = OpenTelemetryPlugin( + tracer_provider=tracer_provider, + runtime=runtime, + ) + config = ConnectConfig(target_host="localhost:7233") + service_client = cast(ServiceClient, Mock(spec=ServiceClient)) + + async def connect(input: ConnectConfig) -> ServiceClient: + assert input.runtime is runtime + return service_client + + assert await plugin.connect_service_client(config, connect) is service_client + + conflicting_config = ConnectConfig( + target_host="localhost:7233", + runtime=cast(Runtime, Mock(spec=Runtime)), + ) + with pytest.raises(ValueError, match="runtime conflicts"): + await plugin.connect_service_client(conflicting_config, connect) + + +def test_force_flush_and_application_owned_shutdown( + tracer_provider: ReplaySafeTracerProvider, monkeypatch: pytest.MonkeyPatch +) -> None: + force_flush = Mock(return_value=True) + shutdown = Mock() + monkeypatch.setattr(tracer_provider, "force_flush", force_flush) + monkeypatch.setattr(tracer_provider, "shutdown", shutdown) + plugin = _application_owned_plugin(tracer_provider) + + assert plugin.force_flush(timedelta(seconds=2)) + assert plugin.shutdown(timedelta(seconds=3)) + assert plugin.shutdown(timedelta(seconds=4)) + + assert force_flush.call_args_list == [call(2000), call(3000)] + shutdown.assert_not_called() + + +def test_plugin_owned_provider_is_shut_down( + tracer_provider: ReplaySafeTracerProvider, monkeypatch: pytest.MonkeyPatch +) -> None: + force_flush = Mock(return_value=True) + shutdown = Mock() + monkeypatch.setattr(tracer_provider, "force_flush", force_flush) + monkeypatch.setattr(tracer_provider, "shutdown", shutdown) + monkeypatch.setattr( + "temporalio.contrib.gcp.cloud_run._opentelemetry._create_tracer_provider", + lambda endpoint, service_name: tracer_provider, + ) + runtime = cast(Runtime, Mock(spec=Runtime)) + runtime_factory = Mock(return_value=runtime) + monkeypatch.setattr( + "temporalio.contrib.gcp.cloud_run._opentelemetry.Runtime", runtime_factory + ) + + plugin = OpenTelemetryPlugin() + + assert plugin.runtime is runtime + telemetry = runtime_factory.call_args.kwargs["telemetry"] + assert isinstance(telemetry.metrics, OpenTelemetryConfig) + assert telemetry.metrics.url == DEFAULT_OTLP_ENDPOINT + assert telemetry.metrics.metric_periodicity == DEFAULT_METRIC_PERIODICITY + assert telemetry.global_tags == {"service_name": DEFAULT_SERVICE_NAME} + assert plugin.shutdown(timedelta(seconds=2)) + force_flush.assert_called_once_with(2000) + shutdown.assert_called_once_with() + + +def test_plugin_owned_provider_is_cleaned_up_on_runtime_failure( + tracer_provider: ReplaySafeTracerProvider, monkeypatch: pytest.MonkeyPatch +) -> None: + shutdown = Mock() + monkeypatch.setattr(tracer_provider, "shutdown", shutdown) + monkeypatch.setattr( + "temporalio.contrib.gcp.cloud_run._opentelemetry._create_tracer_provider", + lambda endpoint, service_name: tracer_provider, + ) + monkeypatch.setattr( + "temporalio.contrib.gcp.cloud_run._opentelemetry.Runtime", + Mock(side_effect=RuntimeError("runtime failed")), + ) + + with pytest.raises(RuntimeError, match="runtime failed"): + OpenTelemetryPlugin() + + shutdown.assert_called_once_with() + + +def test_rejects_conflicting_global_tracer_provider( + tracer_provider: ReplaySafeTracerProvider, monkeypatch: pytest.MonkeyPatch +) -> None: + existing_provider = create_tracer_provider(shutdown_on_exit=False) + set_tracer_provider(existing_provider) + shutdown = Mock() + monkeypatch.setattr(tracer_provider, "shutdown", shutdown) + monkeypatch.setattr( + "temporalio.contrib.gcp.cloud_run._opentelemetry._create_tracer_provider", + lambda endpoint, service_name: tracer_provider, + ) + + with pytest.raises(RuntimeError, match="already configured"): + OpenTelemetryPlugin(runtime=cast(Runtime, Mock(spec=Runtime))) + + shutdown.assert_called_once_with() + + +@pytest.mark.asyncio +@pytest.mark.parametrize("flush_on_worker_stop", [False, True]) +async def test_worker_stop_flush_is_opt_in( + tracer_provider: ReplaySafeTracerProvider, + monkeypatch: pytest.MonkeyPatch, + flush_on_worker_stop: bool, +) -> None: + calls: list[str] = [] + + def record_flush(_timeout_millis: int) -> bool: + calls.append("flush") + return True + + force_flush = Mock(side_effect=record_flush) + monkeypatch.setattr(tracer_provider, "force_flush", force_flush) + plugin = _application_owned_plugin( + tracer_provider, + flush_on_worker_stop=flush_on_worker_stop, + flush_timeout=timedelta(seconds=2), + ) + + async def run(_worker: Worker) -> None: + calls.append("run") + + await plugin.run_worker(cast(Worker, Mock(spec=Worker)), run) + + assert calls == (["run", "flush"] if flush_on_worker_stop else ["run"]) + + +def test_validates_options(tracer_provider: ReplaySafeTracerProvider) -> None: + with pytest.raises(ValueError, match="flush_timeout must be positive"): + _application_owned_plugin(tracer_provider, flush_timeout=timedelta(seconds=-1)) + with pytest.raises(ValueError, match="metric_periodicity cannot be set"): + _application_owned_plugin( + tracer_provider, metric_periodicity=timedelta(seconds=1) + ) + with pytest.raises(ValueError, match="metric_periodicity must be positive"): + build_metrics_telemetry_config(metric_periodicity=timedelta(0)) + + plugin = _application_owned_plugin(tracer_provider) + with pytest.raises(ValueError, match="timeout must be positive"): + plugin.force_flush(timedelta(0)) + + with pytest.raises(TypeError, match="create_tracer_provider"): + OpenTelemetryPlugin( + tracer_provider=NoOpTracerProvider(), + runtime=cast(Runtime, Mock(spec=Runtime)), + ) diff --git a/tests/contrib/opentelemetry/test_serverless.py b/tests/contrib/opentelemetry/test_serverless.py new file mode 100644 index 000000000..a80b53b09 --- /dev/null +++ b/tests/contrib/opentelemetry/test_serverless.py @@ -0,0 +1,179 @@ +"""Tests for temporalio.contrib.opentelemetry._serverless.""" + +from __future__ import annotations + +import sys +from datetime import timedelta + +import pytest +from opentelemetry.sdk.trace.export import BatchSpanProcessor + +from temporalio.contrib.opentelemetry import _serverless +from temporalio.runtime import OpenTelemetryConfig + + +class TestResolveEndpoint: + def test_explicit_wins(self) -> None: + assert ( + _serverless.resolve_endpoint( + "http://explicit:4317", + env={"OTEL_EXPORTER_OTLP_ENDPOINT": "http://env:4317"}, + ) + == "http://explicit:4317" + ) + + def test_env_fallback(self) -> None: + assert ( + _serverless.resolve_endpoint( + None, env={"OTEL_EXPORTER_OTLP_ENDPOINT": "http://env:4317"} + ) + == "http://env:4317" + ) + + def test_default_when_nothing_set(self) -> None: + assert _serverless.resolve_endpoint(None, env={}) == ( + _serverless.DEFAULT_OTLP_ENDPOINT + ) + + def test_empty_explicit_and_env_are_skipped(self) -> None: + assert ( + _serverless.resolve_endpoint("", env={"OTEL_EXPORTER_OTLP_ENDPOINT": ""}) + == _serverless.DEFAULT_OTLP_ENDPOINT + ) + + def test_custom_default(self) -> None: + assert ( + _serverless.resolve_endpoint(None, env={}, default="http://custom:4317") + == "http://custom:4317" + ) + + +class TestResolveServiceName: + def test_explicit_wins(self) -> None: + assert ( + _serverless.resolve_service_name( + "explicit", + ["FALLBACK"], + "default", + env={"OTEL_SERVICE_NAME": "otel", "FALLBACK": "fallback"}, + ) + == "explicit" + ) + + def test_otel_service_name_precedes_fallbacks(self) -> None: + assert ( + _serverless.resolve_service_name( + None, + ["FALLBACK"], + "default", + env={"OTEL_SERVICE_NAME": "otel", "FALLBACK": "fallback"}, + ) + == "otel" + ) + + def test_fallback_env_vars_checked_in_order(self) -> None: + assert ( + _serverless.resolve_service_name( + None, + ["FIRST", "SECOND"], + "default", + env={"SECOND": "second-value"}, + ) + == "second-value" + ) + assert ( + _serverless.resolve_service_name( + None, + ["FIRST", "SECOND"], + "default", + env={"FIRST": "first-value", "SECOND": "second-value"}, + ) + == "first-value" + ) + + def test_default_when_nothing_set(self) -> None: + assert ( + _serverless.resolve_service_name(None, ["FALLBACK"], "default", env={}) + == "default" + ) + + def test_empty_values_are_skipped(self) -> None: + # Empty strings are falsy and are skipped, falling through to the next + # source in the chain. + assert ( + _serverless.resolve_service_name( + "", + ["FIRST", "SECOND"], + "default", + env={"OTEL_SERVICE_NAME": "", "FIRST": "", "SECOND": "second-value"}, + ) + == "second-value" + ) + + def test_whitespace_is_not_stripped(self) -> None: + # No stripping: a whitespace-only value is truthy and is used as-is, + # matching the pre-existing AWS Lambda behavior. + assert ( + _serverless.resolve_service_name( + None, + ["FALLBACK"], + "default", + env={"OTEL_SERVICE_NAME": " ", "FALLBACK": "fallback"}, + ) + == " " + ) + + +class TestBuildMetricsTelemetryConfig: + def test_with_service_name(self) -> None: + config = _serverless.build_metrics_telemetry_config( + endpoint="http://collector:4317", + service_name="my-service", + metric_periodicity=timedelta(seconds=30), + ) + assert isinstance(config.metrics, OpenTelemetryConfig) + assert config.metrics.url == "http://collector:4317" + assert config.metrics.metric_periodicity == timedelta(seconds=30) + assert config.global_tags == {"service_name": "my-service"} + + def test_without_service_name(self) -> None: + config = _serverless.build_metrics_telemetry_config( + endpoint="http://collector:4317", + service_name="", + metric_periodicity=None, + ) + assert isinstance(config.metrics, OpenTelemetryConfig) + assert config.metrics.url == "http://collector:4317" + assert config.metrics.metric_periodicity is None + assert config.global_tags == {} + + def test_empty_endpoint_uses_default(self) -> None: + config = _serverless.build_metrics_telemetry_config( + endpoint="", + service_name="svc", + metric_periodicity=None, + ) + assert isinstance(config.metrics, OpenTelemetryConfig) + assert config.metrics.url == _serverless.DEFAULT_OTLP_ENDPOINT + + +class TestBuildOtlpSpanProcessor: + def test_returns_batch_span_processor(self) -> None: + processor = _serverless.build_otlp_span_processor("http://localhost:4317") + try: + assert isinstance(processor, BatchSpanProcessor) + finally: + processor.shutdown() + + def test_raises_import_error_when_exporter_absent( + self, monkeypatch: pytest.MonkeyPatch + ) -> None: + # Setting the module to None in sys.modules makes the lazy import raise + # ImportError, simulating the exporter not being installed. + monkeypatch.setitem( + sys.modules, + "opentelemetry.exporter.otlp.proto.grpc.trace_exporter", + None, + ) + with pytest.raises(ImportError): + _serverless.build_otlp_span_processor("http://localhost:4317") diff --git a/uv.lock b/uv.lock index 22d2244c6..fe7c202ee 100644 --- a/uv.lock +++ b/uv.lock @@ -4711,6 +4711,12 @@ aioboto3 = [ { name = "aioboto3" }, { name = "types-aioboto3", extra = ["s3"] }, ] +cloud-run-worker-otel = [ + { name = "opentelemetry-api" }, + { name = "opentelemetry-exporter-otlp-proto-grpc" }, + { name = "opentelemetry-sdk" }, + { name = "protobuf" }, +] deepagents = [ { name = "deepagents", marker = "python_full_version >= '3.11'" }, { name = "langchain", marker = "python_full_version >= '3.11'" }, @@ -4817,14 +4823,18 @@ requires-dist = [ { name = "mcp", marker = "extra == 'openai-agents'", specifier = ">=1.9.4,<2" }, { name = "nexus-rpc", specifier = "==1.4.0" }, { name = "openai-agents", marker = "extra == 'openai-agents'", specifier = ">=0.19.2,<0.20" }, + { name = "opentelemetry-api", marker = "extra == 'cloud-run-worker-otel'", specifier = ">=1.26,<2" }, { name = "opentelemetry-api", marker = "extra == 'lambda-worker-otel'", specifier = ">=1.26,<2" }, { name = "opentelemetry-api", marker = "extra == 'opentelemetry'", specifier = ">=1.26,<2" }, + { name = "opentelemetry-exporter-otlp-proto-grpc", marker = "extra == 'cloud-run-worker-otel'", specifier = ">=1.11.1,<2" }, { name = "opentelemetry-exporter-otlp-proto-grpc", marker = "extra == 'lambda-worker-otel'", specifier = ">=1.11.1,<2" }, + { name = "opentelemetry-sdk", marker = "extra == 'cloud-run-worker-otel'", specifier = ">=1.26,<2" }, { name = "opentelemetry-sdk", marker = "extra == 'lambda-worker-otel'", specifier = ">=1.26,<2" }, { name = "opentelemetry-sdk", marker = "extra == 'opentelemetry'", specifier = ">=1.26,<2" }, { name = "opentelemetry-sdk-extension-aws", marker = "extra == 'lambda-worker-otel'", specifier = ">=2.0.0,<3" }, { name = "opentelemetry-semantic-conventions", marker = "extra == 'lambda-worker-otel'", specifier = ">=0.40b0,<1" }, { name = "protobuf", specifier = ">=3.20,<8.0.0" }, + { name = "protobuf", marker = "extra == 'cloud-run-worker-otel'", specifier = "<7" }, { name = "pydantic", marker = "extra == 'pydantic'", specifier = ">=2.0.0,<3" }, { name = "python-dateutil", marker = "python_full_version < '3.11'", specifier = ">=2.8.2,<3" }, { name = "strands-agents", marker = "extra == 'strands-agents'", specifier = ">=1.39.0" }, @@ -4832,7 +4842,7 @@ requires-dist = [ { name = "types-protobuf", specifier = ">=3.20,<8.0.0" }, { name = "typing-extensions", specifier = ">=4.2.0,<5" }, ] -provides-extras = ["grpc", "opentelemetry", "pydantic", "openai-agents", "google-adk", "langgraph", "langsmith", "deepagents", "lambda-worker-otel", "aioboto3", "google-genai", "strands-agents"] +provides-extras = ["grpc", "opentelemetry", "pydantic", "openai-agents", "google-adk", "langgraph", "langsmith", "deepagents", "lambda-worker-otel", "cloud-run-worker-otel", "aioboto3", "google-genai", "strands-agents"] [package.metadata.requires-dev] dev = [