diff --git a/CHANGELOG.md b/CHANGELOG.md index 34b3435..728f9a2 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -4,6 +4,28 @@ All notable changes to this project are documented in this file. ## Unreleased +- Auto-instrumentation, an OpenAI Agents SDK adapter, and MCP trace propagation: + - `logquill.instrument.anthropic(logger)` / `.openai(logger)` / `.litellm(logger)` + patch the Anthropic, OpenAI, and litellm Python SDKs so every LLM call they + make anywhere in the process emits a `logger.llm_call(...)` — no call-site + changes. Each has a matching `.uninstrument()`, lives behind its own + optional extra (`logquill[instrument-anthropic]`, `[instrument-openai]`, + `[instrument-litellm]`), and is imported lazily: `import logquill.instrument` + never imports a provider SDK. Streaming calls are a documented gap in this + release — passed through untouched rather than partially instrumented. + - `OpenAIAgentsAdapter` (`pip install logquill[openai-agents]`) maps the + OpenAI Agents SDK's `RunHooks` the same way the existing adapters map their + frameworks: every agent activation (including a handoff's target) is its + own `invoke_agent` span, and `on_llm_end` becomes a real `.llm_call()` with + token usage, so `OTLPTransport` exports a full run — spans, tokens, cost — + with zero manual logging calls. + - `logquill.mcp` (`propagate()`/`inbound()`) propagates trace context over an + MCP request's `_meta` field and stamps `meta.mcp.*` on records logged while + handling one, with no dependency on the `mcp` package itself. A client's + `run_id` rides along informationally as `meta.mcp.run_id`; it never + overrides the handling process's own `RunPlugin` run id. + - The record contract gained `meta.mcp.run_id`. + - 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 diff --git a/MIGRATING.md b/MIGRATING.md index 3bfa930..c8c424b 100644 --- a/MIGRATING.md +++ b/MIGRATING.md @@ -59,7 +59,7 @@ has them unless you use those. | `llm.cost_usd`, `llm.latency_ms` | `llm` block | number ≥ 0 | | `meta.retry_count` | `meta` | integer ≥ 0 | | `meta.state_diff` | `meta` | object | -| `meta.mcp.server`, `meta.mcp.tool` | `meta` | string | +| `meta.mcp.server`, `meta.mcp.tool`, `meta.mcp.run_id` | `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) | diff --git a/README.md b/README.md index ec2634b..a5d6ca4 100644 --- a/README.md +++ b/README.md @@ -26,8 +26,9 @@ for what's landed so far. - **Pluggable formatters** — `JSONFormatter` (default, machine-readable), `TextFormatter` (human-readable, for terminals), and `LogfmtFormatter` (`key=value`, the Heroku/Go convention); implement `format(record) -> str` for your own — see [Formatters](#formatters) - **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) +- **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]`), `AutoGenAdapter` (`pip install logquill[autogen]`), and `OpenAIAgentsAdapter` (`pip install logquill[openai-agents]`) — 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) +- **Auto-instrumentation & MCP** — `logquill.instrument.anthropic(logger)`/`.openai(logger)`/`.litellm(logger)` patch a provider SDK so every LLM call logs itself, no call-site changes; `OpenAIAgentsAdapter` covers the OpenAI Agents SDK; `logquill.mcp` propagates trace context over MCP requests and stamps `meta.mcp.*` — see [Auto-instrumentation and MCP](#auto-instrumentation-and-mcp) - **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 @@ -695,6 +696,29 @@ LangChain's own `run_id`/`parent_run_id` are written directly onto field renaming, not translation. `langchain-core` is never imported unless you import `logquill.adapters.langchain` yourself. +`OpenAIAgentsAdapter` maps the [OpenAI Agents SDK](https://openai.github.io/openai-agents-python/)'s +`RunHooks` the same way — pass an instance as `Runner.run(..., hooks=...)`: + +```bash +pip install logquill[openai-agents] +``` + +```python +from agents import Agent, Runner +from logquill import Logger, RunPlugin +from logquill.adapters.openai_agents import OpenAIAgentsAdapter + +log = Logger("app") +hooks = OpenAIAgentsAdapter(log.child("agent").use(RunPlugin())) +agent = Agent(name="assistant", instructions="...") +result = await Runner.run(agent, "hello", hooks=hooks) +``` + +Every agent activation (including a handoff's target) gets its own +`invoke_agent {agent.name}` span; `on_llm_end` becomes a real `.llm_call()` +with token usage, so `OTLPTransport` exports the whole run with cost and +token fields attached, with zero manual logging calls. + ### LangGraph LangGraph nodes execute as ordinary LangChain `Runnable`s, so @@ -948,6 +972,68 @@ Three things worth knowing: 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. +## Auto-instrumentation and MCP + +**Auto-instrumentation.** `logquill.instrument.(logger)` patches a +provider SDK so every LLM call it makes anywhere in the process logs itself — +no call-site changes, and no code path forgets to log. Each provider lives +behind its own extra and is imported lazily; `import logquill.instrument` +never imports a provider SDK: + +```bash +pip install logquill[instrument-anthropic] # or [instrument-openai] / [instrument-litellm] +``` + +```python +import logquill.instrument as instrument +from logquill import Logger + +logger = Logger("app") +instrument.anthropic(logger) + +# unchanged call site — no logger, no llm_call(), nothing added here +response = client.messages.create(model="...", max_tokens=100, messages=[...]) + +instrument.anthropic.uninstrument() # each instrumenter has a matching uninstrument() +``` + +`logquill.instrument.litellm` is the broadest one: since litellm itself fans +out to 100+ providers behind one interface, instrumenting it covers all of +them from this one call. Calling an instrumenter twice without an +intervening `.uninstrument()` raises, so a call is never wrapped twice. +**Streaming calls are a known, documented gap** in this release — a +`stream=True` call passes through untouched rather than being partially +instrumented; see each provider module's docstring +(`logquill.instrument.anthropic`, `.openai`, `.litellm`). + +**MCP (Model Context Protocol).** `logquill.mcp` propagates trace context +over an MCP request and stamps `meta.mcp.*` on records logged while handling +one — with no dependency on the `mcp` package itself, since both sides just +build or read a plain dict (exactly what an MCP client's `meta=` argument and +a server's inbound `_meta` are): + +```python +from logquill import Logger, TraceContextPlugin +from logquill.mcp import inbound, propagate + +# client, before calling a tool: +await session.call_tool("search", {"q": "..."}, meta=propagate(run_id=agent_run_id)) + +# server, inside the tool handler: +server_log = Logger("mcp.files", plugins=[TraceContextPlugin()]) +with inbound(ctx.request_context.meta, server="files", tool="search"): + server_log.info("handling search") # meta.trace_id matches the client's own + # call, and meta.mcp = {"server": "files", + # "tool": "search", "run_id": agent_run_id} +``` + +`propagate()`'s keys are namespaced under `logquill/`, an unreserved `_meta` +prefix per the MCP spec (reserved prefixes contain a +`modelcontextprotocol`/`mcp` label). A client's `run_id` never overrides the +server's own `RunPlugin` run id — it rides along informationally, under +`meta.mcp.run_id`, since the server handling a tool call is its own run, not +a continuation of the client's. + ## Async dispatch & serverless safety By default, every log call dispatches to its transports synchronously — a diff --git a/logquill/adapters/openai_agents.py b/logquill/adapters/openai_agents.py new file mode 100644 index 0000000..de89de9 --- /dev/null +++ b/logquill/adapters/openai_agents.py @@ -0,0 +1,156 @@ +"""`OpenAIAgentsAdapter` — translates the OpenAI Agents SDK's `RunHooks` +lifecycle callbacks into `.thought()/.action()/.observation()/.decision()` +and `span()` calls, per `LogQuillAdapter`'s "thin mapping, not a +reimplementation" contract. + +Requires the optional `openai-agents` package (`pip install +logquill[openai-agents]`), imported lazily — importing this module before +it's installed raises an actionable `ImportError`, and importing `logquill` +itself never imports `agents`. +""" + +from __future__ import annotations + +import time +from typing import Any + +try: + from agents.lifecycle import RunHooks +except ImportError as exc: + raise ImportError( + "logquill.adapters.openai_agents requires the optional `openai-agents` " + "package — install with `pip install logquill[openai-agents]`." + ) from exc + +from logquill.adapters.base import LogQuillAdapter +from logquill.logger import Logger +from logquill.span import SpanContext, new_span_id + + +def _agent_name(agent: Any) -> str: + name = getattr(agent, "name", None) + return name if isinstance(name, str) and name else "agent" + + +def _tool_name(tool: Any) -> str: + name = getattr(tool, "name", None) + return name if isinstance(name, str) and name else "tool" + + +def _usage_fields(response: Any) -> dict[str, Any]: + usage = getattr(response, "usage", None) + return { + "tokens_in": getattr(usage, "input_tokens", None), + "tokens_out": getattr(usage, "output_tokens", None), + } + + +# `type: ignore[misc]` — `RunHooks` types as `Any` whenever `openai-agents` +# isn't installed in the environment running mypy (an optional dependency, +# never in this project's `dev` extra — see pyproject.toml), and mypy refuses +# to let a class subclass something typed `Any`. With the real package +# installed, this subclasses the genuine `RunHooks` and the ignore is inert. +class OpenAIAgentsAdapter(LogQuillAdapter, RunHooks): # type: ignore[misc] + """Pass an instance as `Runner.run(..., hooks=...)` — no separate + registration step needed. + + | OpenAI Agents SDK hook | LogQuill call | + |------------------------------------|------------------------------------| + | `on_agent_start` / `on_agent_end` | opens/closes a `span()` | + | `on_llm_end` | `.llm_call()` with token usage and elapsed time | + | `on_tool_start` / `on_tool_end` | `.action(tool=...)` / `.observation(tool=...)` | + | `on_handoff` | `.decision()` naming `from_agent`/`to_agent` | + + Every agent activation gets its own `invoke_agent {agent.name}` span, not + just the outermost — a handoff's target is its own nested span, never a + reuse of its predecessor's. `on_llm_end`'s token counts and `duration_ms` + (timed since the matching `on_llm_start`) come straight from + `ModelResponse.usage`. Tagging `.action()`/`.observation()` with `tool=` + is what makes them show up as `execute_tool` spans under `OTLPTransport` + and get `retry_count` tracked automatically. + + A run's `RunContextWrapper`/`AgentHookContext` has no id of its own to + key spans by (unlike LangChain's `run_id`), so this adapter uses the + context object's identity — one `Runner.run()` call reuses one context + throughout, so nesting (including across a handoff) still comes out + right. The one sharp edge: if the *same* tool is invoked twice + concurrently in one run, their `on_tool_start`/`on_tool_end` pairs can + cross-attribute duration, since the SDK's hooks don't hand this adapter + a per-call id to key on — sequential calls (the common case, and the + SDK's default) are unaffected. + """ + + def __init__(self, agent_log: Logger) -> None: + LogQuillAdapter.__init__(self, agent_log) + self._agent_spans: dict[int, list[SpanContext]] = {} + self._llm_starts: dict[int, float] = {} + self._tool_spans: dict[tuple[int, str], list[float]] = {} + + def _stack(self, context: Any) -> list[SpanContext]: + return self._agent_spans.setdefault(id(context), []) + + async def on_agent_start(self, context: Any, agent: Any) -> None: + """Opens a span for this agent's activation, nested under whichever + agent (if any) is already active in this run — including a handoff's + target, which gets its own nested span rather than reusing its + predecessor's.""" + stack = self._stack(context) + parent_span_id = stack[-1]._span_id if stack else None + span = self.log.span( + _agent_name(agent), + span_id=new_span_id(), + parent_span_id=parent_span_id, + operation="invoke_agent", + agent_name=_agent_name(agent), + ) + span.__enter__() + stack.append(span) + + async def on_agent_end(self, context: Any, agent: Any, output: Any) -> None: + """Closes the span opened by the matching `on_agent_start`.""" + stack = self._stack(context) + if stack: + stack.pop().__exit__(None, None, None) + if not stack: + self._agent_spans.pop(id(context), None) + + async def on_llm_start( + self, context: Any, agent: Any, system_prompt: str | None, input_items: Any + ) -> None: + """Records the call's start time, for the matching `on_llm_end`'s + `duration_ms`.""" + self._llm_starts[id(context)] = time.monotonic() + + async def on_llm_end(self, context: Any, agent: Any, response: Any) -> None: + """Emits `.llm_call()` with token usage and elapsed time.""" + start = self._llm_starts.pop(id(context), None) + duration_ms = round((time.monotonic() - start) * 1000, 3) if start is not None else None + model = getattr(agent, "model", None) + self.log.llm_call( + model=model if isinstance(model, str) else None, + latency_ms=duration_ms, + **_usage_fields(response), + ) + + async def on_tool_start(self, context: Any, agent: Any, tool: Any) -> None: + """Emits `.action(tool=...)` and records the call's start time.""" + name = _tool_name(tool) + self._tool_spans.setdefault((id(context), name), []).append(time.monotonic()) + self.log.action(f"call {name}", tool=name, agent_name=_agent_name(agent)) + + async def on_tool_end(self, context: Any, agent: Any, tool: Any, result: Any) -> None: + """Emits `.observation(tool=...)` with `duration_ms` measured since + the matching `on_tool_start` — see the class docstring for the one + case (concurrent calls of the same tool) this timing can misattribute.""" + name = _tool_name(tool) + starts = self._tool_spans.get((id(context), name)) + start = starts.pop() if starts else None + duration_ms = round((time.monotonic() - start) * 1000, 3) if start is not None else None + self.log.observation(f"{name} done", tool=name, duration_ms=duration_ms) + + async def on_handoff(self, context: Any, from_agent: Any, to_agent: Any) -> None: + """Emits `.decision()` naming the handoff's source and destination + agents.""" + self.log.decision( + "handoff", from_agent=_agent_name(from_agent), to_agent=_agent_name(to_agent) + ) diff --git a/logquill/instrument/__init__.py b/logquill/instrument/__init__.py new file mode 100644 index 0000000..0479a9c --- /dev/null +++ b/logquill/instrument/__init__.py @@ -0,0 +1,30 @@ +"""One-line instrumentation for provider SDKs: patch a client library so every +LLM call it makes emits `logger.llm_call(...)` — token counts, model, finish +reason, latency — with no call-site changes anywhere in your code. + + import logquill.instrument as instrument + + instrument.anthropic(agent_log) + # every client.messages.create(...) anywhere in the process now logs itself + ... + instrument.anthropic.uninstrument() + +Each provider lives behind its own optional extra (`pip install +logquill[instrument-anthropic]`, `[instrument-openai]`, `[instrument-litellm]`) +and is imported lazily, inside the call to `instrument.(logger)` — +**importing `logquill`, or `logquill.instrument`, never imports a provider +SDK.** Calling an instrumenter twice without a matching `.uninstrument()` +raises, so a call is never wrapped twice. + +Every instrumenter currently covers non-streaming calls only; a streaming +call is passed through untouched rather than partially instrumented — see +each provider module's docstring. +""" + +from __future__ import annotations + +from logquill.instrument.anthropic import anthropic +from logquill.instrument.litellm import litellm +from logquill.instrument.openai import openai + +__all__ = ["anthropic", "litellm", "openai"] diff --git a/logquill/instrument/_util.py b/logquill/instrument/_util.py new file mode 100644 index 0000000..22926b6 --- /dev/null +++ b/logquill/instrument/_util.py @@ -0,0 +1,87 @@ +"""Shared plumbing for the `logquill.instrument.*` patchers: swap one method +or function out for a wrapper, remember how to put it back, and never let a +bug in the wrapper break the real call it wraps. +""" + +from __future__ import annotations + +import logging +from dataclasses import dataclass +from typing import Any, Callable, TypeVar + +_logger = logging.getLogger("logquill") + +T = TypeVar("T") + + +@dataclass +class Patch: + """One reversible attribute swap.""" + + target: Any + name: str + original: Any + + def revert(self) -> None: + """Puts `original` back.""" + setattr(self.target, self.name, self.original) + + +def apply(target: Any, name: str, build_wrapper: Callable[[Any], Any]) -> Patch: + """Replaces `target.name` with `build_wrapper(current_value)` and returns + a `Patch` that can put it back.""" + original = getattr(target, name) + setattr(target, name, build_wrapper(original)) + return Patch(target, name, original) + + +def record_safely(do_record: Callable[[], None]) -> None: + """Runs `do_record` (which logs one call's fields), swallowing and + logging any exception. A bug in instrumentation must never surface as a + failure of the API call it's instrumenting — the whole point of + `instrument()` is that call sites don't change, including their error + handling.""" + try: + do_record() + except Exception: + _logger.exception("logquill.instrument: failed to record a call") + + +class Instrumenter: + """One `logquill.instrument.` entry: calling it patches, + `.uninstrument()` puts everything back. Calling it twice without an + intervening `.uninstrument()` raises, so a wrapper is never wrapped + twice — each instrumented method would otherwise fire this provider's + `llm_call` twice per real call, once per layer. + """ + + def __init__(self, name: str, apply_patches: Callable[..., list[Patch]]) -> None: + """`name` is used only in the "already active" error message. + `apply_patches` does the actual patching (typically importing the + provider SDK lazily and calling `apply()` once per method) and + returns the `Patch`es to revert later.""" + self._name = name + self._apply_patches = apply_patches + self._patches: list[Patch] = [] + + def __call__(self, *args: Any, **kwargs: Any) -> None: + """Patches the provider SDK. See the concrete module (e.g. + `logquill.instrument.anthropic`) for what arguments it takes.""" + if self._patches: + raise RuntimeError( + f"logquill.instrument.{self._name} is already active — call " + f".uninstrument() first if you want to change its logger or options" + ) + self._patches = self._apply_patches(*args, **kwargs) + + def uninstrument(self) -> None: + """Restores every method this patched. Safe to call even when not + currently instrumented.""" + while self._patches: + self._patches.pop().revert() + + @property + def active(self) -> bool: + """Whether `instrument()` has been called without a matching + `.uninstrument()` since.""" + return bool(self._patches) diff --git a/logquill/instrument/anthropic.py b/logquill/instrument/anthropic.py new file mode 100644 index 0000000..3ecf1f9 --- /dev/null +++ b/logquill/instrument/anthropic.py @@ -0,0 +1,83 @@ +"""`logquill.instrument.anthropic(logger)` — patches the Anthropic Python +SDK so every `client.messages.create(...)` call (sync and async) emits a +`logger.llm_call(...)` with no call-site changes. + +Requires the optional `anthropic` package (`pip install +logquill[instrument-anthropic]`), imported lazily inside `instrument()` — +importing `logquill.instrument` never imports it. + +**Scope**: only non-streaming calls (`stream` unset or `False`) are +instrumented. A streaming call (`stream=True`, or `.stream()`) is passed +through untouched — accumulating usage across a stream safely, without +disturbing the SDK's own streaming ergonomics, needs more design than this +pass covers; it's a tracked gap, not a silent one (see the module docstring +in `logquill/instrument/__init__.py`). +""" + +from __future__ import annotations + +import functools +import time +from typing import Any + +from logquill.instrument._util import Instrumenter, Patch, apply, record_safely +from logquill.logger import Logger + + +def _record(logger: Logger, kwargs: dict[str, Any], response: Any, duration_ms: float) -> None: + usage = getattr(response, "usage", None) + logger.llm_call( + model=getattr(response, "model", None) or kwargs.get("model"), + tokens_in=getattr(usage, "input_tokens", None), + tokens_out=getattr(usage, "output_tokens", None), + latency_ms=duration_ms, + finish_reason=getattr(response, "stop_reason", None), + provider="anthropic", + ) + + +def _wrap_sync(original: Any, logger: Logger) -> Any: + @functools.wraps(original) + def wrapper(self: Any, *args: Any, **kwargs: Any) -> Any: + if kwargs.get("stream"): + return original(self, *args, **kwargs) + start = time.monotonic() + response = original(self, *args, **kwargs) + record_safely(lambda: _record(logger, kwargs, response, (time.monotonic() - start) * 1000)) + return response + + return wrapper + + +def _wrap_async(original: Any, logger: Logger) -> Any: + @functools.wraps(original) + async def wrapper(self: Any, *args: Any, **kwargs: Any) -> Any: + if kwargs.get("stream"): + return await original(self, *args, **kwargs) + start = time.monotonic() + response = await original(self, *args, **kwargs) + record_safely(lambda: _record(logger, kwargs, response, (time.monotonic() - start) * 1000)) + return response + + return wrapper + + +def _apply(logger: Logger) -> list[Patch]: + try: + from anthropic.resources.messages import messages as messages_module + except ImportError as exc: + raise ImportError( + "logquill.instrument.anthropic requires the optional `anthropic` package — " + "install with `pip install logquill[instrument-anthropic]`." + ) from exc + return [ + apply(messages_module.Messages, "create", lambda original: _wrap_sync(original, logger)), + apply( + messages_module.AsyncMessages, "create", lambda original: _wrap_async(original, logger) + ), + ] + + +anthropic = Instrumenter("anthropic", _apply) +"""Call `anthropic(logger)` to instrument, `anthropic.uninstrument()` to +undo it. See the module docstring for what's covered.""" diff --git a/logquill/instrument/litellm.py b/logquill/instrument/litellm.py new file mode 100644 index 0000000..ff6514c --- /dev/null +++ b/logquill/instrument/litellm.py @@ -0,0 +1,91 @@ +"""`logquill.instrument.litellm(logger)` — patches `litellm.completion`/ +`litellm.acompletion` so every call emits a `logger.llm_call(...)` with no +call-site changes. Since litellm itself fans out to 100+ providers behind one +interface, this one instrumentation covers all of them. + +Requires the optional `litellm` package (`pip install +logquill[instrument-litellm]`), imported lazily inside `instrument()`. + +**Scope**: only non-streaming calls are instrumented; see +`logquill.instrument.anthropic` for why streaming is a tracked gap rather +than a silent one. +""" + +from __future__ import annotations + +import functools +import time +from typing import Any + +from logquill.instrument._util import Instrumenter, Patch, apply, record_safely +from logquill.logger import Logger + + +def _provider(response: Any) -> str | None: + hidden = getattr(response, "_hidden_params", None) + provider = hidden.get("custom_llm_provider") if isinstance(hidden, dict) else None + return provider if isinstance(provider, str) else None + + +def _finish_reason(response: Any) -> str | None: + choices = getattr(response, "choices", None) + if not choices: + return None + return getattr(choices[0], "finish_reason", None) + + +def _record(logger: Logger, kwargs: dict[str, Any], response: Any, duration_ms: float) -> None: + usage = getattr(response, "usage", None) + logger.llm_call( + model=getattr(response, "model", None) or kwargs.get("model"), + tokens_in=getattr(usage, "prompt_tokens", None), + tokens_out=getattr(usage, "completion_tokens", None), + latency_ms=duration_ms, + finish_reason=_finish_reason(response), + provider=_provider(response), + ) + + +def _wrap_sync(original: Any, logger: Logger) -> Any: + @functools.wraps(original) + def wrapper(*args: Any, **kwargs: Any) -> Any: + if kwargs.get("stream"): + return original(*args, **kwargs) + start = time.monotonic() + response = original(*args, **kwargs) + record_safely(lambda: _record(logger, kwargs, response, (time.monotonic() - start) * 1000)) + return response + + return wrapper + + +def _wrap_async(original: Any, logger: Logger) -> Any: + @functools.wraps(original) + async def wrapper(*args: Any, **kwargs: Any) -> Any: + if kwargs.get("stream"): + return await original(*args, **kwargs) + start = time.monotonic() + response = await original(*args, **kwargs) + record_safely(lambda: _record(logger, kwargs, response, (time.monotonic() - start) * 1000)) + return response + + return wrapper + + +def _apply(logger: Logger) -> list[Patch]: + try: + import litellm as litellm_module + except ImportError as exc: + raise ImportError( + "logquill.instrument.litellm requires the optional `litellm` package — " + "install with `pip install logquill[instrument-litellm]`." + ) from exc + return [ + apply(litellm_module, "completion", lambda original: _wrap_sync(original, logger)), + apply(litellm_module, "acompletion", lambda original: _wrap_async(original, logger)), + ] + + +litellm = Instrumenter("litellm", _apply) +"""Call `litellm(logger)` to instrument, `litellm.uninstrument()` to undo it. +See the module docstring for what's covered.""" diff --git a/logquill/instrument/openai.py b/logquill/instrument/openai.py new file mode 100644 index 0000000..821677f --- /dev/null +++ b/logquill/instrument/openai.py @@ -0,0 +1,124 @@ +"""`logquill.instrument.openai(logger)` — patches the OpenAI Python SDK so +every `client.chat.completions.create(...)` and `client.responses.create(...)` +call (sync and async) emits a `logger.llm_call(...)` with no call-site changes. + +Requires the optional `openai` package (`pip install +logquill[instrument-openai]`), imported lazily inside `instrument()`. + +**Scope**: only non-streaming calls are instrumented; see +`logquill.instrument.anthropic` for why streaming is a tracked gap rather +than a silent one. +""" + +from __future__ import annotations + +import functools +import time +from typing import Any + +from logquill.instrument._util import Instrumenter, Patch, apply, record_safely +from logquill.logger import Logger + + +def _chat_finish_reason(response: Any) -> str | None: + choices = getattr(response, "choices", None) + if not choices: + return None + return getattr(choices[0], "finish_reason", None) + + +def _record_chat(logger: Logger, kwargs: dict[str, Any], response: Any, duration_ms: float) -> None: + usage = getattr(response, "usage", None) + logger.llm_call( + model=getattr(response, "model", None) or kwargs.get("model"), + tokens_in=getattr(usage, "prompt_tokens", None), + tokens_out=getattr(usage, "completion_tokens", None), + latency_ms=duration_ms, + finish_reason=_chat_finish_reason(response), + provider="openai", + ) + + +def _responses_finish_reason(response: Any) -> str | None: + incomplete = getattr(response, "incomplete_details", None) + reason = getattr(incomplete, "reason", None) if incomplete is not None else None + return reason if isinstance(reason, str) else getattr(response, "status", None) + + +def _record_responses( + logger: Logger, kwargs: dict[str, Any], response: Any, duration_ms: float +) -> None: + usage = getattr(response, "usage", None) + logger.llm_call( + model=getattr(response, "model", None) or kwargs.get("model"), + tokens_in=getattr(usage, "input_tokens", None), + tokens_out=getattr(usage, "output_tokens", None), + latency_ms=duration_ms, + finish_reason=_responses_finish_reason(response), + provider="openai", + ) + + +def _wrap_sync(original: Any, logger: Logger, record: Any) -> Any: + @functools.wraps(original) + def wrapper(self: Any, *args: Any, **kwargs: Any) -> Any: + if kwargs.get("stream"): + return original(self, *args, **kwargs) + start = time.monotonic() + response = original(self, *args, **kwargs) + record_safely(lambda: record(logger, kwargs, response, (time.monotonic() - start) * 1000)) + return response + + return wrapper + + +def _wrap_async(original: Any, logger: Logger, record: Any) -> Any: + @functools.wraps(original) + async def wrapper(self: Any, *args: Any, **kwargs: Any) -> Any: + if kwargs.get("stream"): + return await original(self, *args, **kwargs) + start = time.monotonic() + response = await original(self, *args, **kwargs) + record_safely(lambda: record(logger, kwargs, response, (time.monotonic() - start) * 1000)) + return response + + return wrapper + + +def _apply(logger: Logger) -> list[Patch]: + try: + from openai.resources import responses as responses_module + from openai.resources.chat import completions as completions_module + except ImportError as exc: + raise ImportError( + "logquill.instrument.openai requires the optional `openai` package — " + "install with `pip install logquill[instrument-openai]`." + ) from exc + return [ + apply( + completions_module.Completions, + "create", + lambda original: _wrap_sync(original, logger, _record_chat), + ), + apply( + completions_module.AsyncCompletions, + "create", + lambda original: _wrap_async(original, logger, _record_chat), + ), + apply( + responses_module.Responses, + "create", + lambda original: _wrap_sync(original, logger, _record_responses), + ), + apply( + responses_module.AsyncResponses, + "create", + lambda original: _wrap_async(original, logger, _record_responses), + ), + ] + + +openai = Instrumenter("openai", _apply) +"""Call `openai(logger)` to instrument, `openai.uninstrument()` to undo it. +Covers both the Chat Completions and Responses APIs. See the module +docstring for what's covered.""" diff --git a/logquill/mcp.py b/logquill/mcp.py new file mode 100644 index 0000000..b8a6eef --- /dev/null +++ b/logquill/mcp.py @@ -0,0 +1,84 @@ +"""Trace-context propagation over MCP (Model Context Protocol) requests, plus +stamping `meta.mcp.*` on records logged while handling one. + +No dependency on the `mcp` package: both helpers just build or read a plain +`dict[str, str]`, since that is exactly what an MCP client's `meta=` argument +and a server handler's inbound request `_meta` are — the MCP spec's `_meta` +field is a generic metadata bag any request, notification or result may +carry (see the spec's "General fields" / `_meta` section). The client passes +`propagate()`'s result as `meta=` on its call; the server passes the request's +inbound `_meta` (in the Python SDK, `ctx.request_context.meta`) to `inbound()`. + +Keys are namespaced under `logquill/`, a plain, unreserved `_meta` prefix — +the spec reserves only prefixes containing a `modelcontextprotocol`/`mcp` +label, and `logquill/traceparent` doesn't. +""" + +from __future__ import annotations + +from collections.abc import Iterator, Mapping +from contextlib import contextmanager +from typing import Any + +from logquill.context import bind_context +from logquill.plugins.trace_context_plugin import ( + generate_trace_id, + reset_traceparent, + set_traceparent, +) + +TRACEPARENT_KEY = "logquill/traceparent" +RUN_ID_KEY = "logquill/run-id" + + +def propagate(*, trace_id: str | None = None, run_id: str | None = None) -> dict[str, str]: + """Client side: build the `_meta` dict to pass on an outbound MCP request + so the server shares this call's trace — e.g. + + await session.call_tool("search", args, meta=propagate(run_id=agent_run_id)) + + Carries a W3C `traceparent` for `trace_id` (an active one if you don't + pass one and none is given, else freshly generated — the same resolution + `TraceContextPlugin` itself uses for an outbound call), so a server whose + own logger has a `TraceContextPlugin` picks up the *same* `trace_id`. + `run_id`, if given, rides along too — the server's `inbound()` records it + under `meta.mcp.run_id`, informationally; it never overrides the server's + own `RunPlugin` run id, since that would conflate two different runs. + """ + span_id = generate_trace_id()[:16] + meta = {TRACEPARENT_KEY: f"00-{trace_id or generate_trace_id()}-{span_id}-01"} + if run_id: + meta[RUN_ID_KEY] = run_id + return meta + + +@contextmanager +def inbound(meta: Mapping[str, Any] | None, *, server: str, tool: str) -> Iterator[None]: + """Server side: wraps handling one MCP request (typically a whole tool + call). `meta` is the inbound request's `_meta` field exactly as the SDK + hands it to you — in the official Python SDK, `ctx.request_context.meta` + from inside a tool handler; `None` (no `_meta` sent) is fine. + + For the duration of the block: + + - if `meta` carries a `propagate()`-built `traceparent`, `TraceContextPlugin` + on any logger used inside resolves to that same `trace_id` — the + mechanism is `set_traceparent()`, the same one HTTP middleware uses for + an inbound header. + - every record logged inside picks up `meta.mcp.server`/`meta.mcp.tool` + (this call's own `server`/`tool`) and, if the client sent one, + `meta.mcp.run_id` — via `bind_context()`, so nothing needs passing + through call signatures by hand. + """ + meta = meta or {} + traceparent = meta.get(TRACEPARENT_KEY) + token = set_traceparent(traceparent if isinstance(traceparent, str) else None) + mcp_context: dict[str, Any] = {"server": server, "tool": tool} + run_id = meta.get(RUN_ID_KEY) + if isinstance(run_id, str) and run_id: + mcp_context["run_id"] = run_id + try: + with bind_context(mcp=mcp_context): + yield + finally: + reset_traceparent(token) diff --git a/pyproject.toml b/pyproject.toml index 37104e4..e8c18dd 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -45,6 +45,10 @@ otel = [ "opentelemetry-sdk>=1.39", "opentelemetry-exporter-otlp-proto-http>=1.39", ] +instrument-anthropic = ["anthropic>=0.40"] +instrument-openai = ["openai>=1.50"] +instrument-litellm = ["litellm>=1.50"] +openai-agents = ["openai-agents>=0.0.3"] postgres = ["psycopg2-binary>=2.9"] mysql = ["pymysql>=1.1"] mongodb = ["pymongo>=4.6"] @@ -160,6 +164,13 @@ module = ["opentelemetry", "opentelemetry.*"] ignore_missing_imports = true follow_imports = "skip" +# Every provider SDK / agent framework `logquill.instrument`/`logquill.adapters` +# touches is optional and imported lazily, for the same reason as above. +[[tool.mypy.overrides]] +module = ["anthropic", "anthropic.*", "openai", "openai.*", "litellm", "agents", "agents.*"] +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 fb0b415..e43c265 100644 --- a/schema/golden_records.json +++ b/schema/golden_records.json @@ -396,6 +396,23 @@ "model": "m" } } + }, + { + "name": "MCP fields with a client-side run_id", + "record": { + "schema_version": "2.0", + "timestamp": "2026-01-01T00:00:00.000Z", + "level": "INFO", + "logger": "mcp.files", + "message": "handling search", + "meta": { + "mcp": { + "server": "files", + "tool": "search", + "run_id": "run-1" + } + } + } } ], "invalid": [ @@ -898,6 +915,22 @@ } }, "why": "input_messages is an array" + }, + { + "name": "mcp.run_id is a number", + "record": { + "schema_version": "2.0", + "timestamp": "2026-01-01T00:00:00.000Z", + "level": "INFO", + "logger": "a", + "message": "x", + "meta": { + "mcp": { + "run_id": 1 + } + } + }, + "why": "mcp.run_id is a string" } ], "legacy": [ diff --git a/schema/record.schema.json b/schema/record.schema.json index 8e904a6..339b7f5 100644 --- a/schema/record.schema.json +++ b/schema/record.schema.json @@ -160,6 +160,10 @@ }, "tool": { "type": "string" + }, + "run_id": { + "type": "string", + "description": "Informational: which client-side agent run triggered this MCP request. Never overrides this process's own run_id." } } }, diff --git a/tests/test_adapters/test_openai_agents_adapter.py b/tests/test_adapters/test_openai_agents_adapter.py new file mode 100644 index 0000000..ed532d5 --- /dev/null +++ b/tests/test_adapters/test_openai_agents_adapter.py @@ -0,0 +1,189 @@ +from __future__ import annotations + +import asyncio +import importlib +import sys +import types +from types import ModuleType + +import pytest + +from logquill.logger import Logger +from logquill.transports.transport import CollectingTransport + + +class FakeUsage: + def __init__(self, input_tokens: int, output_tokens: int) -> None: + self.input_tokens = input_tokens + self.output_tokens = output_tokens + + +class FakeResponse: + def __init__(self, input_tokens: int = 10, output_tokens: int = 4) -> None: + self.usage = FakeUsage(input_tokens, output_tokens) + + +class FakeAgent: + def __init__(self, name: str, model: str = "gpt-fake") -> None: + self.name = name + self.model = model + + +class FakeTool: + def __init__(self, name: str) -> None: + self.name = name + + +def _install_fake_agents(monkeypatch: pytest.MonkeyPatch) -> None: + # Same `sys.modules` injection pattern as the other framework adapters + # (see `tests/test_adapters/test_langchain_adapter.py`), so this + # exercises the real adapter against a stand-in `RunHooks` without + # requiring the actual, heavier `openai-agents` package. + class FakeRunHooks: + pass + + lifecycle_module = types.ModuleType("agents.lifecycle") + lifecycle_module.RunHooks = FakeRunHooks # type: ignore[attr-defined] + agents_module = types.ModuleType("agents") + agents_module.lifecycle = lifecycle_module # type: ignore[attr-defined] + + monkeypatch.setitem(sys.modules, "agents", agents_module) + monkeypatch.setitem(sys.modules, "agents.lifecycle", lifecycle_module) + + +def _load_adapter_module(monkeypatch: pytest.MonkeyPatch) -> ModuleType: + _install_fake_agents(monkeypatch) + monkeypatch.delitem(sys.modules, "logquill.adapters.openai_agents", raising=False) + return importlib.import_module("logquill.adapters.openai_agents") + + +def test_raises_an_actionable_error_without_openai_agents( + monkeypatch: pytest.MonkeyPatch, +) -> None: + monkeypatch.setitem(sys.modules, "agents", None) + monkeypatch.delitem(sys.modules, "logquill.adapters.openai_agents", raising=False) + + with pytest.raises(ImportError, match=r"logquill\[openai-agents\]"): + importlib.import_module("logquill.adapters.openai_agents") + + +def test_a_run_with_a_tool_call_and_an_llm_call_produces_a_full_span_tree( + monkeypatch: pytest.MonkeyPatch, +) -> None: + module = _load_adapter_module(monkeypatch) + OpenAIAgentsAdapter = module.OpenAIAgentsAdapter + + sink = CollectingTransport() + logger = Logger("app.agent", transports=[sink]) + adapter = OpenAIAgentsAdapter(logger) + context = object() + agent = FakeAgent("planner") + + async def run() -> None: + await adapter.on_agent_start(context, agent) + await adapter.on_tool_start(context, agent, FakeTool("search")) + await adapter.on_tool_end(context, agent, FakeTool("search"), "result") + await adapter.on_llm_start(context, agent, None, []) + await adapter.on_llm_end(context, agent, FakeResponse(input_tokens=100, output_tokens=25)) + await adapter.on_agent_end(context, agent, "done") + + asyncio.run(run()) + + kinds = [r["meta"]["kind"] for r in sink.records] + assert kinds == ["action", "observation", "action", "span"] + + action, observation, llm_action, span = sink.records + assert action["meta"]["tool"] == "search" + assert isinstance(observation["meta"]["duration_ms"], float) + assert observation["meta"]["tool"] == "search" + assert llm_action["llm"] == { + "model": "gpt-fake", + "tokens_in": 100, + "tokens_out": 25, + "latency_ms": llm_action["llm"]["latency_ms"], + } + assert span["message"] == "planner" + assert span["meta"]["operation"] == "invoke_agent" + assert span["meta"]["agent_name"] == "planner" + assert isinstance(span["meta"]["duration_ms"], float) + + # zero-effort nesting: every event happened inside the agent's span + assert action["meta"]["parent_span_id"] == span["meta"]["span_id"] + assert observation["meta"]["parent_span_id"] == span["meta"]["span_id"] + assert llm_action["meta"]["parent_span_id"] == span["meta"]["span_id"] + + +def test_a_handoff_opens_a_nested_span_for_the_new_agent(monkeypatch: pytest.MonkeyPatch) -> None: + module = _load_adapter_module(monkeypatch) + OpenAIAgentsAdapter = module.OpenAIAgentsAdapter + + sink = CollectingTransport() + logger = Logger("app.agent", transports=[sink]) + adapter = OpenAIAgentsAdapter(logger) + context = object() + planner, specialist = FakeAgent("planner"), FakeAgent("specialist") + + async def run() -> None: + await adapter.on_agent_start(context, planner) + await adapter.on_handoff(context, planner, specialist) + await adapter.on_agent_start(context, specialist) + await adapter.on_agent_end(context, specialist, "done") + await adapter.on_agent_end(context, planner, "done") + + asyncio.run(run()) + + decision, inner_span, outer_span = sink.records + assert decision["meta"]["kind"] == "decision" + assert decision["meta"]["from_agent"] == "planner" + assert decision["meta"]["to_agent"] == "specialist" + assert inner_span["message"] == "specialist" + assert outer_span["message"] == "planner" + assert inner_span["meta"]["parent_span_id"] == outer_span["meta"]["span_id"] + assert "parent_span_id" not in outer_span["meta"] + + +def test_two_runs_gathered_concurrently_do_not_share_span_state( + monkeypatch: pytest.MonkeyPatch, +) -> None: + module = _load_adapter_module(monkeypatch) + OpenAIAgentsAdapter = module.OpenAIAgentsAdapter + + sink = CollectingTransport() + logger = Logger("app.agent", transports=[sink]) + adapter = OpenAIAgentsAdapter(logger) + + async def one_run(context: object, name: str) -> None: + await adapter.on_agent_start(context, FakeAgent(name)) + await adapter.on_tool_start(context, FakeAgent(name), FakeTool("search")) + await asyncio.sleep(0) # yield control, so the two runs genuinely interleave + await adapter.on_tool_end(context, FakeAgent(name), FakeTool("search"), "r") + await adapter.on_agent_end(context, FakeAgent(name), "done") + + async def run_both() -> None: + await asyncio.gather(one_run(object(), "a"), one_run(object(), "b")) + + asyncio.run(run_both()) + + by_message: dict[str, list[dict]] = {} + for record in sink.records: + by_message.setdefault(record["message"], []).append(record) + a_action, b_action = by_message["call search"] + a_span, b_span = by_message["a"][0], by_message["b"][0] + + # each tool call is parented to its own run's span, never the other run's + assert a_action["meta"]["parent_span_id"] == a_span["meta"]["span_id"] + assert b_action["meta"]["parent_span_id"] == b_span["meta"]["span_id"] + assert "parent_span_id" not in a_span["meta"] + assert "parent_span_id" not in b_span["meta"] + assert adapter._agent_spans == {} + + +def test_llm_call_without_a_matching_start_has_no_duration(monkeypatch: pytest.MonkeyPatch) -> None: + module = _load_adapter_module(monkeypatch) + + sink = CollectingTransport() + adapter = module.OpenAIAgentsAdapter(Logger("app", transports=[sink])) + + asyncio.run(adapter.on_llm_end(object(), FakeAgent("a"), FakeResponse())) + + assert "latency_ms" not in sink.records[0]["llm"] diff --git a/tests/test_instrument/__init__.py b/tests/test_instrument/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/tests/test_instrument/fakes.py b/tests/test_instrument/fakes.py new file mode 100644 index 0000000..95e0df1 --- /dev/null +++ b/tests/test_instrument/fakes.py @@ -0,0 +1,179 @@ +"""Fake stand-ins for the provider SDKs `logquill.instrument.*` patches — the +same `sys.modules` injection pattern the framework adapter tests use (see +`tests/test_adapters/test_langchain_adapter.py`), so these tests exercise the +real patching code against objects shaped like the real SDKs' without +requiring the actual (heavier) packages to be installed. +""" + +from __future__ import annotations + +import types +from dataclasses import dataclass, field +from typing import Any, Callable + + +@dataclass +class Usage: + input_tokens: int | None = None + output_tokens: int | None = None + prompt_tokens: int | None = None + completion_tokens: int | None = None + + +@dataclass +class Message: + """Stands in for `anthropic.types.Message`.""" + + model: str = "claude-fake" + stop_reason: str | None = "end_turn" + usage: Usage = field(default_factory=lambda: Usage(input_tokens=10, output_tokens=4)) + + +def install_fake_anthropic(monkeypatch: Any, *, calls: list[str] | None = None) -> None: + def create(self: Any, **kwargs: Any) -> Message: + if calls is not None: + calls.append("sync") + return Message() + + async def async_create(self: Any, **kwargs: Any) -> Message: + if calls is not None: + calls.append("async") + return Message() + + messages_module = types.ModuleType("anthropic.resources.messages.messages") + messages_module.Messages = type("Messages", (), {"create": create}) # type: ignore[attr-defined] + messages_module.AsyncMessages = type( # type: ignore[attr-defined] + "AsyncMessages", (), {"create": async_create} + ) + package = types.ModuleType("anthropic.resources.messages") + package.messages = messages_module # type: ignore[attr-defined] + anthropic_module = types.ModuleType("anthropic") + resources_module = types.ModuleType("anthropic.resources") + resources_module.messages = package # type: ignore[attr-defined] + anthropic_module.resources = resources_module # type: ignore[attr-defined] + + import sys + + monkeypatch.setitem(sys.modules, "anthropic", anthropic_module) + monkeypatch.setitem(sys.modules, "anthropic.resources", resources_module) + monkeypatch.setitem(sys.modules, "anthropic.resources.messages", package) + monkeypatch.setitem(sys.modules, "anthropic.resources.messages.messages", messages_module) + + +@dataclass +class Choice: + finish_reason: str | None = "stop" + + +@dataclass +class ChatCompletion: + """Stands in for `openai.types.chat.ChatCompletion`.""" + + model: str = "gpt-fake" + usage: Usage = field(default_factory=lambda: Usage(prompt_tokens=8, completion_tokens=3)) + choices: list[Choice] = field(default_factory=lambda: [Choice()]) + + +@dataclass +class IncompleteDetails: + reason: str | None = None + + +@dataclass +class Response: + """Stands in for `openai.types.responses.Response`.""" + + model: str = "gpt-fake" + status: str = "completed" + incomplete_details: IncompleteDetails | None = None + usage: Usage = field(default_factory=lambda: Usage(input_tokens=6, output_tokens=2)) + + +def install_fake_openai(monkeypatch: Any, *, calls: list[str] | None = None) -> None: + def chat_create(self: Any, **kwargs: Any) -> ChatCompletion: + if calls is not None: + calls.append("chat.sync") + return ChatCompletion() + + async def async_chat_create(self: Any, **kwargs: Any) -> ChatCompletion: + if calls is not None: + calls.append("chat.async") + return ChatCompletion() + + def responses_create(self: Any, **kwargs: Any) -> Response: + if calls is not None: + calls.append("responses.sync") + return Response() + + async def async_responses_create(self: Any, **kwargs: Any) -> Response: + if calls is not None: + calls.append("responses.async") + return Response() + + completions_module = types.ModuleType("openai.resources.chat.completions") + completions_module.Completions = type( # type: ignore[attr-defined] + "Completions", (), {"create": chat_create} + ) + completions_module.AsyncCompletions = type( # type: ignore[attr-defined] + "AsyncCompletions", (), {"create": async_chat_create} + ) + responses_module = types.ModuleType("openai.resources.responses") + responses_module.Responses = type( # type: ignore[attr-defined] + "Responses", (), {"create": responses_create} + ) + responses_module.AsyncResponses = type( # type: ignore[attr-defined] + "AsyncResponses", (), {"create": async_responses_create} + ) + chat_module = types.ModuleType("openai.resources.chat") + chat_module.completions = completions_module # type: ignore[attr-defined] + resources_module = types.ModuleType("openai.resources") + resources_module.chat = chat_module # type: ignore[attr-defined] + resources_module.responses = responses_module # type: ignore[attr-defined] + openai_module = types.ModuleType("openai") + openai_module.resources = resources_module # type: ignore[attr-defined] + + import sys + + monkeypatch.setitem(sys.modules, "openai", openai_module) + monkeypatch.setitem(sys.modules, "openai.resources", resources_module) + monkeypatch.setitem(sys.modules, "openai.resources.chat", chat_module) + monkeypatch.setitem(sys.modules, "openai.resources.chat.completions", completions_module) + monkeypatch.setitem(sys.modules, "openai.resources.responses", responses_module) + + +@dataclass +class LitellmResponse: + """Stands in for `litellm.types.utils.ModelResponse`.""" + + model: str = "litellm-fake" + usage: Usage = field(default_factory=lambda: Usage(prompt_tokens=5, completion_tokens=1)) + choices: list[Choice] = field(default_factory=lambda: [Choice(finish_reason="stop")]) + _hidden_params: dict[str, Any] = field( + default_factory=lambda: {"custom_llm_provider": "openai"} + ) + + +def install_fake_litellm( + monkeypatch: Any, + *, + calls: list[str] | None = None, + completion: Callable[..., Any] | None = None, + acompletion: Callable[..., Any] | None = None, +) -> None: + def default_completion(**kwargs: Any) -> LitellmResponse: + if calls is not None: + calls.append("sync") + return LitellmResponse() + + async def default_acompletion(**kwargs: Any) -> LitellmResponse: + if calls is not None: + calls.append("async") + return LitellmResponse() + + litellm_module = types.ModuleType("litellm") + litellm_module.completion = completion or default_completion # type: ignore[attr-defined] + litellm_module.acompletion = acompletion or default_acompletion # type: ignore[attr-defined] + + import sys + + monkeypatch.setitem(sys.modules, "litellm", litellm_module) diff --git a/tests/test_instrument/test_anthropic.py b/tests/test_instrument/test_anthropic.py new file mode 100644 index 0000000..ea3cf31 --- /dev/null +++ b/tests/test_instrument/test_anthropic.py @@ -0,0 +1,152 @@ +from __future__ import annotations + +import asyncio +import sys + +import pytest + +from logquill import Logger +from logquill.transports.transport import CollectingTransport +from tests.test_instrument.fakes import install_fake_anthropic + + +def _reload() -> None: + sys.modules.pop("logquill.instrument.anthropic", None) + import logquill.instrument.anthropic # noqa: F401 + + +@pytest.fixture(autouse=True) +def _fresh_module(monkeypatch: pytest.MonkeyPatch) -> None: + _reload() + yield + from logquill.instrument.anthropic import anthropic + + if anthropic.active: + anthropic.uninstrument() + + +def test_a_sync_call_emits_llm_call_with_no_code_changes(monkeypatch: pytest.MonkeyPatch) -> None: + install_fake_anthropic(monkeypatch) + from anthropic.resources.messages.messages import Messages + + from logquill.instrument.anthropic import anthropic + + sink = CollectingTransport() + logger = Logger("app", transports=[sink]) + anthropic(logger) + + response = Messages().create(model="claude-fake", max_tokens=100, messages=[]) + + assert response.model == "claude-fake" # the real return value still reaches the caller + assert len(sink.records) == 1 + record = sink.records[0] + assert record["llm"] == { + "model": "claude-fake", + "tokens_in": 10, + "tokens_out": 4, + "finish_reason": "end_turn", + "latency_ms": record["llm"]["latency_ms"], + } + assert record["meta"]["provider"] == "anthropic" + assert isinstance(record["llm"]["latency_ms"], float) + + +def test_an_async_call_is_instrumented_too(monkeypatch: pytest.MonkeyPatch) -> None: + install_fake_anthropic(monkeypatch) + from anthropic.resources.messages.messages import AsyncMessages + + from logquill.instrument.anthropic import anthropic + + sink = CollectingTransport() + logger = Logger("app", transports=[sink]) + anthropic(logger) + + asyncio.run(AsyncMessages().create(model="claude-fake", max_tokens=100, messages=[])) + + assert len(sink.records) == 1 + assert sink.records[0]["llm"]["tokens_in"] == 10 + + +def test_a_streaming_call_passes_through_untouched(monkeypatch: pytest.MonkeyPatch) -> None: + calls: list[str] = [] + install_fake_anthropic(monkeypatch, calls=calls) + from anthropic.resources.messages.messages import Messages + + from logquill.instrument.anthropic import anthropic + + sink = CollectingTransport() + anthropic(Logger("app", transports=[sink])) + + Messages().create(model="claude-fake", max_tokens=100, messages=[], stream=True) + + assert calls == ["sync"] # the real (fake) create() still ran + assert sink.records == [] # but nothing was logged — streaming is a known gap + + +def test_uninstrument_restores_the_original_method(monkeypatch: pytest.MonkeyPatch) -> None: + install_fake_anthropic(monkeypatch) + from anthropic.resources.messages.messages import Messages + + from logquill.instrument.anthropic import anthropic + + original = Messages.create + anthropic(Logger("app")) + assert Messages.create is not original + + anthropic.uninstrument() + + assert Messages.create is original + anthropic.uninstrument() # idempotent + + +def test_instrumenting_twice_without_uninstrument_raises(monkeypatch: pytest.MonkeyPatch) -> None: + install_fake_anthropic(monkeypatch) + from logquill.instrument.anthropic import anthropic + + anthropic(Logger("app")) + + with pytest.raises(RuntimeError, match="already active"): + anthropic(Logger("app")) + + +def test_a_broken_logger_does_not_break_the_real_call(monkeypatch: pytest.MonkeyPatch) -> None: + install_fake_anthropic(monkeypatch) + from anthropic.resources.messages.messages import Messages + + from logquill.instrument.anthropic import anthropic + + class ExplodingLogger(Logger): + def llm_call(self, *args: object, **kwargs: object) -> None: # type: ignore[override] + raise RuntimeError("logging is broken") + + anthropic(ExplodingLogger("app")) + + response = Messages().create(model="claude-fake", max_tokens=100, messages=[]) + + assert response.model == "claude-fake" # the API call's result still comes back + + +def test_a_missing_anthropic_package_gives_an_install_hint(monkeypatch: pytest.MonkeyPatch) -> None: + monkeypatch.setitem(sys.modules, "anthropic.resources.messages", None) + from logquill.instrument.anthropic import anthropic + + with pytest.raises(ImportError, match=r"pip install logquill\[instrument-anthropic\]"): + anthropic(Logger("app")) + + +def test_importing_logquill_instrument_does_not_import_anthropic() -> None: + import subprocess + + result = subprocess.run( + [ + sys.executable, + "-c", + "import sys, logquill.instrument; " + "print(any(m.startswith('anthropic') for m in sys.modules))", + ], + capture_output=True, + text=True, + check=True, + ) + + assert result.stdout.strip() == "False" diff --git a/tests/test_instrument/test_litellm.py b/tests/test_instrument/test_litellm.py new file mode 100644 index 0000000..0fb0a5d --- /dev/null +++ b/tests/test_instrument/test_litellm.py @@ -0,0 +1,97 @@ +from __future__ import annotations + +import asyncio +import sys + +import pytest + +from logquill import Logger +from logquill.transports.transport import CollectingTransport +from tests.test_instrument.fakes import install_fake_litellm + + +@pytest.fixture(autouse=True) +def _fresh_module() -> None: + sys.modules.pop("logquill.instrument.litellm", None) + import logquill.instrument.litellm # noqa: F401 + + yield + from logquill.instrument.litellm import litellm + + if litellm.active: + litellm.uninstrument() + + +def test_a_sync_completion_is_instrumented(monkeypatch: pytest.MonkeyPatch) -> None: + install_fake_litellm(monkeypatch) + import litellm as litellm_module + + from logquill.instrument.litellm import litellm + + sink = CollectingTransport() + litellm(Logger("app", transports=[sink])) + + litellm_module.completion(model="litellm-fake", messages=[]) + + assert sink.records[0]["llm"] == { + "model": "litellm-fake", + "tokens_in": 5, + "tokens_out": 1, + "finish_reason": "stop", + "latency_ms": sink.records[0]["llm"]["latency_ms"], + } + # litellm resolves the real underlying provider; we surface it as-is + assert sink.records[0]["meta"]["provider"] == "openai" + + +def test_an_async_completion_is_instrumented(monkeypatch: pytest.MonkeyPatch) -> None: + install_fake_litellm(monkeypatch) + import litellm as litellm_module + + from logquill.instrument.litellm import litellm + + sink = CollectingTransport() + litellm(Logger("app", transports=[sink])) + + asyncio.run(litellm_module.acompletion(model="litellm-fake", messages=[])) + + assert sink.records[0]["llm"]["tokens_in"] == 5 + + +def test_a_streaming_completion_passes_through_untouched(monkeypatch: pytest.MonkeyPatch) -> None: + calls: list[str] = [] + install_fake_litellm(monkeypatch, calls=calls) + import litellm as litellm_module + + from logquill.instrument.litellm import litellm + + sink = CollectingTransport() + litellm(Logger("app", transports=[sink])) + + litellm_module.completion(model="litellm-fake", messages=[], stream=True) + + assert calls == ["sync"] + assert sink.records == [] + + +def test_uninstrument_restores_both_functions(monkeypatch: pytest.MonkeyPatch) -> None: + install_fake_litellm(monkeypatch) + import litellm as litellm_module + + from logquill.instrument.litellm import litellm + + original_completion = litellm_module.completion + original_acompletion = litellm_module.acompletion + litellm(Logger("app")) + litellm.uninstrument() + + assert litellm_module.completion is original_completion + assert litellm_module.acompletion is original_acompletion + + +def test_a_missing_litellm_package_gives_an_install_hint(monkeypatch: pytest.MonkeyPatch) -> None: + monkeypatch.setitem(sys.modules, "litellm", None) + from logquill.instrument.litellm import litellm + + with pytest.raises(ImportError, match=r"pip install logquill\[instrument-litellm\]"): + litellm(Logger("app")) diff --git a/tests/test_instrument/test_openai.py b/tests/test_instrument/test_openai.py new file mode 100644 index 0000000..eee6a16 --- /dev/null +++ b/tests/test_instrument/test_openai.py @@ -0,0 +1,143 @@ +from __future__ import annotations + +import asyncio +import sys + +import pytest + +from logquill import Logger +from logquill.transports.transport import CollectingTransport +from tests.test_instrument.fakes import IncompleteDetails, Response, install_fake_openai + + +@pytest.fixture(autouse=True) +def _fresh_module() -> None: + sys.modules.pop("logquill.instrument.openai", None) + import logquill.instrument.openai # noqa: F401 + + yield + from logquill.instrument.openai import openai + + if openai.active: + openai.uninstrument() + + +def test_chat_completions_are_instrumented(monkeypatch: pytest.MonkeyPatch) -> None: + install_fake_openai(monkeypatch) + from openai.resources.chat.completions import Completions + + from logquill.instrument.openai import openai + + sink = CollectingTransport() + openai(Logger("app", transports=[sink])) + + Completions().create(model="gpt-fake", messages=[]) + + assert sink.records[0]["llm"] == { + "model": "gpt-fake", + "tokens_in": 8, + "tokens_out": 3, + "finish_reason": "stop", + "latency_ms": sink.records[0]["llm"]["latency_ms"], + } + assert sink.records[0]["meta"]["provider"] == "openai" + + +def test_async_chat_completions_are_instrumented(monkeypatch: pytest.MonkeyPatch) -> None: + install_fake_openai(monkeypatch) + from openai.resources.chat.completions import AsyncCompletions + + from logquill.instrument.openai import openai + + sink = CollectingTransport() + openai(Logger("app", transports=[sink])) + + asyncio.run(AsyncCompletions().create(model="gpt-fake", messages=[])) + + assert sink.records[0]["llm"]["tokens_in"] == 8 + + +def test_the_responses_api_is_instrumented_too(monkeypatch: pytest.MonkeyPatch) -> None: + install_fake_openai(monkeypatch) + from openai.resources.responses import Responses + + from logquill.instrument.openai import openai + + sink = CollectingTransport() + openai(Logger("app", transports=[sink])) + + Responses().create(model="gpt-fake", input="hi") + + assert sink.records[0]["llm"] == { + "model": "gpt-fake", + "tokens_in": 6, + "tokens_out": 2, + "finish_reason": "completed", + "latency_ms": sink.records[0]["llm"]["latency_ms"], + } + + +def test_an_incomplete_response_reports_its_own_reason(monkeypatch: pytest.MonkeyPatch) -> None: + install_fake_openai(monkeypatch) + from openai.resources.responses import Responses + + from logquill.instrument.openai import openai + + def create(self: object, **kwargs: object) -> Response: + return Response( + status="incomplete", incomplete_details=IncompleteDetails(reason="max_output_tokens") + ) + + Responses.create = create # type: ignore[method-assign] + sink = CollectingTransport() + openai(Logger("app", transports=[sink])) + + Responses().create(model="gpt-fake", input="hi") + + assert sink.records[0]["llm"]["finish_reason"] == "max_output_tokens" + + +def test_a_streaming_chat_call_passes_through_untouched(monkeypatch: pytest.MonkeyPatch) -> None: + install_fake_openai(monkeypatch) + from openai.resources.chat.completions import Completions + + from logquill.instrument.openai import openai + + sink = CollectingTransport() + openai(Logger("app", transports=[sink])) + + Completions().create(model="gpt-fake", messages=[], stream=True) + + assert sink.records == [] + + +def test_uninstrument_restores_all_four_methods(monkeypatch: pytest.MonkeyPatch) -> None: + install_fake_openai(monkeypatch) + from openai.resources.chat.completions import AsyncCompletions, Completions + from openai.resources.responses import AsyncResponses, Responses + + from logquill.instrument.openai import openai + + originals = ( + Completions.create, + AsyncCompletions.create, + Responses.create, + AsyncResponses.create, + ) + openai(Logger("app")) + openai.uninstrument() + + assert ( + Completions.create, + AsyncCompletions.create, + Responses.create, + AsyncResponses.create, + ) == originals + + +def test_a_missing_openai_package_gives_an_install_hint(monkeypatch: pytest.MonkeyPatch) -> None: + monkeypatch.setitem(sys.modules, "openai.resources.chat", None) + from logquill.instrument.openai import openai + + with pytest.raises(ImportError, match=r"pip install logquill\[instrument-openai\]"): + openai(Logger("app")) diff --git a/tests/test_mcp.py b/tests/test_mcp.py new file mode 100644 index 0000000..5d29cf4 --- /dev/null +++ b/tests/test_mcp.py @@ -0,0 +1,130 @@ +from __future__ import annotations + +import json + +from logquill import Logger, RunPlugin, TraceContextPlugin +from logquill.mcp import RUN_ID_KEY, TRACEPARENT_KEY, inbound, propagate +from logquill.transports.transport import CollectingTransport + + +def test_propagate_carries_a_w3c_traceparent() -> None: + meta = propagate() + + parts = meta[TRACEPARENT_KEY].split("-") + assert len(parts) == 4 + assert parts[0] == "00" + assert len(parts[1]) == 32 # trace id + assert len(parts[2]) == 16 # span id + assert RUN_ID_KEY not in meta + + +def test_propagate_uses_a_given_trace_id() -> None: + meta = propagate(trace_id="4bf92f3577b34da6a3ce929d0e0e4736") + + assert meta[TRACEPARENT_KEY].split("-")[1] == "4bf92f3577b34da6a3ce929d0e0e4736" + + +def test_propagate_includes_run_id_only_when_given() -> None: + assert RUN_ID_KEY not in propagate() + assert propagate(run_id="run-1")[RUN_ID_KEY] == "run-1" + + +def test_the_client_server_call_shares_one_trace_id_end_to_end() -> None: + client_sink = CollectingTransport() + client_logger = Logger("client", transports=[client_sink], plugins=[TraceContextPlugin()]) + call_record = client_logger.info("calling the search tool") + assert call_record is not None + client_trace_id = call_record["meta"]["trace_id"] + + # the client sends `propagate(trace_id=...)`'s result as `meta=` on its MCP + # request — simulated here as a JSON round trip, since that's what + # actually crosses the wire + outbound_meta = json.loads(json.dumps(propagate(trace_id=client_trace_id))) + + server_sink = CollectingTransport() + server_logger = Logger("server", transports=[server_sink], plugins=[TraceContextPlugin()]) + with inbound(outbound_meta, server="files", tool="search"): + server_logger.info("handling the call") + + server_trace_id = server_sink.records[0]["meta"]["trace_id"] + assert outbound_meta[TRACEPARENT_KEY].split("-")[1] == client_trace_id + assert server_trace_id == client_trace_id + + +def test_inbound_stamps_mcp_server_and_tool_on_every_record_in_the_block() -> None: + sink = CollectingTransport() + logger = Logger("server", transports=[sink]) + + with inbound(None, server="files", tool="read_file"): + logger.info("first") + logger.warn("second") + logger.error("outside the block") + + assert sink.records[0]["meta"]["mcp"] == {"server": "files", "tool": "read_file"} + assert sink.records[1]["meta"]["mcp"] == {"server": "files", "tool": "read_file"} + assert "mcp" not in sink.records[2]["meta"] + + +def test_inbound_with_no_meta_at_all_still_stamps_mcp_but_has_no_trace_to_inherit() -> None: + sink = CollectingTransport() + logger = Logger("server", transports=[sink], plugins=[TraceContextPlugin()]) + + with inbound(None, server="files", tool="search"): + logger.info("handled") + + assert sink.records[0]["meta"]["mcp"] == {"server": "files", "tool": "search"} + assert isinstance(sink.records[0]["meta"]["trace_id"], str) # generated, not inherited + + +def test_a_client_run_id_lands_under_meta_mcp_and_never_shadows_the_servers_own_run_plugin() -> ( + None +): + sink = CollectingTransport() + logger = Logger("server", transports=[sink], plugins=[RunPlugin(run_id="server-run")]) + meta = propagate(run_id="client-run") + + with inbound(meta, server="files", tool="search"): + logger.info("handled") + + assert sink.records[0]["meta"]["mcp"]["run_id"] == "client-run" + assert sink.records[0]["meta"]["run_id"] == "server-run" + + +def test_inbound_restores_the_previous_traceparent_after_the_block() -> None: + from logquill.plugins.trace_context_plugin import reset_traceparent, set_traceparent + + outer_token = set_traceparent("00-" + "a" * 32 + "-" + "b" * 16 + "-01") + sink = CollectingTransport() + logger = Logger("server", transports=[sink], plugins=[TraceContextPlugin()]) + try: + with inbound(propagate(), server="files", tool="search"): + pass + logger.info("after the block") + finally: + reset_traceparent(outer_token) + + assert sink.records[0]["meta"]["trace_id"] == "a" * 32 + + +def test_inbound_restores_the_traceparent_even_if_the_block_raises() -> None: + from logquill.plugins.trace_context_plugin import _current_traceparent + + before = _current_traceparent.get() + try: + with inbound(propagate(), server="files", tool="search"): + raise ValueError("tool failed") + except ValueError: + pass + + assert _current_traceparent.get() == before + + +def test_a_malformed_traceparent_falls_back_to_generating_a_fresh_trace_id() -> None: + sink = CollectingTransport() + logger = Logger("server", transports=[sink], plugins=[TraceContextPlugin()]) + + with inbound({TRACEPARENT_KEY: "not-a-real-traceparent"}, server="files", tool="search"): + logger.info("handled") + + assert isinstance(sink.records[0]["meta"]["trace_id"], str) + assert sink.records[0]["meta"]["trace_id"] != "not-a-real-traceparent"