From d83982c33505bac0fe093f911c91eb669cce0cda Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E6=B5=81=E5=B1=BF?= Date: Mon, 20 Jul 2026 20:54:48 +0800 Subject: [PATCH 1/3] fix(hermes): preserve provider response IDs --- .../CHANGELOG.md | 6 + .../instrumentation/hermes_agent/helpers.py | 47 +- .../hermes_agent/instrumentor.py | 7 + .../instrumentation/hermes_agent/wrappers.py | 149 +++++- .../tests/test_live_llm_calls.py | 13 +- .../tests/test_telemetry_spec.py | 425 ++++++++++++++++++ 6 files changed, 644 insertions(+), 3 deletions(-) diff --git a/instrumentation-loongsuite/loongsuite-instrumentation-hermes-agent/CHANGELOG.md b/instrumentation-loongsuite/loongsuite-instrumentation-hermes-agent/CHANGELOG.md index 74194d067..0b45db9ee 100644 --- a/instrumentation-loongsuite/loongsuite-instrumentation-hermes-agent/CHANGELOG.md +++ b/instrumentation-loongsuite/loongsuite-instrumentation-hermes-agent/CHANGELOG.md @@ -7,6 +7,12 @@ 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. + ## Version 0.7.0 (2026-07-03) ### Fixed 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..e39128ccd 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 @@ -69,6 +69,38 @@ def obj_get(value: Any, field: str, default: Any = None) -> Any: return getattr(value, field, default) +def response_identifier( + value: Any, + *, + fields: tuple[str, ...] = ( + "id", + "request_id", + "response_id", + "_request_id", + ), +) -> str | None: + """Return a normalized provider/framework response identifier. + + OpenAI-compatible providers expose ``id`` while native DashScope-style + responses commonly call the same correlation value ``request_id``. Keep + the extraction intentionally narrow so unrelated object identifiers are + never promoted to ``gen_ai.response.id``. + """ + + if isinstance(value, str): + normalized_value = value.strip() + return normalized_value or None + + for field in fields: + candidate = obj_get(value, field) + if not isinstance(candidate, (str, int)): + continue + normalized = str(candidate).strip() + if normalized: + return normalized + return None + + def _normalize_platform(value: Any) -> str: platform = getattr(value, "value", value) return str(platform or "").strip().lower() @@ -674,6 +706,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 +715,18 @@ def update_llm_invocation_from_response( else: invocation.response_model_name = invocation.request_model - response_id = getattr(response, "id", None) + # Prefer the ID captured from the provider SDK stream. Hermes currently + # synthesizes ``stream-*`` IDs while aggregating OpenAI-compatible chunks, + # so the final framework response is only a fallback. Native-style + # ``request_id`` fields are checked before the framework ``id`` for the + # same reason. + response_id = response_identifier( + provider_response_id, + fields=("id",), + ) or response_identifier( + response, + fields=("request_id", "id", "response_id", "_request_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..68beb281d 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 @@ -37,6 +39,7 @@ push_state, reset_state, resolve_entry_platform, + response_identifier, start_step, state, step_finish_reason, @@ -59,6 +62,131 @@ "embedding ", "embeddings ", ) +_PROVIDER_CREATE_WRAPPED = "_otel_hermes_response_id_capture_wrapped" +_PROVIDER_ATTEMPT = threading.local() +_MISSING_PROVIDER_ATTEMPT = 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 + + @property + def response_id(self) -> str | None: + with self._lock: + return self._response_id + + def record(self, value: Any) -> None: + response_id = response_identifier(value) + if response_id is None: + return + with self._lock: + if self._response_id is None: + self._response_id = response_id + + +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 + 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: + if resource is None or getattr(resource, _PROVIDER_CREATE_WRAPPED, False): + return + create = getattr(resource, "create", None) + if not callable(create): + return + + @functools.wraps(create) + def _capturing_create(*args, **kwargs): + attempt = _ProviderResponseAttempt() + _PROVIDER_ATTEMPT.value = 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 + + try: + setattr(resource, "create", _capturing_create) + setattr(resource, _PROVIDER_CREATE_WRAPPED, True) + 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) + resources = ( + getattr(getattr(client, "chat", None), "completions", None), + getattr(client, "responses", None), + getattr(client, "completions", None), + ) + for resource in resources: + _wrap_provider_create(resource) + return client def _resolve_handler( @@ -307,16 +435,24 @@ 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_PROVIDER_ATTEMPT, + ) + _PROVIDER_ATTEMPT.value = None 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 +463,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 +521,11 @@ def _wrapped_first_delta(): finish_step(instance, "error", exc=exc) raise finally: + if previous_provider_attempt is _MISSING_PROVIDER_ATTEMPT: + with suppress(AttributeError): + del _PROVIDER_ATTEMPT.value + else: + _PROVIDER_ATTEMPT.value = previous_provider_attempt current_state["active_llm_depth"] = max( 0, current_state["active_llm_depth"] - 1 ) 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..b4f861beb 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,423 @@ 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(), + _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" + + +def test_streaming_retry_uses_successful_provider_attempt_id( + instrumentation_module, + tracer_provider, + meter_provider, + span_exporter, +): + runtime = _runtime(instrumentation_module, tracer_provider, meter_provider) + agent = _FakeAgent(session_id="session-provider-retry") + 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("chatcmpl-successful-attempt")]), + ] + ) + + 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() + for _chunk in request_client.chat.completions.create(stream=True): + on_first_delta() + return _response( + content="恢复成功", + response_id="stream-retry-framework-id", + ) + + 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"] + == "chatcmpl-successful-attempt" + ) + + +def test_streaming_retry_does_not_reuse_failed_provider_attempt_id( + instrumentation_module, + tracer_provider, + meter_provider, + span_exporter, +): + 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()]), + ] + ) + + 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() + for _chunk in request_client.chat.completions.create(stream=True): + on_first_delta() + 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"] + == "stream-retry-framework-fallback" + ) + + +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_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, From 93411f96cfcb985acd953318e1842165e56869a4 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E6=B5=81=E5=B1=BF?= Date: Tue, 21 Jul 2026 19:55:09 +0800 Subject: [PATCH 2/3] refactor(genai): share provider response ID resolution --- .../instrumentation/hermes_agent/helpers.py | 44 +---- .../instrumentation/hermes_agent/wrappers.py | 52 ++++-- .../tests/test_telemetry_spec.py | 152 ++++++++++-------- .../CHANGELOG-loongsuite.md | 5 + .../opentelemetry/util/genai/response_id.py | 79 +++++++++ .../tests/test_response_id.py | 90 +++++++++++ 6 files changed, 300 insertions(+), 122 deletions(-) create mode 100644 util/opentelemetry-util-genai/src/opentelemetry/util/genai/response_id.py create mode 100644 util/opentelemetry-util-genai/tests/test_response_id.py 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 e39128ccd..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, @@ -69,38 +70,6 @@ def obj_get(value: Any, field: str, default: Any = None) -> Any: return getattr(value, field, default) -def response_identifier( - value: Any, - *, - fields: tuple[str, ...] = ( - "id", - "request_id", - "response_id", - "_request_id", - ), -) -> str | None: - """Return a normalized provider/framework response identifier. - - OpenAI-compatible providers expose ``id`` while native DashScope-style - responses commonly call the same correlation value ``request_id``. Keep - the extraction intentionally narrow so unrelated object identifiers are - never promoted to ``gen_ai.response.id``. - """ - - if isinstance(value, str): - normalized_value = value.strip() - return normalized_value or None - - for field in fields: - candidate = obj_get(value, field) - if not isinstance(candidate, (str, int)): - continue - normalized = str(candidate).strip() - if normalized: - return normalized - return None - - def _normalize_platform(value: Any) -> str: platform = getattr(value, "value", value) return str(platform or "").strip().lower() @@ -715,17 +684,10 @@ def update_llm_invocation_from_response( else: invocation.response_model_name = invocation.request_model - # Prefer the ID captured from the provider SDK stream. Hermes currently - # synthesizes ``stream-*`` IDs while aggregating OpenAI-compatible chunks, - # so the final framework response is only a fallback. Native-style - # ``request_id`` fields are checked before the framework ``id`` for the - # same reason. - response_id = response_identifier( + response_id = resolve_response_id( provider_response_id, - fields=("id",), - ) or response_identifier( response, - fields=("request_id", "id", "response_id", "_request_id"), + 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/wrappers.py b/instrumentation-loongsuite/loongsuite-instrumentation-hermes-agent/src/opentelemetry/instrumentation/hermes_agent/wrappers.py index 68beb281d..223ca18e8 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 @@ -26,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 ( @@ -39,7 +40,6 @@ push_state, reset_state, resolve_entry_platform, - response_identifier, start_step, state, step_finish_reason, @@ -64,7 +64,8 @@ ) _PROVIDER_CREATE_WRAPPED = "_otel_hermes_response_id_capture_wrapped" _PROVIDER_ATTEMPT = threading.local() -_MISSING_PROVIDER_ATTEMPT = object() +_PROVIDER_CAPTURE = threading.local() +_MISSING_THREAD_LOCAL = object() class _ProviderResponseAttempt: @@ -80,7 +81,10 @@ def response_id(self) -> str | None: return self._response_id def record(self, value: Any) -> None: - response_id = response_identifier(value) + response_id = extract_response_id( + value, + fields=("request_id", "id", "response_id"), + ) if response_id is None: return with self._lock: @@ -99,6 +103,10 @@ 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 @@ -143,16 +151,21 @@ def __getattr__(self, name: str) -> Any: def _wrap_provider_create(resource: Any) -> None: - if resource is None or getattr(resource, _PROVIDER_CREATE_WRAPPED, False): - return - create = getattr(resource, "create", None) - if not callable(create): + 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) @@ -164,9 +177,9 @@ def _capturing_create(*args, **kwargs): return _ProviderStreamProxy(response, attempt) return response + setattr(_capturing_create, _PROVIDER_CREATE_WRAPPED, True) try: setattr(resource, "create", _capturing_create) - setattr(resource, _PROVIDER_CREATE_WRAPPED, True) except (AttributeError, TypeError): # A provider resource may use slots/read-only descriptors. Telemetry # must never make the model call fail; the framework response remains @@ -179,13 +192,9 @@ class ProviderClientWrapper: def __call__(self, wrapped, instance, args, kwargs): client = wrapped(*args, **kwargs) - resources = ( - getattr(getattr(client, "chat", None), "completions", None), - getattr(client, "responses", None), - getattr(client, "completions", None), + _wrap_provider_create( + getattr(getattr(client, "chat", None), "completions", None) ) - for resource in resources: - _wrap_provider_create(resource) return client @@ -443,9 +452,15 @@ def __call__(self, wrapped, instance, args, kwargs): previous_provider_attempt = getattr( _PROVIDER_ATTEMPT, "value", - _MISSING_PROVIDER_ATTEMPT, + _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: @@ -521,11 +536,16 @@ def _wrapped_first_delta(): finish_step(instance, "error", exc=exc) raise finally: - if previous_provider_attempt is _MISSING_PROVIDER_ATTEMPT: + 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/test_telemetry_spec.py b/instrumentation-loongsuite/loongsuite-instrumentation-hermes-agent/tests/test_telemetry_spec.py index b4f861beb..a3fe5f592 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 @@ -523,60 +523,20 @@ def test_streaming_request_id_from_usage_trailer_is_preserved( assert llm_span.attributes["gen_ai.response.id"] == "dashscope-request-456" -def test_streaming_retry_uses_successful_provider_attempt_id( - instrumentation_module, - tracer_provider, - meter_provider, - span_exporter, -): - runtime = _runtime(instrumentation_module, tracer_provider, meter_provider) - agent = _FakeAgent(session_id="session-provider-retry") - 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("chatcmpl-successful-attempt")]), - ] - ) - - 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() - for _chunk in request_client.chat.completions.create(stream=True): - on_first_delta() - return _response( - content="恢复成功", - response_id="stream-retry-framework-id", - ) - - 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"] - == "chatcmpl-successful-attempt" - ) - - -def test_streaming_retry_does_not_reuse_failed_provider_attempt_id( +@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") @@ -587,21 +547,39 @@ def test_streaming_retry_does_not_reuse_failed_provider_attempt_id( client = _provider_client( [ iter([_stream_chunk("chatcmpl-failed-attempt")]), - iter([_stream_chunk()]), + iter([_stream_chunk(retry_response_id)]), ] ) 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() - for _chunk in request_client.chat.completions.create(stream=True): - 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", @@ -615,10 +593,7 @@ def streaming_call(_api_kwargs, *, on_first_delta): ) llm_span = _spans_by_kind(span_exporter, "LLM")[0] - assert ( - llm_span.attributes["gen_ai.response.id"] - == "stream-retry-framework-fallback" - ) + assert llm_span.attributes["gen_ai.response.id"] == expected_response_id def test_streaming_without_provider_id_falls_back_to_hermes_response_id( @@ -692,6 +667,53 @@ def streaming_call(_api_kwargs, *, on_first_delta): ) +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, 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..d295e6b11 --- /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 + + +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: + candidate = ( + response.get(field) + if isinstance(response, Mapping) + else 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" + ) From 89e6c5f68295f43fa2be60892f2c7d4a260b4c43 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E6=B5=81=E5=B1=BF?= Date: Wed, 22 Jul 2026 11:31:14 +0800 Subject: [PATCH 3/3] fix(hermes): align response id dependencies --- .../CHANGELOG.md | 2 ++ .../pyproject.toml | 10 +++++----- .../instrumentation/hermes_agent/wrappers.py | 18 +++++++++++++----- .../tests/requirements.txt | 10 +++++----- .../tests/test_telemetry_spec.py | 2 +- .../opentelemetry/util/genai/response_id.py | 12 ++++++------ 6 files changed, 32 insertions(+), 22 deletions(-) diff --git a/instrumentation-loongsuite/loongsuite-instrumentation-hermes-agent/CHANGELOG.md b/instrumentation-loongsuite/loongsuite-instrumentation-hermes-agent/CHANGELOG.md index 0b45db9ee..2c87fcfe0 100644 --- a/instrumentation-loongsuite/loongsuite-instrumentation-hermes-agent/CHANGELOG.md +++ b/instrumentation-loongsuite/loongsuite-instrumentation-hermes-agent/CHANGELOG.md @@ -12,6 +12,8 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 - 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) 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/wrappers.py b/instrumentation-loongsuite/loongsuite-instrumentation-hermes-agent/src/opentelemetry/instrumentation/hermes_agent/wrappers.py index 223ca18e8..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 @@ -74,6 +74,7 @@ class _ProviderResponseAttempt: 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: @@ -81,15 +82,22 @@ def response_id(self) -> str | None: return self._response_id def record(self, value: Any) -> None: - response_id = extract_response_id( - value, - fields=("request_id", "id", "response_id"), - ) + # 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 self._response_id is None: + if priority > self._response_id_priority: self._response_id = response_id + self._response_id_priority = priority class _ProviderResponseIdCapture: 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_telemetry_spec.py b/instrumentation-loongsuite/loongsuite-instrumentation-hermes-agent/tests/test_telemetry_spec.py index a3fe5f592..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 @@ -512,7 +512,7 @@ def test_streaming_request_id_from_usage_trailer_is_preserved( [ iter( [ - _stream_chunk(), + _stream_chunk("chatcmpl-chunk-456"), _stream_chunk(request_id="dashscope-request-456"), ] ) 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 index d295e6b11..eac3467da 100644 --- a/util/opentelemetry-util-genai/src/opentelemetry/util/genai/response_id.py +++ b/util/opentelemetry-util-genai/src/opentelemetry/util/genai/response_id.py @@ -17,7 +17,7 @@ from __future__ import annotations from collections.abc import Mapping, Sequence -from typing import Any +from typing import Any, cast def _normalize_response_id(value: Any) -> str | None: @@ -46,11 +46,11 @@ def extract_response_id( for field in fields: try: - candidate = ( - response.get(field) - if isinstance(response, Mapping) - else getattr(response, field, None) - ) + 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.