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
89 changes: 89 additions & 0 deletions src/agentex/lib/core/observability/event_metrics.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,89 @@
"""OTel metrics for event delivery (task/event_send -> live workflow).

Records whether an EVENT_SEND signal actually reached a *running* Temporal
workflow, vs. only being accepted by the ACP server. ``event/send`` is async:
the ACP server returns a 200 ack before the signal is even attempted (see
``base_acp_server._handle_jsonrpc`` background branch), so HTTP status is not a
delivery signal — this counter is. Under load, tasks that have already
completed or idle-timed-out have no workflow to receive the event, which
Temporal surfaces as an ``RPCError`` with status ``NOT_FOUND``.

The meter is no-op when the application hasn't configured a ``MeterProvider``,
so importing this module is safe for runtimes that don't use OTel. Instruments
are created lazily on first ``get_event_metrics()`` call so a ``MeterProvider``
configured *after* this module is imported still binds correctly.

Recording is gated on ``AGENTEX_EVENT_METRICS`` (default on) and every record
is best-effort — it never raises into the business path.

Cardinality is bounded: the only attribute is ``outcome``, drawn from a small
fixed set (see the ``OUTCOME_*`` constants). Resource attributes
(``service.name``, ``k8s.*``, etc.) come from the application's OTel resource
configuration and are added to every series automatically.
"""

from __future__ import annotations

import os
from typing import Optional

from opentelemetry import metrics

# Outcome label values (bounded cardinality — do NOT add task IDs etc.).
OUTCOME_DELIVERED = "delivered"
OUTCOME_NO_LIVE_WORKFLOW = "no_live_workflow"
OUTCOME_ERROR = "error"


class EventMetrics:
"""Lazily-created OTel instruments for event-delivery telemetry."""

def __init__(self) -> None:
meter = metrics.get_meter("agentex.events")
self.delivery = meter.create_counter(
name="agentex.events.delivery",
unit="1",
description=(
"task/event_send deliveries tagged with outcome (delivered / "
"no_live_workflow / error). 'no_live_workflow' means the signal "
"had no running workflow to receive it (task already completed "
"or idle-timed-out). Use delivered / (delivered + "
"no_live_workflow) as the true event-delivery rate."
),
)


_event_metrics: Optional[EventMetrics] = None


def get_event_metrics() -> EventMetrics:
"""Return the event metrics singleton, creating it on first use."""
global _event_metrics
if _event_metrics is None:
_event_metrics = EventMetrics()
return _event_metrics


_metrics_enabled: Optional[bool] = None


def is_event_metrics_enabled() -> bool:
"""Whether event-delivery metric recording is enabled (AGENTEX_EVENT_METRICS)."""
global _metrics_enabled
if _metrics_enabled is None:
raw = os.environ.get("AGENTEX_EVENT_METRICS", "1").strip().lower()
_metrics_enabled = raw not in ("0", "false", "no", "off")
return _metrics_enabled


def record_event_delivery(outcome: str) -> None:
"""Best-effort bump of the event-delivery counter for the given outcome.

``outcome`` should be one of the ``OUTCOME_*`` constants.
"""
if not is_event_metrics_enabled():
return
try:
get_event_metrics().delivery.add(1, {"outcome": outcome})
except Exception:
pass
48 changes: 38 additions & 10 deletions src/agentex/lib/core/temporal/services/temporal_task_service.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,8 @@
from typing import Any
from datetime import timedelta

from temporalio.service import RPCError, RPCStatusCode

from agentex.types.task import Task
from agentex.types.agent import Agent
from agentex.types.event import Event
Expand All @@ -11,6 +13,12 @@
from agentex.lib.core.clients.temporal.types import WorkflowState
from agentex.lib.core.temporal.types.workflow import SignalName
from agentex.lib.core.clients.temporal.temporal_client import TemporalClient
from agentex.lib.core.observability.event_metrics import (
OUTCOME_DELIVERED,
OUTCOME_ERROR,
OUTCOME_NO_LIVE_WORKFLOW,
record_event_delivery,
)


class TemporalTaskService:
Expand Down Expand Up @@ -63,16 +71,36 @@ async def get_state(self, task_id: str) -> WorkflowState:
)

async def send_event(self, agent: Agent, task: Task, event: Event, request: dict | None = None) -> None:
return await self._temporal_client.send_signal(
workflow_id=task.id,
signal=SignalName.RECEIVE_EVENT.value,
payload=SendEventParams(
agent=agent,
task=task,
event=event,
request=request,
).model_dump(),
)
# event/send is accepted+acked by the ACP server before the signal is
# attempted, so the only place we learn whether the event actually
# reached a *running* workflow is here. Record the true delivery outcome
# (see agentex.events.delivery). Behaviour is unchanged — we still raise.
try:
result = await self._temporal_client.send_signal(
workflow_id=task.id,
signal=SignalName.RECEIVE_EVENT.value,
payload=SendEventParams(
agent=agent,
task=task,
event=event,
request=request,
).model_dump(),
)
except RPCError as e:
# NOT_FOUND == no running workflow to receive the event (task already
# completed or idle-timed-out). Distinct from unexpected errors so we
# can measure the "arrived too late" rate separately.
record_event_delivery(
OUTCOME_NO_LIVE_WORKFLOW
if e.status == RPCStatusCode.NOT_FOUND
else OUTCOME_ERROR
)
raise
except Exception:
record_event_delivery(OUTCOME_ERROR)
raise
record_event_delivery(OUTCOME_DELIVERED)
return result

async def interrupt(self, agent: Agent, task: Task, request: dict | None = None) -> None:
"""Forward a task/interrupt to the running workflow as a dedicated signal.
Expand Down
3 changes: 3 additions & 0 deletions src/agentex/lib/core/temporal/workers/worker.py
Original file line number Diff line number Diff line change
Expand Up @@ -151,6 +151,7 @@ def __init__(
task_queue,
max_workers: int = 10,
max_concurrent_activities: int = 10,
max_cached_workflows: int = 1000,
health_check_port: int | None = None,
plugins: list = [],
interceptors: list = [],
Expand All @@ -162,6 +163,7 @@ def __init__(
self.activity_handles = []
self.max_workers = max_workers
self.max_concurrent_activities = max_concurrent_activities
self.max_cached_workflows = max_cached_workflows
self.health_check_server_running = False
self.healthy = False
self.health_check_port = (
Expand Down Expand Up @@ -227,6 +229,7 @@ async def run(
activities=activities,
workflow_runner=UnsandboxedWorkflowRunner(),
max_concurrent_activities=self.max_concurrent_activities,
max_cached_workflows=self.max_cached_workflows,
build_id=str(uuid.uuid4()),
debug_mode=debug_enabled, # Disable deadlock detection in debug mode
interceptors=self.interceptors, # Pass interceptors to Worker
Expand Down
Loading