From bd26f5a80f6f85c3e4074cbc4efb43d14452822d Mon Sep 17 00:00:00 2001 From: zhanghui Date: Thu, 23 Jul 2026 20:47:39 +0800 Subject: [PATCH] feat(observability): capture LLM + embedding token usage MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Surface gen_ai.* model + token attributes onto the active span so Langfuse can compute cost — without touching everalgo: - UsageRecordingClient wraps the LLM client and records response.usage after each chat(); get_llm_client composes it over the existing _LoggingLLMClient only when observability is enabled (disabled default stays overhead-free). - OpenAIEmbeddingProvider records its response.usage (input tokens) onto the active span too. Tokens land on the everos.extract / everos.reflect.consolidate generation spans and the search embedding recall; no-op when tracing is off. Co-Authored-By: Claude Opus 4.8 (1M context) --- .../component/embedding/openai_provider.py | 9 ++ src/everos/component/llm/_usage_client.py | 62 +++++++++++ src/everos/component/llm/client.py | 12 ++- .../test_embedding/test_usage_span.py | 95 ++++++++++++++++ .../test_component/test_llm/test_client.py | 42 +++++++- .../test_llm/test_usage_client.py | 101 ++++++++++++++++++ 6 files changed, 318 insertions(+), 3 deletions(-) create mode 100644 src/everos/component/llm/_usage_client.py create mode 100644 tests/unit/test_component/test_embedding/test_usage_span.py create mode 100644 tests/unit/test_component/test_llm/test_usage_client.py diff --git a/src/everos/component/embedding/openai_provider.py b/src/everos/component/embedding/openai_provider.py index f756127..08d131d 100644 --- a/src/everos/component/embedding/openai_provider.py +++ b/src/everos/component/embedding/openai_provider.py @@ -23,6 +23,8 @@ from collections.abc import Sequence import openai +from everos.core.observability.tracing import set_generation_usage + from .protocol import EmbeddingServiceError @@ -94,5 +96,12 @@ class OpenAIEmbeddingProvider: ) except openai.OpenAIError as exc: raise EmbeddingServiceError(str(exc)) from exc + # Surface token usage onto the active span (e.g. everos.search.embed_query). + # No-op when tracing is off; embeddings report only input (prompt) tokens. + usage = getattr(response, "usage", None) + set_generation_usage( + model=self._model, + input_tokens=usage.prompt_tokens if usage else None, + ) # OpenAI returns ``data`` indexed by request order; truncate to ``dim``. return [list(item.embedding[: self.dim]) for item in response.data] diff --git a/src/everos/component/llm/_usage_client.py b/src/everos/component/llm/_usage_client.py new file mode 100644 index 0000000..a4bb7b7 --- /dev/null +++ b/src/everos/component/llm/_usage_client.py @@ -0,0 +1,62 @@ +"""Token-usage-recording LLM client wrapper. + +Wraps any :class:`everalgo.llm.LLMClient` and, after each ``chat`` call, +writes ``response.usage`` (+ model) onto the current OpenTelemetry span via +``set_generation_usage``. The wrapped client's behaviour is otherwise +untouched — the response is returned verbatim and all other attributes +delegate through. + +This is how token counts reach the ``everos.extract`` generation span: +everalgo's extractors own the ``chat`` call and discard the ``ChatResponse`` +(keeping only ``.content``), so usage is captured here at the client +boundary — no everalgo change required. When no span is active (e.g. the +search path, or tracing disabled), ``set_generation_usage`` no-ops. +""" + +from __future__ import annotations + +from typing import TYPE_CHECKING, Any + +from everos.core.observability.tracing import set_generation_usage + +if TYPE_CHECKING: + from everalgo.llm import ChatMessage, ChatResponse + from everalgo.llm.protocols import LLMClient + from pydantic import BaseModel + + +class UsageRecordingClient: + """LLM client proxy that records token usage onto the active span.""" + + def __init__(self, inner: LLMClient) -> None: + self._inner = inner + + async def chat( + self, + messages: list[ChatMessage], + *, + model: str | None = None, + temperature: float | None = None, + max_tokens: int | None = None, + response_format: type[BaseModel] | None = None, + **extra: Any, + ) -> ChatResponse: + response = await self._inner.chat( + messages, + model=model, + temperature=temperature, + max_tokens=max_tokens, + response_format=response_format, + **extra, + ) + usage = response.usage + set_generation_usage( + model=response.model, + input_tokens=usage.prompt_tokens if usage else None, + output_tokens=usage.completion_tokens if usage else None, + ) + return response + + def __getattr__(self, name: str) -> Any: + # Delegate everything else to the wrapped client. + return getattr(self._inner, name) diff --git a/src/everos/component/llm/client.py b/src/everos/component/llm/client.py index 846dcf1..5b0de8c 100644 --- a/src/everos/component/llm/client.py +++ b/src/everos/component/llm/client.py @@ -16,6 +16,8 @@ from everalgo.llm.protocols import LLMClient from everos.config import load_settings from everos.core.observability.logging import get_logger +from ._usage_client import UsageRecordingClient + logger = get_logger(__name__) @@ -38,7 +40,8 @@ def get_llm_client() -> LLMClient: if _llm_client is not None: return _llm_client - llm_cfg = load_settings().llm + settings = load_settings() + llm_cfg = settings.llm api_key = ( llm_cfg.api_key.get_secret_value() if llm_cfg.api_key is not None else None ) @@ -46,13 +49,18 @@ def get_llm_client() -> LLMClient: raise LLMNotConfiguredError( "LLM is required; set EVEROS_LLM__API_KEY + EVEROS_LLM__BASE_URL" ) - _llm_client = build_client( + client: LLMClient = build_client( LLMConfig( model=llm_cfg.model, api_key=api_key, base_url=llm_cfg.base_url, ) ) + # Wrap for OTel token capture only when tracing is on — keeps the + # disabled path (the default) allocation- and overhead-free. + if settings.observability.enabled: + client = UsageRecordingClient(client) + _llm_client = client logger.info("llm_client_built", model=llm_cfg.model) return _llm_client diff --git a/tests/unit/test_component/test_embedding/test_usage_span.py b/tests/unit/test_component/test_embedding/test_usage_span.py new file mode 100644 index 0000000..6005607 --- /dev/null +++ b/tests/unit/test_component/test_embedding/test_usage_span.py @@ -0,0 +1,95 @@ +"""``OpenAIEmbeddingProvider`` surfaces token usage onto the active span. + +The OpenAI embeddings response carries ``usage``; the provider records it +via ``set_generation_usage`` (no contract change) so the +``everos.search.embed_query`` embedding span shows model + input tokens. +Captured via an in-memory exporter. +""" + +from __future__ import annotations + +from collections.abc import Iterator +from types import SimpleNamespace + +import pytest +from opentelemetry.sdk.trace.export import SimpleSpanProcessor +from opentelemetry.sdk.trace.export.in_memory_span_exporter import ( + InMemorySpanExporter, +) + +from everos.component.embedding.openai_provider import OpenAIEmbeddingProvider +from everos.config.settings import ObservabilitySettings +from everos.core.observability.tracing import ( + force_flush, + init_tracing, + memory_span, + shutdown_tracing, +) + + +class _FakeEmbeddings: + def __init__(self, response: object) -> None: + self._response = response + + async def create(self, *, model: str, input: list[str]) -> object: + return self._response + + +class _FakeClient: + def __init__(self, response: object) -> None: + self.embeddings = _FakeEmbeddings(response) + + +@pytest.fixture(autouse=True) +def _reset() -> Iterator[None]: + shutdown_tracing() + yield + shutdown_tracing() + + +@pytest.fixture +def captured() -> Iterator[InMemorySpanExporter]: + exporter = InMemorySpanExporter() + init_tracing( + ObservabilitySettings(enabled=True, endpoint="http://collector.invalid"), + span_processor=SimpleSpanProcessor(exporter), + ) + yield exporter + + +def _provider(response: object) -> OpenAIEmbeddingProvider: + provider = OpenAIEmbeddingProvider( + model="emb-m", api_key="k", base_url="http://x", dim=8 + ) + provider._client = _FakeClient(response) # type: ignore[assignment] + return provider + + +async def test_embed_records_usage_to_active_span( + captured: InMemorySpanExporter, +) -> None: + resp = SimpleNamespace( + data=[SimpleNamespace(embedding=[0.1] * 8)], + usage=SimpleNamespace(prompt_tokens=5, total_tokens=5), + ) + provider = _provider(resp) + with memory_span("everos.search.embed_query", observation_type="embedding"): + vec = await provider.embed("hello") + assert len(vec) == 8 + force_flush() + attrs = captured.get_finished_spans()[0].attributes + assert attrs["gen_ai.request.model"] == "emb-m" + assert attrs["gen_ai.usage.input_tokens"] == 5 + + +async def test_embed_without_usage_does_not_raise( + captured: InMemorySpanExporter, +) -> None: + resp = SimpleNamespace(data=[SimpleNamespace(embedding=[0.1] * 8)], usage=None) + provider = _provider(resp) + with memory_span("everos.search.embed_query", observation_type="embedding"): + vec = await provider.embed("hello") + assert len(vec) == 8 + force_flush() + attrs = captured.get_finished_spans()[0].attributes + assert "gen_ai.usage.input_tokens" not in attrs diff --git a/tests/unit/test_component/test_llm/test_client.py b/tests/unit/test_component/test_llm/test_client.py index eb7427e..46098e8 100644 --- a/tests/unit/test_component/test_llm/test_client.py +++ b/tests/unit/test_component/test_llm/test_client.py @@ -8,8 +8,9 @@ import pytest from pydantic import SecretStr from everos.component.llm import LLMNotConfiguredError +from everos.component.llm._usage_client import UsageRecordingClient from everos.config import Settings -from everos.config.settings import LLMSettings +from everos.config.settings import LLMSettings, ObservabilitySettings _client_mod = importlib.import_module("everos.component.llm.client") @@ -62,3 +63,42 @@ def test_returns_singleton_when_configured(monkeypatch: pytest.MonkeyPatch) -> N assert first is sentinel assert first is second + + +def _patch_settings_with_observability( + monkeypatch: pytest.MonkeyPatch, *, enabled: bool +) -> None: + cfg = Settings( + llm=LLMSettings( + model="gpt-4.1-mini", + api_key=SecretStr("sk-test"), + base_url="https://example.test", + ), + observability=ObservabilitySettings(enabled=enabled), + ) + monkeypatch.setattr(_client_mod, "load_settings", lambda: cfg) + + +def test_wraps_client_when_observability_enabled( + monkeypatch: pytest.MonkeyPatch, +) -> None: + _reset_singleton(monkeypatch) + _patch_settings_with_observability(monkeypatch, enabled=True) + sentinel = object() + monkeypatch.setattr(_client_mod, "build_client", lambda cfg: sentinel) + + client = _client_mod.get_llm_client() + + assert isinstance(client, UsageRecordingClient) + assert client._inner is sentinel + + +def test_does_not_wrap_client_when_observability_disabled( + monkeypatch: pytest.MonkeyPatch, +) -> None: + _reset_singleton(monkeypatch) + _patch_settings_with_observability(monkeypatch, enabled=False) + sentinel = object() + monkeypatch.setattr(_client_mod, "build_client", lambda cfg: sentinel) + + assert _client_mod.get_llm_client() is sentinel diff --git a/tests/unit/test_component/test_llm/test_usage_client.py b/tests/unit/test_component/test_llm/test_usage_client.py new file mode 100644 index 0000000..6583bc8 --- /dev/null +++ b/tests/unit/test_component/test_llm/test_usage_client.py @@ -0,0 +1,101 @@ +"""``UsageRecordingClient`` — wraps an LLM client and records token usage +onto the active span, without altering the response or the call. + +The extractor calls happen inside an ``everos.extract`` generation span; the +wrapper's job is to surface ``response.usage`` there (Langfuse computes cost +from model + tokens). Everything else delegates to the wrapped client. +""" + +from __future__ import annotations + +from collections.abc import Iterator +from typing import Any + +import pytest +from everalgo.llm import ChatMessage, ChatResponse, Usage +from opentelemetry.sdk.trace.export import SimpleSpanProcessor +from opentelemetry.sdk.trace.export.in_memory_span_exporter import ( + InMemorySpanExporter, +) + +from everos.component.llm._usage_client import UsageRecordingClient +from everos.config.settings import ObservabilitySettings +from everos.core.observability.tracing import ( + force_flush, + init_tracing, + memory_span, + shutdown_tracing, +) + + +class _FakeInner: + """Minimal stand-in LLM client that returns a preset response.""" + + def __init__(self, response: ChatResponse | None) -> None: + self._response = response + self.calls: list[Any] = [] + self.some_attr = 42 + + async def chat(self, messages: list[ChatMessage], **kwargs: Any) -> ChatResponse: + self.calls.append((messages, kwargs)) + assert self._response is not None + return self._response + + +@pytest.fixture(autouse=True) +def _reset() -> Iterator[None]: + shutdown_tracing() + yield + shutdown_tracing() + + +@pytest.fixture +def captured() -> Iterator[InMemorySpanExporter]: + exporter = InMemorySpanExporter() + init_tracing( + ObservabilitySettings(enabled=True, endpoint="http://collector.invalid"), + span_processor=SimpleSpanProcessor(exporter), + ) + yield exporter + + +async def test_records_usage_to_active_span(captured: InMemorySpanExporter) -> None: + resp = ChatResponse( + content="x", model="gpt-x", usage=Usage(prompt_tokens=7, completion_tokens=3) + ) + client = UsageRecordingClient(_FakeInner(resp)) + with memory_span("everos.extract", observation_type="generation"): + out = await client.chat([ChatMessage(role="user", content="hi")], model="gpt-x") + assert out is resp + force_flush() + attrs = captured.get_finished_spans()[0].attributes + assert attrs["gen_ai.request.model"] == "gpt-x" + assert attrs["gen_ai.usage.input_tokens"] == 7 + assert attrs["gen_ai.usage.output_tokens"] == 3 + + +async def test_passes_through_when_usage_absent( + captured: InMemorySpanExporter, +) -> None: + resp = ChatResponse(content="x", model="gpt-x", usage=None) + client = UsageRecordingClient(_FakeInner(resp)) + with memory_span("everos.extract", observation_type="generation"): + out = await client.chat([ChatMessage(role="user", content="hi")]) + assert out is resp + force_flush() + attrs = captured.get_finished_spans()[0].attributes + assert attrs["gen_ai.request.model"] == "gpt-x" + assert "gen_ai.usage.input_tokens" not in attrs + + +async def test_delegates_unknown_attributes() -> None: + client = UsageRecordingClient(_FakeInner(None)) + assert client.some_attr == 42 + + +async def test_chat_forwards_kwargs_to_inner() -> None: + resp = ChatResponse(content="x", model="m") + inner = _FakeInner(resp) + client = UsageRecordingClient(inner) + await client.chat([ChatMessage(role="user", content="hi")], temperature=0.5) + assert inner.calls[0][1]["temperature"] == 0.5