Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
24 changes: 24 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
110 changes: 110 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 <run_id> --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

Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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 <pid> 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
Expand Down
69 changes: 69 additions & 0 deletions benchmarks/test_flight_recorder_memory.py
Original file line number Diff line number Diff line change
@@ -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
14 changes: 14 additions & 0 deletions logquill/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -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 (
Expand All @@ -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
Expand Down Expand Up @@ -65,6 +73,7 @@
__version__ = "1.0.0"

__all__ = [
"AdaptiveSamplingPlugin",
"AlertingPlugin",
"AppriseAlertPlugin",
"AppInsightsTransport",
Expand All @@ -87,12 +96,15 @@
"FIELD_CLASSES",
"FieldClass",
"FileTransport",
"FlightRecorderPlugin",
"Formatter",
"FunctionPlugin",
"HTTPTransport",
"JSONFormatter",
"KafkaTransport",
"Level",
"LevelEnvWatcher",
"LevelFileWatcher",
"LogfmtFormatter",
"LLMBlock",
"LogQuillAdapter",
Expand All @@ -115,6 +127,7 @@
"RedactPlugin",
"RedisTransport",
"RunPlugin",
"RunSummaryPlugin",
"SCHEMA_VERSION",
"SQLLogRow",
"SQLiteTransport",
Expand All @@ -134,6 +147,7 @@
"disable",
"enable",
"format_exc_info",
"install_signal_level_handler",
"is_enabled",
"load_config",
"logger_from_env",
Expand Down
Loading
Loading