Repository navigation
feat(tracing): stamp sgp evals run and row ids on SGP spans #545
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Open
mohammadatallah-scale
wants to merge
8
commits into
next
Choose a base branch
from
mohammad/ove-1253-sgp-evals-span-metadata
base: next
Could not load branches
Branch not found: {{ refName }}
Loading
Could not load tags
Nothing to show
Loading
Are you sure you want to change the base?
Some commits from the old base branch may be removed from the timeline,
and old review comments may become outdated.
+631
−7
Open
Changes from all commits
Commits
Show all changes
8 commits
Select commit
Hold shift + click to select a range
d718f78
feat(tracing): stamp sgp evals run and row ids on SGP spans
mohammadatallah-scale 43211ba
fix(tracing): keep list-shaped eval spans searchable and drop stale e…
mohammadatallah-scale a348297
fix(tracing): stamp sgp evals ids on local activity spans too
mohammadatallah-scale 618f226
fix(tracing): drop a worker's eval attrs when a task's activities car…
mohammadatallah-scale ce0dcfa
fix(tracing): pin eval ids on a span when it starts so a reused task …
mohammadatallah-scale f1d9eb4
fix(tracing): pin eval ids for plain spans and keep them until the up…
mohammadatallah-scale 3e31b85
chore: restore uv.lock
mohammadatallah-scale 848212a
fix(tracing): keep eval span id captures apart from plain span pins s…
mohammadatallah-scale File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
There are no files selected for viewing
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,119 @@ | ||
| """Run/row attribution for spans of SGP evals generation-unit tasks. | ||
|
|
||
| The evals service tags each unit's task with ``task_metadata`` carrying | ||
| ``sgp_evals`` plus the run, row and attempt it belongs to. The ids are copied | ||
| onto the SGP copy of every span as flat ``sgp_evals_*`` keys so a spans search | ||
| on ``extra_metadata`` finds a unit's spans. Flat because sgp-traces treats | ||
| dotted keys differently on ClickHouse and Postgres. | ||
| """ | ||
|
|
||
| from __future__ import annotations | ||
|
|
||
| import threading | ||
| from typing import Any | ||
| from collections import OrderedDict | ||
|
|
||
| from agentex.types.span import Span | ||
|
|
||
| __all__ = ( | ||
| "MEMO_KEY", | ||
| "register_task", | ||
| "register_task_metadata", | ||
| "span_attrs_from_task_metadata", | ||
| "attrs_for_span", | ||
| "capture_for_span", | ||
| "release_span", | ||
| ) | ||
|
|
||
| TASK_METADATA_MARKER = "sgp_evals" | ||
| SPAN_KEY_PREFIX = "sgp_evals_" | ||
| _TASK_METADATA_KEYS = ("generation_run_id", "row_id", "attempt_idx") | ||
| # Temporal workflow memo key the ACP server sets so the worker process can stamp the same attrs. | ||
| MEMO_KEY = "sgp_evals_span_attrs" | ||
|
|
||
| # Only eval tasks are stored, so this stays tiny. The LRU bound caps a long-lived agent process. | ||
| _MAX_TASKS = 10_000 | ||
| _attrs_by_task: OrderedDict[str, dict[str, Any]] = OrderedDict() | ||
| # Attrs captured when a span starts, so a later registry change cannot alter a span still queued for export. | ||
| _attrs_by_span: OrderedDict[str, dict[str, Any]] = OrderedDict() | ||
| # Spans pinned with no attrs live apart so a flood of plain spans cannot evict an eval span's pinned ids. | ||
| _plain_spans: OrderedDict[str, None] = OrderedDict() | ||
| _lock = threading.Lock() | ||
|
|
||
|
|
||
| def span_attrs_from_task_metadata(task_metadata: Any) -> dict[str, Any] | None: | ||
| """The flat span attrs for an evals generation-unit task, or None for any other task.""" | ||
| if not isinstance(task_metadata, dict) or task_metadata.get(TASK_METADATA_MARKER) is None: | ||
| return None | ||
| attrs = { | ||
| f"{SPAN_KEY_PREFIX}{key}": task_metadata[key] for key in _TASK_METADATA_KEYS if task_metadata.get(key) is not None | ||
| } | ||
| return attrs or None | ||
|
|
||
|
|
||
| def register_task(task_id: str, attrs: dict[str, Any]) -> None: | ||
| with _lock: | ||
| _attrs_by_task[task_id] = attrs | ||
| _attrs_by_task.move_to_end(task_id) | ||
| while len(_attrs_by_task) > _MAX_TASKS: | ||
| _attrs_by_task.popitem(last=False) | ||
|
|
||
|
|
||
| def register_task_metadata(task_id: str, task_metadata: Any) -> dict[str, Any] | None: | ||
| """Remember an eval task's span attrs. No-op (and no lookup cost) for every other task.""" | ||
| attrs = span_attrs_from_task_metadata(task_metadata) | ||
| if attrs is not None: | ||
| register_task(task_id, attrs) | ||
| else: | ||
| unregister_task(task_id) | ||
| return attrs | ||
|
|
||
|
|
||
| def unregister_task(task_id: str) -> None: | ||
| if _attrs_by_task: | ||
| with _lock: | ||
| _attrs_by_task.pop(task_id, None) | ||
|
|
||
|
|
||
| def _lookup(span: Span) -> dict[str, Any]: | ||
| for key in (span.task_id, span.trace_id): | ||
| if key and key in _attrs_by_task: | ||
| return dict(_attrs_by_task[key]) | ||
| return {} | ||
|
|
||
|
|
||
| def capture_for_span(span: Span) -> None: | ||
| """Pin the span's attrs at start, empty included, so an eval task registered later cannot claim it.""" | ||
| with _lock: | ||
| attrs = _lookup(span) | ||
| store: OrderedDict[str, Any] = _attrs_by_span if attrs else _plain_spans | ||
| store[span.id] = attrs or None | ||
| while len(store) > _MAX_TASKS: | ||
| store.popitem(last=False) | ||
|
|
||
|
|
||
| def release_span(span_id: str) -> None: | ||
| if _attrs_by_span or _plain_spans: | ||
| with _lock: | ||
| _attrs_by_span.pop(span_id, None) | ||
| _plain_spans.pop(span_id, None) | ||
|
|
||
|
|
||
| def attrs_for_span(span: Span) -> dict[str, Any]: | ||
| """Attrs captured at span start, else those of the task found by ``span.task_id`` then ``span.trace_id``.""" | ||
| if not _attrs_by_task and not _attrs_by_span and not _plain_spans: | ||
| return {} | ||
| with _lock: | ||
| if span.id in _attrs_by_span: | ||
| return dict(_attrs_by_span[span.id]) | ||
| if span.id in _plain_spans: | ||
| return {} | ||
| return _lookup(span) | ||
|
|
||
|
|
||
| def clear() -> None: | ||
| """Reset the registry (test isolation).""" | ||
| with _lock: | ||
| _attrs_by_task.clear() | ||
| _attrs_by_span.clear() | ||
| _plain_spans.clear() |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,83 @@ | ||
| """Carry an eval task's span attrs from its Temporal workflow memo into the worker's activities. | ||
|
|
||
| Spans are emitted by activities in the worker process, which never sees the ACP | ||
| server's task. The ACP server puts the attrs in the workflow memo, the outbound | ||
| interceptor copies them onto each activity's headers, and the activity | ||
| interceptor registers them for the task so the SGP processor can stamp spans. | ||
| Non-eval workflows have no memo entry and get no header. | ||
| """ | ||
|
|
||
| from __future__ import annotations | ||
|
|
||
| from typing import Any, override | ||
|
|
||
| from temporalio import activity, workflow | ||
| from temporalio.worker import ( | ||
| Interceptor, | ||
| StartActivityInput, | ||
| ExecuteActivityInput, | ||
| StartLocalActivityInput, | ||
| ActivityInboundInterceptor, | ||
| WorkflowInboundInterceptor, | ||
| WorkflowOutboundInterceptor, | ||
| ) | ||
| from temporalio.converter import default | ||
|
|
||
| from agentex.lib.utils.logging import make_logger | ||
| from agentex.lib.core.tracing.sgp_evals import MEMO_KEY, register_task, unregister_task | ||
|
|
||
| logger = make_logger(__name__) | ||
|
|
||
| ATTRS_HEADER = "sgp-evals-span-attrs" | ||
| _converter = default().payload_converter | ||
|
|
||
|
|
||
| class SGPEvalsInterceptor(Interceptor): | ||
| @override | ||
| def intercept_activity(self, next: ActivityInboundInterceptor) -> ActivityInboundInterceptor: | ||
| return _ActivityInbound(next) | ||
|
|
||
| @override | ||
| def workflow_interceptor_class(self, input: Any) -> type[WorkflowInboundInterceptor] | None: | ||
| return _WorkflowInbound | ||
|
|
||
|
|
||
| class _WorkflowInbound(WorkflowInboundInterceptor): | ||
| @override | ||
| def init(self, outbound: WorkflowOutboundInterceptor) -> None: | ||
| super().init(_WorkflowOutbound(outbound)) | ||
|
|
||
|
|
||
| def _add_attrs_header(input: StartActivityInput | StartLocalActivityInput) -> None: | ||
| attrs = workflow.memo_value(MEMO_KEY, default=None) | ||
| if isinstance(attrs, dict) and attrs: | ||
| input.headers = {**input.headers, ATTRS_HEADER: _converter.to_payload(attrs)} | ||
|
|
||
|
|
||
| class _WorkflowOutbound(WorkflowOutboundInterceptor): | ||
| @override | ||
| def start_activity(self, input: StartActivityInput) -> workflow.ActivityHandle[Any]: | ||
| _add_attrs_header(input) | ||
| return super().start_activity(input) | ||
|
|
||
| @override | ||
| def start_local_activity(self, input: StartLocalActivityInput) -> workflow.ActivityHandle[Any]: | ||
| _add_attrs_header(input) | ||
| return super().start_local_activity(input) | ||
|
|
||
|
|
||
| class _ActivityInbound(ActivityInboundInterceptor): | ||
| @override | ||
| async def execute_activity(self, input: ExecuteActivityInput) -> Any: | ||
| payload = input.headers.get(ATTRS_HEADER) | ||
| # The workflow id is the task id (see TemporalTaskService.submit_task). | ||
| task_id = activity.info().workflow_id | ||
| if task_id and payload is None: | ||
| # A reused workflow id must not inherit a previous eval run's ids. | ||
| unregister_task(task_id) | ||
|
mohammadatallah-scale marked this conversation as resolved.
|
||
| elif task_id and payload is not None: | ||
| try: | ||
| register_task(task_id, _converter.from_payload(payload, dict)) | ||
| except Exception: | ||
| logger.warning("failed to read sgp evals span attrs from activity headers", exc_info=True) | ||
| return await super().execute_activity(input) | ||
|
mohammadatallah-scale marked this conversation as resolved.
|
||
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Oops, something went wrong.
Oops, something went wrong.
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
Uh oh!
There was an error while loading. Please reload this page.