diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index ac3cad9..08835d6 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -20,7 +20,7 @@ jobs: with: python-version: ${{ matrix.python-version }} - - run: pip install -e ".[dev,http,apprise]" + - run: pip install -e ".[dev,http,apprise,otel]" - run: ruff check . diff --git a/CHANGELOG.md b/CHANGELOG.md index f221763..34b3435 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -4,6 +4,32 @@ All notable changes to this project are documented in this file. ## Unreleased +- LLM calls and OpenTelemetry export: + - `logger.llm_call(model=, tokens_in=, tokens_out=, cost_usd=, latency_ms=, + finish_reason=)` records an LLM call with its numbers in the record's + first-class `llm` block. A value that breaks the contract is dropped with a + warning instead of raising. + - `OTLPTransport` (`pip install logquill[otel]`) exports agent tracing as real + OpenTelemetry spans through the SDK: `span()` blocks, tool `.action()`s and + LLM calls, with the ids, parents and timings the records carry, so a + collector shows exactly the tree that was logged. Spans are named and + attributed per the OpenTelemetry GenAI conventions (`invoke_agent`, + `execute_tool`, `chat`, with token counts, model and finish reason). + - `OTelLogsTransport` sends every record as an OTLP log record, with the same + trace and span ids so logs and spans join up. + - The convention names live in one file, `logquill/semconv.py`, pinned to a + named release (semantic-conventions 1.44.0), because the GenAI conventions + are still experimental and have already renamed attributes. + `OTEL_SEMCONV_STABILITY_OPT_IN` is honored, `semconv_version="legacy"` picks + the older names, and old and new names are never emitted together. Prompt + and completion text is opt-in only. + - `span(name, capture_state=...)` records what changed during a block as + `meta.state_diff`, and a tool `.action()` that is reopened before it + succeeded gets `meta.retry_count` automatically. + - The record contract gained optional `meta` fields for this: `tool`, + `tool_call_id`, `provider`, `agent_name`, `agent_id`, `response_model`, + `operation`, and opt-in `input_messages` / `output_messages`. + - Type-checking no longer depends on whether OpenTelemetry is installed. - **Breaking: the record shape and the Python floor changed.** See [MIGRATING.md](MIGRATING.md). - logquill now requires **Python 3.10 or newer** and is tested on 3.10–3.14. diff --git a/MIGRATING.md b/MIGRATING.md index 7a3d89b..3bfa930 100644 --- a/MIGRATING.md +++ b/MIGRATING.md @@ -47,8 +47,10 @@ LogQuill record, or one written by a newer major version. ## 3. New optional fields -These are defined now so both languages agree on them. Nothing emits them on its -own yet, and no record has them unless you add them. +These are defined so both languages agree on them. `logger.llm_call()` writes +the `llm` block, `span(capture_state=...)` writes `state_diff`, and repeated tool +calls get `retry_count` automatically; nothing else writes them, and no record +has them unless you use those. | Field | Where | Type | |---|---|---| @@ -58,6 +60,9 @@ own yet, and no record has them unless you add them. | `meta.retry_count` | `meta` | integer ≥ 0 | | `meta.state_diff` | `meta` | object | | `meta.mcp.server`, `meta.mcp.tool` | `meta` | string | +| `meta.tool`, `meta.tool_call_id`, `meta.provider`, `meta.agent_name`, `meta.agent_id`, `meta.response_model` | `meta` | string | +| `meta.operation` | `meta` | `chat`, `text_completion` or `invoke_agent` | +| `meta.input_messages`, `meta.output_messages` | `meta` | array (opt-in prompt/completion content) | `llm` is its own block, not part of `meta`, because cost and latency dashboards need stable names for it. A record that isn't an LLM call has no diff --git a/README.md b/README.md index 036cc95..ec2634b 100644 --- a/README.md +++ b/README.md @@ -27,6 +27,7 @@ for what's landed so far. - **Config from file/env** — `load_config(dict)`, `logger_from_file(path)` (JSON/YAML), `logger_from_env()` build a `Logger` from one config shape — see [Config](#config) - **Plugin pipeline** — `ContextPlugin`, `RedactPlugin` (by key), `PIIRedactPlugin` (by pattern), `SamplingPlugin` (with tail-based elevation), `TamperEvidentPlugin` (hash-chained logs), `TraceContextPlugin` (cross-service trace correlation), and `AlertingPlugin` (`SlackAlertPlugin`/`PagerDutyAlertPlugin`/`EmailAlertPlugin`, deduplicated) out of the box; a broken plugin can't crash logging; `.use()` also accepts a plain function, no subclassing required (see [Plugins](#plugins)) - **Agentic & harness tracing** — `.child()` loggers, `RunPlugin`, `.thought()/.action()/.observation()/.decision()`, `with agent_log.span(...)`, and framework adapters — `LangChainAdapter` (`pip install logquill[langchain]`), `LangGraphAdapter` (`pip install logquill[langgraph]`, adds checkpoint interrupt/resume events on top), `CrewAIAdapter` (`pip install logquill[crewai]`), `LlamaIndexAdapter` (`pip install logquill[llamaindex]`), and `AutoGenAdapter` (`pip install logquill[autogen]`) — see [Agentic & harness tracing](#agentic--harness-tracing) +- **LLM calls & OpenTelemetry** — `logger.llm_call(model=, tokens_in=, tokens_out=, cost_usd=, ...)` records a first-class `llm` block, `span(capture_state=...)` records what state changed, repeated tool calls get `meta.retry_count` automatically, and `OTLPTransport` (`pip install logquill[otel]`) exports it all as real OpenTelemetry spans named and attributed per the GenAI semantic conventions; `OTelLogsTransport` sends records as OTLP logs — see [LLM calls & OpenTelemetry](#llm-calls--opentelemetry) - **Non-blocking async dispatch** — `Logger(async_dispatch=True)` moves transport writes onto a background thread with a bounded queue and a configurable backpressure policy (`drop_oldest`/`drop_newest`/`block`); `flush()`/`flush_async()` and a `with_lambda`/`with_cloud_function`/`with_azure_function` decorator make serverless shutdown safe — see [Async dispatch & serverless safety](#async-dispatch--serverless-safety) - **Zero required runtime dependencies** — stdlib only; `aiohttp` (`logquill[http]`) is opt-in, for keep-alive HTTP delivery - **Typed throughout** — `mypy --strict` clean on the public API @@ -817,6 +818,136 @@ event classes; that's a real divergence, not just a detail, so it needs its own adapter rather than reusing this one. `autogen-core` is never imported unless you import `logquill.adapters.autogen` yourself. +## LLM calls & OpenTelemetry + +**Recording an LLM call.** `logger.llm_call()` writes one record with the LLM +call's numbers in their own `llm` block — fixed names and types, so a cost or +latency dashboard can rely on them: + +```python +from logquill import Logger + +logger = Logger("app.agent") + +record = logger.llm_call( + "chat", + model="example-model", + tokens_in=1200, + tokens_out=340, + cost_usd=0.0123, + latency_ms=2150.5, + finish_reason="stop", + provider="anthropic", # extra keywords go in `meta` +) + +assert record["llm"]["tokens_in"] == 1200 +assert record["meta"] == {"kind": "action", "provider": "anthropic"} +``` + +Leave out what you don't have; a value of the wrong type (a negative count, a +string where a number belongs) is dropped with a warning rather than raising. +Don't put prompt or completion text in `meta` unless you mean every transport +to receive it. + +**What changed during a step.** `span(name, capture_state=...)` calls your +function on entering and leaving the block and records what differs as +`meta.state_diff` — only the keys that changed, deep-copied so in-place +mutation is seen, and omitted if nothing changed: + +```python +from logquill import CollectingTransport, Logger + +sink = CollectingTransport() +logger = Logger("app.agent", transports=[sink]) +state = {"items": 1, "user": "ada"} + +with logger.span("add_item", capture_state=lambda: state): + state["items"] = 2 + +assert sink.records[0]["meta"]["state_diff"] == {"before": {"items": 1}, "after": {"items": 2}} +``` + +`state_diff` lands in `meta`, so `PIIRedactPlugin` sees it; `RedactPlugin` only +matches top-level keys, so don't capture secrets. + +**Retries.** An `.action()` that names its tool (`tool="search"`) is tracked: +if the same call — same tool, same enclosing span, same `tool_call_id` if you +give one — is reopened before it succeeded, the record gets `meta.retry_count` +(1, 2, 3, ...). A successful `.observation(tool=...)` ends the chain, so a loop +that legitimately calls one tool many times isn't reported as retries: + +```python +from logquill import Logger + +logger = Logger("app.agent") + +first = logger.action("look it up", tool="search") +logger.observation("timed out", tool="search", error="TimeoutError: slow") +second = logger.action("look it up", tool="search") + +assert "retry_count" not in first["meta"] +assert second["meta"]["retry_count"] == 1 +``` + +**Exporting to OpenTelemetry.** `pip install logquill[otel]`, then attach +`OTLPTransport`. It turns the records that describe work with a start and an +end into real spans, with the ids, parents and timings the records carry: + +- a `span()` marked `operation="invoke_agent"` (or given an `agent_name`) is + `invoke_agent {agent}`; any other span keeps its own name, +- an `.action()` naming its `tool` is `execute_tool {tool}`, +- a record with an `llm` block is `chat {model}` with the model, token counts + and finish reason as the standard GenAI attributes. + +```python +from opentelemetry.sdk.trace.export.in_memory_span_exporter import InMemorySpanExporter + +from logquill import Logger, OTLPTransport, RunPlugin + +exporter = InMemorySpanExporter() # in real use: leave `span_exporter` out and set `endpoint` +logger = Logger( + "app.agent", + transports=[OTLPTransport(span_exporter=exporter, processor="simple")], + plugins=[RunPlugin()], +) + +with logger.span("run", operation="invoke_agent", agent_name="planner"): + logger.llm_call("chat", model="example-model", tokens_in=1200, tokens_out=340, provider="anthropic") + logger.action("look it up", tool="search", duration_ms=40) +logger.close() + +spans = {span.name: span for span in exporter.get_finished_spans()} +assert set(spans) == {"invoke_agent planner", "chat example-model", "execute_tool search"} +chat = spans["chat example-model"] +assert chat.attributes["gen_ai.usage.input_tokens"] == 1200 +assert chat.parent.span_id == spans["invoke_agent planner"].context.span_id +``` + +To send to a collector, give it an endpoint instead: +`OTLPTransport(endpoint="http://localhost:4318/v1/traces", service_name="my-agent")`. +Export is batched off the calling thread. Records that aren't spans are +ignored here; `OTelLogsTransport(endpoint=".../v1/logs")` sends every record +as an OTLP log record — level as severity, `meta` as attributes, and the same +trace and span ids, so a collector can join a log line to its span. + +Three things worth knowing: + +- **The attribute names are pinned to one release of the GenAI conventions** + (semantic-conventions 1.44.0) and live in a single file, + `logquill/semconv.py`. Those conventions are still marked *Development* + upstream and have already renamed attributes (`gen_ai.system` became + `gen_ai.provider.name`), so a future release is an edit to that file. + `semconv_version="legacy"` selects the older provider and token-count names + for a backend that hasn't caught up; `OTEL_SEMCONV_STABILITY_OPT_IN=gen_ai_latest_experimental` + selects the latest. Old and new names are never emitted together. +- **Prompt and completion text is never exported by default.** It's only + included — from `meta.input_messages` / `meta.output_messages`, as JSON — if + you pass `capture_content=True` or set + `OTEL_INSTRUMENTATION_GENAI_CAPTURE_MESSAGE_CONTENT=true`. +- Give records a `run_id` (`RunPlugin`) or a `trace_id` (`TraceContextPlugin`) + and one run is one trace. Without either, spans that name a parent still share + a trace with their siblings, but not with their grandparents. + ## Async dispatch & serverless safety By default, every log call dispatches to its transports synchronously — a diff --git a/logquill/__init__.py b/logquill/__init__.py index 6178484..09cf5a3 100644 --- a/logquill/__init__.py +++ b/logquill/__init__.py @@ -39,6 +39,8 @@ from logquill.transports.nosql.dynamodb_transport import DynamoDBTransport from logquill.transports.nosql.mongodb_transport import MongoDBTransport from logquill.transports.nosql.redis_transport import RedisTransport +from logquill.transports.otel.logs_transport import OTelLogsTransport +from logquill.transports.otel.otlp_transport import OTLPTransport from logquill.transports.queue.base_queue_transport import BaseQueueTransport from logquill.transports.queue.kafka_transport import KafkaTransport from logquill.transports.queue.pubsub_transport import PubSubTransport @@ -86,6 +88,8 @@ "MongoDBTransport", "MySQLTransport", "NewRelicTransport", + "OTLPTransport", + "OTelLogsTransport", "OptLogger", "PIIRedactPlugin", "PagerDutyAlertPlugin", diff --git a/logquill/logger.py b/logquill/logger.py index 725ca77..08e74e5 100644 --- a/logquill/logger.py +++ b/logquill/logger.py @@ -2,7 +2,7 @@ import contextlib import logging -from typing import Any +from typing import Any, Callable from logquill import shutdown from logquill.context import current_context @@ -11,7 +11,8 @@ from logquill.opt import OptLogger, caller_info, resolve_lazy from logquill.plugins.context_plugin import ContextPlugin from logquill.plugins.plugin import FunctionPlugin, MiddlewareFunc, Plugin -from logquill.records import LogRecord, create_record +from logquill.records import LLMBlock, LogRecord, build_llm_block, create_record +from logquill.retry import RetryTracker from logquill.span import SpanContext, current_span_id from logquill.toggle import is_enabled from logquill.transports.transport import Transport @@ -72,6 +73,7 @@ def __init__( if async_dispatch else None ) + self._retries = RetryTracker() if flush_at_exit: shutdown.register(self) @@ -166,6 +168,24 @@ def opt(self, *, lazy: bool = False, depth: int | None = None) -> OptLogger: """ return OptLogger(self, lazy=lazy, depth=depth) + def _track_retries(self, level: Level, meta: dict[str, Any]) -> None: + """Stamp `meta.retry_count` on a tool `.action()` that reopens a call + which hasn't succeeded yet, and end the chain on a successful + `.observation()` for it. Only records naming their tool in `meta.tool` + take part.""" + tool = meta.get("tool") + if not isinstance(tool, str) or not tool: + return + call_id = meta.get("tool_call_id") + call_id = call_id if isinstance(call_id, str) else None + kind = meta.get("kind") + if kind == "action": + count = self._retries.opened(current_span_id(), tool, call_id) + if count > 0: + meta.setdefault("retry_count", count) + elif kind == "observation" and level < Level.ERROR and "error" not in meta: + self._retries.succeeded(current_span_id(), tool, call_id) + def _local_redactor(self) -> LocalRedactor: """Chain every plugin's `redact_local` into one function for `diagnose` mode. A plugin whose hook raises fails closed: the value @@ -216,6 +236,7 @@ def _log( *, lazy: bool = False, depth: int | None = None, + llm: LLMBlock | None = None, ) -> LogRecord | None: if level < self._level or not is_enabled(self.name): return None @@ -223,6 +244,8 @@ def _log( if lazy: meta = resolve_lazy(meta) + self._track_retries(level, meta) + if depth is not None: caller = caller_info(depth) if caller is not None: @@ -250,7 +273,7 @@ def _log( if stack is not None: meta["stack"] = stack - record = create_record(level=level, logger=self.name, message=message, meta=meta) + record = create_record(level=level, logger=self.name, message=message, meta=meta, llm=llm) bound_context = current_context() if bound_context: @@ -347,6 +370,40 @@ def decision(self, message: str, /, **meta: Any) -> LogRecord | None: decision for a step or run, for harness/agentic tracing.""" return self._log(Level.INFO, message, {"kind": "decision", **meta}) + def llm_call( + self, + message: str = "llm_call", + /, + *, + model: str | None = None, + tokens_in: int | None = None, + tokens_out: int | None = None, + cost_usd: float | None = None, + latency_ms: float | None = None, + finish_reason: str | None = None, + **meta: Any, + ) -> LogRecord | None: + """Log one LLM call as an `.action()` carrying the record's first-class + `llm` block (`model`, `tokens_in`, `tokens_out`, `cost_usd`, + `latency_ms`, `finish_reason`). Those fields are what cost and latency + dashboards read, and what `OTLPTransport` exports as the standard token + and model attributes. Any you leave out are simply absent; one of the + wrong type is dropped with a warning rather than raising. + + Extra keyword arguments go in `meta` as usual — e.g. `provider="openai"`. + Don't put prompt or completion text in `meta` unless you mean every + transport to receive it. + """ + llm = build_llm_block( + model=model, + tokens_in=tokens_in, + tokens_out=tokens_out, + cost_usd=cost_usd, + latency_ms=latency_ms, + finish_reason=finish_reason, + ) + return self._log(Level.INFO, message, {"kind": "action", **meta}, llm=llm) + def span( self, name: str, @@ -354,6 +411,7 @@ def span( *, span_id: str | None = None, parent_span_id: str | None = None, + capture_state: Callable[[], Any] | None = None, **meta: Any, ) -> SpanContext: """`with agent_log.span("call_llm"):` — on exit, emits one record @@ -368,5 +426,21 @@ def span( `span_id`/`parent_span_id` normally auto-generate/auto-nest; pass them explicitly to adopt an id handed in from elsewhere (see `logquill.adapters.langchain.LangChainAdapter` for an example). + + `capture_state=lambda: {...}` is called on entering and on leaving the + block, and what changed between the two is recorded as + `meta.state_diff` (`{"before": ..., "after": ...}`, only the keys that + changed when both are dicts; omitted if nothing did). The values are + deep-copied, so an in-place mutation is seen; a state that can't be + copied, or a callable that raises, just means no `state_diff`. It's + stored in `meta`, so a redaction plugin that recurses (`PIIRedactPlugin`) + sees it, but `RedactPlugin` only matches top-level keys. """ - return SpanContext(self, name, span_id=span_id, parent_span_id=parent_span_id, **meta) + return SpanContext( + self, + name, + span_id=span_id, + parent_span_id=parent_span_id, + capture_state=capture_state, + **meta, + ) diff --git a/logquill/opt.py b/logquill/opt.py index 145d9d0..4bf27f4 100644 --- a/logquill/opt.py +++ b/logquill/opt.py @@ -4,7 +4,7 @@ from typing import TYPE_CHECKING, Any from logquill.levels import Level -from logquill.records import LogRecord +from logquill.records import LogRecord, build_llm_block if TYPE_CHECKING: from logquill.logger import Logger @@ -123,3 +123,34 @@ def decision(self, message: str, /, **meta: Any) -> LogRecord | None: return self._logger._log( Level.INFO, message, {"kind": "decision", **meta}, lazy=self._lazy, depth=self._depth ) + + def llm_call( + self, + message: str = "llm_call", + /, + *, + model: str | None = None, + tokens_in: int | None = None, + tokens_out: int | None = None, + cost_usd: float | None = None, + latency_ms: float | None = None, + finish_reason: str | None = None, + **meta: Any, + ) -> LogRecord | None: + """`Logger.llm_call` with this view's options.""" + llm = build_llm_block( + model=model, + tokens_in=tokens_in, + tokens_out=tokens_out, + cost_usd=cost_usd, + latency_ms=latency_ms, + finish_reason=finish_reason, + ) + return self._logger._log( + Level.INFO, + message, + {"kind": "action", **meta}, + lazy=self._lazy, + depth=self._depth, + llm=llm, + ) diff --git a/logquill/plugins/trace_context_plugin.py b/logquill/plugins/trace_context_plugin.py index 91e0064..675f656 100644 --- a/logquill/plugins/trace_context_plugin.py +++ b/logquill/plugins/trace_context_plugin.py @@ -123,7 +123,7 @@ def _resolve_trace_id(self) -> str: @staticmethod def _from_active_otel_span() -> str | None: try: - from opentelemetry import trace as otel_trace # type: ignore[import-not-found] + from opentelemetry import trace as otel_trace except ImportError: return None span = otel_trace.get_current_span() diff --git a/logquill/records.py b/logquill/records.py index 0c8bbfa..01848f4 100644 --- a/logquill/records.py +++ b/logquill/records.py @@ -1,5 +1,6 @@ from __future__ import annotations +import logging import re from collections.abc import Mapping from datetime import datetime, timezone @@ -7,6 +8,8 @@ from logquill.levels import Level +_logger = logging.getLogger("logquill") + #: The record shape this version writes. Shared with logquill-js: a record #: carries it as `schema_version` so a reader can tell which shape it has. SCHEMA_VERSION = "2.0" @@ -131,3 +134,46 @@ def parse_record(raw: Mapping[str, Any]) -> LogRecord: raise ValueError(f"record field 'llm' must be an object, got {type(llm).__name__}") return cast(LogRecord, record) + + +_LLM_INT_FIELDS = ("tokens_in", "tokens_out") +_LLM_NUMBER_FIELDS = ("cost_usd", "latency_ms") +_LLM_STR_FIELDS = ("model", "finish_reason") + + +def build_llm_block(**fields: Any) -> LLMBlock | None: + """Assemble an `LLMBlock` from keyword fields, leaving out anything that is + `None` and anything that doesn't fit the contract (a negative count, a + string where a number belongs) — with a warning naming the field, since a + log call must not raise over a bad value, and a record that breaks the + schema is worse than one missing a field. Returns `None` if nothing is + left, so a record with no usable LLM data has no `llm` key at all.""" + block: dict[str, Any] = {} + for name, value in fields.items(): + if value is None: + continue + if name in _LLM_INT_FIELDS: + valid = isinstance(value, int) and not isinstance(value, bool) and value >= 0 + elif name in _LLM_NUMBER_FIELDS: + valid = isinstance(value, (int, float)) and not isinstance(value, bool) and value >= 0 + elif name in _LLM_STR_FIELDS: + valid = isinstance(value, str) + else: + valid = False + if valid: + block[name] = value + else: + _logger.warning( + "llm_call: ignoring %s=%r — expected %s", name, value, _expected_llm_type(name) + ) + return cast(LLMBlock, block) if block else None + + +def _expected_llm_type(name: str) -> str: + if name in _LLM_INT_FIELDS: + return "a non-negative integer" + if name in _LLM_NUMBER_FIELDS: + return "a non-negative number" + if name in _LLM_STR_FIELDS: + return "a string" + return "one of model, tokens_in, tokens_out, cost_usd, latency_ms, finish_reason" diff --git a/logquill/retry.py b/logquill/retry.py new file mode 100644 index 0000000..d29528c --- /dev/null +++ b/logquill/retry.py @@ -0,0 +1,47 @@ +from __future__ import annotations + +import threading +from collections import OrderedDict + + +class RetryTracker: + """Counts how many times the same tool call has been reopened, so + `meta.retry_count` can be stamped on `.action()` records automatically — a + retry loop is one of the commonest agent failure modes and is otherwise + invisible in the logs. + + A tool call is identified by `(enclosing span, tool name, tool call id)`. + The first `.action()` for it is attempt 0; each further `.action()` for the + same call before it succeeds is a retry, so it comes back as 1, 2, 3, ... A + successful `.observation()` for the call ends the chain: the next action + for that tool starts again at 0, which is what stops a loop that + legitimately calls one tool ten times from being reported as ten retries. + + Bounded: at most `max_entries` calls are remembered, the oldest forgotten + first, so a long-running process can't grow it without limit. + """ + + def __init__(self, max_entries: int = 1024) -> None: + """`max_entries` is how many distinct in-flight tool calls to remember.""" + if max_entries < 1: + raise ValueError(f"max_entries must be >= 1, got {max_entries}") + self.max_entries = max_entries + self._attempts: OrderedDict[tuple[str, str, str], int] = OrderedDict() + self._lock = threading.Lock() + + def opened(self, span_id: str | None, tool: str, call_id: str | None = None) -> int: + """Record another action for this call; returns its retry count (0 for + the first attempt).""" + key = (span_id or "", tool, call_id or "") + with self._lock: + count = self._attempts.get(key, -1) + 1 + self._attempts[key] = count + self._attempts.move_to_end(key) + while len(self._attempts) > self.max_entries: + self._attempts.popitem(last=False) + return count + + def succeeded(self, span_id: str | None, tool: str, call_id: str | None = None) -> None: + """The call finished successfully; forget it.""" + with self._lock: + self._attempts.pop((span_id or "", tool, call_id or ""), None) diff --git a/logquill/semconv.py b/logquill/semconv.py new file mode 100644 index 0000000..8d95250 --- /dev/null +++ b/logquill/semconv.py @@ -0,0 +1,376 @@ +"""LogQuill records -> OpenTelemetry GenAI semantic-convention names. + +This is the only module in logquill that contains GenAI attribute names, span +names or operation names. Everything else asks it: a transport passes a +record in and gets a span name, a span kind and a dict of attributes back. + +The GenAI conventions are still marked *Development* upstream, and the names +have already changed more than once (`gen_ai.system` became +`gen_ai.provider.name`; `prompt_tokens`/`completion_tokens` became +`input_tokens`/`output_tokens`). So the names live in one table per convention +generation, pinned to a named release, and moving to a newer release is an edit +to this file and nothing else. It never imports OpenTelemetry, so it can be +tested — and read — on its own. + +Pinned to the GenAI conventions as published with semantic-conventions +`CONVENTION_VERSION`. +""" + +from __future__ import annotations + +import json +import logging +import os +from collections.abc import Mapping +from dataclasses import dataclass, field +from datetime import datetime +from typing import Any, Literal + +from logquill.records import LogRecord + +_logger = logging.getLogger("logquill") + +#: The semantic-conventions release the `latest` names below come from. +CONVENTION_VERSION = "1.44.0" + +#: The value of `OTEL_SEMCONV_STABILITY_OPT_IN` that asks for the latest +#: experimental GenAI conventions. +OPT_IN_ENV = "OTEL_SEMCONV_STABILITY_OPT_IN" +OPT_IN_LATEST = "gen_ai_latest_experimental" + +#: `OTEL_INSTRUMENTATION_GENAI_CAPTURE_MESSAGE_CONTENT=true` turns on prompt and +#: completion capture, the same switch other GenAI instrumentation uses. +CAPTURE_CONTENT_ENV = "OTEL_INSTRUMENTATION_GENAI_CAPTURE_MESSAGE_CONTENT" + +#: Attributes LogQuill adds that no convention defines. Namespaced so they can't +#: collide with anything upstream adds later. +_EXT = "logquill." + +SpanKindName = Literal["client", "internal"] + + +@dataclass(frozen=True) +class Convention: + """One generation of GenAI attribute names. `names` maps LogQuill's own + keys (the ones used by the functions below) to the attribute a backend + receives.""" + + version: str + names: Mapping[str, str] = field(repr=False) + + +_LATEST = Convention( + version=CONVENTION_VERSION, + names={ + "operation": "gen_ai.operation.name", + "provider": "gen_ai.provider.name", + "request_model": "gen_ai.request.model", + "response_model": "gen_ai.response.model", + "input_tokens": "gen_ai.usage.input_tokens", + "output_tokens": "gen_ai.usage.output_tokens", + "finish_reasons": "gen_ai.response.finish_reasons", + "agent_name": "gen_ai.agent.name", + "agent_id": "gen_ai.agent.id", + "conversation_id": "gen_ai.conversation.id", + "tool_name": "gen_ai.tool.name", + "tool_call_id": "gen_ai.tool.call.id", + "input_messages": "gen_ai.input.messages", + "output_messages": "gen_ai.output.messages", + }, +) + +#: The names in use before the provider and token-count renames — for a backend +#: or dashboard that hasn't caught up. Only ever chosen explicitly. +_LEGACY = Convention( + version="legacy", + names={ + **_LATEST.names, + "provider": "gen_ai.system", + "input_tokens": "gen_ai.usage.prompt_tokens", + "output_tokens": "gen_ai.usage.completion_tokens", + }, +) + +_CONVENTIONS: dict[str, Convention] = {"latest": _LATEST, "legacy": _LEGACY} + +#: What is used when neither the caller nor the environment says. Kept separate +#: from `_LATEST` on purpose: when a newer generation is added, the default can +#: stay on the one people already export while `gen_ai_latest_experimental` +#: opts in to the new one, which is how OpenTelemetry asks for this to work. +_DEFAULT = _LATEST + +_ERROR_TYPE = "error.type" + +# operation names +INVOKE_AGENT = "invoke_agent" +EXECUTE_TOOL = "execute_tool" +CHAT = "chat" +TEXT_COMPLETION = "text_completion" + +#: `provider` when nothing says which one — the attribute is required upstream. +UNKNOWN_PROVIDER = "unknown" + +_warned_opt_in: set[str] = set() + + +def resolve_convention( + explicit: str | None = None, environ: Mapping[str, str] | None = None +) -> Convention: + """Choose the naming generation. An explicit `explicit` ("latest" or + "legacy") wins; otherwise `OTEL_SEMCONV_STABILITY_OPT_IN` is read (a + comma-separated list, of which `gen_ai_latest_experimental` selects the + latest names); otherwise the default is used. + + Old and new names are never emitted together: an opt-in value asking for + that (anything ending `/dup`) is ignored with a one-time warning, and + exactly one generation is used. + """ + if explicit is not None: + try: + return _CONVENTIONS[explicit] + except KeyError: + raise ValueError( + f"semconv must be one of {', '.join(sorted(_CONVENTIONS))}, got {explicit!r}" + ) from None + + env = os.environ if environ is None else environ + tokens = [token.strip() for token in env.get(OPT_IN_ENV, "").split(",") if token.strip()] + for token in tokens: + if token.endswith("/dup") and token not in _warned_opt_in: + _warned_opt_in.add(token) + _logger.warning( + "%s=%s: LogQuill never emits old and new GenAI attribute names together — " + "using one generation only", + OPT_IN_ENV, + token, + ) + return _CONVENTIONS["latest"] if OPT_IN_LATEST in tokens else _DEFAULT + + +def capture_content_enabled(environ: Mapping[str, str] | None = None) -> bool: + """Whether the environment switches on prompt/completion capture.""" + env = os.environ if environ is None else environ + return env.get(CAPTURE_CONTENT_ENV, "").strip().lower() == "true" + + +@dataclass(frozen=True) +class SpanSpec: + """A span a record describes: its name, kind, attributes, and — because a + record is written when the work *ends* — how long it lasted.""" + + name: str + kind: SpanKindName + attributes: dict[str, Any] + duration_ms: float | None = None + is_error: bool = False + events: tuple[tuple[str, dict[str, Any]], ...] = () + + +def _meta(record: LogRecord) -> Mapping[str, Any]: + meta = record.get("meta") + return meta if isinstance(meta, dict) else {} + + +def _string(value: Any) -> str | None: + return value if isinstance(value, str) and value else None + + +def _number(value: Any) -> float | None: + if isinstance(value, bool) or not isinstance(value, (int, float)): + return None + return float(value) + + +def _json(value: Any) -> str: + try: + return json.dumps(value, default=str, separators=(",", ":")) + except (TypeError, ValueError): + return repr(value) + + +def _error_type(meta: Mapping[str, Any]) -> str | None: + error = _string(meta.get("error")) + if error is None: + return None + head = error.split(":", 1)[0].strip() + return head if head and " " not in head else "_OTHER" + + +def is_span_record(record: LogRecord) -> bool: + """A record `Logger.span()` wrote on exit.""" + meta = _meta(record) + return meta.get("kind") == "span" and isinstance(meta.get("span_id"), str) + + +def is_llm_record(record: LogRecord) -> bool: + """A record carrying an `llm` block.""" + return isinstance(record.get("llm"), dict) + + +def is_tool_record(record: LogRecord) -> bool: + """An `.action()` naming the tool it called, in `meta.tool`.""" + meta = _meta(record) + return meta.get("kind") == "action" and _string(meta.get("tool")) is not None + + +def _common(convention: Convention, meta: Mapping[str, Any], attributes: dict[str, Any]) -> None: + names = convention.names + thread_id = _string(meta.get("thread_id")) + if thread_id: + attributes[names["conversation_id"]] = thread_id + agent_name = _string(meta.get("agent_name")) + if agent_name: + attributes[names["agent_name"]] = agent_name + if run_id := _string(meta.get("run_id")): + attributes[_EXT + "run_id"] = run_id + if node := _string(meta.get("node_name")): + attributes[_EXT + "node_name"] = node + retry = meta.get("retry_count") + if isinstance(retry, int) and not isinstance(retry, bool): + attributes[_EXT + "retry_count"] = retry + error_type = _error_type(meta) + if error_type: + attributes[_ERROR_TYPE] = error_type + + +def _events(meta: Mapping[str, Any]) -> tuple[tuple[str, dict[str, Any]], ...]: + diff = meta.get("state_diff") + if isinstance(diff, dict): + return ((_EXT + "state_diff", {_EXT + "state": _json(diff)}),) + return () + + +def llm_spec( + record: LogRecord, + convention: Convention, + *, + default_provider: str = UNKNOWN_PROVIDER, + capture_content: bool = False, +) -> SpanSpec: + """The inference span for a record with an `llm` block: `chat {model}`, or + `text_completion {model}` when `meta.operation` says so. Prompt and + completion text (`meta.input_messages`/`output_messages`) are added only if + `capture_content` is true.""" + names = convention.names + llm: Mapping[str, Any] = record.get("llm") or {} + meta = _meta(record) + + operation = TEXT_COMPLETION if meta.get("operation") == TEXT_COMPLETION else CHAT + model = _string(llm.get("model")) + attributes: dict[str, Any] = { + names["operation"]: operation, + names["provider"]: _string(meta.get("provider")) or default_provider, + } + if model: + attributes[names["request_model"]] = model + if response_model := _string(meta.get("response_model")): + attributes[names["response_model"]] = response_model + if isinstance(llm.get("tokens_in"), int): + attributes[names["input_tokens"]] = llm["tokens_in"] + if isinstance(llm.get("tokens_out"), int): + attributes[names["output_tokens"]] = llm["tokens_out"] + if finish_reason := _string(llm.get("finish_reason")): + attributes[names["finish_reasons"]] = [finish_reason] + if isinstance(llm.get("cost_usd"), (int, float)): + attributes[_EXT + "cost_usd"] = float(llm["cost_usd"]) + if capture_content: + for key in ("input_messages", "output_messages"): + if meta.get(key) is not None: + attributes[names[key]] = _json(meta[key]) + _common(convention, meta, attributes) + + return SpanSpec( + name=f"{operation} {model}" if model else operation, + kind="client", + attributes=attributes, + duration_ms=_number(llm.get("latency_ms")), + is_error=_error_type(meta) is not None or record.get("level") in ("ERROR", "FATAL"), + events=_events(meta), + ) + + +def tool_spec(record: LogRecord, convention: Convention) -> SpanSpec: + """The `execute_tool {tool}` span for an `.action()` that names its tool.""" + names = convention.names + meta = _meta(record) + tool = str(meta["tool"]) + attributes: dict[str, Any] = {names["operation"]: EXECUTE_TOOL, names["tool_name"]: tool} + if call_id := _string(meta.get("tool_call_id")): + attributes[names["tool_call_id"]] = call_id + _common(convention, meta, attributes) + return SpanSpec( + name=f"{EXECUTE_TOOL} {tool}", + kind="internal", + attributes=attributes, + duration_ms=_number(meta.get("duration_ms")), + is_error=_error_type(meta) is not None, + events=_events(meta), + ) + + +def span_spec( + record: LogRecord, + convention: Convention, + *, + default_provider: str = UNKNOWN_PROVIDER, +) -> SpanSpec: + """The span for a `Logger.span()` record. A span that says it is an agent + run — `meta.operation="invoke_agent"`, or it names an `agent_name` — is an + `invoke_agent {agent}` span; any other span keeps its own name and carries + no GenAI operation. (A span isn't assumed to be an agent run just because + it is the outermost one: that would rename ordinary spans.)""" + names = convention.names + meta = _meta(record) + attributes: dict[str, Any] = {} + _common(convention, meta, attributes) + + is_agent_run = meta.get("operation") == INVOKE_AGENT or _string(meta.get("agent_name")) + if is_agent_run: + agent = _string(meta.get("agent_name")) or str(record["logger"]) + attributes[names["operation"]] = INVOKE_AGENT + attributes[names["provider"]] = _string(meta.get("provider")) or default_provider + attributes[names["agent_name"]] = agent + if agent_id := _string(meta.get("agent_id")): + attributes[names["agent_id"]] = agent_id + name = f"{INVOKE_AGENT} {agent}" + kind: SpanKindName = "client" + else: + name = str(record["message"]) + kind = "internal" + + return SpanSpec( + name=name, + kind=kind, + attributes=attributes, + duration_ms=_number(meta.get("duration_ms")), + is_error=record.get("level") in ("ERROR", "FATAL"), + events=_events(meta), + ) + + +def log_attributes(record: LogRecord, convention: Convention) -> dict[str, Any]: + """The GenAI attributes to put on a *log* record: the `llm` block's usage + and model, so a log line about an LLM call is queryable by the same names + as its span.""" + if not is_llm_record(record): + return {} + llm: Mapping[str, Any] = record["llm"] + names = convention.names + attributes: dict[str, Any] = {names["operation"]: CHAT} + if model := _string(llm.get("model")): + attributes[names["request_model"]] = model + if isinstance(llm.get("tokens_in"), int): + attributes[names["input_tokens"]] = llm["tokens_in"] + if isinstance(llm.get("tokens_out"), int): + attributes[names["output_tokens"]] = llm["tokens_out"] + if finish_reason := _string(llm.get("finish_reason")): + attributes[names["finish_reasons"]] = [finish_reason] + if isinstance(llm.get("cost_usd"), (int, float)): + attributes[_EXT + "cost_usd"] = float(llm["cost_usd"]) + return attributes + + +def epoch_ns(timestamp: str) -> int: + """A record's ISO 8601 timestamp as nanoseconds since the Unix epoch.""" + parsed = datetime.fromisoformat(timestamp.replace("Z", "+00:00")) + return int(parsed.timestamp() * 1_000_000_000) diff --git a/logquill/span.py b/logquill/span.py index 3beb44a..8dd0fd9 100644 --- a/logquill/span.py +++ b/logquill/span.py @@ -1,16 +1,20 @@ from __future__ import annotations +import copy +import logging import time import uuid from contextvars import ContextVar, Token from types import TracebackType -from typing import TYPE_CHECKING, Any +from typing import TYPE_CHECKING, Any, Callable from logquill.levels import Level if TYPE_CHECKING: from logquill.logger import Logger +_logger = logging.getLogger("logquill") + _current_span_id: ContextVar[str | None] = ContextVar("logquill_span_id", default=None) @@ -21,6 +25,43 @@ def current_span_id() -> str | None: return _current_span_id.get() +_UNCAPTURED = object() + + +def _snapshot(capture: Callable[[], Any]) -> Any: + """A deep copy of the captured state — a copy, because the state usually + mutates in place and a reference would show "after" for both. Returns + `_UNCAPTURED` if the callable raises or the value can't be copied, so a + broken snapshot never breaks the code being traced.""" + try: + return copy.deepcopy(capture()) + except Exception: + _logger.debug("span: capture_state failed; recording no state_diff", exc_info=True) + return _UNCAPTURED + + +def _state_diff(before: Any, after: Any) -> dict[str, Any] | None: + """What changed between two snapshots, or `None` if nothing did. For two + dicts, only the keys that differ (a key added shows only under `after`, one + removed only under `before`); otherwise the two whole values.""" + try: + if isinstance(before, dict) and isinstance(after, dict): + changed = [ + key + for key in before.keys() | after.keys() + if key not in before or key not in after or before[key] != after[key] + ] + if not changed: + return None + return { + "before": {key: before[key] for key in changed if key in before}, + "after": {key: after[key] for key in changed if key in after}, + } + return None if before == after else {"before": before, "after": after} + except Exception: + return None + + def new_span_id() -> str: """A 16-hex-char id, matching the shape of an OTel span id.""" return uuid.uuid4().hex[:16] @@ -55,9 +96,12 @@ def __init__( *, span_id: str | None = None, parent_span_id: str | None = None, + capture_state: Callable[[], Any] | None = None, **meta: Any, ) -> None: - """`span_id` defaults to a freshly generated id; `parent_span_id` + """`capture_state`, if given, is called on entering and on leaving the + block; what changed between the two lands in `meta.state_diff` (see + `Logger.span`). `span_id` defaults to a freshly generated id; `parent_span_id` overrides the auto-nesting that would otherwise come from any enclosing span active in this execution context — see the class docstring for why a caller would pass either explicitly.""" @@ -66,12 +110,16 @@ def __init__( self._meta = meta self._span_id = span_id or new_span_id() self._explicit_parent_span_id = parent_span_id + self._capture_state = capture_state + self._state_before: Any = _UNCAPTURED self._token: Token[str | None] | None = None self._start = 0.0 def __enter__(self) -> SpanContext: """Push this span's id as the current span for this execution context and start its duration timer.""" + if self._capture_state is not None: + self._state_before = _snapshot(self._capture_state) self._token = _current_span_id.set(self._span_id) self._start = time.monotonic() return self @@ -98,6 +146,12 @@ def __exit__( if self._explicit_parent_span_id is not None: meta["parent_span_id"] = self._explicit_parent_span_id meta.setdefault("kind", "span") + if self._capture_state is not None and self._state_before is not _UNCAPTURED: + state_after = _snapshot(self._capture_state) + if state_after is not _UNCAPTURED: + diff = _state_diff(self._state_before, state_after) + if diff is not None: + meta.setdefault("state_diff", diff) if exc_type is not None: meta["error"] = f"{exc_type.__name__}: {exc}" diff --git a/logquill/transports/otel/__init__.py b/logquill/transports/otel/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/logquill/transports/otel/_ids.py b/logquill/transports/otel/_ids.py new file mode 100644 index 0000000..edec9d7 --- /dev/null +++ b/logquill/transports/otel/_ids.py @@ -0,0 +1,44 @@ +from __future__ import annotations + +import hashlib +import re +from collections.abc import Mapping +from typing import Any + +_TRACE_ID = re.compile(r"^[0-9a-f]{32}$") +_SPAN_ID = re.compile(r"^[0-9a-f]{16}$") + + +def _digest(text: str, nbytes: int) -> int: + return int.from_bytes(hashlib.sha256(text.encode("utf-8")).digest()[:nbytes], "big") + + +def span_id_for(value: str) -> int: + """A 64-bit OpenTelemetry span id for a LogQuill span id: the id itself if + it already is 16 hex characters (the shape `Logger.span()` generates), + otherwise a stable hash of it — LangChain's run ids are UUIDs, for example. + The same input always gives the same id, so a child's `parent_span_id` and + its parent's own `span_id` still meet.""" + if _SPAN_ID.match(value) and int(value, 16) != 0: + return int(value, 16) + return _digest(value, 8) or 1 + + +def trace_id_for(meta: Mapping[str, Any]) -> int | None: + """The 128-bit trace id a record belongs to: `meta.trace_id` if it's a valid + 32-hex id (what `TraceContextPlugin` writes), else derived from + `meta.run_id` so every span of one agent run shares a trace. `None` if the + record names neither.""" + trace_id = meta.get("trace_id") + if isinstance(trace_id, str) and _TRACE_ID.match(trace_id) and int(trace_id, 16) != 0: + return int(trace_id, 16) + run_id = meta.get("run_id") + if isinstance(run_id, str) and run_id: + return _digest("run:" + run_id, 16) or 1 + return None + + +def trace_id_from_parent(parent_span_id: str) -> int: + """A trace id for spans that name a parent but no run or trace: derived from + the parent's id, so siblings share a trace.""" + return _digest("parent:" + parent_span_id, 16) or 1 diff --git a/logquill/transports/otel/logs_transport.py b/logquill/transports/otel/logs_transport.py new file mode 100644 index 0000000..eba6401 --- /dev/null +++ b/logquill/transports/otel/logs_transport.py @@ -0,0 +1,139 @@ +from __future__ import annotations + +import json +from typing import Any + +from logquill import semconv +from logquill.formatters import Formatter +from logquill.records import LogRecord +from logquill.transports.otel import _ids +from logquill.transports.transport import Transport + +_SEVERITY = { + "TRACE": "TRACE", + "DEBUG": "DEBUG", + "INFO": "INFO", + "WARN": "WARN", + "ERROR": "ERROR", + "FATAL": "FATAL", +} + + +def _attribute_value(value: Any) -> Any: + """An OpenTelemetry attribute value: scalars and homogeneous lists of them + pass through; anything else (a nested dict, a mixed list, an object) becomes + a JSON string, since attributes can't hold structure.""" + if isinstance(value, (str, bool, int, float)): + return value + if isinstance(value, (list, tuple)) and value: + kinds = {type(item) for item in value} + if len(kinds) == 1 and kinds <= {str, bool, int, float}: + return list(value) + try: + return json.dumps(value, default=str, separators=(",", ":")) + except (TypeError, ValueError): + return repr(value) + + +class OTelLogsTransport(Transport): + """Sends every LogQuill record to an OpenTelemetry collector as an OTLP + *log record*, without going through the span path — for a backend that + takes logs, or to ship the records `OTLPTransport` doesn't turn into spans. + + The message is the body, the level maps to OTLP severity, and `meta` + becomes attributes (nested values as JSON strings; a traceback in + `meta.stack` goes to `exception.stacktrace`). A record with a `trace_id` or + `run_id`, or a `span_id`, carries the same ids `OTLPTransport` gives its + spans, so a collector can join a log line to its span. An `llm` block adds + the standard model and token attributes. + + Requires `pip install logquill[otel]`, imported lazily. + """ + + def __init__( + self, + *, + endpoint: str | None = None, + headers: dict[str, str] | None = None, + service_name: str = "logquill", + log_exporter: Any = None, + processor: str = "batch", + semconv_version: str | None = None, + timeout: float = 10.0, + formatter: Formatter | None = None, + ) -> None: + """Arguments mirror `OTLPTransport`: `endpoint` is the OTLP/HTTP logs + URL, `log_exporter` swaps in any OpenTelemetry log exporter, and + `processor` is `"batch"` or `"simple"`.""" + super().__init__(formatter) + try: + from opentelemetry.sdk._logs import LoggerProvider + from opentelemetry.sdk._logs.export import ( + BatchLogRecordProcessor, + SimpleLogRecordProcessor, + ) + from opentelemetry.sdk.resources import Resource + except ImportError as exc: + raise ImportError( + "OTelLogsTransport requires the OpenTelemetry SDK — install with " + "`pip install logquill[otel]`." + ) from exc + if processor not in ("batch", "simple"): + raise ValueError(f"processor must be 'batch' or 'simple', got {processor!r}") + if log_exporter is None: + try: + from opentelemetry.exporter.otlp.proto.http._log_exporter import OTLPLogExporter + except ImportError as exc: + raise ImportError( + "OTelLogsTransport needs the OTLP/HTTP exporter — install with " + "`pip install logquill[otel]`." + ) from exc + log_exporter = OTLPLogExporter(endpoint=endpoint, headers=headers, timeout=timeout) + + self._convention = semconv.resolve_convention(semconv_version) + self._provider = LoggerProvider(resource=Resource.create({"service.name": service_name})) + self._provider.add_log_record_processor( + BatchLogRecordProcessor(log_exporter) + if processor == "batch" + else SimpleLogRecordProcessor(log_exporter) + ) + self._logger = self._provider.get_logger("logquill") + + def write(self, formatted: str, record: LogRecord) -> None: + """Emits `record` as an OTLP log record.""" + from opentelemetry._logs import LogRecord as OTelLogRecord + from opentelemetry._logs import SeverityNumber + + meta = record["meta"] + attributes: dict[str, Any] = {} + for key, value in meta.items(): + if key == "stack": + attributes["exception.stacktrace"] = _attribute_value(value) + elif value is not None: + attributes[key] = _attribute_value(value) + attributes["logquill.logger"] = record["logger"] + attributes["logquill.schema_version"] = record["schema_version"] + attributes.update(semconv.log_attributes(record, self._convention)) + + span_id = meta.get("span_id") + timestamp = semconv.epoch_ns(record["timestamp"]) + self._logger.emit( + OTelLogRecord( + timestamp=timestamp, + observed_timestamp=timestamp, + trace_id=_ids.trace_id_for(meta), + span_id=_ids.span_id_for(span_id) if isinstance(span_id, str) and span_id else None, + severity_text=record["level"], + severity_number=getattr(SeverityNumber, _SEVERITY[record["level"]]), + body=record["message"], + attributes=attributes, + ) + ) + + def flush(self) -> None: + """Exports every log record still waiting in the batch processor.""" + self._provider.force_flush() + + def close(self) -> None: + """Flushes, then shuts the exporter down.""" + self._provider.shutdown() diff --git a/logquill/transports/otel/otlp_transport.py b/logquill/transports/otel/otlp_transport.py new file mode 100644 index 0000000..1b83f8c --- /dev/null +++ b/logquill/transports/otel/otlp_transport.py @@ -0,0 +1,213 @@ +from __future__ import annotations + +import logging +import threading +from typing import Any + +from logquill import semconv +from logquill.formatters import Formatter +from logquill.records import LogRecord +from logquill.transports.otel import _ids +from logquill.transports.transport import Transport + +_logger = logging.getLogger("logquill") + + +def _require_otel() -> None: + try: + import opentelemetry.sdk.trace # noqa: F401 + except ImportError as exc: + raise ImportError( + "OTLPTransport requires the OpenTelemetry SDK — install with " + "`pip install logquill[otel]`." + ) from exc + + +class OTLPTransport(Transport): + """Exports LogQuill's agent tracing as real OpenTelemetry spans. + + A log line and a span are different things, so this doesn't format text: it + reads each record's structure and, for the records that describe work with a + beginning and an end, creates an OpenTelemetry span through the SDK, so it + goes through the usual processors and exporters to any OTLP collector + (Jaeger, Tempo, Honeycomb, Datadog, Grafana, ...): + + - a `Logger.span()` block becomes a span. One that marks itself as an agent + run — `logger.span("run", operation="invoke_agent", agent_name="planner")` — + is `invoke_agent planner`; any other keeps its own name. + - an `.action()` naming its tool in `meta.tool` becomes `execute_tool {tool}`. + - a record with an `llm` block (`Logger.llm_call()`) becomes `chat {model}` + with the model, token counts and finish reason as the standard GenAI + attributes. + + Every other record is ignored here; send those to a collector too with + `OTelLogsTransport`. The span's real ids, start and end are taken from the + record (`span_id`, `parent_span_id`, `duration_ms`, the timestamp), so the + tree a collector shows is exactly the tree that was logged. Spans of one run + share a trace when the records carry a `run_id` or `trace_id`. + + Attribute names come from `logquill.semconv`, pinned to one release of the + (still experimental) GenAI conventions; see `semconv` for choosing the + naming generation. Prompt and completion text is never exported unless you + turn on `capture_content`. + + Requires `pip install logquill[otel]`, imported lazily. Export is batched + and happens off the caller's thread; `flush()`/`close()` send what's pending. + """ + + def __init__( + self, + *, + endpoint: str | None = None, + headers: dict[str, str] | None = None, + service_name: str = "logquill", + span_exporter: Any = None, + processor: str = "batch", + semconv_version: str | None = None, + default_provider: str = semconv.UNKNOWN_PROVIDER, + capture_content: bool | None = None, + timeout: float = 10.0, + formatter: Formatter | None = None, + ) -> None: + """`endpoint` is the OTLP/HTTP traces URL (default: the SDK's, i.e. + `OTEL_EXPORTER_OTLP_*` environment variables, else localhost). + `span_exporter` replaces the OTLP exporter with any OpenTelemetry + `SpanExporter` — for tests, or a different protocol; `processor` is + `"batch"` (default) or `"simple"` (export each span immediately). + `semconv_version` is `"latest"` or `"legacy"`; unset, it follows + `OTEL_SEMCONV_STABILITY_OPT_IN`. `default_provider` names the provider + for records without `meta.provider`. `capture_content` defaults to the + `OTEL_INSTRUMENTATION_GENAI_CAPTURE_MESSAGE_CONTENT` environment + variable, else off.""" + super().__init__(formatter) + _require_otel() + if processor not in ("batch", "simple"): + raise ValueError(f"processor must be 'batch' or 'simple', got {processor!r}") + from opentelemetry.sdk.resources import Resource + from opentelemetry.sdk.trace import TracerProvider + from opentelemetry.sdk.trace.export import BatchSpanProcessor, SimpleSpanProcessor + from opentelemetry.sdk.trace.id_generator import RandomIdGenerator + + self._convention = semconv.resolve_convention(semconv_version) + self._default_provider = default_provider + self._capture_content = ( + semconv.capture_content_enabled() if capture_content is None else capture_content + ) + + pending = threading.local() + + class _RecordIds(RandomIdGenerator): # type: ignore[misc] + """Uses the ids the record being exported names, when it names any.""" + + def generate_span_id(self) -> int: + return int(getattr(pending, "span_id", None) or super().generate_span_id()) + + def generate_trace_id(self) -> int: + return int(getattr(pending, "trace_id", None) or super().generate_trace_id()) + + self._pending = pending + if span_exporter is None: + try: + from opentelemetry.exporter.otlp.proto.http.trace_exporter import OTLPSpanExporter + except ImportError as exc: + raise ImportError( + "OTLPTransport needs the OTLP/HTTP exporter — install with " + "`pip install logquill[otel]`." + ) from exc + span_exporter = OTLPSpanExporter(endpoint=endpoint, headers=headers, timeout=timeout) + + self._provider = TracerProvider( + resource=Resource.create({"service.name": service_name}), + id_generator=_RecordIds(), + ) + span_processor = ( + BatchSpanProcessor(span_exporter) + if processor == "batch" + else SimpleSpanProcessor(span_exporter) + ) + self._provider.add_span_processor(span_processor) + self._tracer = self._provider.get_tracer("logquill") + + def _spec_for(self, record: LogRecord) -> semconv.SpanSpec | None: + if semconv.is_llm_record(record): + return semconv.llm_spec( + record, + self._convention, + default_provider=self._default_provider, + capture_content=self._capture_content, + ) + if semconv.is_tool_record(record): + return semconv.tool_spec(record, self._convention) + if semconv.is_span_record(record): + return semconv.span_spec( + record, self._convention, default_provider=self._default_provider + ) + return None + + def write(self, formatted: str, record: LogRecord) -> None: + """Creates and ends the span this record describes, if it describes one + (see the class docstring); ignores the record otherwise.""" + spec = self._spec_for(record) + if spec is None: + return + + from opentelemetry import trace + from opentelemetry.trace import ( + NonRecordingSpan, + SpanContext, + SpanKind, + Status, + StatusCode, + TraceFlags, + ) + + meta = record["meta"] + end_ns = semconv.epoch_ns(record["timestamp"]) + start_ns = end_ns - int(spec.duration_ms * 1_000_000) if spec.duration_ms else end_ns + + own = meta.get("span_id") + parent = meta.get("parent_span_id") + trace_id = _ids.trace_id_for(meta) + parent_context = None + if isinstance(parent, str) and parent: + if trace_id is None: + trace_id = _ids.trace_id_from_parent(parent) + parent_context = trace.set_span_in_context( + NonRecordingSpan( + SpanContext( + trace_id=trace_id, + span_id=_ids.span_id_for(parent), + is_remote=True, + trace_flags=TraceFlags(TraceFlags.SAMPLED), + ) + ) + ) + + self._pending.trace_id = trace_id + self._pending.span_id = _ids.span_id_for(own) if isinstance(own, str) and own else None + try: + span = self._tracer.start_span( + spec.name, + context=parent_context, + kind=SpanKind.CLIENT if spec.kind == "client" else SpanKind.INTERNAL, + attributes=spec.attributes, + start_time=start_ns, + ) + finally: + self._pending.trace_id = self._pending.span_id = None + for name, attributes in spec.events: + span.add_event(name, attributes, timestamp=end_ns) + if spec.is_error: + description = meta.get("error") + span.set_status( + Status(StatusCode.ERROR, description if isinstance(description, str) else None) + ) + span.end(end_time=end_ns) + + def flush(self) -> None: + """Exports every span still waiting in the batch processor.""" + self._provider.force_flush() + + def close(self) -> None: + """Flushes, then shuts the exporter down.""" + self._provider.shutdown() diff --git a/pyproject.toml b/pyproject.toml index cb594db..37104e4 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -40,6 +40,11 @@ dependencies = [] [project.optional-dependencies] http = ["aiohttp>=3.9"] apprise = ["apprise>=1.8"] +otel = [ + "opentelemetry-api>=1.39", + "opentelemetry-sdk>=1.39", + "opentelemetry-exporter-otlp-proto-http>=1.39", +] postgres = ["psycopg2-binary>=2.9"] mysql = ["pymysql>=1.1"] mongodb = ["pymongo>=4.6"] @@ -112,6 +117,10 @@ include = [ line-length = 100 target-version = "py310" +[tool.ruff.lint.per-file-ignores] +# Test fixtures read better as one record per line than wrapped to 100 columns. +"tests/**" = ["E501"] + [tool.ruff.lint] select = ["E", "F", "W", "I", "UP", "B", "C4", "SIM"] # Importing `Callable`/`Sequence`/... from `typing` still works and is used @@ -144,6 +153,13 @@ warn_unused_ignores = false module = ["apprise", "aiohttp"] ignore_missing_imports = true +# OpenTelemetry is optional too, and is imported lazily; treating it as untyped +# keeps `mypy --strict` giving the same answer with or without it installed. +[[tool.mypy.overrides]] +module = ["opentelemetry", "opentelemetry.*"] +ignore_missing_imports = true +follow_imports = "skip" + [tool.pytest.ini_options] testpaths = ["tests"] asyncio_mode = "auto" diff --git a/schema/golden_records.json b/schema/golden_records.json index 72055dc..fb0b415 100644 --- a/schema/golden_records.json +++ b/schema/golden_records.json @@ -310,6 +310,92 @@ "message": "hello", "meta": {} } + }, + { + "name": "tool call with id and retry", + "record": { + "schema_version": "2.0", + "timestamp": "2026-01-01T00:00:00.000Z", + "level": "INFO", + "logger": "app.agent", + "message": "search", + "meta": { + "kind": "action", + "tool": "search", + "tool_call_id": "call_1", + "retry_count": 1, + "duration_ms": 120.5 + } + } + }, + { + "name": "LLM call with provider and response model", + "record": { + "schema_version": "2.0", + "timestamp": "2026-01-01T00:00:00.000Z", + "level": "INFO", + "logger": "app.agent", + "message": "chat", + "meta": { + "kind": "action", + "provider": "openai", + "response_model": "example-model-2026", + "operation": "chat" + }, + "llm": { + "model": "example-model", + "tokens_in": 10, + "tokens_out": 4, + "finish_reason": "stop" + } + } + }, + { + "name": "agent run span with agent identity", + "record": { + "schema_version": "2.0", + "timestamp": "2026-01-01T00:00:00.000Z", + "level": "INFO", + "logger": "app.agent", + "message": "run", + "meta": { + "kind": "span", + "run_id": "run-1", + "span_id": "00f067aa0ba902b7", + "duration_ms": 40, + "agent_name": "planner", + "agent_id": "agent-7", + "operation": "invoke_agent" + } + } + }, + { + "name": "opt-in message content", + "record": { + "schema_version": "2.0", + "timestamp": "2026-01-01T00:00:00.000Z", + "level": "INFO", + "logger": "app.agent", + "message": "chat", + "meta": { + "kind": "action", + "input_messages": [ + { + "role": "user", + "parts": [ + { + "type": "text", + "content": "hi" + } + ] + } + ], + "output_messages": [] + }, + "llm": { + "model": "m" + } + } } ], "invalid": [ @@ -728,6 +814,90 @@ "hash": "abc" } } + }, + { + "name": "tool is a number", + "record": { + "schema_version": "2.0", + "timestamp": "2026-01-01T00:00:00.000Z", + "level": "INFO", + "logger": "app.agent", + "message": "x", + "meta": { + "tool": 7 + } + }, + "why": "tool is a string" + }, + { + "name": "tool_call_id is a number", + "record": { + "schema_version": "2.0", + "timestamp": "2026-01-01T00:00:00.000Z", + "level": "INFO", + "logger": "app.agent", + "message": "x", + "meta": { + "tool_call_id": 7 + } + }, + "why": "tool_call_id is a string" + }, + { + "name": "provider is a number", + "record": { + "schema_version": "2.0", + "timestamp": "2026-01-01T00:00:00.000Z", + "level": "INFO", + "logger": "app.agent", + "message": "x", + "meta": { + "provider": 1 + } + }, + "why": "provider is a string" + }, + { + "name": "agent_name is a number", + "record": { + "schema_version": "2.0", + "timestamp": "2026-01-01T00:00:00.000Z", + "level": "INFO", + "logger": "app.agent", + "message": "x", + "meta": { + "agent_name": 1 + } + }, + "why": "agent_name is a string" + }, + { + "name": "unknown operation", + "record": { + "schema_version": "2.0", + "timestamp": "2026-01-01T00:00:00.000Z", + "level": "INFO", + "logger": "app.agent", + "message": "x", + "meta": { + "operation": "embed" + } + }, + "why": "operation is chat, text_completion or invoke_agent" + }, + { + "name": "input_messages is a string", + "record": { + "schema_version": "2.0", + "timestamp": "2026-01-01T00:00:00.000Z", + "level": "INFO", + "logger": "app.agent", + "message": "x", + "meta": { + "input_messages": "hello" + } + }, + "why": "input_messages is an array" } ], "legacy": [ diff --git a/schema/record.schema.json b/schema/record.schema.json index 1877c0a..8e904a6 100644 --- a/schema/record.schema.json +++ b/schema/record.schema.json @@ -108,6 +108,45 @@ "prev_hash": { "type": "string", "pattern": "^[0-9a-f]{64}$" + }, + "tool": { + "type": "string", + "description": "Name of the tool an action called." + }, + "tool_call_id": { + "type": "string", + "description": "Identifies one tool call across its retries." + }, + "provider": { + "type": "string", + "description": "The LLM provider, e.g. openai, anthropic." + }, + "agent_name": { + "type": "string", + "description": "Human-readable name of the agent." + }, + "agent_id": { + "type": "string" + }, + "operation": { + "enum": [ + "chat", + "text_completion", + "invoke_agent" + ], + "description": "Overrides the operation inferred for OpenTelemetry export." + }, + "response_model": { + "type": "string", + "description": "The model that actually answered, if it differs from llm.model." + }, + "input_messages": { + "type": "array", + "description": "Opt-in prompt content; only exported to OpenTelemetry when content capture is on." + }, + "output_messages": { + "type": "array", + "description": "Opt-in completion content." } } }, diff --git a/tests/test_llm_call.py b/tests/test_llm_call.py new file mode 100644 index 0000000..702b3f6 --- /dev/null +++ b/tests/test_llm_call.py @@ -0,0 +1,133 @@ +from __future__ import annotations + +import json +import logging +from typing import Any + +import pytest + +from logquill import JSONFormatter, Logger +from logquill.transports.transport import CollectingTransport + + +def _logger() -> tuple[Logger, CollectingTransport]: + sink = CollectingTransport() + return Logger("app.agent", transports=[sink]), sink + + +def test_llm_call_writes_the_first_class_llm_block() -> None: + logger, sink = _logger() + + record = logger.llm_call( + "chat", + model="example-model", + tokens_in=1200, + tokens_out=340, + cost_usd=0.0123, + latency_ms=2150.5, + finish_reason="stop", + provider="openai", + ) + + assert record is sink.records[0] + assert record["message"] == "chat" + assert record["llm"] == { + "model": "example-model", + "tokens_in": 1200, + "tokens_out": 340, + "cost_usd": 0.0123, + "latency_ms": 2150.5, + "finish_reason": "stop", + } + assert record["meta"] == {"kind": "action", "provider": "openai"} + assert record["level"] == "INFO" + + +def test_llm_call_leaves_out_fields_you_did_not_give() -> None: + logger, _ = _logger() + + record = logger.llm_call(model="m", tokens_in=0) + + assert record is not None + assert record["llm"] == {"model": "m", "tokens_in": 0} + assert record["message"] == "llm_call" + + +def test_llm_call_with_nothing_to_report_has_no_llm_key() -> None: + logger, _ = _logger() + + record = logger.llm_call() + + assert record is not None + assert "llm" not in record + + +@pytest.mark.parametrize( + "bad", + [ + {"tokens_in": -1}, + {"tokens_in": 1.5}, + {"tokens_in": "12"}, + {"tokens_out": True}, + {"cost_usd": "0.1"}, + {"latency_ms": -0.5}, + {"model": 7}, + {"finish_reason": ["stop"]}, + ], +) +def test_a_value_that_breaks_the_contract_is_dropped_with_a_warning( + bad: dict[str, Any], caplog: pytest.LogCaptureFixture +) -> None: + logger, _ = _logger() + good = {"tokens_out": 3, "finish_reason": "stop"} + + with caplog.at_level(logging.WARNING, logger="logquill"): + record = logger.llm_call(**{**good, **bad}) + + assert record is not None + assert record.get("llm") == ({k: v for k, v in good.items() if k not in bad} or None) + assert "llm_call: ignoring" in caplog.text + + +def test_llm_call_respects_the_level_and_disable() -> None: + import logquill + + logger = Logger("app.agent", level="ERROR") + assert logger.llm_call(model="m") is None + + loud = Logger("mylib", level="INFO") + logquill.disable("mylib") + assert loud.llm_call(model="m") is None + + +def test_opt_views_have_llm_call_too() -> None: + logger, sink = _logger() + + record = logger.opt(depth=0).llm_call("chat", model="m", tokens_in=3) + + assert record is not None + assert record["llm"] == {"model": "m", "tokens_in": 3} + assert record["meta"]["caller"]["function"] == "test_opt_views_have_llm_call_too" + assert sink.records == [record] + + +def test_the_record_validates_against_the_published_schema() -> None: + jsonschema = pytest.importorskip("jsonschema") + from pathlib import Path + + schema = json.loads( + (Path(__file__).parent.parent / "schema" / "record.schema.json").read_text() + ) + logger, _ = _logger() + record = logger.llm_call( + "chat", + model="m", + tokens_in=1, + tokens_out=2, + cost_usd=0.1, + latency_ms=3, + finish_reason="stop", + provider="p", + ) + + jsonschema.validate(json.loads(JSONFormatter().format(record)), schema) # type: ignore[arg-type] diff --git a/tests/test_otel_logs_transport.py b/tests/test_otel_logs_transport.py new file mode 100644 index 0000000..a141f43 --- /dev/null +++ b/tests/test_otel_logs_transport.py @@ -0,0 +1,145 @@ +from __future__ import annotations + +import sys +from typing import Any + +import pytest + +pytest.importorskip("opentelemetry.sdk") + +from opentelemetry._logs import SeverityNumber # noqa: E402 +from opentelemetry.sdk._logs import export as _log_export # noqa: E402 + +# renamed in newer SDKs; the old name warns +InMemoryLogExporter = ( + getattr(_log_export, "InMemoryLogRecordExporter", None) or _log_export.InMemoryLogExporter +) + +from logquill import Logger, OTelLogsTransport, RunPlugin # noqa: E402 + + +def _setup(**kwargs: Any) -> tuple[Logger, InMemoryLogExporter]: + exporter = InMemoryLogExporter() + transport = OTelLogsTransport(log_exporter=exporter, processor="simple", **kwargs) + return Logger("app.api", level="TRACE", transports=[transport], plugins=[RunPlugin()]), exporter + + +def _emitted(exporter: InMemoryLogExporter) -> list[Any]: + return [item.log_record for item in exporter.get_finished_logs()] + + +def test_every_record_becomes_a_log_record_with_body_and_severity() -> None: + logger, exporter = _setup() + + logger.trace("t") + logger.debug("d") + logger.info("i") + logger.warn("w") + logger.error("e") + logger.fatal("f") + + emitted = _emitted(exporter) + assert [r.body for r in emitted] == ["t", "d", "i", "w", "e", "f"] + assert [r.severity_text for r in emitted] == [ + "TRACE", + "DEBUG", + "INFO", + "WARN", + "ERROR", + "FATAL", + ] + assert [r.severity_number for r in emitted] == [ + SeverityNumber.TRACE, + SeverityNumber.DEBUG, + SeverityNumber.INFO, + SeverityNumber.WARN, + SeverityNumber.ERROR, + SeverityNumber.FATAL, + ] + + +def test_meta_becomes_attributes_with_structure_flattened_to_json() -> None: + logger, exporter = _setup() + + logger.info( + "hello", user_id=42, ok=True, ratio=0.5, tags=["a", "b"], nested={"a": [1]}, none=None + ) + + attributes = dict(_emitted(exporter)[0].attributes) + assert attributes["user_id"] == 42 + assert attributes["ok"] is True + assert attributes["ratio"] == 0.5 + assert attributes["tags"] == ("a", "b") + assert attributes["nested"] == '{"a":[1]}' + assert "none" not in attributes + assert attributes["logquill.logger"] == "app.api" + assert attributes["logquill.schema_version"] == "2.0" + + +def test_a_mixed_list_becomes_a_json_string() -> None: + logger, exporter = _setup() + + logger.info("hello", mixed=[1, "a"]) + + assert dict(_emitted(exporter)[0].attributes)["mixed"] == '[1,"a"]' + + +def test_a_traceback_goes_to_the_standard_exception_attribute() -> None: + logger, exporter = _setup() + + try: + _ = 1 / 0 + except ZeroDivisionError as exc: + logger.error("failed", exc_info=exc) + + attributes = dict(_emitted(exporter)[0].attributes) + assert "ZeroDivisionError" in attributes["exception.stacktrace"] + assert "stack" not in attributes + + +def test_the_ids_match_the_ones_the_span_exporter_uses() -> None: + logger, exporter = _setup() + + with logger.span("work", span_id="00f067aa0ba902b7"): + pass + + record = _emitted(exporter)[0] + assert format(record.span_id, "016x") == "00f067aa0ba902b7" + assert record.trace_id is not None and record.trace_id != 0 + + +def test_an_llm_block_adds_the_standard_usage_attributes() -> None: + logger, exporter = _setup() + + logger.llm_call( + "chat", model="m", tokens_in=7, tokens_out=3, cost_usd=0.25, finish_reason="stop" + ) + + attributes = dict(_emitted(exporter)[0].attributes) + assert attributes["gen_ai.usage.input_tokens"] == 7 + assert attributes["gen_ai.usage.output_tokens"] == 3 + assert attributes["gen_ai.request.model"] == "m" + assert attributes["logquill.cost_usd"] == 0.25 + + +def test_flush_and_close_with_the_batch_processor() -> None: + exporter = InMemoryLogExporter() + logger = Logger("app", transports=[OTelLogsTransport(log_exporter=exporter)]) + + logger.info("hello") + logger.flush() + + assert len(exporter.get_finished_logs()) == 1 + logger.close() + + +def test_a_bad_processor_is_rejected_up_front() -> None: + with pytest.raises(ValueError, match="processor must be"): + OTelLogsTransport(log_exporter=InMemoryLogExporter(), processor="later") + + +def test_a_missing_sdk_gives_an_install_hint(monkeypatch: pytest.MonkeyPatch) -> None: + monkeypatch.setitem(sys.modules, "opentelemetry.sdk._logs", None) + + with pytest.raises(ImportError, match=r"pip install logquill\[otel\]"): + OTelLogsTransport(log_exporter=InMemoryLogExporter()) diff --git a/tests/test_otel_wire.py b/tests/test_otel_wire.py new file mode 100644 index 0000000..7aa896f --- /dev/null +++ b/tests/test_otel_wire.py @@ -0,0 +1,151 @@ +"""OTLP over HTTP end to end: a real exporter posts to a local collector +stand-in, and the protobuf it received is decoded and checked.""" + +from __future__ import annotations + +import gzip +import threading +from collections.abc import Iterator +from http.server import BaseHTTPRequestHandler, HTTPServer +from typing import Any + +import pytest + +pytest.importorskip("opentelemetry.exporter.otlp.proto.http") + +from opentelemetry.proto.collector.logs.v1.logs_service_pb2 import ( # noqa: E402 + ExportLogsServiceRequest, +) +from opentelemetry.proto.collector.trace.v1.trace_service_pb2 import ( # noqa: E402 + ExportTraceServiceRequest, +) + +from logquill import Logger, OTelLogsTransport, OTLPTransport, RunPlugin # noqa: E402 + + +class Collector: + def __init__(self) -> None: + self.requests: dict[str, list[bytes]] = {"/v1/traces": [], "/v1/logs": []} + self.received = threading.Event() + + +@pytest.fixture() +def collector() -> Iterator[tuple[str, Collector]]: + state = Collector() + + class Handler(BaseHTTPRequestHandler): + def do_POST(self) -> None: # noqa: N802 + body = self.rfile.read(int(self.headers["Content-Length"])) + if self.headers.get("Content-Encoding") == "gzip": + body = gzip.decompress(body) + state.requests.setdefault(self.path, []).append(body) + self.send_response(200) + self.send_header("Content-Length", "0") + self.end_headers() + state.received.set() + + def log_message(self, format: str, *args: Any) -> None: + pass + + httpd = HTTPServer(("127.0.0.1", 0), Handler) + thread = threading.Thread( + target=httpd.serve_forever, kwargs={"poll_interval": 0.01}, daemon=True + ) + thread.start() + try: + yield f"http://127.0.0.1:{httpd.server_port}", state + finally: + httpd.shutdown() + httpd.server_close() + + +def _attribute(attributes: Any, key: str) -> Any: + for item in attributes: + if item.key == key: + return item.value + raise KeyError(key) + + +def test_an_agent_run_reaches_a_collector_with_the_right_tree_and_usage_values( + collector: tuple[str, Collector], +) -> None: + base, state = collector + logger = Logger( + "app.agent", + transports=[OTLPTransport(endpoint=f"{base}/v1/traces", service_name="agent-service")], + plugins=[RunPlugin()], + ) + + with logger.span("run", operation="invoke_agent", agent_name="planner"): + logger.llm_call( + "chat", + model="example-model", + tokens_in=1200, + tokens_out=340, + finish_reason="stop", + provider="anthropic", + ) + logger.close() + + (body,) = state.requests["/v1/traces"] + request = ExportTraceServiceRequest.FromString(body) + (resource_spans,) = request.resource_spans + assert ( + _attribute(resource_spans.resource.attributes, "service.name").string_value + == "agent-service" + ) + spans = {span.name: span for scope in resource_spans.scope_spans for span in scope.spans} + assert set(spans) == {"invoke_agent planner", "chat example-model"} + + run, chat = spans["invoke_agent planner"], spans["chat example-model"] + assert run.parent_span_id == b"" + assert chat.parent_span_id == run.span_id + assert chat.trace_id == run.trace_id + assert _attribute(chat.attributes, "gen_ai.usage.input_tokens").int_value == 1200 + assert _attribute(chat.attributes, "gen_ai.usage.output_tokens").int_value == 340 + assert _attribute(chat.attributes, "gen_ai.request.model").string_value == "example-model" + assert _attribute(chat.attributes, "gen_ai.operation.name").string_value == "chat" + assert ( + list(_attribute(chat.attributes, "gen_ai.response.finish_reasons").array_value.values)[ + 0 + ].string_value + == "stop" + ) + + +def test_logs_reach_a_collector_and_join_their_span_by_id( + collector: tuple[str, Collector], +) -> None: + base, state = collector + logger = Logger( + "app.agent", + transports=[ + OTLPTransport(endpoint=f"{base}/v1/traces"), + OTelLogsTransport(endpoint=f"{base}/v1/logs"), + ], + plugins=[RunPlugin()], + ) + + with logger.span("work", span_id="00f067aa0ba902b7"): + logger.warn("something odd", user_id=42) + logger.close() + + span = next( + span + for body in state.requests["/v1/traces"] + for scope in ExportTraceServiceRequest.FromString(body).resource_spans[0].scope_spans + for span in scope.spans + ) + log_records = [ + record + for body in state.requests["/v1/logs"] + for scope in ExportLogsServiceRequest.FromString(body).resource_logs[0].scope_logs + for record in scope.log_records + ] + warning = next(r for r in log_records if r.body.string_value == "something odd") + assert warning.severity_text == "WARN" + assert _attribute(warning.attributes, "user_id").int_value == 42 + # the span's own log line carries the span's id, so the two can be joined + span_log = next(r for r in log_records if r.body.string_value == "work") + assert span_log.span_id == span.span_id + assert span_log.trace_id == span.trace_id diff --git a/tests/test_otlp_transport.py b/tests/test_otlp_transport.py new file mode 100644 index 0000000..2ec9af3 --- /dev/null +++ b/tests/test_otlp_transport.py @@ -0,0 +1,314 @@ +from __future__ import annotations + +import subprocess +import sys +from typing import Any + +import pytest + +pytest.importorskip("opentelemetry.sdk") + +from opentelemetry.sdk.trace import ReadableSpan # noqa: E402 +from opentelemetry.sdk.trace.export.in_memory_span_exporter import ( # noqa: E402 + InMemorySpanExporter, +) +from opentelemetry.trace import SpanKind, StatusCode # noqa: E402 + +from logquill import Logger, OTLPTransport, RunPlugin, TraceContextPlugin # noqa: E402 +from logquill.plugins.trace_context_plugin import ( # noqa: E402 + reset_traceparent, + set_traceparent, +) + + +def _setup(**kwargs: Any) -> tuple[Logger, OTLPTransport, InMemorySpanExporter]: + exporter = InMemorySpanExporter() + transport = OTLPTransport(span_exporter=exporter, processor="simple", **kwargs) + logger = Logger("app.agent", level="TRACE", transports=[transport], plugins=[RunPlugin()]) + return logger, transport, exporter + + +def _by_name(spans: list[ReadableSpan]) -> dict[str, ReadableSpan]: + return {span.name: span for span in spans} + + +def _attrs(span: ReadableSpan) -> dict[str, Any]: + return dict(span.attributes or {}) + + +def _record_an_agent_run(logger: Logger) -> None: + with logger.span("run", operation="invoke_agent", agent_name="planner"): + logger.thought("plan the work") # not a span: ignored by the span exporter + with logger.span("plan_step"): + logger.llm_call( + "chat", + model="example-model", + tokens_in=1200, + tokens_out=340, + cost_usd=0.0123, + latency_ms=250, + finish_reason="stop", + provider="anthropic", + ) + logger.action("look it up", tool="search", tool_call_id="call_1", duration_ms=40) + logger.close() + + +def test_a_recorded_agent_run_exports_the_correct_span_tree() -> None: + logger, _, exporter = _setup() + + _record_an_agent_run(logger) + + spans = _by_name(list(exporter.get_finished_spans())) + assert set(spans) == { + "invoke_agent planner", + "plan_step", + "chat example-model", + "execute_tool search", + } + run, step = spans["invoke_agent planner"], spans["plan_step"] + chat, tool = spans["chat example-model"], spans["execute_tool search"] + + assert run.parent is None + assert step.parent is not None and step.parent.span_id == run.context.span_id + assert chat.parent is not None and chat.parent.span_id == step.context.span_id + assert tool.parent is not None and tool.parent.span_id == run.context.span_id + # one run, one trace + assert len({s.context.trace_id for s in spans.values()}) == 1 + + +def test_the_llm_span_carries_the_usage_and_model_attributes() -> None: + logger, _, exporter = _setup() + + _record_an_agent_run(logger) + + chat = _by_name(list(exporter.get_finished_spans()))["chat example-model"] + assert chat.kind == SpanKind.CLIENT + assert _attrs(chat) == { + "gen_ai.operation.name": "chat", + "gen_ai.provider.name": "anthropic", + "gen_ai.request.model": "example-model", + "gen_ai.usage.input_tokens": 1200, + "gen_ai.usage.output_tokens": 340, + "gen_ai.response.finish_reasons": ("stop",), + "logquill.cost_usd": 0.0123, + "logquill.run_id": chat.attributes["logquill.run_id"], # type: ignore[index] + } + + +def test_agent_and_tool_spans_use_the_conventional_names_and_attributes() -> None: + logger, _, exporter = _setup() + + _record_an_agent_run(logger) + + spans = _by_name(list(exporter.get_finished_spans())) + run, tool = spans["invoke_agent planner"], spans["execute_tool search"] + assert run.kind == SpanKind.CLIENT + assert _attrs(run)["gen_ai.operation.name"] == "invoke_agent" + assert _attrs(run)["gen_ai.agent.name"] == "planner" + assert tool.kind == SpanKind.INTERNAL + assert _attrs(tool)["gen_ai.operation.name"] == "execute_tool" + assert _attrs(tool)["gen_ai.tool.name"] == "search" + assert _attrs(tool)["gen_ai.tool.call.id"] == "call_1" + assert "gen_ai.operation.name" not in _attrs(spans["plan_step"]) + + +def test_span_start_and_end_come_from_the_record() -> None: + logger, _, exporter = _setup() + + _record_an_agent_run(logger) + + spans = _by_name(list(exporter.get_finished_spans())) + chat, tool = spans["chat example-model"], spans["execute_tool search"] + assert chat.end_time - chat.start_time == pytest.approx(250_000_000, abs=1_000) + assert tool.end_time - tool.start_time == pytest.approx(40_000_000, abs=1_000) + + +def test_a_trace_id_from_trace_context_is_used_as_the_otel_trace_id() -> None: + exporter = InMemorySpanExporter() + logger = Logger( + "app", + transports=[OTLPTransport(span_exporter=exporter, processor="simple")], + plugins=[TraceContextPlugin()], + ) + + token = set_traceparent("00-4bf92f3577b34da6a3ce929d0e0e4736-00f067aa0ba902b7-01") + try: + with logger.span("work"): + pass + finally: + reset_traceparent(token) + + (span,) = exporter.get_finished_spans() + assert format(span.context.trace_id, "032x") == "4bf92f3577b34da6a3ce929d0e0e4736" + + +def test_a_span_keeps_its_own_span_id_so_children_and_logs_can_point_at_it() -> None: + logger, _, exporter = _setup() + with logger.span("outer", span_id="00f067aa0ba902b7"): + pass + logger.close() + + (span,) = exporter.get_finished_spans() + assert format(span.context.span_id, "016x") == "00f067aa0ba902b7" + + +def test_a_span_id_that_is_not_hex_is_mapped_consistently() -> None: + logger, _, exporter = _setup() + parent_id = "3f0a6f4e-5c1b-4d0b-9d7a-0a1b2c3d4e5f" # e.g. a LangChain run id + with logger.span("parent", span_id=parent_id), logger.span("child"): + pass + logger.close() + + spans = _by_name(list(exporter.get_finished_spans())) + assert spans["child"].parent is not None + assert spans["child"].parent.span_id == spans["parent"].context.span_id + + +def test_failures_become_error_spans_with_an_error_type() -> None: + logger, _, exporter = _setup() + + with pytest.raises(ValueError), logger.span("risky"): + raise ValueError("bad input") + logger.action("look it up", tool="search", error="TimeoutError: too slow") + logger.close() + + spans = _by_name(list(exporter.get_finished_spans())) + risky, tool = spans["risky"], spans["execute_tool search"] + assert risky.status.status_code == StatusCode.ERROR + assert risky.status.description == "ValueError: bad input" + assert _attrs(risky)["error.type"] == "ValueError" + assert tool.status.status_code == StatusCode.ERROR + assert _attrs(tool)["error.type"] == "TimeoutError" + + +def test_state_diff_and_retry_count_are_exported() -> None: + logger, _, exporter = _setup() + state = {"n": 1} + + with logger.span("act", capture_state=lambda: state): + state["n"] = 2 + logger.action("go", tool="search", duration_ms=1) + logger.action("go", tool="search", duration_ms=1) + logger.close() + + spans = list(exporter.get_finished_spans()) + act = _by_name(spans)["act"] + (event,) = act.events + assert event.name == "logquill.state_diff" + assert event.attributes["logquill.state"] == '{"before":{"n":1},"after":{"n":2}}' # type: ignore[index] + retries = sorted( + _attrs(s).get("logquill.retry_count", 0) for s in spans if s.name == "execute_tool search" + ) + assert retries == [0, 1] + + +def test_records_that_are_not_spans_are_ignored() -> None: + logger, _, exporter = _setup() + + logger.info("just a log line") + logger.thought("thinking") + logger.observation("saw something") + logger.error("oops") + + assert exporter.get_finished_spans() == () + + +def test_prompt_and_completion_text_is_not_exported_unless_asked_for() -> None: + messages = [{"role": "user", "parts": [{"type": "text", "content": "secret prompt"}]}] + + off_logger, _, off = _setup() + off_logger.llm_call("chat", model="m", input_messages=messages) + on_logger, _, on = _setup(capture_content=True) + on_logger.llm_call("chat", model="m", input_messages=messages) + + (off_span,) = off.get_finished_spans() + (on_span,) = on.get_finished_spans() + assert "gen_ai.input.messages" not in _attrs(off_span) + assert "secret prompt" in _attrs(on_span)["gen_ai.input.messages"] + + +def test_content_capture_can_be_switched_on_from_the_environment( + monkeypatch: pytest.MonkeyPatch, +) -> None: + monkeypatch.setenv("OTEL_INSTRUMENTATION_GENAI_CAPTURE_MESSAGE_CONTENT", "true") + logger, _, exporter = _setup() + + logger.llm_call("chat", model="m", output_messages=[{"role": "assistant", "parts": []}]) + + (span,) = exporter.get_finished_spans() + assert "gen_ai.output.messages" in _attrs(span) + + +def test_the_legacy_naming_generation_is_selectable_and_never_mixed_with_the_new() -> None: + logger, _, exporter = _setup(semconv_version="legacy") + + logger.llm_call("chat", model="m", tokens_in=1, tokens_out=2, provider="openai") + + (span,) = exporter.get_finished_spans() + attributes = _attrs(span) + assert attributes["gen_ai.system"] == "openai" + assert attributes["gen_ai.usage.prompt_tokens"] == 1 + assert attributes["gen_ai.usage.completion_tokens"] == 2 + assert "gen_ai.provider.name" not in attributes + assert "gen_ai.usage.input_tokens" not in attributes + + +def test_the_environment_opt_in_selects_the_names(monkeypatch: pytest.MonkeyPatch) -> None: + monkeypatch.setenv("OTEL_SEMCONV_STABILITY_OPT_IN", "gen_ai_latest_experimental") + logger, _, exporter = _setup() + + logger.llm_call("chat", model="m", tokens_in=1) + + (span,) = exporter.get_finished_spans() + assert _attrs(span)["gen_ai.usage.input_tokens"] == 1 + + +def test_the_service_name_becomes_a_resource_attribute() -> None: + logger, _, exporter = _setup(service_name="checkout-agent") + + with logger.span("work"): + pass + + (span,) = exporter.get_finished_spans() + assert span.resource.attributes["service.name"] == "checkout-agent" + + +def test_flush_and_close_work_with_the_batch_processor() -> None: + exporter = InMemorySpanExporter() + transport = OTLPTransport(span_exporter=exporter) # batch, the default + logger = Logger("app", transports=[transport]) + + with logger.span("work"): + pass + logger.flush() + + assert len(exporter.get_finished_spans()) == 1 + logger.close() + + +def test_a_bad_processor_is_rejected_up_front() -> None: + with pytest.raises(ValueError, match="processor must be"): + OTLPTransport(span_exporter=InMemorySpanExporter(), processor="eventually") + + +def test_importing_logquill_does_not_import_opentelemetry() -> None: + result = subprocess.run( + [ + sys.executable, + "-c", + "import sys, logquill; print(any(m.startswith('opentelemetry') for m in sys.modules))", + ], + capture_output=True, + text=True, + check=True, + ) + + assert result.stdout.strip() == "False" + + +def test_a_missing_sdk_gives_an_install_hint(monkeypatch: pytest.MonkeyPatch) -> None: + monkeypatch.setitem(sys.modules, "opentelemetry.sdk.trace", None) + + with pytest.raises(ImportError, match=r"pip install logquill\[otel\]"): + OTLPTransport(span_exporter=InMemorySpanExporter()) diff --git a/tests/test_retry.py b/tests/test_retry.py new file mode 100644 index 0000000..b804af9 --- /dev/null +++ b/tests/test_retry.py @@ -0,0 +1,129 @@ +from __future__ import annotations + +import pytest + +from logquill import Logger +from logquill.retry import RetryTracker +from logquill.transports.transport import CollectingTransport + + +def test_the_first_attempt_is_zero_and_each_reopening_counts_up() -> None: + tracker = RetryTracker() + + assert [tracker.opened(None, "search") for _ in range(3)] == [0, 1, 2] + + +def test_a_success_ends_the_chain() -> None: + tracker = RetryTracker() + tracker.opened(None, "search") + tracker.opened(None, "search") + + tracker.succeeded(None, "search") + + assert tracker.opened(None, "search") == 0 + + +def test_calls_are_told_apart_by_span_tool_and_call_id() -> None: + tracker = RetryTracker() + + assert tracker.opened("s1", "search") == 0 + assert tracker.opened("s2", "search") == 0 + assert tracker.opened("s1", "fetch") == 0 + assert tracker.opened("s1", "search", "c1") == 0 + assert tracker.opened("s1", "search") == 1 + + +def test_the_tracker_is_bounded_and_forgets_the_oldest_calls() -> None: + tracker = RetryTracker(max_entries=3) + + for i in range(10): + tracker.opened(None, f"tool-{i}") + + assert len(tracker._attempts) == 3 + assert tracker.opened(None, "tool-0") == 0 # long forgotten, so not a retry + + +def test_max_entries_must_be_positive() -> None: + with pytest.raises(ValueError, match="max_entries"): + RetryTracker(max_entries=0) + + +def _logger() -> tuple[Logger, CollectingTransport]: + sink = CollectingTransport() + return Logger("app.agent", transports=[sink]), sink + + +def test_a_repeated_tool_action_gets_retry_count_stamped() -> None: + logger, sink = _logger() + + logger.action("call", tool="search") + logger.error("failed", kind="observation", tool="search", error="boom") + logger.action("call", tool="search") + logger.action("call", tool="search") + + assert [ + r["meta"].get("retry_count") for r in sink.records if r["meta"]["kind"] == "action" + ] == [ + None, + 1, + 2, + ] + + +def test_a_successful_observation_stops_later_calls_counting_as_retries() -> None: + logger, sink = _logger() + + for _ in range(3): + logger.action("call", tool="search") + logger.observation("done", tool="search") + + assert all("retry_count" not in r["meta"] for r in sink.records) + + +def test_an_observation_carrying_an_error_does_not_end_the_chain() -> None: + logger, sink = _logger() + + logger.action("call", tool="search") + logger.observation("nope", tool="search", error="TimeoutError: slow") + logger.action("call", tool="search") + + assert sink.records[-1]["meta"]["retry_count"] == 1 + + +def test_an_explicit_retry_count_is_kept() -> None: + logger, sink = _logger() + + logger.action("call", tool="search") + logger.action("call", tool="search", retry_count=7) + + assert sink.records[-1]["meta"]["retry_count"] == 7 + + +def test_actions_without_a_tool_are_never_counted() -> None: + logger, sink = _logger() + + logger.action("think") + logger.action("think") + + assert all("retry_count" not in r["meta"] for r in sink.records) + + +def test_the_same_tool_in_different_spans_is_not_a_retry() -> None: + logger, sink = _logger() + + with logger.span("step-1"): + logger.action("call", tool="search") + with logger.span("step-2"): + logger.action("call", tool="search") + + assert all("retry_count" not in r["meta"] for r in sink.records) + + +def test_tool_call_id_keeps_two_concurrent_calls_apart() -> None: + logger, sink = _logger() + + logger.action("call", tool="search", tool_call_id="a") + logger.action("call", tool="search", tool_call_id="b") + logger.action("call", tool="search", tool_call_id="a") + + assert [r["meta"].get("retry_count") for r in sink.records] == [None, None, 1] diff --git a/tests/test_semconv.py b/tests/test_semconv.py new file mode 100644 index 0000000..5e4d0cb --- /dev/null +++ b/tests/test_semconv.py @@ -0,0 +1,276 @@ +from __future__ import annotations + +import logging +import re +from pathlib import Path +from typing import Any + +import pytest + +from logquill import Level, semconv +from logquill.records import create_record + + +def _record(message: str = "m", *, llm: dict[str, Any] | None = None, **meta: Any) -> Any: + record = create_record( + level=Level.INFO, logger="app.agent", message=message, meta=meta, llm=llm + ) + record["timestamp"] = "2026-01-01T00:00:10.000Z" + return record + + +LATEST = semconv.resolve_convention("latest") +LEGACY = semconv.resolve_convention("legacy") + + +# --- the mapping --------------------------------------------------------------- + + +def test_an_llm_record_maps_to_a_chat_span_with_usage_attributes() -> None: + record = _record( + "chat", + kind="action", + provider="openai", + llm={ + "model": "m1", + "tokens_in": 12, + "tokens_out": 5, + "cost_usd": 0.5, + "latency_ms": 250, + "finish_reason": "stop", + }, + ) + + spec = semconv.llm_spec(record, LATEST) + + assert spec.name == "chat m1" + assert spec.kind == "client" + assert spec.duration_ms == 250 + assert spec.attributes == { + "gen_ai.operation.name": "chat", + "gen_ai.provider.name": "openai", + "gen_ai.request.model": "m1", + "gen_ai.usage.input_tokens": 12, + "gen_ai.usage.output_tokens": 5, + "gen_ai.response.finish_reasons": ["stop"], + "logquill.cost_usd": 0.5, + } + + +def test_text_completion_is_chosen_by_meta_operation() -> None: + spec = semconv.llm_spec(_record("c", operation="text_completion", llm={"model": "m1"}), LATEST) + + assert spec.name == "text_completion m1" + assert spec.attributes["gen_ai.operation.name"] == "text_completion" + + +def test_an_llm_call_without_a_model_or_provider_still_has_the_required_attributes() -> None: + spec = semconv.llm_spec(_record(llm={"tokens_in": 1}), LATEST, default_provider="acme") + + assert spec.name == "chat" + assert spec.attributes["gen_ai.provider.name"] == "acme" + assert "gen_ai.request.model" not in spec.attributes + + +def test_a_tool_action_maps_to_execute_tool() -> None: + spec = semconv.tool_spec( + _record( + "s", kind="action", tool="search", tool_call_id="c1", retry_count=2, duration_ms=30 + ), + LATEST, + ) + + assert spec.name == "execute_tool search" + assert spec.kind == "internal" + assert spec.duration_ms == 30 + assert spec.attributes == { + "gen_ai.operation.name": "execute_tool", + "gen_ai.tool.name": "search", + "gen_ai.tool.call.id": "c1", + "logquill.retry_count": 2, + } + + +def test_a_span_marked_as_an_agent_run_is_invoke_agent() -> None: + run = _record( + "run", kind="span", span_id="a" * 16, operation="invoke_agent", thread_id="conv-1" + ) + + spec = semconv.span_spec(run, LATEST) + + assert spec.name == "invoke_agent app.agent" + assert spec.kind == "client" + assert spec.attributes["gen_ai.operation.name"] == "invoke_agent" + assert spec.attributes["gen_ai.conversation.id"] == "conv-1" + + +def test_an_ordinary_span_keeps_its_name_even_when_it_is_the_outermost_of_a_run() -> None: + root = _record("run", kind="span", span_id="a" * 16, run_id="run-1") + nested = _record("step", kind="span", span_id="b" * 16, run_id="run-1", parent_span_id="a" * 16) + + for record, name in ((root, "run"), (nested, "step")): + spec = semconv.span_spec(record, LATEST) + assert spec.name == name + assert spec.kind == "internal" + assert "gen_ai.operation.name" not in spec.attributes + + +def test_naming_an_agent_makes_a_span_an_agent_run() -> None: + record = _record( + "go", kind="span", span_id="a" * 16, operation="invoke_agent", agent_name="planner" + ) + + spec = semconv.span_spec(record, LATEST) + + assert spec.name == "invoke_agent planner" + assert spec.attributes["gen_ai.agent.name"] == "planner" + + +def test_error_spans_carry_error_type() -> None: + spec = semconv.tool_spec( + _record("s", kind="action", tool="t", error="TimeoutError: too slow"), LATEST + ) + + assert spec.is_error is True + assert spec.attributes["error.type"] == "TimeoutError" + + +def test_state_diff_becomes_a_span_event() -> None: + spec = semconv.span_spec( + _record( + "s", kind="span", span_id="a" * 16, state_diff={"before": {"n": 1}, "after": {"n": 2}} + ), + LATEST, + ) + + ((name, attributes),) = spec.events + assert name == "logquill.state_diff" + assert attributes == {"logquill.state": '{"before":{"n":1},"after":{"n":2}}'} + + +def test_message_content_is_only_included_when_captured() -> None: + messages = [{"role": "user", "parts": [{"type": "text", "content": "hi"}]}] + record = _record("c", input_messages=messages, output_messages=[], llm={"model": "m"}) + + off = semconv.llm_spec(record, LATEST) + on = semconv.llm_spec(record, LATEST, capture_content=True) + + assert "gen_ai.input.messages" not in off.attributes + assert ( + on.attributes["gen_ai.input.messages"] + == '[{"role":"user","parts":[{"type":"text","content":"hi"}]}]' + ) + assert on.attributes["gen_ai.output.messages"] == "[]" + + +def test_only_the_named_record_kinds_are_recognized() -> None: + assert not semconv.is_span_record(_record("plain")) + assert not semconv.is_tool_record(_record("a", kind="action")) # no tool named + assert not semconv.is_tool_record(_record("a", kind="observation", tool="t")) + assert not semconv.is_llm_record(_record("plain")) + + +def test_log_attributes_for_an_llm_record() -> None: + record = _record(llm={"model": "m", "tokens_in": 3, "tokens_out": 4, "finish_reason": "length"}) + + assert semconv.log_attributes(record, LATEST) == { + "gen_ai.operation.name": "chat", + "gen_ai.request.model": "m", + "gen_ai.usage.input_tokens": 3, + "gen_ai.usage.output_tokens": 4, + "gen_ai.response.finish_reasons": ["length"], + } + assert semconv.log_attributes(_record("plain"), LATEST) == {} + + +def test_epoch_ns_reads_the_record_timestamp() -> None: + assert semconv.epoch_ns("1970-01-01T00:00:01.500Z") == 1_500_000_000 + + +# --- choosing the naming generation --------------------------------------------- + + +def test_the_legacy_generation_uses_the_older_names_for_provider_and_tokens() -> None: + spec = semconv.llm_spec( + _record(provider="openai", llm={"model": "m", "tokens_in": 1, "tokens_out": 2}), LEGACY + ) + + assert spec.attributes["gen_ai.system"] == "openai" + assert spec.attributes["gen_ai.usage.prompt_tokens"] == 1 + assert spec.attributes["gen_ai.usage.completion_tokens"] == 2 + assert "gen_ai.provider.name" not in spec.attributes + assert "gen_ai.usage.input_tokens" not in spec.attributes + + +def test_explicit_choice_wins_and_unknown_choice_lists_the_valid_ones() -> None: + assert ( + semconv.resolve_convention( + "legacy", {"OTEL_SEMCONV_STABILITY_OPT_IN": "gen_ai_latest_experimental"} + ) + is LEGACY + ) + with pytest.raises(ValueError, match="latest, legacy"): + semconv.resolve_convention("v9") + + +def test_the_environment_opt_in_selects_the_latest_names() -> None: + env = {"OTEL_SEMCONV_STABILITY_OPT_IN": "http, gen_ai_latest_experimental"} + + assert semconv.resolve_convention(None, env) is LATEST + assert semconv.resolve_convention(None, {}) is semconv._DEFAULT + + +def test_a_dual_emission_request_is_refused_with_a_warning( + caplog: pytest.LogCaptureFixture, +) -> None: + semconv._warned_opt_in.clear() + + with caplog.at_level(logging.WARNING, logger="logquill"): + convention = semconv.resolve_convention( + None, {"OTEL_SEMCONV_STABILITY_OPT_IN": "gen_ai/dup"} + ) + semconv.resolve_convention(None, {"OTEL_SEMCONV_STABILITY_OPT_IN": "gen_ai/dup"}) + + assert convention is semconv._DEFAULT + assert len([r for r in caplog.records if "never emits old and new" in r.getMessage()]) == 1 + + +@pytest.mark.parametrize( + "value, expected", [("true", True), ("TRUE", True), ("false", False), ("", False)] +) +def test_content_capture_follows_the_environment(value: str, expected: bool) -> None: + assert semconv.capture_content_enabled({semconv.CAPTURE_CONTENT_ENV: value}) is expected + assert semconv.capture_content_enabled({}) is False + + +# --- the convention strings live in exactly one file ------------------------------ + + +def test_no_genai_attribute_name_appears_outside_the_mapping_module() -> None: + package = Path(__file__).resolve().parent.parent / "logquill" + pattern = re.compile(r"gen_ai[._]|\bgen_ai\b") + offenders = [ + f"{path.relative_to(package)}:{number}" + for path in package.rglob("*.py") + if path.name != "semconv.py" + for number, line in enumerate(path.read_text(encoding="utf-8").splitlines(), start=1) + if pattern.search(line) + ] + + assert offenders == [] + + +def test_moving_to_a_newer_convention_is_an_edit_to_the_mapping_module_only( + monkeypatch: pytest.MonkeyPatch, +) -> None: + renamed = semconv.Convention( + version="9.9.9", + names={**semconv._LATEST.names, "input_tokens": "gen_ai.usage.renamed_input_tokens"}, + ) + monkeypatch.setitem(semconv._CONVENTIONS, "latest", renamed) + + spec = semconv.llm_spec( + _record(llm={"model": "m", "tokens_in": 4}), semconv.resolve_convention("latest") + ) + + assert spec.attributes["gen_ai.usage.renamed_input_tokens"] == 4 diff --git a/tests/test_state_diff.py b/tests/test_state_diff.py new file mode 100644 index 0000000..581d299 --- /dev/null +++ b/tests/test_state_diff.py @@ -0,0 +1,125 @@ +from __future__ import annotations + +from typing import Any + +from logquill import Logger +from logquill.transports.transport import CollectingTransport + + +def _logger() -> tuple[Logger, CollectingTransport]: + sink = CollectingTransport() + return Logger("app.agent", transports=[sink]), sink + + +def _span_meta(sink: CollectingTransport) -> dict[str, Any]: + return sink.records[-1]["meta"] + + +def test_state_diff_records_only_the_keys_that_changed() -> None: + logger, sink = _logger() + state = {"items": 1, "user": "ada", "step": "plan"} + + with logger.span("act", capture_state=lambda: state): + state["items"] = 2 + state["step"] = "act" + + assert _span_meta(sink)["state_diff"] == { + "before": {"items": 1, "step": "plan"}, + "after": {"items": 2, "step": "act"}, + } + + +def test_added_and_removed_keys_show_on_one_side_only() -> None: + logger, sink = _logger() + state: dict[str, Any] = {"gone": 1, "kept": 1} + + with logger.span("act", capture_state=lambda: state): + del state["gone"] + state["new"] = 2 + + assert _span_meta(sink)["state_diff"] == {"before": {"gone": 1}, "after": {"new": 2}} + + +def test_an_in_place_mutation_of_nested_state_is_seen() -> None: + logger, sink = _logger() + state = {"messages": ["a"]} + + with logger.span("act", capture_state=lambda: state): + state["messages"].append("b") + + assert _span_meta(sink)["state_diff"] == { + "before": {"messages": ["a"]}, + "after": {"messages": ["a", "b"]}, + } + + +def test_no_change_means_no_state_diff() -> None: + logger, sink = _logger() + state = {"n": 1} + + with logger.span("act", capture_state=lambda: state): + pass + + assert "state_diff" not in _span_meta(sink) + + +def test_non_dict_state_records_the_whole_values() -> None: + logger, sink = _logger() + counter = [0] + + with logger.span("act", capture_state=lambda: counter[0]): + counter[0] = 5 + + assert _span_meta(sink)["state_diff"] == {"before": 0, "after": 5} + + +def test_a_capture_function_that_raises_never_breaks_the_traced_code() -> None: + logger, sink = _logger() + + def broken() -> Any: + raise RuntimeError("cannot snapshot") + + with logger.span("act", capture_state=broken): + pass + + assert "state_diff" not in _span_meta(sink) + assert _span_meta(sink)["kind"] == "span" + + +def test_state_that_cannot_be_copied_is_skipped() -> None: + logger, sink = _logger() + + class Uncopyable: + def __deepcopy__(self, memo: Any) -> Any: + raise TypeError("no copies") + + with logger.span("act", capture_state=lambda: {"handle": Uncopyable()}): + pass + + assert "state_diff" not in _span_meta(sink) + + +def test_the_diff_is_recorded_even_if_the_block_raises() -> None: + logger, sink = _logger() + state = {"n": 1} + + try: + with logger.span("act", capture_state=lambda: state): + state["n"] = 2 + raise ValueError("boom") + except ValueError: + pass + + meta = _span_meta(sink) + assert meta["state_diff"] == {"before": {"n": 1}, "after": {"n": 2}} + assert meta["error"] == "ValueError: boom" + + +def test_an_explicit_state_diff_in_meta_wins() -> None: + logger, sink = _logger() + state = {"n": 1} + + with logger.span("act", capture_state=lambda: state, state_diff={"mine": True}): + state["n"] = 2 + + assert _span_meta(sink)["state_diff"] == {"mine": True}