diff --git a/CHANGELOG.md b/CHANGELOG.md index 53cf2c2..cab347e 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -4,6 +4,30 @@ All notable changes to this project are documented in this file. ## Unreleased +- Runtime intelligence: + - `FlightRecorderPlugin` generalizes `SamplingPlugin`'s tail-based + elevation from "a record sampling happened to drop" to "a record below + the level you actually want shipped" — buffers low-level records + (`DEBUG`/`TRACE`) per run and ships the whole buffered trail only if + that run later reaches an error level, so a run at `DEBUG` costs + nothing to ship unless it actually fails. Bounded the same way tail + elevation is (`max_buffered_records`/`max_runs`), and verified to stay + memory-bounded under a sustained multi-run burst. + - `RunSummaryPlugin` emits one record — total tokens, cost, tool calls, + retries, and errors — when a run's outermost span closes, so a + dashboard doesn't need to re-walk every record for a run's totals. + - `AdaptiveSamplingPlugin` always keeps errors and slow spans, samples + everything else, and caps ordinary traffic to a bytes-per-second + budget — adjusting its own rate up when there's room to spare and down + when it's saturated each window, based on demand rather than bytes + actually emitted (which stop growing once the budget saturates, and + would otherwise read a saturated gate as "plenty of room"). + - `install_signal_level_handler`, `LevelFileWatcher`, and + `LevelEnvWatcher` change a `Logger`'s level without restarting the + process — a signal that re-reads an environment variable, a polled + file, or a polled environment variable (only useful when something + inside the same process is the one changing it — documented plainly, + since that's a real limit worth knowing about, not a gap to paper over). - Privacy, content policy, and an audit trail: - Every `Logger` now has a `content_policy` (`"off"` by default \| `"hash"` \| `"truncate"` \| `"full"`), governing the content fields an LLM call or diff --git a/README.md b/README.md index 1152242..aadd4e6 100644 --- a/README.md +++ b/README.md @@ -37,6 +37,7 @@ for what's landed so far. - **Parse any log file** — `parse()` pulls structured fields out of a log file (LogQuill's own or a legacy format) with a regex, streaming line by line — see [Parsing log files](#parsing-log-files) - **CLI** — `logquill tail app.log --level=warn --json -f` for filtering/following a JSONL log file in local dev, no extra install — see [CLI](#cli) - **Local-first trace viewer** — `logquill trace --file logs.jsonl` prints an annotated span tree (streamed, bounded memory at gigabyte scale); `logquill serve` runs a small offline web UI (stdlib only — run list, span tree/waterfall, search, level filter) reading JSONL or a `SQLiteTransport` database; `logquill dev` live-renders the current run as it happens — no account, nothing leaves your machine +- **Runtime intelligence** — `FlightRecorderPlugin` ships a run's buffered DEBUG trail only if it errors; `RunSummaryPlugin` emits one record with a run's total tokens/cost/tool calls/retries/errors; `AdaptiveSamplingPlugin` always keeps errors and slow spans while capping ordinary traffic to a bytes-per-second budget; `install_signal_level_handler`/`LevelFileWatcher`/`LevelEnvWatcher` change a logger's level without a restart — see [Runtime level control](#runtime-level-control) ## Install @@ -518,6 +519,84 @@ Buffering is bounded by `max_buffered_records` and `max_traces` — the oldest buffered trace is evicted once either limit is hit, so a single high-cardinality or long-lived trace can't grow memory without limit. +### The flight recorder: DEBUG context, only for runs that fail + +`FlightRecorderPlugin` generalizes the same idea from "a record sampling +happened to drop" to "a record below the level you actually want shipped, +period": run a logger at `DEBUG` so fine-grained context always exists, but +only pay to ship it for the runs that turn out to matter. + +```python +from logquill import CollectingTransport, FlightRecorderPlugin, Logger + +sink = CollectingTransport() +recorder = FlightRecorderPlugin(transports=[sink]) # ship_at=INFO, flush_at=ERROR +logger = Logger("app.agent", level="DEBUG", transports=[sink], plugins=[recorder]) + +# a run that finishes cleanly never ships its DEBUG trail +logger.debug("fetched 40 candidates", run_id="run-1") # buffered, not shipped +logger.debug("filtered to top 3", run_id="run-1") # buffered, not shipped +logger.info("run finished", run_id="run-1") # ships on its own +assert [r["message"] for r in sink.records] == ["run finished"] + +# a run that errors ships its full DEBUG trail, flushed by the error itself +logger.debug("fetched 40 candidates", run_id="run-2") # buffered for now +logger.error("ranking model timed out", run_id="run-2") # flushes the buffer +assert [r["message"] for r in sink.records[-2:]] == [ + "fetched 40 candidates", + "ranking model timed out", +] +``` + +Bounded the same way tail-based elevation is: `max_buffered_records` and +`max_runs` cap the buffer, the oldest run evicted (and lost, not shipped) +once either is hit. + +### One record per run: `RunSummaryPlugin` + +Emits a single summary record — total tokens, cost, tool calls, retries, +and errors — when a run's outermost span closes, so a dashboard doesn't +need to re-walk every record to answer "how much did this run cost": + +```python +from logquill import CollectingTransport, Logger, RunPlugin, RunSummaryPlugin + +sink = CollectingTransport() +summary = RunSummaryPlugin(transports=[sink]) +logger = Logger("app.agent", transports=[sink], plugins=[RunPlugin(), summary]) + +with logger.span("run", operation="invoke_agent", agent_name="planner"): + logger.llm_call("chat", model="m", tokens_in=1200, tokens_out=340, cost_usd=0.02) + logger.action("search", tool="search") + +run_summary = sink.records[-1] +assert run_summary["message"] == "run summary" +assert run_summary["meta"]["tokens_in"] == 1200 +assert run_summary["meta"]["tool_calls"] == 1 +assert run_summary["meta"]["errors"] == 0 +``` + +### Adaptive sampling with a byte-rate budget + +`AdaptiveSamplingPlugin` always keeps errors and slow spans, samples +everything else, and caps how many bytes of that ordinary traffic pass +through per second — adjusting its own rate up when there's room to spare +and down when it's saturated, instead of one fixed rate tuned by hand: + +```python +from logquill import AdaptiveSamplingPlugin, Logger + +logger = Logger( + "app", + plugins=[AdaptiveSamplingPlugin(base_rate=0.1, slow_ms=1000, max_bytes_per_second=100_000)], +) + +logger.error("always kept") # errors bypass sampling entirely +with logger.span("slow_call", duration_ms=1500): # slow spans bypass it too + pass +logger.info("ordinary traffic") # sampled at the current rate +``` + ### PII redaction by pattern, not just key `RedactPlugin` redacts by exact key match. `PIIRedactPlugin` complements it by @@ -1483,6 +1562,37 @@ To read back what `TextFormatter` wrote, use the ready-made `TEXT_LOG_PATTERN` with `cast=TEXT_LOG_CASTS` (which decodes `meta` from JSON); for `LogfmtFormatter` output, `parse_logfmt(line)` returns a dict of strings. +## Runtime level control + +Change a `Logger`'s level without restarting the process — three +independent mechanisms; pick whichever fits your deployment: + +```python +import os, signal +from logquill import Logger, install_signal_level_handler + +logger = Logger("app") +install_signal_level_handler(logger, env_var="APP_LOG_LEVEL") +# kill -USR1 now re-reads APP_LOG_LEVEL and applies it +``` + +```python +from logquill import Logger, LevelFileWatcher + +logger = Logger("app") +watcher = LevelFileWatcher(logger, "/var/run/app.level", poll_interval=2.0) +watcher.start() # `echo DEBUG > /var/run/app.level` to turn on DEBUG, no restart +# watcher.stop() on shutdown +``` + +`LevelEnvWatcher(logger, "APP_LOG_LEVEL")` polls an environment variable the +same way — but only helps when something *inside* this same process is the +one changing it (a config-reload callback, `python-dotenv`'s +`override=True`); an OS-level environment change made from outside the +process, the way a shell normally would, is never visible to code already +running inside it. Use the signal or file mechanism for anything triggered +externally. + ## CLI Installing `logquill` also installs a `logquill` command for local diff --git a/benchmarks/test_flight_recorder_memory.py b/benchmarks/test_flight_recorder_memory.py new file mode 100644 index 0000000..1b9be2a --- /dev/null +++ b/benchmarks/test_flight_recorder_memory.py @@ -0,0 +1,69 @@ +"""The exit criterion for the flight recorder: a sustained burst of +low-level records across many runs, almost none of which error, stays +memory-bounded — the buffer's own bookkeeping never exceeds its configured +limits, and peak memory during the burst doesn't grow with the burst size. +Slow-ish and memory-instrumented, so it lives here rather than in the +default `pytest` run — see `measure.py`'s module docstring for why these +are a separate CI job. +""" + +from __future__ import annotations + +import tracemalloc + +from logquill import Logger +from logquill.plugins.flight_recorder_plugin import FlightRecorderPlugin +from logquill.transports.transport import CollectingTransport + +BURST = 100_000 +MAX_BUFFERED_RECORDS = 2_000 +MAX_RUNS = 200 + +#: However large the burst, buffering low-level records for thousands of +#: concurrent runs must stay nowhere near proportional to the burst size. +MEMORY_BUDGET_BYTES = 15_000_000 + + +def test_a_sustained_burst_of_mostly_passing_runs_stays_memory_bounded() -> None: + sink = CollectingTransport() + recorder = FlightRecorderPlugin( + transports=[sink], max_buffered_records=MAX_BUFFERED_RECORDS, max_runs=MAX_RUNS + ) + logger = Logger("app.agent", level="DEBUG", transports=[sink], plugins=[recorder]) + + tracemalloc.start() + try: + baseline, _ = tracemalloc.get_traced_memory() + for i in range(BURST): + run_id = f"run-{i % 10_000}" + logger.debug("step", run_id=run_id, i=i, payload="x" * 80) + if i % 50_000 == 0: # a rare run actually errors + logger.error("oops", run_id=run_id) + _current, peak = tracemalloc.get_traced_memory() + finally: + tracemalloc.stop() + + assert recorder._buffered_count <= MAX_BUFFERED_RECORDS + assert len(recorder._buffer) <= MAX_RUNS + + peak_bytes = peak - baseline + assert peak_bytes < MEMORY_BUDGET_BYTES, ( + f"buffering a {BURST:,}-record burst used {peak_bytes:,} bytes against a " + f"budget of {MEMORY_BUDGET_BYTES:,} — memory grew with the burst, not the bound" + ) + + +def test_a_single_run_that_never_errors_is_dropped_once_evicted_not_leaked() -> None: + """A burst entirely under one run_id that never reaches flush_at: once + the record cap is hit, the *oldest* buffered records for that run are + evicted — this run's own buffer never grows past the configured cap + even though it alone accounts for the whole burst.""" + sink = CollectingTransport() + recorder = FlightRecorderPlugin(transports=[sink], max_buffered_records=500, max_runs=10) + logger = Logger("app.agent", level="DEBUG", transports=[sink], plugins=[recorder]) + + for i in range(20_000): + logger.debug("step", run_id="single-run", i=i) + + assert recorder._buffered_count <= 500 + assert sink.records == [] # never errored: nothing shipped diff --git a/logquill/__init__.py b/logquill/__init__.py index 05872a3..4bf7115 100644 --- a/logquill/__init__.py +++ b/logquill/__init__.py @@ -9,16 +9,19 @@ from logquill.logger import Logger from logquill.opt import OptLogger from logquill.parsing import TEXT_LOG_CASTS, TEXT_LOG_PATTERN, parse, parse_logfmt +from logquill.plugins.adaptive_sampling_plugin import AdaptiveSamplingPlugin from logquill.plugins.alerting_plugin import AlertingPlugin from logquill.plugins.apprise_alert_plugin import AppriseAlertPlugin from logquill.plugins.context_plugin import ContextPlugin from logquill.plugins.email_alert_plugin import EmailAlertPlugin +from logquill.plugins.flight_recorder_plugin import FlightRecorderPlugin from logquill.plugins.pagerduty_alert_plugin import PagerDutyAlertPlugin from logquill.plugins.pii_redact_plugin import PIIRedactPlugin from logquill.plugins.plugin import FunctionPlugin, Plugin from logquill.plugins.rate_limit_plugin import RateLimitPlugin from logquill.plugins.redact_plugin import RedactPlugin from logquill.plugins.run_plugin import RunPlugin +from logquill.plugins.run_summary_plugin import RunSummaryPlugin from logquill.plugins.sampling_plugin import SamplingPlugin from logquill.plugins.slack_alert_plugin import SlackAlertPlugin from logquill.plugins.tamper_evident_plugin import ( @@ -32,6 +35,11 @@ from logquill.plugins.trace_context_plugin import TraceContextPlugin from logquill.privacy import CONTENT_FIELDS, FIELD_CLASSES, ContentCapturePolicy, FieldClass from logquill.records import SCHEMA_VERSION, LLMBlock, LogRecord, parse_record +from logquill.runtime_level import ( + LevelEnvWatcher, + LevelFileWatcher, + install_signal_level_handler, +) from logquill.serverless import with_azure_function, with_cloud_function, with_lambda from logquill.toggle import disable, enable, is_enabled from logquill.transports.batching_transport import BatchingTransport @@ -65,6 +73,7 @@ __version__ = "1.0.0" __all__ = [ + "AdaptiveSamplingPlugin", "AlertingPlugin", "AppriseAlertPlugin", "AppInsightsTransport", @@ -87,12 +96,15 @@ "FIELD_CLASSES", "FieldClass", "FileTransport", + "FlightRecorderPlugin", "Formatter", "FunctionPlugin", "HTTPTransport", "JSONFormatter", "KafkaTransport", "Level", + "LevelEnvWatcher", + "LevelFileWatcher", "LogfmtFormatter", "LLMBlock", "LogQuillAdapter", @@ -115,6 +127,7 @@ "RedactPlugin", "RedisTransport", "RunPlugin", + "RunSummaryPlugin", "SCHEMA_VERSION", "SQLLogRow", "SQLiteTransport", @@ -134,6 +147,7 @@ "disable", "enable", "format_exc_info", + "install_signal_level_handler", "is_enabled", "load_config", "logger_from_env", diff --git a/logquill/plugins/adaptive_sampling_plugin.py b/logquill/plugins/adaptive_sampling_plugin.py new file mode 100644 index 0000000..a38bcc1 --- /dev/null +++ b/logquill/plugins/adaptive_sampling_plugin.py @@ -0,0 +1,146 @@ +from __future__ import annotations + +import json +import random +import time +from typing import Callable + +from logquill.levels import Level, parse_level +from logquill.plugins.plugin import Plugin +from logquill.records import LogRecord + + +class AdaptiveSamplingPlugin(Plugin): + """Keeps every error and every "slow" span unconditionally, samples + everything else, and caps how many bytes of that ordinary traffic pass + through in any one second — adjusting how much it keeps as traffic + rises and falls, rather than a single fixed rate. + + - **Always kept**: a record at `keep_at` (default `ERROR`) or above, and + a span record (`meta.kind == "span"`) whose `meta.duration_ms` reaches + `slow_ms` (default 1000) — the two categories you can't afford to + lose to sampling, independent of the budget below. + - **Everything else** is rate-sampled at this plugin's *current* rate + (starting at `base_rate`), then checked against a rolling + one-second byte budget (`max_bytes_per_second`, estimated the same + way `BatchingTransport` estimates buffered size — JSON-encoded + length). A record that passes the rate check but would push the + current second over budget is dropped anyway. + - **Adaptive**: at the end of each one-second window, the rate adjusts + based on *demand* — the byte size of every record that passed the + rate check that window, whether or not the budget gate then also let + it through — against the budget: well under it raises the rate (more + gets through), over it lowers the rate (bounded to `[min_rate, + max_rate]`). Demand, not bytes actually emitted, is what the + adjustment has to track: once the budget is saturated, emitted bytes + stop growing no matter how much more is being asked for, so using + emitted bytes as the signal would read a saturated gate as "plenty of + room" and keep raising the rate into an already-over-budget window. + + `rng`/`clock` are injectable for deterministic tests; they default to + `random.random`/`time.monotonic`. + """ + + def __init__( + self, + *, + base_rate: float = 0.1, + min_rate: float = 0.01, + max_rate: float = 1.0, + keep_at: int | str | Level = Level.ERROR, + slow_ms: float = 1000.0, + max_bytes_per_second: int = 100_000, + window_seconds: float = 1.0, + rng: Callable[[], float] | None = None, + clock: Callable[[], float] | None = None, + ) -> None: + """`base_rate` is the starting (and, absent any traffic, steady- + state) sample rate for ordinary records; it adapts within + `[min_rate, max_rate]` from there. Raises `ValueError` if any rate + is outside `[0.0, 1.0]`, or `min_rate > max_rate`.""" + for name, value in ( + ("base_rate", base_rate), + ("min_rate", min_rate), + ("max_rate", max_rate), + ): + if not 0.0 <= value <= 1.0: + raise ValueError(f"{name} must be between 0 and 1, got {value!r}") + if min_rate > max_rate: + raise ValueError(f"min_rate ({min_rate!r}) must be <= max_rate ({max_rate!r})") + + self.rate = base_rate + self.min_rate = min_rate + self.max_rate = max_rate + self.keep_at = parse_level(keep_at) + self.slow_ms = slow_ms + self.max_bytes_per_second = max_bytes_per_second + self.window_seconds = window_seconds + self._rng = rng or random.random + self._clock = clock or time.monotonic + + self._window_start = self._clock() + self._window_demand_bytes = 0 + self._window_emitted_bytes = 0 + + def _is_always_kept(self, record: LogRecord) -> bool: + if Level[record["level"]] >= self.keep_at: + return True + meta = record["meta"] + if meta.get("kind") != "span": + return False + duration = meta.get("duration_ms") + return ( + isinstance(duration, (int, float)) + and not isinstance(duration, bool) + and duration >= self.slow_ms + ) + + def _estimate_size(self, record: LogRecord) -> int: + try: + return len(json.dumps(record, separators=(",", ":"), default=str).encode("utf-8")) + except (TypeError, ValueError, RecursionError): + return len(str(record)) + + def _roll_window_if_due(self) -> None: + now = self._clock() + elapsed = now - self._window_start + if elapsed < self.window_seconds: + return + + demand = ( + self._window_demand_bytes / self.max_bytes_per_second + if self.max_bytes_per_second + else 0.0 + ) + if demand < 0.5: + self.rate = min(self.rate * 1.2, self.max_rate) + elif demand > 0.9: + self.rate = max(self.rate * 0.5, self.min_rate) + + # Each window starts fresh from `now`, not `now - elapsed`, so a long + # gap between records (e.g. idle traffic) doesn't retroactively + # "owe" several missed windows' worth of rate adjustments at once. + self._window_start = now + self._window_demand_bytes = 0 + self._window_emitted_bytes = 0 + + def before_log(self, record: LogRecord) -> LogRecord | None: + """Keeps an error/slow-span record unconditionally; otherwise + samples at the current rate and checks the byte budget, dropping + the record if either fails.""" + self._roll_window_if_due() + + if self._is_always_kept(record): + return record + + if self._rng() >= self.rate: + return None + + size = self._estimate_size(record) + self._window_demand_bytes += size + + if self._window_emitted_bytes + size > self.max_bytes_per_second: + return None + + self._window_emitted_bytes += size + return record diff --git a/logquill/plugins/flight_recorder_plugin.py b/logquill/plugins/flight_recorder_plugin.py new file mode 100644 index 0000000..a85b050 --- /dev/null +++ b/logquill/plugins/flight_recorder_plugin.py @@ -0,0 +1,115 @@ +from __future__ import annotations + +from collections import OrderedDict + +from logquill.levels import Level, parse_level +from logquill.plugins.plugin import Plugin +from logquill.records import LogRecord +from logquill.transports.transport import Transport + + +class FlightRecorderPlugin(Plugin): + """A bounded, per-run ring buffer of low-level records, shipped only if + that run later errors — generalizing `SamplingPlugin`'s tail-based + elevation from "a record sampling happened to drop" to "a record below + the level you actually want shipped, period." + + The idea: run a logger at `DEBUG`/`TRACE` so fine-grained context is + always *available*, but don't pay to ship that volume to a transport + for every run — only for the runs that turn out to matter. A record at + or above `ship_at` (default `INFO`) goes straight through, same as + without this plugin. A record below `ship_at` is instead buffered under + its `meta[run_key]`. If any later record for that same run reaches + `flush_at` (default `ERROR`), every buffered record for that run is + flushed straight to `transports` and the run ships unconditionally from + then on — so a failing run ships its full `DEBUG` trail, and a passing + one never pays to ship any of it. + + Flushing writes directly to `transports` (the same list given to the + `Logger`), bypassing `before_log`/`after_log` for any plugin *after* + this one in the pipeline, the same tradeoff `SamplingPlugin` makes — + put `FlightRecorderPlugin` last if that matters for your pipeline. + + Buffering is bounded: at most `max_buffered_records` records total and + `max_runs` distinct run ids are held at once. Once either limit is hit, + the oldest buffered run is evicted (and lost, not flushed) — a + deliberate bounded-memory trade-off: an unbounded per-run buffer would + let one pathologically long or high-cardinality run grow memory without + limit, which is exactly the failure mode this plugin exists to avoid + while still keeping DEBUG-level detail available for runs that error. + A record with no `meta[run_key]` isn't part of any run this plugin can + buffer, so it's just passed through as-is regardless of level. + """ + + def __init__( + self, + *, + transports: list[Transport], + run_key: str = "run_id", + ship_at: int | str | Level = Level.INFO, + flush_at: int | str | Level = Level.ERROR, + max_buffered_records: int = 1000, + max_runs: int = 200, + ) -> None: + """`transports` must be the same list given to the `Logger` — this + is how a flushed run's buffered records actually reach a sink. + `max_buffered_records`/`max_runs` bound the buffer's memory; the + oldest run is evicted (unflushed) once either is exceeded.""" + self.transports = transports + self.run_key = run_key + self.ship_at = parse_level(ship_at) + self.flush_at = parse_level(flush_at) + self.max_buffered_records = max_buffered_records + self.max_runs = max_runs + self._buffer: OrderedDict[object, list[LogRecord]] = OrderedDict() + self._buffered_count = 0 + self._flushed: OrderedDict[object, None] = OrderedDict() + + def before_log(self, record: LogRecord) -> LogRecord | None: + """Ships `record` unconditionally if it has no run id, if its run + already flushed, or if its own level reaches `flush_at` (triggering + the flush). Otherwise ships it if its level reaches `ship_at`, or + buffers it under its run id and drops it for now.""" + run_id = record["meta"].get(self.run_key) + if run_id is None: + return record + + if run_id in self._flushed: + return record + + level = Level[record["level"]] + if level >= self.flush_at: + self._flush(run_id) + return record + + if level >= self.ship_at: + return record + + self._buffer_record(run_id, record) + return None + + def _flush(self, run_id: object) -> None: + self._flushed[run_id] = None + buffered = self._buffer.pop(run_id, []) + self._buffered_count -= len(buffered) + for buffered_record in buffered: + for transport in self.transports: + transport.write(transport.format(buffered_record), buffered_record) + + def _buffer_record(self, run_id: object, record: LogRecord) -> None: + if run_id in self._buffer: + self._buffer.move_to_end(run_id) + else: + if len(self._buffer) >= self.max_runs: + self._evict_oldest_run() + self._buffer[run_id] = [] + + self._buffer[run_id].append(record) + self._buffered_count += 1 + + while self._buffered_count > self.max_buffered_records and self._buffer: + self._evict_oldest_run() + + def _evict_oldest_run(self) -> None: + _, oldest_records = self._buffer.popitem(last=False) + self._buffered_count -= len(oldest_records) diff --git a/logquill/plugins/run_summary_plugin.py b/logquill/plugins/run_summary_plugin.py new file mode 100644 index 0000000..3a237e6 --- /dev/null +++ b/logquill/plugins/run_summary_plugin.py @@ -0,0 +1,142 @@ +from __future__ import annotations + +from collections import OrderedDict +from dataclasses import dataclass, field +from typing import Any + +from logquill.levels import Level +from logquill.plugins.plugin import Plugin +from logquill.records import LogRecord, create_record +from logquill.transports.transport import Transport + + +@dataclass +class _RunStats: + tokens_in: int = 0 + tokens_out: int = 0 + cost_usd: float = 0.0 + tool_calls: int = 0 + retries: int = 0 + errors: int = 0 + record_count: int = 0 + root_span_id: str | None = None + duration_ms: float | None = None + _extra: dict[str, Any] = field(default_factory=dict) + + +class RunSummaryPlugin(Plugin): + """Emits one summary record — total tokens, cost, tool calls, retries, + and errors — when a run's outermost span closes, so a dashboard or the + local viewer can read a run's totals without re-walking every one of + its records. + + A "run" is any distinct `meta[run_key]` value (default `run_id`, what + `RunPlugin` stamps). Its close is detected as a span record (`meta.kind + == "span"`) for that run with no `parent_span_id` — the outermost span + `Logger.span()` produces. Aggregation happens in `after_log`, so it sees + every record exactly as it was actually dispatched. + + Like `SamplingPlugin`'s tail-based elevation and `FlightRecorderPlugin`, + the summary record is written directly to `transports` (the same list + given to the `Logger`), bypassing `before_log`/`after_log` for any + plugin *after* this one — put `RunSummaryPlugin` last if that matters + for your pipeline (e.g. after `TamperEvidentPlugin`, the summary record + itself won't be hash-chained). + + Bounded: at most `max_runs` runs' aggregates are held in memory at + once — the oldest in-progress run's aggregate is dropped (not + summarized) if a new run's first record arrives once that limit is hit, + the same bounded-memory trade-off `SamplingPlugin`/`FlightRecorderPlugin` + make for their own per-run/per-trace state. + """ + + def __init__( + self, + *, + transports: list[Transport], + run_key: str = "run_id", + logger_name: str = "app.run_summary", + max_runs: int = 1000, + ) -> None: + """`transports` must be the same list given to the `Logger` — see + the class docstring for why. `logger_name` is the `logger` field the + emitted summary record carries.""" + self.transports = transports + self.run_key = run_key + self.logger_name = logger_name + self.max_runs = max_runs + self._stats: OrderedDict[object, _RunStats] = OrderedDict() + + def after_log(self, record: LogRecord) -> None: + """Aggregates `record` into its run's running totals, and emits (and + forgets) that run's summary once its outermost span closes.""" + meta = record["meta"] + run_id = meta.get(self.run_key) + if run_id is None: + return + + stats = self._stats.get(run_id) + if stats is None: + if len(self._stats) >= self.max_runs: + self._stats.popitem(last=False) + stats = _RunStats() + self._stats[run_id] = stats + else: + self._stats.move_to_end(run_id) + + stats.record_count += 1 + if Level[record["level"]] >= Level.ERROR: + stats.errors += 1 + retry_count = meta.get("retry_count") + if isinstance(retry_count, int) and not isinstance(retry_count, bool): + stats.retries += retry_count + if meta.get("kind") == "action" and isinstance(meta.get("tool"), str): + stats.tool_calls += 1 + + llm = record.get("llm") + if isinstance(llm, dict): + if isinstance(llm.get("tokens_in"), int) and not isinstance(llm.get("tokens_in"), bool): + stats.tokens_in += llm["tokens_in"] + if isinstance(llm.get("tokens_out"), int) and not isinstance( + llm.get("tokens_out"), bool + ): + stats.tokens_out += llm["tokens_out"] + cost = llm.get("cost_usd") + if isinstance(cost, (int, float)) and not isinstance(cost, bool): + stats.cost_usd += cost + + is_root_span = ( + meta.get("kind") == "span" + and isinstance(meta.get("span_id"), str) + and "parent_span_id" not in meta + ) + if is_root_span: + stats.root_span_id = meta.get("span_id") + duration = meta.get("duration_ms") + if isinstance(duration, (int, float)) and not isinstance(duration, bool): + stats.duration_ms = float(duration) + del self._stats[run_id] + self._emit_summary(run_id, stats) + + def _emit_summary(self, run_id: object, stats: _RunStats) -> None: + summary_meta: dict[str, Any] = { + "kind": "run_summary", + self.run_key: run_id, + "tokens_in": stats.tokens_in, + "tokens_out": stats.tokens_out, + "cost_usd": stats.cost_usd, + "tool_calls": stats.tool_calls, + "retries": stats.retries, + "errors": stats.errors, + "record_count": stats.record_count, + } + if stats.duration_ms is not None: + summary_meta["duration_ms"] = stats.duration_ms + if stats.root_span_id is not None: + summary_meta["parent_span_id"] = stats.root_span_id + + summary = create_record( + level=Level.INFO, logger=self.logger_name, message="run summary", meta=summary_meta + ) + for transport in self.transports: + transport.write(transport.format(summary), summary) diff --git a/logquill/runtime_level.py b/logquill/runtime_level.py new file mode 100644 index 0000000..8896b88 --- /dev/null +++ b/logquill/runtime_level.py @@ -0,0 +1,197 @@ +"""Changing a `Logger`'s level at runtime, without restarting the process. + +Three independent mechanisms, matching what a deployment commonly has +available — pick whichever fits; none of them needs the other two: + +- `install_signal_level_handler` — an OS signal (`kill -USR1 `) + re-reads an environment variable and applies it. +- `LevelFileWatcher` — a background thread polls a file an orchestrator + (or you, by hand) writes a level name into. +- `LevelEnvWatcher` — a background thread polls an environment variable + something *inside* this same process updates (a config-reload callback, + `python-dotenv` with `override=True`, ...). An OS-level environment + change made from *outside* the process is never visible to code already + running inside it — that's a real, unavoidable limit of what an + environment variable is, not a gap in this watcher; use the file or + signal mechanism for anything triggered externally. + +An unrecognized or missing value is ignored with a warning in every case — +none of these may ever raise out of a signal handler or a background +thread. +""" + +from __future__ import annotations + +import logging +import os +import signal +import threading +from pathlib import Path +from typing import Callable + +from logquill.logger import Logger + +_logger = logging.getLogger("logquill") + + +def _apply_level(logger: Logger, raw: str | None, *, source: str) -> None: + if raw is None: + return + value = raw.strip() + if not value: + return + try: + logger.set_level(value) + except (TypeError, ValueError): + _logger.warning( + "logquill: ignoring %s=%r for logger %r — not a valid level " + "name (TRACE, DEBUG, INFO, WARN, ERROR, FATAL) or number", + source, + raw, + logger.name, + ) + + +def install_signal_level_handler( + logger: Logger, + *, + signum: int | None = None, + env_var: str = "LOGQUILL_LEVEL", +) -> Callable[[], None]: + """Registers a handler so sending `signum` to this process (default + `SIGUSR1` — `kill -USR1 `) re-reads `env_var` and applies it via + `logger.set_level()`. + + Returns a function that uninstalls the handler, restoring whatever + handler was registered for `signum` before this call — safe to install + for more than one logger and uninstall in any order, since each call + only remembers and restores its own previous handler. + + Python signal handlers only run on the main thread, and only the main + thread may call `signal.signal()` — call this from your process's main + thread, typically near startup. Raises `RuntimeError` on a platform + with no `SIGUSR1` (Windows) unless you pass an explicit `signum`. + """ + if signum is None: + try: + signum = signal.SIGUSR1 + except AttributeError: + raise RuntimeError( + "install_signal_level_handler(): this platform has no SIGUSR1 " + "(Windows) — pass an explicit signum, or use LevelFileWatcher/" + "LevelEnvWatcher instead, which don't need a POSIX signal" + ) from None + + previous = signal.getsignal(signum) + + def handler(received_signum: int, frame: object) -> None: + _apply_level(logger, os.environ.get(env_var), source=env_var) + + signal.signal(signum, handler) + + def uninstall() -> None: + signal.signal(signum, previous) + + return uninstall + + +class _PollingLevelWatcher: + """Shared polling loop behind `LevelFileWatcher`/`LevelEnvWatcher`: a + daemon thread that calls `read()` every `poll_interval` seconds and + applies the result if it changed since the last check.""" + + def __init__( + self, + logger: Logger, + *, + poll_interval: float, + source_name: str, + read: Callable[[], str | None], + ) -> None: + self._logger = logger + self._poll_interval = poll_interval + self._source_name = source_name + self._read = read + self._stop_event = threading.Event() + self._thread: threading.Thread | None = None + self._last_seen: str | None = None + + def poll_once(self) -> None: + """Checks the source once and applies it if it changed since the + last check — the single unit of work `start()`'s background loop + repeats. Exposed directly so a test (or a caller on its own + schedule) can drive this deterministically, with no thread or + sleep involved.""" + raw = self._read() + if raw != self._last_seen: + self._last_seen = raw + _apply_level(self._logger, raw, source=self._source_name) + + def start(self) -> None: + """Starts the background polling thread. Idempotent: calling this + again while already running does nothing.""" + if self._thread is not None: + return + self._stop_event.clear() + self._thread = threading.Thread(target=self._run, name="logquill-level-watch", daemon=True) + self._thread.start() + + def stop(self, timeout: float | None = 5.0) -> None: + """Stops the background thread, waiting up to `timeout` seconds for + it to actually exit. Safe to call even if never started, or more + than once.""" + self._stop_event.set() + if self._thread is not None: + self._thread.join(timeout=timeout) + self._thread = None + + def _run(self) -> None: + while not self._stop_event.is_set(): + self.poll_once() + self._stop_event.wait(self._poll_interval) + + +class LevelFileWatcher(_PollingLevelWatcher): + """Polls `path`'s content (expected to be a bare level name like + `"DEBUG"`, or a number) every `poll_interval` seconds, applying it via + `logger.set_level()` whenever it changes — for an orchestrator (or you, + by hand: `echo DEBUG > /tmp/app.level`) that can write a file but not + send a signal. + + A missing file reads as `None` — not an error — so starting the + watcher before the file exists, or deleting it later, doesn't raise; + it's simply treated as "no change" until the file reappears with new + content. + """ + + def __init__(self, logger: Logger, path: str | Path, *, poll_interval: float = 2.0) -> None: + self._path = Path(path) + + def read() -> str | None: + try: + return self._path.read_text(encoding="utf-8") + except OSError: + return None + + super().__init__(logger, poll_interval=poll_interval, source_name=str(path), read=read) + + +class LevelEnvWatcher(_PollingLevelWatcher): + """Polls `os.environ[env_var]` every `poll_interval` seconds, applying + it via `logger.set_level()` whenever it changes. + + Only helps when something *inside this same process* is the one + changing the environment (a config-reload callback, + `python-dotenv.load_dotenv(override=True)`, your own code doing + `os.environ[...] = ...`) — an environment variable changed from + *outside* the process, the way a shell or orchestrator normally would, + is never visible to code already running inside it. Use + `LevelFileWatcher` or `install_signal_level_handler` for anything + triggered externally. + """ + + def __init__(self, logger: Logger, env_var: str, *, poll_interval: float = 2.0) -> None: + def read() -> str | None: + return os.environ.get(env_var) + + super().__init__(logger, poll_interval=poll_interval, source_name=env_var, read=read) diff --git a/tests/test_plugins/test_adaptive_sampling_plugin.py b/tests/test_plugins/test_adaptive_sampling_plugin.py new file mode 100644 index 0000000..2459330 --- /dev/null +++ b/tests/test_plugins/test_adaptive_sampling_plugin.py @@ -0,0 +1,237 @@ +from __future__ import annotations + +import pytest + +from logquill import Logger +from logquill.plugins.adaptive_sampling_plugin import AdaptiveSamplingPlugin +from logquill.transports.transport import CollectingTransport + + +def _logger(**kwargs: object) -> tuple[Logger, CollectingTransport, AdaptiveSamplingPlugin]: + sink = CollectingTransport() + plugin = AdaptiveSamplingPlugin(**kwargs) # type: ignore[arg-type] + return Logger("app.test", transports=[sink], plugins=[plugin]), sink, plugin + + +# --- construction validation -------------------------------------------------- + + +@pytest.mark.parametrize("field", ["base_rate", "min_rate", "max_rate"]) +def test_a_rate_outside_0_1_raises(field: str) -> None: + with pytest.raises(ValueError, match=field): + AdaptiveSamplingPlugin(**{field: 1.5}) # type: ignore[arg-type] + + +def test_min_rate_above_max_rate_raises() -> None: + with pytest.raises(ValueError, match="min_rate"): + AdaptiveSamplingPlugin(min_rate=0.9, max_rate=0.1) + + +# --- always-kept categories --------------------------------------------------- + + +def test_errors_are_always_kept_regardless_of_sampling() -> None: + logger, sink, _ = _logger(base_rate=0.0, rng=lambda: 0.99) + + logger.error("always kept") + + assert len(sink.records) == 1 + + +def test_fatal_is_also_always_kept() -> None: + logger, sink, _ = _logger(base_rate=0.0, rng=lambda: 0.99) + + logger.fatal("always kept") + + assert len(sink.records) == 1 + + +def test_ordinary_info_is_dropped_when_the_rate_misses() -> None: + logger, sink, _ = _logger(base_rate=0.0, rng=lambda: 0.99) + + logger.info("dropped") + + assert sink.records == [] + + +def test_a_slow_span_is_always_kept() -> None: + logger, sink, _ = _logger(base_rate=0.0, rng=lambda: 0.99, slow_ms=10) + + # an explicit duration_ms in **meta overrides the span's own computed + # timing (SpanContext.__exit__ unpacks it last), making this deterministic + with logger.span("work", duration_ms=50): + pass + + assert [r["message"] for r in sink.records] == ["work"] + + +def test_a_fast_span_is_not_automatically_kept() -> None: + sink = CollectingTransport() + plugin = AdaptiveSamplingPlugin(base_rate=0.0, rng=lambda: 0.99, slow_ms=10_000) + logger = Logger("app.test", transports=[sink], plugins=[plugin]) + + with logger.span("fast"): + pass + + assert sink.records == [] + + +def test_a_non_span_record_is_never_treated_as_slow_even_with_a_duration_ms() -> None: + logger, sink, _ = _logger(base_rate=0.0, rng=lambda: 0.99, slow_ms=10) + + logger.info("not a span", duration_ms=9999) + + assert sink.records == [] + + +def test_keep_at_is_configurable() -> None: + logger, sink, _ = _logger(base_rate=0.0, rng=lambda: 0.99, keep_at="WARN") + + logger.warn("kept at the lower threshold") + logger.info("still dropped") + + assert [r["message"] for r in sink.records] == ["kept at the lower threshold"] + + +# --- byte budget --------------------------------------------------------------- + + +def test_the_byte_budget_caps_records_within_one_window() -> None: + logger, sink, _ = _logger( + base_rate=1.0, rng=lambda: 0.0, max_bytes_per_second=300, clock=lambda: 0.0 + ) + + for i in range(20): + logger.info("x" * 50, i=i) + + assert 0 < len(sink.records) < 20 + + +def test_a_record_too_big_for_an_empty_budget_is_dropped() -> None: + logger, sink, _ = _logger( + base_rate=1.0, rng=lambda: 0.0, max_bytes_per_second=5, clock=lambda: 0.0 + ) + + logger.info("way more than five bytes of content here") + + assert sink.records == [] + + +def test_the_byte_budget_resets_on_a_new_window() -> None: + clock = {"t": 0.0} + logger, sink, _ = _logger( + base_rate=1.0, rng=lambda: 0.0, max_bytes_per_second=200, clock=lambda: clock["t"] + ) + + logger.info("x" * 40) # ~161 bytes: fits the 200-byte budget + logger.info("y" * 40) # a second one this window would not: dropped + assert len(sink.records) == 1 + + clock["t"] = 10.0 # well past window_seconds=1.0: budget resets + logger.info("z" * 40) + + assert len(sink.records) == 2 + + +# --- adaptivity ------------------------------------------------------------ + + +def test_rate_rises_after_a_window_with_low_demand() -> None: + clock = {"t": 0.0} + logger, _, plugin = _logger( + base_rate=0.1, + rng=lambda: 0.0, + clock=lambda: clock["t"], + max_bytes_per_second=1_000_000, + ) + + logger.info("a") + clock["t"] = 2.0 + logger.info("b") # rolls the window; last window's demand was tiny + + assert plugin.rate > 0.1 + + +def test_rate_falls_after_a_window_with_high_demand() -> None: + clock = {"t": 0.0} + logger, _, plugin = _logger( + base_rate=0.5, rng=lambda: 0.0, clock=lambda: clock["t"], max_bytes_per_second=10 + ) + + logger.info("x" * 200) # demands far more than the 10-byte budget + clock["t"] = 2.0 + logger.info("y") + + assert plugin.rate < 0.5 + + +def test_rate_never_exceeds_max_rate() -> None: + clock = {"t": 0.0} + logger, _, plugin = _logger( + base_rate=0.9, + max_rate=0.95, + rng=lambda: 0.0, + clock=lambda: clock["t"], + max_bytes_per_second=1_000_000, + ) + + for t in range(1, 20): + clock["t"] = float(t) + logger.info("a") + + assert plugin.rate <= 0.95 + + +def test_rate_never_drops_below_min_rate() -> None: + clock = {"t": 0.0} + logger, _, plugin = _logger( + base_rate=0.1, + min_rate=0.05, + rng=lambda: 0.0, + clock=lambda: clock["t"], + max_bytes_per_second=1, + ) + + for t in range(1, 20): + clock["t"] = float(t) + logger.info("x" * 200) + + assert plugin.rate >= 0.05 + + +def test_demand_that_is_dropped_by_the_budget_still_lowers_the_rate() -> None: + """The bug this guards against: using *emitted* bytes (which stop + growing once the budget saturates) instead of *demanded* bytes as the + adaptivity signal would read a saturated gate as "plenty of room" and + raise the rate, even though far more was being asked for than the + budget allows.""" + clock = {"t": 0.0} + logger, sink, plugin = _logger( + base_rate=0.5, rng=lambda: 0.0, clock=lambda: clock["t"], max_bytes_per_second=10 + ) + + # the very first record already exceeds the 10-byte budget, so nothing + # is ever actually emitted this window — only demanded + logger.info("x" * 200) + assert sink.records == [] + clock["t"] = 2.0 + logger.info("y") + + assert plugin.rate < 0.5 + + +def test_a_long_idle_gap_does_not_retroactively_apply_several_windows_worth() -> None: + clock = {"t": 0.0} + logger, _, plugin = _logger( + base_rate=0.1, rng=lambda: 0.0, clock=lambda: clock["t"], max_bytes_per_second=1_000_000 + ) + + logger.info("a") + clock["t"] = 100.0 # a long idle gap + logger.info("b") + rate_after_first_roll = plugin.rate + clock["t"] = 101.0 + logger.info("c") + + # only one more adjustment step happened, not ~100 compounding ones + assert rate_after_first_roll < plugin.rate <= rate_after_first_roll * 1.21 diff --git a/tests/test_plugins/test_flight_recorder_plugin.py b/tests/test_plugins/test_flight_recorder_plugin.py new file mode 100644 index 0000000..0a7bc7a --- /dev/null +++ b/tests/test_plugins/test_flight_recorder_plugin.py @@ -0,0 +1,147 @@ +from __future__ import annotations + +from logquill.logger import Logger +from logquill.plugins.flight_recorder_plugin import FlightRecorderPlugin +from logquill.transports.transport import CollectingTransport + + +def test_a_debug_record_below_ship_at_is_buffered_not_shipped() -> None: + sink = CollectingTransport() + recorder = FlightRecorderPlugin(transports=[sink]) + logger = Logger("app.test", level="DEBUG", transports=[sink], plugins=[recorder]) + + assert logger.debug("context", run_id="r1") is None + assert sink.records == [] + + +def test_an_info_record_ships_immediately_unbuffered() -> None: + sink = CollectingTransport() + recorder = FlightRecorderPlugin(transports=[sink]) + logger = Logger("app.test", level="DEBUG", transports=[sink], plugins=[recorder]) + + record = logger.info("normal", run_id="r1") + + assert record is not None + assert sink.records == [record] + + +def test_a_record_with_no_run_id_passes_through_regardless_of_level() -> None: + sink = CollectingTransport() + recorder = FlightRecorderPlugin(transports=[sink]) + logger = Logger("app.test", level="DEBUG", transports=[sink], plugins=[recorder]) + + record = logger.debug("loose") + + assert record is not None + assert sink.records == [record] + + +def test_a_failing_run_ships_its_full_debug_trail() -> None: + sink = CollectingTransport() + recorder = FlightRecorderPlugin(transports=[sink]) + logger = Logger("app.test", level="DEBUG", transports=[sink], plugins=[recorder]) + + logger.debug("step 1", run_id="r1") + logger.debug("step 2", run_id="r1") + assert sink.records == [] + + record = logger.error("step 3", run_id="r1") + + assert record is not None + messages = [r["message"] for r in sink.records] + assert messages == ["step 1", "step 2", "step 3"] + + +def test_a_passing_run_never_ships_its_debug_trail() -> None: + sink = CollectingTransport() + recorder = FlightRecorderPlugin(transports=[sink]) + logger = Logger("app.test", level="DEBUG", transports=[sink], plugins=[recorder]) + + logger.debug("step 1", run_id="r1") + logger.debug("step 2", run_id="r1") + logger.info("step 3", run_id="r1") # ships on its own; run never errors + + messages = [r["message"] for r in sink.records] + assert messages == ["step 3"] + + +def test_flushing_only_affects_the_matching_run() -> None: + sink = CollectingTransport() + recorder = FlightRecorderPlugin(transports=[sink]) + logger = Logger("app.test", level="DEBUG", transports=[sink], plugins=[recorder]) + + logger.debug("other run", run_id="r2") + logger.debug("step 1", run_id="r1") + logger.error("step 2", run_id="r1") + + messages = [r["message"] for r in sink.records] + assert "other run" not in messages + assert messages == ["step 1", "step 2"] + + +def test_records_after_flushing_ship_unconditionally_even_below_ship_at() -> None: + sink = CollectingTransport() + recorder = FlightRecorderPlugin(transports=[sink]) + logger = Logger("app.test", level="DEBUG", transports=[sink], plugins=[recorder]) + + logger.error("triggers flush", run_id="r1") + record = logger.debug("after flush", run_id="r1") + + assert record is not None + assert sink.records[-1]["message"] == "after flush" + + +def test_flush_at_can_be_reached_without_meeting_ship_at() -> None: + # a WARN below the default ship_at (INFO is below WARN, so this is + # actually above — use a custom, lower ship_at/flush_at pairing instead + # to exercise flush_at triggering independently of ship_at. + sink = CollectingTransport() + recorder = FlightRecorderPlugin(transports=[sink], ship_at="ERROR", flush_at="WARN") + logger = Logger("app.test", level="DEBUG", transports=[sink], plugins=[recorder]) + + logger.info("buffered", run_id="r1") # below ship_at=ERROR: buffered + record = logger.warn("flush trigger", run_id="r1") # below ship_at, but reaches flush_at + + assert record is not None + messages = [r["message"] for r in sink.records] + assert messages == ["buffered", "flush trigger"] + + +def test_buffer_is_bounded_by_max_runs() -> None: + sink = CollectingTransport() + recorder = FlightRecorderPlugin(transports=[sink], max_runs=1) + logger = Logger("app.test", level="DEBUG", transports=[sink], plugins=[recorder]) + + logger.debug("run one", run_id="r1") + logger.debug("run two", run_id="r2") # evicts r1's buffer (max_runs=1) + logger.error("errors r1", run_id="r1") + + messages = [r["message"] for r in sink.records] + assert "run one" not in messages + assert "errors r1" in messages + + +def test_buffer_is_bounded_by_max_buffered_records() -> None: + sink = CollectingTransport() + recorder = FlightRecorderPlugin(transports=[sink], max_buffered_records=1) + logger = Logger("app.test", level="DEBUG", transports=[sink], plugins=[recorder]) + + logger.debug("run one, record one", run_id="r1") + logger.debug("run one, record two", run_id="r1") # evicts r1's first record + logger.error("errors r1", run_id="r1") + + messages = [r["message"] for r in sink.records] + assert "run one, record one" not in messages + assert "errors r1" in messages + + +def test_a_broken_plugin_pipeline_bounded_memory_under_a_long_passing_run() -> None: + sink = CollectingTransport() + recorder = FlightRecorderPlugin(transports=[sink], max_buffered_records=50, max_runs=10) + logger = Logger("app.test", level="DEBUG", transports=[sink], plugins=[recorder]) + + for i in range(5000): + logger.debug(f"step {i}", run_id=f"r{i % 20}") + + assert recorder._buffered_count <= 50 + assert len(recorder._buffer) <= 10 diff --git a/tests/test_plugins/test_run_summary_plugin.py b/tests/test_plugins/test_run_summary_plugin.py new file mode 100644 index 0000000..ea48785 --- /dev/null +++ b/tests/test_plugins/test_run_summary_plugin.py @@ -0,0 +1,182 @@ +from __future__ import annotations + +from logquill import Logger, RunPlugin +from logquill.plugins.run_summary_plugin import RunSummaryPlugin +from logquill.transports.transport import CollectingTransport + + +def _logger(**kwargs: object) -> tuple[Logger, CollectingTransport, RunSummaryPlugin]: + sink = CollectingTransport() + summary = RunSummaryPlugin(transports=[sink], **kwargs) # type: ignore[arg-type] + logger = Logger( + "app.agent", + transports=[sink], + plugins=[RunPlugin(run_id="run-1"), summary], + content_policy="full", + ) + return logger, sink, summary + + +def _run_summary(sink: CollectingTransport) -> dict: + (summary,) = (r for r in sink.records if r["meta"].get("kind") == "run_summary") + return summary + + +def test_no_summary_until_the_outermost_span_closes() -> None: + logger, sink, _ = _logger() + + with logger.span("run"): + logger.info("step") + assert all(r["meta"].get("kind") != "run_summary" for r in sink.records) + + +def test_summary_emitted_when_the_outermost_span_closes() -> None: + logger, sink, _ = _logger() + + with logger.span("run"): + logger.info("step") + + summary = _run_summary(sink) + assert summary["message"] == "run summary" + assert summary["logger"] == "app.run_summary" + assert summary["level"] == "INFO" + + +def test_tokens_and_cost_are_summed_across_every_llm_call() -> None: + logger, sink, _ = _logger() + + with logger.span("run"): + logger.llm_call("chat", model="m", tokens_in=100, tokens_out=20, cost_usd=0.01) + logger.llm_call("chat", model="m", tokens_in=50, tokens_out=10, cost_usd=0.02) + + summary = _run_summary(sink) + assert summary["meta"]["tokens_in"] == 150 + assert summary["meta"]["tokens_out"] == 30 + assert summary["meta"]["cost_usd"] == 0.03 + + +def test_tool_calls_are_counted() -> None: + logger, sink, _ = _logger() + + with logger.span("run"): + logger.action("call", tool="search") + logger.action("call", tool="search") # retry of the same tool + logger.action("think") # not a tool call — no `tool=` + logger.observation("done", tool="search") + + summary = _run_summary(sink) + assert summary["meta"]["tool_calls"] == 2 + + +def test_retry_counts_are_summed() -> None: + logger, sink, _ = _logger() + + with logger.span("run"): + logger.action("call", tool="search") # retry_count=None + logger.action("call", tool="search") # auto-tracked: retry_count=1 + + summary = _run_summary(sink) + assert summary["meta"]["retries"] == 1 + + +def test_errors_are_counted() -> None: + logger, sink, _ = _logger() + + with logger.span("run"): + logger.error("first") + logger.fatal("second") + logger.warn("not counted") + logger.info("not counted either") + + summary = _run_summary(sink) + assert summary["meta"]["errors"] == 2 + + +def test_record_count_and_duration_are_reported() -> None: + logger, sink, _ = _logger() + + with logger.span("run"): + logger.info("one") + logger.info("two") + + summary = _run_summary(sink) + # one + two + the root span's own closing record = 3 + assert summary["meta"]["record_count"] == 3 + assert isinstance(summary["meta"]["duration_ms"], float) + + +def test_summary_is_nested_under_the_runs_root_span() -> None: + logger, sink, _ = _logger() + + with logger.span("run"): + pass + + summary = _run_summary(sink) + root = next(r for r in sink.records if r["message"] == "run") + assert summary["meta"]["parent_span_id"] == root["meta"]["span_id"] + + +def test_a_run_with_no_errors_still_gets_a_summary() -> None: + logger, sink, _ = _logger() + + with logger.span("run"): + logger.info("fine") + + summary = _run_summary(sink) + assert summary["meta"]["errors"] == 0 + + +def test_two_runs_are_summarized_independently() -> None: + sink = CollectingTransport() + summary_plugin = RunSummaryPlugin(transports=[sink]) + logger_a = Logger( + "app.agent", + transports=[sink], + plugins=[RunPlugin(run_id="run-a"), summary_plugin], + ) + logger_b = Logger( + "app.agent", + transports=[sink], + plugins=[RunPlugin(run_id="run-b"), summary_plugin], + ) + + with logger_a.span("run"): + logger_a.action("call", tool="search") + with logger_b.span("run"): + pass + + summaries = { + r["meta"]["run_id"]: r for r in sink.records if r["meta"].get("kind") == "run_summary" + } + assert summaries["run-a"]["meta"]["tool_calls"] == 1 + assert summaries["run-b"]["meta"]["tool_calls"] == 0 + + +def test_a_record_with_no_run_id_is_ignored() -> None: + sink = CollectingTransport() + summary_plugin = RunSummaryPlugin(transports=[sink]) + logger = Logger("app.agent", transports=[sink], plugins=[summary_plugin]) + + logger.info("loose, no run_id") + + assert sink.records == [sink.records[0]] # just the one record, no summary + + +def test_a_nested_span_closing_does_not_trigger_a_summary() -> None: + logger, sink, _ = _logger() + + with logger.span("run"): + with logger.span("step"): + pass + assert all(r["meta"].get("kind") != "run_summary" for r in sink.records) + + +def test_aggregation_state_is_bounded_by_max_runs() -> None: + sink = CollectingTransport() + summary_plugin = RunSummaryPlugin(transports=[sink], max_runs=2) + logger = Logger("app.agent", transports=[sink], plugins=[summary_plugin]) + + for i in range(10): + logger.info("step", run_id=f"r{i}") + + assert len(summary_plugin._stats) <= 2 diff --git a/tests/test_runtime_level.py b/tests/test_runtime_level.py new file mode 100644 index 0000000..32b93e7 --- /dev/null +++ b/tests/test_runtime_level.py @@ -0,0 +1,229 @@ +from __future__ import annotations + +import logging +import os +import signal +import time +from pathlib import Path + +import pytest + +from logquill import Level, Logger +from logquill.runtime_level import ( + LevelEnvWatcher, + LevelFileWatcher, + install_signal_level_handler, +) + +# --- install_signal_level_handler -------------------------------------------- + + +def test_sending_the_signal_applies_the_env_vars_level(monkeypatch: pytest.MonkeyPatch) -> None: + logger = Logger("app", level="INFO") + monkeypatch.setenv("TEST_LOGQUILL_LEVEL", "DEBUG") + uninstall = install_signal_level_handler( + logger, signum=signal.SIGUSR1, env_var="TEST_LOGQUILL_LEVEL" + ) + try: + os.kill(os.getpid(), signal.SIGUSR1) + finally: + uninstall() + + assert logger.level == Level.DEBUG + + +def test_uninstall_restores_the_previous_handler() -> None: + logger = Logger("app") + sentinel_calls = [] + original = signal.signal(signal.SIGUSR1, lambda *_: sentinel_calls.append(1)) + try: + uninstall = install_signal_level_handler(logger, signum=signal.SIGUSR1) + uninstall() + + os.kill(os.getpid(), signal.SIGUSR1) + time.sleep(0.05) + + assert sentinel_calls == [1] # the original handler ran, not logquill's + finally: + signal.signal(signal.SIGUSR1, original) + + +def test_missing_env_var_on_signal_leaves_the_level_unchanged( + monkeypatch: pytest.MonkeyPatch, +) -> None: + logger = Logger("app", level="WARN") + monkeypatch.delenv("TEST_LOGQUILL_LEVEL_UNSET", raising=False) + uninstall = install_signal_level_handler( + logger, signum=signal.SIGUSR1, env_var="TEST_LOGQUILL_LEVEL_UNSET" + ) + try: + os.kill(os.getpid(), signal.SIGUSR1) + finally: + uninstall() + + assert logger.level == Level.WARN + + +def test_an_invalid_level_value_warns_and_leaves_the_level_unchanged( + monkeypatch: pytest.MonkeyPatch, caplog: pytest.LogCaptureFixture +) -> None: + logger = Logger("app", level="WARN") + monkeypatch.setenv("TEST_LOGQUILL_LEVEL", "NOT_A_LEVEL") + uninstall = install_signal_level_handler( + logger, signum=signal.SIGUSR1, env_var="TEST_LOGQUILL_LEVEL" + ) + try: + with caplog.at_level(logging.WARNING, logger="logquill"): + os.kill(os.getpid(), signal.SIGUSR1) + finally: + uninstall() + + assert logger.level == Level.WARN + assert any("not a valid level" in r.getMessage() for r in caplog.records) + + +def test_without_sigusr1_on_this_platform_raises_unless_an_explicit_signum_is_given( + monkeypatch: pytest.MonkeyPatch, +) -> None: + import logquill.runtime_level as module + + class NoSigUsr1: + def getsignal(self, *a: object) -> object: + return signal.SIG_DFL + + def signal(self, *a: object) -> object: + return signal.SIG_DFL + + monkeypatch.setattr(module, "signal", NoSigUsr1()) + logger = Logger("app") + + with pytest.raises(RuntimeError, match="SIGUSR1"): + install_signal_level_handler(logger) + + # an explicit signum sidesteps the missing SIGUSR1 entirely + uninstall = install_signal_level_handler(logger, signum=signal.SIGUSR1) + uninstall() + + +# --- _PollingLevelWatcher / LevelFileWatcher --------------------------------- + + +def test_file_watcher_applies_the_files_content(tmp_path: Path) -> None: + logger = Logger("app", level="INFO") + path = tmp_path / "level.txt" + path.write_text("DEBUG") + watcher = LevelFileWatcher(logger, path) + + watcher.poll_once() + + assert logger.level == Level.DEBUG + + +def test_file_watcher_before_the_file_exists_is_a_no_op() -> None: + logger = Logger("app", level="WARN") + watcher = LevelFileWatcher(logger, "/no/such/file/at/all.txt") + + watcher.poll_once() + + assert logger.level == Level.WARN + + +def test_file_watcher_only_reapplies_on_an_actual_change(tmp_path: Path) -> None: + logger = Logger("app", level="INFO") + path = tmp_path / "level.txt" + path.write_text("DEBUG") + watcher = LevelFileWatcher(logger, path) + watcher.poll_once() + logger.set_level("WARN") # someone else changes it in between + + watcher.poll_once() # file content unchanged since last poll + + assert logger.level == Level.WARN # not stomped back to DEBUG + + +def test_file_watcher_picks_up_a_later_change(tmp_path: Path) -> None: + logger = Logger("app", level="INFO") + path = tmp_path / "level.txt" + path.write_text("DEBUG") + watcher = LevelFileWatcher(logger, path) + watcher.poll_once() + + path.write_text("ERROR") + watcher.poll_once() + + assert logger.level == Level.ERROR + + +def test_file_watcher_tolerates_a_trailing_newline(tmp_path: Path) -> None: + logger = Logger("app") + path = tmp_path / "level.txt" + path.write_text("DEBUG\n") + watcher = LevelFileWatcher(logger, path) + + watcher.poll_once() + + assert logger.level == Level.DEBUG + + +def test_file_watcher_start_and_stop_run_a_real_background_thread(tmp_path: Path) -> None: + logger = Logger("app", level="INFO") + path = tmp_path / "level.txt" + watcher = LevelFileWatcher(logger, path, poll_interval=0.02) + watcher.start() + try: + path.write_text("TRACE") + for _ in range(100): + if logger.level == Level.TRACE: + break + time.sleep(0.01) + assert logger.level == Level.TRACE + finally: + watcher.stop() + + +def test_start_is_idempotent_and_stop_is_safe_without_start(tmp_path: Path) -> None: + logger = Logger("app") + watcher = LevelFileWatcher(logger, tmp_path / "level.txt", poll_interval=0.02) + + watcher.start() + watcher.start() # no second thread, no error + watcher.stop() + watcher.stop() # safe to call again + + never_started = LevelFileWatcher(logger, tmp_path / "level.txt") + never_started.stop() # safe even though start() was never called + + +# --- LevelEnvWatcher ---------------------------------------------------------- + + +def test_env_watcher_applies_a_changed_value(monkeypatch: pytest.MonkeyPatch) -> None: + logger = Logger("app", level="INFO") + monkeypatch.setenv("TEST_LOGQUILL_ENV_LEVEL", "ERROR") + watcher = LevelEnvWatcher(logger, "TEST_LOGQUILL_ENV_LEVEL") + + watcher.poll_once() + + assert logger.level == Level.ERROR + + +def test_env_watcher_with_no_var_set_is_a_no_op(monkeypatch: pytest.MonkeyPatch) -> None: + logger = Logger("app", level="WARN") + monkeypatch.delenv("TEST_LOGQUILL_ENV_LEVEL_UNSET", raising=False) + watcher = LevelEnvWatcher(logger, "TEST_LOGQUILL_ENV_LEVEL_UNSET") + + watcher.poll_once() + + assert logger.level == Level.WARN + + +def test_env_watcher_only_reapplies_on_an_actual_change(monkeypatch: pytest.MonkeyPatch) -> None: + logger = Logger("app", level="INFO") + monkeypatch.setenv("TEST_LOGQUILL_ENV_LEVEL", "DEBUG") + watcher = LevelEnvWatcher(logger, "TEST_LOGQUILL_ENV_LEVEL") + watcher.poll_once() + logger.set_level("ERROR") + + watcher.poll_once() # unchanged since last poll + + assert logger.level == Level.ERROR