feat(observability): capture LLM + embedding token usage
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) <noreply@anthropic.com>
This commit is contained in:
parent
ab5cf0447e
commit
bd26f5a80f
|
|
@ -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]
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
|
@ -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
|
||||
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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
|
||||
Loading…
Reference in New Issue