Skip to content

Commit 3dca81f

Browse files
authored
feat(worker): add OTLP metrics exporter options (headers, HTTP transport, delta temporality) (#501)
1 parent 941bbc4 commit 3dca81f

2 files changed

Lines changed: 129 additions & 3 deletions

File tree

src/agentex/lib/core/temporal/workers/worker.py

Lines changed: 25 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -16,7 +16,7 @@
1616
Interceptor,
1717
UnsandboxedWorkflowRunner,
1818
)
19-
from temporalio.runtime import Runtime, TelemetryConfig, OpenTelemetryConfig
19+
from temporalio.runtime import Runtime, TelemetryConfig, OpenTelemetryConfig, OpenTelemetryMetricTemporality
2020
from temporalio.converter import (
2121
PayloadCodec,
2222
DataConverter,
@@ -98,6 +98,10 @@ async def get_temporal_client(
9898
plugins: list = [],
9999
payload_codec: PayloadCodec | None = None,
100100
data_converter: DataConverter | None = None,
101+
*,
102+
metrics_headers: dict[str, str] | None = None,
103+
metrics_use_http: bool = False,
104+
metrics_temporality_delta: bool = False,
101105
) -> Client:
102106
if plugins != []: # We don't need to validate the plugins if they are empty
103107
_validate_plugins(plugins)
@@ -143,7 +147,16 @@ async def get_temporal_client(
143147
if not metrics_url:
144148
client = await Client.connect(**connect_kwargs)
145149
else:
146-
runtime = Runtime(telemetry=TelemetryConfig(metrics=OpenTelemetryConfig(url=metrics_url)))
150+
runtime = Runtime(telemetry=TelemetryConfig(metrics=OpenTelemetryConfig(
151+
url=metrics_url,
152+
headers=metrics_headers or {},
153+
http=metrics_use_http,
154+
metric_temporality=(
155+
OpenTelemetryMetricTemporality.DELTA
156+
if metrics_temporality_delta
157+
else OpenTelemetryMetricTemporality.CUMULATIVE
158+
),
159+
)))
147160
connect_kwargs["runtime"] = runtime
148161
client = await Client.connect(**connect_kwargs)
149162
return client
@@ -161,6 +174,10 @@ def __init__(
161174
metrics_url: str | None = None,
162175
payload_codec: PayloadCodec | None = None,
163176
data_converter: DataConverter | None = None,
177+
*,
178+
metrics_headers: dict[str, str] | None = None,
179+
metrics_use_http: bool = False,
180+
metrics_temporality_delta: bool = False,
164181
):
165182
self.task_queue = task_queue
166183
self.activity_handles = []
@@ -174,6 +191,9 @@ def __init__(
174191
self.plugins = plugins
175192
self.interceptors = interceptors
176193
self.metrics_url = metrics_url
194+
self.metrics_headers = metrics_headers
195+
self.metrics_use_http = metrics_use_http
196+
self.metrics_temporality_delta = metrics_temporality_delta
177197
self.payload_codec = payload_codec
178198
self.data_converter = data_converter
179199

@@ -211,6 +231,9 @@ async def run(
211231
temporal_address=os.environ.get("TEMPORAL_ADDRESS", "localhost:7233"),
212232
plugins=self.plugins,
213233
metrics_url=self.metrics_url,
234+
metrics_headers=self.metrics_headers,
235+
metrics_use_http=self.metrics_use_http,
236+
metrics_temporality_delta=self.metrics_temporality_delta,
214237
payload_codec=self.payload_codec,
215238
data_converter=self.data_converter,
216239
)

tests/lib/test_agentex_worker.py

Lines changed: 104 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,5 @@
11
import os
2-
from unittest.mock import patch
2+
from unittest.mock import AsyncMock, MagicMock, patch
33

44
import pytest
55

@@ -88,3 +88,106 @@ def test_worker_init_basic_attributes(self):
8888
assert worker.health_check_server_running is False
8989
assert worker.healthy is False
9090
assert worker.plugins == []
91+
92+
def test_worker_stores_metrics_params(self):
93+
from agentex.lib.core.temporal.workers.worker import AgentexWorker
94+
95+
worker = AgentexWorker(
96+
task_queue="test-queue",
97+
health_check_port=8080,
98+
metrics_url="http://example.com/v1/metrics",
99+
metrics_headers={"Authorization": "Api-Token tok"},
100+
metrics_use_http=True,
101+
metrics_temporality_delta=True,
102+
)
103+
104+
assert worker.metrics_url == "http://example.com/v1/metrics"
105+
assert worker.metrics_headers == {"Authorization": "Api-Token tok"}
106+
assert worker.metrics_use_http is True
107+
assert worker.metrics_temporality_delta is True
108+
109+
def test_worker_metrics_params_default_to_none_and_false(self):
110+
from agentex.lib.core.temporal.workers.worker import AgentexWorker
111+
112+
worker = AgentexWorker(task_queue="test-queue", health_check_port=8080)
113+
114+
assert worker.metrics_url is None
115+
assert worker.metrics_headers is None
116+
assert worker.metrics_use_http is False
117+
assert worker.metrics_temporality_delta is False
118+
119+
120+
class TestGetTemporalClientMetricsConfig:
121+
"""Tests that metrics params reach OpenTelemetryConfig correctly."""
122+
123+
async def test_metrics_params_reach_otel_config(self):
124+
from temporalio.client import Client
125+
from temporalio.runtime import OpenTelemetryMetricTemporality
126+
127+
from agentex.lib.core.temporal.workers.worker import get_temporal_client
128+
129+
with patch.object(Client, "connect", new=AsyncMock(return_value=MagicMock())), \
130+
patch("agentex.lib.core.temporal.workers.worker.Runtime"), \
131+
patch("agentex.lib.core.temporal.workers.worker.TelemetryConfig"), \
132+
patch("agentex.lib.core.temporal.workers.worker.OpenTelemetryConfig") as mock_otel:
133+
await get_temporal_client(
134+
"localhost:7233",
135+
metrics_url="http://example.com/v1/metrics",
136+
metrics_headers={"Authorization": "Api-Token tok"},
137+
metrics_use_http=True,
138+
metrics_temporality_delta=True,
139+
)
140+
141+
mock_otel.assert_called_once_with(
142+
url="http://example.com/v1/metrics",
143+
headers={"Authorization": "Api-Token tok"},
144+
http=True,
145+
metric_temporality=OpenTelemetryMetricTemporality.DELTA,
146+
)
147+
148+
async def test_delta_false_maps_to_cumulative(self):
149+
from temporalio.client import Client
150+
from temporalio.runtime import OpenTelemetryMetricTemporality
151+
152+
from agentex.lib.core.temporal.workers.worker import get_temporal_client
153+
154+
with patch.object(Client, "connect", new=AsyncMock(return_value=MagicMock())), \
155+
patch("agentex.lib.core.temporal.workers.worker.Runtime"), \
156+
patch("agentex.lib.core.temporal.workers.worker.TelemetryConfig"), \
157+
patch("agentex.lib.core.temporal.workers.worker.OpenTelemetryConfig") as mock_otel:
158+
await get_temporal_client(
159+
"localhost:7233",
160+
metrics_url="http://example.com/v1/metrics",
161+
metrics_temporality_delta=False,
162+
)
163+
164+
_, kwargs = mock_otel.call_args
165+
assert kwargs["metric_temporality"] == OpenTelemetryMetricTemporality.CUMULATIVE
166+
167+
async def test_none_headers_defaults_to_empty_dict(self):
168+
from temporalio.client import Client
169+
170+
from agentex.lib.core.temporal.workers.worker import get_temporal_client
171+
172+
with patch.object(Client, "connect", new=AsyncMock(return_value=MagicMock())), \
173+
patch("agentex.lib.core.temporal.workers.worker.Runtime"), \
174+
patch("agentex.lib.core.temporal.workers.worker.TelemetryConfig"), \
175+
patch("agentex.lib.core.temporal.workers.worker.OpenTelemetryConfig") as mock_otel:
176+
await get_temporal_client(
177+
"localhost:7233",
178+
metrics_url="http://example.com/v1/metrics",
179+
)
180+
181+
_, kwargs = mock_otel.call_args
182+
assert kwargs["headers"] == {}
183+
184+
async def test_no_metrics_url_skips_runtime(self):
185+
from temporalio.client import Client
186+
187+
from agentex.lib.core.temporal.workers.worker import get_temporal_client
188+
189+
with patch.object(Client, "connect", new=AsyncMock(return_value=MagicMock())), \
190+
patch("agentex.lib.core.temporal.workers.worker.Runtime") as mock_runtime:
191+
await get_temporal_client("localhost:7233")
192+
193+
mock_runtime.assert_not_called()

0 commit comments

Comments
 (0)