Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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",
]
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,7 @@
InvokeAgentInvocation,
ReactStepInvocation,
)
from opentelemetry.util.genai.response_id import resolve_response_id
from opentelemetry.util.genai.types import (
FunctionToolDefinition,
GenericToolDefinition,
Expand Down Expand Up @@ -674,14 +675,20 @@ 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:
invocation.response_model_name = response_model
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

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,7 @@
from .metrics import HermesMetrics
from .wrappers import (
LLMCallWrapper,
ProviderClientWrapper,
RunConversationWrapper,
ToolBatchWrapper,
ToolCallWrapper,
Expand Down Expand Up @@ -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",
Expand Down Expand Up @@ -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")
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,13 +17,16 @@
from __future__ import annotations

import contextvars
import functools
import threading
import timeit
from collections.abc import Mapping
from contextlib import suppress
from typing import Any

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 (
Expand Down Expand Up @@ -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(
Expand Down Expand Up @@ -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:
Expand All @@ -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,
)
)

Expand Down Expand Up @@ -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
)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Original file line number Diff line number Diff line change
Expand Up @@ -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"]
Expand Down Expand Up @@ -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])
Expand Down
Loading
Loading