diff --git a/instrumentation-loongsuite/loongsuite-instrumentation-hermes-agent/CHANGELOG.md b/instrumentation-loongsuite/loongsuite-instrumentation-hermes-agent/CHANGELOG.md index 74194d067..2c87fcfe0 100644 --- a/instrumentation-loongsuite/loongsuite-instrumentation-hermes-agent/CHANGELOG.md +++ b/instrumentation-loongsuite/loongsuite-instrumentation-hermes-agent/CHANGELOG.md @@ -7,6 +7,14 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## Unreleased +### Fixed + +- Prefer provider-supplied OpenAI-compatible response IDs (including + DashScope-style `request_id` values) for `gen_ai.response.id` on Hermes LLM + spans, while retaining Hermes's response ID as a fallback. +- Align the Hermes plugin with the OpenTelemetry 1.39.1/0.60b1 release set used + by LoongSuite releases so its standard PyPI dependencies resolve consistently. + ## Version 0.7.0 (2026-07-03) ### Fixed diff --git a/instrumentation-loongsuite/loongsuite-instrumentation-hermes-agent/pyproject.toml b/instrumentation-loongsuite/loongsuite-instrumentation-hermes-agent/pyproject.toml index 863af507d..8306528fb 100644 --- a/instrumentation-loongsuite/loongsuite-instrumentation-hermes-agent/pyproject.toml +++ b/instrumentation-loongsuite/loongsuite-instrumentation-hermes-agent/pyproject.toml @@ -24,11 +24,11 @@ classifiers = [ "Programming Language :: Python :: 3.13", ] dependencies = [ - "opentelemetry-api >= 1.37.0, < 1.40", - "opentelemetry-instrumentation >= 0.60b1, < 0.61", - "opentelemetry-instrumentation-threading >= 0.60b1, < 0.61", - "opentelemetry-sdk >= 1.37.0, < 1.40", - "opentelemetry-semantic-conventions >= 0.60b1, < 0.61", + "opentelemetry-api ~= 1.39.1", + "opentelemetry-instrumentation ~= 0.60b1", + "opentelemetry-instrumentation-threading ~= 0.60b1", + "opentelemetry-sdk ~= 1.39.1", + "opentelemetry-semantic-conventions ~= 0.60b1", "opentelemetry-util-genai", "wrapt >= 1.0.0, < 2.0.0", ] diff --git a/instrumentation-loongsuite/loongsuite-instrumentation-hermes-agent/src/opentelemetry/instrumentation/hermes_agent/helpers.py b/instrumentation-loongsuite/loongsuite-instrumentation-hermes-agent/src/opentelemetry/instrumentation/hermes_agent/helpers.py index 3744478fd..39af8ae4b 100644 --- a/instrumentation-loongsuite/loongsuite-instrumentation-hermes-agent/src/opentelemetry/instrumentation/hermes_agent/helpers.py +++ b/instrumentation-loongsuite/loongsuite-instrumentation-hermes-agent/src/opentelemetry/instrumentation/hermes_agent/helpers.py @@ -35,6 +35,7 @@ InvokeAgentInvocation, ReactStepInvocation, ) +from opentelemetry.util.genai.response_id import resolve_response_id from opentelemetry.util.genai.types import ( FunctionToolDefinition, GenericToolDefinition, @@ -674,6 +675,8 @@ def update_llm_invocation_from_response( invocation: LLMInvocation, instance: Any, response: Any, + *, + provider_response_id: str | None = None, ) -> tuple[int, int, int]: response_model = getattr(response, "model", None) if response_model: @@ -681,7 +684,11 @@ def update_llm_invocation_from_response( else: invocation.response_model_name = invocation.request_model - response_id = getattr(response, "id", None) + response_id = resolve_response_id( + provider_response_id, + response, + framework_fields=("request_id", "id", "response_id"), + ) if response_id: invocation.response_id = response_id diff --git a/instrumentation-loongsuite/loongsuite-instrumentation-hermes-agent/src/opentelemetry/instrumentation/hermes_agent/instrumentor.py b/instrumentation-loongsuite/loongsuite-instrumentation-hermes-agent/src/opentelemetry/instrumentation/hermes_agent/instrumentor.py index cc22683af..17ed96c04 100644 --- a/instrumentation-loongsuite/loongsuite-instrumentation-hermes-agent/src/opentelemetry/instrumentation/hermes_agent/instrumentor.py +++ b/instrumentation-loongsuite/loongsuite-instrumentation-hermes-agent/src/opentelemetry/instrumentation/hermes_agent/instrumentor.py @@ -29,6 +29,7 @@ from .metrics import HermesMetrics from .wrappers import ( LLMCallWrapper, + ProviderClientWrapper, RunConversationWrapper, ToolBatchWrapper, ToolCallWrapper, @@ -127,6 +128,11 @@ def _instrument(self, **kwargs: Any) -> None: "AIAgent._interruptible_streaming_api_call", LLMCallWrapper(handler, metrics, streaming=True), ) + _safe_wrap_function_wrapper( + "run_agent", + "AIAgent._create_request_openai_client", + ProviderClientWrapper(), + ) _safe_wrap_function_wrapper( "run_agent", "AIAgent._invoke_tool", @@ -180,6 +186,7 @@ def _uninstrument(self, **kwargs: Any) -> None: _safe_unwrap(ai_agent, "run_conversation") _safe_unwrap(ai_agent, "_interruptible_api_call") _safe_unwrap(ai_agent, "_interruptible_streaming_api_call") + _safe_unwrap(ai_agent, "_create_request_openai_client") _safe_unwrap(ai_agent, "_invoke_tool") _safe_unwrap(ai_agent, "_execute_tool_calls") _safe_unwrap(model_tools, "handle_function_call") diff --git a/instrumentation-loongsuite/loongsuite-instrumentation-hermes-agent/src/opentelemetry/instrumentation/hermes_agent/wrappers.py b/instrumentation-loongsuite/loongsuite-instrumentation-hermes-agent/src/opentelemetry/instrumentation/hermes_agent/wrappers.py index 73427e2c3..af7b41063 100644 --- a/instrumentation-loongsuite/loongsuite-instrumentation-hermes-agent/src/opentelemetry/instrumentation/hermes_agent/wrappers.py +++ b/instrumentation-loongsuite/loongsuite-instrumentation-hermes-agent/src/opentelemetry/instrumentation/hermes_agent/wrappers.py @@ -17,6 +17,8 @@ from __future__ import annotations import contextvars +import functools +import threading import timeit from collections.abc import Mapping from contextlib import suppress @@ -24,6 +26,7 @@ from opentelemetry import trace as trace_api from opentelemetry.util.genai.extended_handler import ExtendedTelemetryHandler +from opentelemetry.util.genai.response_id import extract_response_id from opentelemetry.util.genai.types import Error from .helpers import ( @@ -59,6 +62,148 @@ "embedding ", "embeddings ", ) +_PROVIDER_CREATE_WRAPPED = "_otel_hermes_response_id_capture_wrapped" +_PROVIDER_ATTEMPT = threading.local() +_PROVIDER_CAPTURE = threading.local() +_MISSING_THREAD_LOCAL = object() + + +class _ProviderResponseAttempt: + """Response-ID state for one provider SDK ``create`` attempt.""" + + def __init__(self) -> None: + self._lock = threading.Lock() + self._response_id: str | None = None + self._response_id_priority = 0 + + @property + def response_id(self) -> str | None: + with self._lock: + return self._response_id + + def record(self, value: Any) -> None: + # DashScope exposes its provider correlation ID as ``request_id`` even + # when an OpenAI-compatible completion ``id`` is also present. + response_id = extract_response_id(value, fields=("request_id",)) + priority = 2 + if response_id is None: + response_id = extract_response_id( + value, + fields=("id", "response_id"), + ) + priority = 1 + if response_id is None: + return + with self._lock: + if priority > self._response_id_priority: + self._response_id = response_id + self._response_id_priority = priority + + +class _ProviderResponseIdCapture: + """Invocation-local pointer to the active provider attempt.""" + + def __init__(self) -> None: + self._lock = threading.Lock() + self._attempt: _ProviderResponseAttempt | None = None + + def adopt_current_attempt(self) -> None: + attempt = getattr(_PROVIDER_ATTEMPT, "value", None) + if attempt is None: + return + _PROVIDER_CAPTURE.value = self + self.adopt(attempt) + + def adopt(self, attempt: _ProviderResponseAttempt) -> None: + with self._lock: + self._attempt = attempt + + @property + def response_id(self) -> str | None: + with self._lock: + attempt = self._attempt + return attempt.response_id if attempt is not None else None + + +class _ProviderStreamProxy: + """Transparent iterator that observes provider chunks without changing them.""" + + def __init__(self, stream: Any, attempt: _ProviderResponseAttempt) -> None: + self._stream = stream + self._attempt = attempt + + def __iter__(self): + return self + + def __next__(self): + chunk = next(self._stream) + self._attempt.record(chunk) + return chunk + + def __enter__(self): + enter = getattr(self._stream, "__enter__", None) + if enter is not None: + entered = enter() + if entered is not self._stream: + self._stream = entered + return self + + def __exit__(self, exc_type, exc_value, traceback): + exit_method = getattr(self._stream, "__exit__", None) + if exit_method is None: + return False + return exit_method(exc_type, exc_value, traceback) + + def __getattr__(self, name: str) -> Any: + return getattr(self._stream, name) + + +def _wrap_provider_create(resource: Any) -> None: + create = ( + getattr(resource, "create", None) if resource is not None else None + ) + if not callable(create) or getattr( + create, _PROVIDER_CREATE_WRAPPED, False + ): + return + + @functools.wraps(create) + def _capturing_create(*args, **kwargs): + attempt = _ProviderResponseAttempt() + _PROVIDER_ATTEMPT.value = attempt + capture = getattr(_PROVIDER_CAPTURE, "value", None) + if capture is not None: + capture.adopt(attempt) + response = create(*args, **kwargs) + attempt.record(response) + + if ( + kwargs.get("stream") is True + and not hasattr(response, "choices") + and hasattr(response, "__iter__") + ): + return _ProviderStreamProxy(response, attempt) + return response + + setattr(_capturing_create, _PROVIDER_CREATE_WRAPPED, True) + try: + setattr(resource, "create", _capturing_create) + except (AttributeError, TypeError): + # A provider resource may use slots/read-only descriptors. Telemetry + # must never make the model call fail; the framework response remains + # available as the fallback ID in that case. + return + + +class ProviderClientWrapper: + """Decorate Hermes request-local SDK clients for response-ID capture.""" + + def __call__(self, wrapped, instance, args, kwargs): + client = wrapped(*args, **kwargs) + _wrap_provider_create( + getattr(getattr(client, "chat", None), "completions", None) + ) + return client def _resolve_handler( @@ -307,16 +452,30 @@ def __call__(self, wrapped, instance, args, kwargs): api_kwargs = args[0] if args else kwargs.get("api_kwargs", {}) invocation = create_llm_invocation(instance, api_kwargs) + provider_response_capture = _ProviderResponseIdCapture() provider = invocation.provider or provider_name(instance) model = invocation.request_model or "" self._handler.start_llm(invocation) started_at = timeit.default_timer() + previous_provider_attempt = getattr( + _PROVIDER_ATTEMPT, + "value", + _MISSING_THREAD_LOCAL, + ) + previous_provider_capture = getattr( + _PROVIDER_CAPTURE, + "value", + _MISSING_THREAD_LOCAL, + ) + _PROVIDER_ATTEMPT.value = None + _PROVIDER_CAPTURE.value = provider_response_capture try: if self._streaming: original_first_delta = kwargs.get("on_first_delta") def _wrapped_first_delta(): + provider_response_capture.adopt_current_attempt() now = timeit.default_timer() invocation.monotonic_first_token_s = now if current_state["first_token_monotonic_s"] is None: @@ -327,9 +486,15 @@ def _wrapped_first_delta(): kwargs["on_first_delta"] = _wrapped_first_delta response = wrapped(*args, **kwargs) + # Covers direct/synchronous streaming implementations that iterate + # on the caller thread instead of Hermes's worker thread. + provider_response_capture.adopt_current_attempt() input_tokens, output_tokens, total_tokens = ( update_llm_invocation_from_response( - invocation, instance, response + invocation, + instance, + response, + provider_response_id=provider_response_capture.response_id, ) ) @@ -379,6 +544,16 @@ def _wrapped_first_delta(): finish_step(instance, "error", exc=exc) raise finally: + if previous_provider_attempt is _MISSING_THREAD_LOCAL: + with suppress(AttributeError): + del _PROVIDER_ATTEMPT.value + else: + _PROVIDER_ATTEMPT.value = previous_provider_attempt + if previous_provider_capture is _MISSING_THREAD_LOCAL: + with suppress(AttributeError): + del _PROVIDER_CAPTURE.value + else: + _PROVIDER_CAPTURE.value = previous_provider_capture current_state["active_llm_depth"] = max( 0, current_state["active_llm_depth"] - 1 ) diff --git a/instrumentation-loongsuite/loongsuite-instrumentation-hermes-agent/tests/requirements.txt b/instrumentation-loongsuite/loongsuite-instrumentation-hermes-agent/tests/requirements.txt index f2c9499dc..e5d978198 100644 --- a/instrumentation-loongsuite/loongsuite-instrumentation-hermes-agent/tests/requirements.txt +++ b/instrumentation-loongsuite/loongsuite-instrumentation-hermes-agent/tests/requirements.txt @@ -19,11 +19,11 @@ pyyaml>=6.0 openai>=1.0.0 wrapt>=1.0.0,<2.0.0 -opentelemetry-api>=1.37.0,<1.40 -opentelemetry-sdk>=1.37.0,<1.40 -opentelemetry-instrumentation>=0.60b1,<0.61 -opentelemetry-instrumentation-threading>=0.60b1,<0.61 -opentelemetry-semantic-conventions>=0.60b1,<0.61 +opentelemetry-api==1.39.1 +opentelemetry-sdk==1.39.1 +opentelemetry-instrumentation==0.60b1 +opentelemetry-instrumentation-threading==0.60b1 +opentelemetry-semantic-conventions==0.60b1 -e util/opentelemetry-util-genai -e instrumentation-loongsuite/loongsuite-instrumentation-hermes-agent diff --git a/instrumentation-loongsuite/loongsuite-instrumentation-hermes-agent/tests/test_live_llm_calls.py b/instrumentation-loongsuite/loongsuite-instrumentation-hermes-agent/tests/test_live_llm_calls.py index b49c13be4..3577329ee 100644 --- a/instrumentation-loongsuite/loongsuite-instrumentation-hermes-agent/tests/test_live_llm_calls.py +++ b/instrumentation-loongsuite/loongsuite-instrumentation-hermes-agent/tests/test_live_llm_calls.py @@ -102,7 +102,10 @@ def test_sync_llm_call_records_single_llm_span_and_metric( assert llm_input_messages[0]["role"] in {"system", "user"} assert all("parts" in message for message in llm_input_messages) assert span.attributes["gen_ai.response.model"] - assert span.attributes["gen_ai.response.id"] + assert ( + span.attributes["gen_ai.response.id"] + == "chatcmpl-16b2e5f0-ebac-9b6a-b1f8-6cdf99efa597" + ) assert span.attributes["gen_ai.usage.input_tokens"] > 0 assert span.attributes["gen_ai.usage.output_tokens"] > 0 assert span.attributes["gen_ai.response.finish_reasons"] @@ -169,6 +172,14 @@ def test_streaming_llm_call_records_ttft( assert len(llm_spans) == 1 assert step_spans[0].attributes["gen_ai.react.finish_reason"] == "stop" assert llm_spans[0].attributes["gen_ai.response.time_to_first_token"] > 0 + assert ( + llm_spans[0].attributes["gen_ai.response.id"] + == "chatcmpl-2d34edbc-ab63-9129-944b-df41352dfa4b" + ) + assert ( + agent_spans[0].attributes["gen_ai.response.id"] + == "chatcmpl-2d34edbc-ab63-9129-944b-df41352dfa4b" + ) assert agent_spans[0].attributes["gen_ai.response.time_to_first_token"] > 0 assert agent_spans[0].parent is None _assert_parent(step_spans[0], agent_spans[0]) diff --git a/instrumentation-loongsuite/loongsuite-instrumentation-hermes-agent/tests/test_telemetry_spec.py b/instrumentation-loongsuite/loongsuite-instrumentation-hermes-agent/tests/test_telemetry_spec.py index 759e5bad8..a90377134 100644 --- a/instrumentation-loongsuite/loongsuite-instrumentation-hermes-agent/tests/test_telemetry_spec.py +++ b/instrumentation-loongsuite/loongsuite-instrumentation-hermes-agent/tests/test_telemetry_spec.py @@ -141,6 +141,10 @@ def fake_wrap_function_wrapper(module_name, name, wrapper): ) assert ("run_agent", "AIAgent.run_conversation") in wrapped_targets + assert ( + "run_agent", + "AIAgent._create_request_openai_client", + ) in wrapped_targets assert ("tools.memory_tool", "memory_tool") in wrapped_targets assert not any( module_name == "tools.session_search_tool" @@ -186,6 +190,10 @@ def fake_unwrap(parent, attribute): instrumentation_module.HermesAgentInstrumentor()._uninstrument() assert (ai_agent, "run_conversation") in unwrapped_targets + assert ( + ai_agent, + "_create_request_openai_client", + ) in unwrapped_targets assert (modules["tools.memory_tool"], "memory_tool") in unwrapped_targets assert not any( parent is modules["tools.delegate_tool"] @@ -389,6 +397,445 @@ def _runtime(instrumentation_module, tracer_provider, meter_provider): ) +class _FakeProviderResource: + def __init__(self, responses): + self._responses = iter(responses) + + def create(self, **_kwargs): + return next(self._responses) + + +class _ReadOnlyProviderResource: + __slots__ = ("_response",) + + def __init__(self, response): + self._response = response + + def create(self, **_kwargs): + return self._response + + +def _provider_client(responses): + resource = _FakeProviderResource(responses) + return SimpleNamespace( + chat=SimpleNamespace(completions=resource), + responses=None, + completions=None, + ) + + +def _stream_chunk(response_id=None, *, request_id=None): + return SimpleNamespace(id=response_id, request_id=request_id) + + +def _run_streaming_capture( + runtime, + agent, + provider_responses, + *, + framework_response_id="stream-framework-id", +): + wrappers_module = importlib.import_module( + "opentelemetry.instrumentation.hermes_agent.wrappers" + ) + provider_client_wrapper = wrappers_module.ProviderClientWrapper() + client = _provider_client(provider_responses) + + def streaming_call(_api_kwargs, *, on_first_delta): + request_client = provider_client_wrapper( + lambda: client, + agent, + (), + {}, + ) + first_delta = True + for _chunk in request_client.chat.completions.create(stream=True): + if first_delta: + first_delta = False + on_first_delta() + return _response( + content="完成", + response_id=framework_response_id, + ) + + runtime.streaming_llm_wrapper( + streaming_call, + agent, + ( + { + "model": agent.model, + "messages": [{"role": "user", "content": "hello"}], + }, + ), + {}, + ) + + +def test_streaming_provider_response_id_overrides_hermes_framework_id( + instrumentation_module, + tracer_provider, + meter_provider, + span_exporter, +): + runtime = _runtime(instrumentation_module, tracer_provider, meter_provider) + agent = _FakeAgent(session_id="session-provider-id") + + _run_streaming_capture( + runtime, + agent, + [ + iter( + [ + _stream_chunk("chatcmpl-provider-123"), + _stream_chunk("chatcmpl-provider-123"), + ] + ) + ], + ) + + llm_span = _spans_by_kind(span_exporter, "LLM")[0] + assert llm_span.attributes["gen_ai.response.id"] == "chatcmpl-provider-123" + + +def test_streaming_request_id_from_usage_trailer_is_preserved( + instrumentation_module, + tracer_provider, + meter_provider, + span_exporter, +): + runtime = _runtime(instrumentation_module, tracer_provider, meter_provider) + agent = _FakeAgent(session_id="session-request-id") + + _run_streaming_capture( + runtime, + agent, + [ + iter( + [ + _stream_chunk("chatcmpl-chunk-456"), + _stream_chunk(request_id="dashscope-request-456"), + ] + ) + ], + ) + + llm_span = _spans_by_kind(span_exporter, "LLM")[0] + assert llm_span.attributes["gen_ai.response.id"] == "dashscope-request-456" + + +@pytest.mark.parametrize( + ("retry_response_id", "expected_response_id"), + [ + ("chatcmpl-successful-attempt", "chatcmpl-successful-attempt"), + (None, "stream-retry-framework-fallback"), + ], +) +def test_streaming_retry_uses_only_final_provider_attempt( + instrumentation_module, + tracer_provider, + meter_provider, + span_exporter, + retry_response_id, + expected_response_id, +): + runtime = _runtime(instrumentation_module, tracer_provider, meter_provider) + agent = _FakeAgent(session_id="session-provider-retry-fallback") + wrappers_module = importlib.import_module( + "opentelemetry.instrumentation.hermes_agent.wrappers" + ) + provider_client_wrapper = wrappers_module.ProviderClientWrapper() + client = _provider_client( + [ + iter([_stream_chunk("chatcmpl-failed-attempt")]), + iter([_stream_chunk(retry_response_id)]), + ] + ) + + def streaming_call(_api_kwargs, *, on_first_delta): + errors = [] + + def provider_worker(): + try: + request_client = provider_client_wrapper( + lambda: client, + agent, + (), + {}, + ) + for _chunk in request_client.chat.completions.create( + stream=True + ): + # Hermes invokes this callback once for the whole logical + # call, even when it starts another provider attempt. + on_first_delta() + for _chunk in request_client.chat.completions.create( + stream=True + ): + pass + except BaseException as exc: # pragma: no cover - defensive + errors.append(exc) + + worker = threading.Thread(target=provider_worker) + worker.start() + worker.join(timeout=5) + assert not worker.is_alive() + assert errors == [] + return _response( + content="恢复成功", + response_id="stream-retry-framework-fallback", + ) + + runtime.streaming_llm_wrapper( + streaming_call, + agent, + ({"model": agent.model, "messages": []},), + {}, + ) + + llm_span = _spans_by_kind(span_exporter, "LLM")[0] + assert llm_span.attributes["gen_ai.response.id"] == expected_response_id + + +def test_streaming_without_provider_id_falls_back_to_hermes_response_id( + instrumentation_module, + tracer_provider, + meter_provider, + span_exporter, +): + runtime = _runtime(instrumentation_module, tracer_provider, meter_provider) + agent = _FakeAgent(session_id="session-provider-fallback") + + _run_streaming_capture( + runtime, + agent, + [iter([_stream_chunk()])], + framework_response_id="stream-hermes-fallback", + ) + + llm_span = _spans_by_kind(span_exporter, "LLM")[0] + assert ( + llm_span.attributes["gen_ai.response.id"] == "stream-hermes-fallback" + ) + + +def test_read_only_provider_resource_falls_back_without_breaking_call( + instrumentation_module, + tracer_provider, + meter_provider, + span_exporter, +): + runtime = _runtime(instrumentation_module, tracer_provider, meter_provider) + agent = _FakeAgent(session_id="session-provider-read-only") + wrappers_module = importlib.import_module( + "opentelemetry.instrumentation.hermes_agent.wrappers" + ) + provider_client_wrapper = wrappers_module.ProviderClientWrapper() + resource = _ReadOnlyProviderResource( + iter([_stream_chunk("chatcmpl-not-captured")]) + ) + client = SimpleNamespace( + chat=SimpleNamespace(completions=resource), + responses=None, + completions=None, + ) + + def streaming_call(_api_kwargs, *, on_first_delta): + request_client = provider_client_wrapper( + lambda: client, + agent, + (), + {}, + ) + for _chunk in request_client.chat.completions.create(stream=True): + on_first_delta() + return _response( + content="完成", + response_id="stream-read-only-fallback", + ) + + runtime.streaming_llm_wrapper( + streaming_call, + agent, + ({"model": agent.model, "messages": []},), + {}, + ) + + llm_span = _spans_by_kind(span_exporter, "LLM")[0] + assert ( + llm_span.attributes["gen_ai.response.id"] + == "stream-read-only-fallback" + ) + + +def test_provider_error_is_reported_without_masking_original_exception( + instrumentation_module, + tracer_provider, + meter_provider, + span_exporter, +): + runtime = _runtime(instrumentation_module, tracer_provider, meter_provider) + agent = _FakeAgent(session_id="session-provider-error") + wrappers_module = importlib.import_module( + "opentelemetry.instrumentation.hermes_agent.wrappers" + ) + provider_client_wrapper = wrappers_module.ProviderClientWrapper() + + class ProviderError(RuntimeError): + pass + + class FailingResource: + @staticmethod + def create(**_kwargs): + raise ProviderError("provider unavailable") + + client = SimpleNamespace( + chat=SimpleNamespace(completions=FailingResource()) + ) + + def streaming_call(_api_kwargs, *, on_first_delta): + del on_first_delta + request_client = provider_client_wrapper( + lambda: client, + agent, + (), + {}, + ) + return request_client.chat.completions.create(stream=True) + + with pytest.raises(ProviderError, match="provider unavailable"): + runtime.streaming_llm_wrapper( + streaming_call, + agent, + ({"model": agent.model, "messages": []},), + {}, + ) + + llm_span = _spans_by_kind(span_exporter, "LLM")[0] + assert llm_span.status.status_code == StatusCode.ERROR + + +def test_provider_attempt_from_previous_call_does_not_leak( + instrumentation_module, + tracer_provider, + meter_provider, + span_exporter, +): + runtime = _runtime(instrumentation_module, tracer_provider, meter_provider) + agent = _FakeAgent(session_id="session-provider-stale-attempt") + wrappers_module = importlib.import_module( + "opentelemetry.instrumentation.hermes_agent.wrappers" + ) + provider_client_wrapper = wrappers_module.ProviderClientWrapper() + client = _provider_client( + [iter([_stream_chunk("chatcmpl-previous-call")])] + ) + request_client = provider_client_wrapper(lambda: client, agent, (), {}) + list(request_client.chat.completions.create(stream=True)) + + runtime.llm_wrapper( + lambda _api_kwargs: _response( + content="新调用", + response_id="framework-current-call", + ), + agent, + ({"model": agent.model, "messages": []},), + {}, + ) + + llm_span = _spans_by_kind(span_exporter, "LLM")[0] + assert ( + llm_span.attributes["gen_ai.response.id"] == "framework-current-call" + ) + + +def test_native_request_id_precedes_framework_response_id( + instrumentation_module, + tracer_provider, + meter_provider, + span_exporter, +): + runtime = _runtime(instrumentation_module, tracer_provider, meter_provider) + agent = _FakeAgent(session_id="session-native-request-id") + response = _response(content="完成", response_id="framework-response-id") + response.request_id = "dashscope-native-request-id" + + runtime.llm_wrapper( + lambda api_kwargs: response, + agent, + ({"model": agent.model, "messages": []},), + {}, + ) + + llm_span = _spans_by_kind(span_exporter, "LLM")[0] + assert ( + llm_span.attributes["gen_ai.response.id"] + == "dashscope-native-request-id" + ) + + +def test_streaming_provider_response_ids_are_isolated_across_threads( + instrumentation_module, + tracer_provider, + meter_provider, + span_exporter, +): + runtime = _runtime(instrumentation_module, tracer_provider, meter_provider) + shared_agent = _FakeAgent(session_id="session-provider-concurrent") + barrier = threading.Barrier(2) + errors = [] + + def run_stream(thread_id): + try: + wrappers_module = importlib.import_module( + "opentelemetry.instrumentation.hermes_agent.wrappers" + ) + provider_client_wrapper = wrappers_module.ProviderClientWrapper() + client = _provider_client( + [iter([_stream_chunk(f"chatcmpl-thread-{thread_id}")])] + ) + + def streaming_call(_api_kwargs, *, on_first_delta): + request_client = provider_client_wrapper( + lambda: client, + shared_agent, + (), + {}, + ) + stream = request_client.chat.completions.create(stream=True) + next(stream) + barrier.wait(timeout=5) + on_first_delta() + return _response( + content=f"完成-{thread_id}", + response_id=f"stream-framework-{thread_id}", + ) + + runtime.streaming_llm_wrapper( + streaming_call, + shared_agent, + ({"model": shared_agent.model, "messages": []},), + {}, + ) + except BaseException as exc: # pragma: no cover - defensive + errors.append(exc) + + threads = [ + threading.Thread(target=run_stream, args=(thread_id,)) + for thread_id in (1, 2) + ] + for thread in threads: + thread.start() + for thread in threads: + thread.join(timeout=10) + + assert errors == [] + assert all(not thread.is_alive() for thread in threads) + assert { + span.attributes["gen_ai.response.id"] + for span in _spans_by_kind(span_exporter, "LLM") + } == {"chatcmpl-thread-1", "chatcmpl-thread-2"} + + def test_agent_layer_reuses_existing_entry_instead_of_creating_a_new_one( instrumentation_module, tracer_provider, diff --git a/util/opentelemetry-util-genai/CHANGELOG-loongsuite.md b/util/opentelemetry-util-genai/CHANGELOG-loongsuite.md index 3812b57b4..e829b1d7e 100644 --- a/util/opentelemetry-util-genai/CHANGELOG-loongsuite.md +++ b/util/opentelemetry-util-genai/CHANGELOG-loongsuite.md @@ -7,6 +7,11 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## Unreleased +### Added + +- Add shared response ID extraction and provider-first fallback helpers for + GenAI instrumentations that receive both provider and framework responses. + ## Version 0.7.0 (2026-07-03) ### Added diff --git a/util/opentelemetry-util-genai/src/opentelemetry/util/genai/response_id.py b/util/opentelemetry-util-genai/src/opentelemetry/util/genai/response_id.py new file mode 100644 index 000000000..eac3467da --- /dev/null +++ b/util/opentelemetry-util-genai/src/opentelemetry/util/genai/response_id.py @@ -0,0 +1,79 @@ +# Copyright The OpenTelemetry Authors +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +"""Helpers for resolving provider and framework response identifiers.""" + +from __future__ import annotations + +from collections.abc import Mapping, Sequence +from typing import Any, cast + + +def _normalize_response_id(value: Any) -> str | None: + if isinstance(value, bool) or not isinstance(value, (str, int)): + return None + normalized = str(value).strip() + return normalized or None + + +def extract_response_id( + response: Any, + *, + fields: Sequence[str] = ("id",), +) -> str | None: + """Extract the first non-empty identifier from ``response``. + + ``response`` may be an identifier itself, a mapping, or an SDK response + object. Field order is deliberately caller-controlled because providers + use different names for the same operation identifier. Transport-only + fields such as OpenAI's ``_request_id`` are not considered implicitly. + """ + + direct_identifier = _normalize_response_id(response) + if direct_identifier is not None: + return direct_identifier + + for field in fields: + try: + if isinstance(response, Mapping): + mapping_response = cast(Mapping[str, Any], response) + candidate = mapping_response.get(field) + else: + candidate = getattr(response, field, None) + except Exception: # pylint: disable=broad-exception-caught + # Provider SDK properties may raise arbitrary lazy-load errors; + # response-ID telemetry must never break the model call. + continue + identifier = _normalize_response_id(candidate) + if identifier is not None: + return identifier + return None + + +def resolve_response_id( + provider_response: Any = None, + framework_response: Any = None, + *, + provider_fields: Sequence[str] = ("id",), + framework_fields: Sequence[str] = ("id",), +) -> str | None: + """Prefer a provider identifier and fall back to the framework response.""" + + return extract_response_id( + provider_response, + fields=provider_fields, + ) or extract_response_id( + framework_response, + fields=framework_fields, + ) diff --git a/util/opentelemetry-util-genai/tests/test_response_id.py b/util/opentelemetry-util-genai/tests/test_response_id.py new file mode 100644 index 000000000..2841f5045 --- /dev/null +++ b/util/opentelemetry-util-genai/tests/test_response_id.py @@ -0,0 +1,90 @@ +# Copyright The OpenTelemetry Authors +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +from types import SimpleNamespace + +import pytest + +from opentelemetry.util.genai.response_id import ( + extract_response_id, + resolve_response_id, +) + + +@pytest.mark.parametrize( + ("response", "fields", "expected"), + [ + (" chatcmpl-123 ", ("id",), "chatcmpl-123"), + (123, ("id",), "123"), + ({"request_id": "dashscope-123"}, ("request_id",), "dashscope-123"), + (SimpleNamespace(id="msg-123"), ("id",), "msg-123"), + ( + {"id": "", "response_id": "resp-123"}, + ("id", "response_id"), + "resp-123", + ), + (SimpleNamespace(id=None), ("id",), None), + (True, ("id",), None), + ], +) +def test_extract_response_id(response, fields, expected): + assert extract_response_id(response, fields=fields) == expected + + +def test_extract_response_id_uses_caller_field_order(): + response = {"id": "completion-123", "request_id": "request-123"} + + assert ( + extract_response_id(response, fields=("request_id", "id")) + == "request-123" + ) + assert ( + extract_response_id(response, fields=("id", "request_id")) + == "completion-123" + ) + + +def test_extract_response_id_does_not_use_transport_id_implicitly(): + response = SimpleNamespace(_request_id="transport-request-123") + + assert extract_response_id(response) is None + + +def test_resolve_response_id_prefers_provider_and_falls_back_to_framework(): + framework_response = SimpleNamespace(id="stream-framework-123") + + assert ( + resolve_response_id("chatcmpl-provider-123", framework_response) + == "chatcmpl-provider-123" + ) + assert ( + resolve_response_id(None, framework_response) == "stream-framework-123" + ) + + +def test_extract_response_id_ignores_raising_sdk_property(): + class RaisingResponse: + @property + def request_id(self): + raise RuntimeError("not loaded") + + id = "chatcmpl-after-error" + + assert ( + extract_response_id( + RaisingResponse(), + fields=("request_id", "id"), + ) + == "chatcmpl-after-error" + )