Skip to content
Open
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
9 changes: 8 additions & 1 deletion src/agentex/lib/core/temporal/workers/worker.py
Original file line number Diff line number Diff line change
Expand Up @@ -98,6 +98,7 @@ async def get_temporal_client(
plugins: list = [],
payload_codec: PayloadCodec | None = None,
data_converter: DataConverter | None = None,
metrics_headers: dict[str, str] | None = None,
) -> Client:
if plugins != []: # We don't need to validate the plugins if they are empty
_validate_plugins(plugins)
Expand Down Expand Up @@ -143,7 +144,10 @@ async def get_temporal_client(
if not metrics_url:
client = await Client.connect(**connect_kwargs)
else:
runtime = Runtime(telemetry=TelemetryConfig(metrics=OpenTelemetryConfig(url=metrics_url)))
runtime = Runtime(telemetry=TelemetryConfig(metrics=OpenTelemetryConfig(
url=metrics_url,
headers=metrics_headers or {},
)))
connect_kwargs["runtime"] = runtime
client = await Client.connect(**connect_kwargs)
return client
Expand All @@ -161,6 +165,7 @@ def __init__(
metrics_url: str | None = None,
payload_codec: PayloadCodec | None = None,
data_converter: DataConverter | None = None,
metrics_headers: dict[str, str] | None = None,
):
self.task_queue = task_queue
self.activity_handles = []
Expand All @@ -174,6 +179,7 @@ def __init__(
self.plugins = plugins
self.interceptors = interceptors
self.metrics_url = metrics_url
self.metrics_headers = metrics_headers
self.payload_codec = payload_codec
self.data_converter = data_converter

Expand Down Expand Up @@ -211,6 +217,7 @@ async def run(
temporal_address=os.environ.get("TEMPORAL_ADDRESS", "localhost:7233"),
plugins=self.plugins,
metrics_url=self.metrics_url,
metrics_headers=self.metrics_headers,
payload_codec=self.payload_codec,
data_converter=self.data_converter,
)
Expand Down
Loading