diff --git a/CHANGELOG.md b/CHANGELOG.md index 728f9a2..0d4f343 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -4,6 +4,24 @@ All notable changes to this project are documented in this file. ## Unreleased +- Local-first trace viewer: + - `logquill trace --file logs.jsonl` reconstructs and prints one + agent run's span tree, annotated with each span's own duration and the + token/cost totals rolled up from everything nested under it. It streams + the file line by line, so tracing one run out of a multi-gigabyte log + costs memory proportional to that run, not the file (a gigabyte-scale + test asserts this). `--json` prints the same tree as nested data instead. + - `logquill serve --file logs.jsonl` (or `--db logs.sqlite` for a + `SQLiteTransport` database) runs a small local web UI — run list, a + combined span-tree/waterfall view, search, and a level filter — built + entirely on the stdlib (`http.server`, `sqlite3`): no new dependency, no + account, nothing leaves the machine. Reading from a SQLite database never + shows token/cost annotations, since that transport's fixed schema doesn't + store the `llm` block. + - `logquill dev logs.jsonl` follows a file like `tail -f`, but live-renders + the current run's span tree (colorized, screen-cleared between redraws) + instead of flat lines, following whichever run is most recently active + unless `--run-id` pins it to one. - 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 diff --git a/README.md b/README.md index a5d6ca4..db18b3e 100644 --- a/README.md +++ b/README.md @@ -36,6 +36,7 @@ for what's landed so far. - **Cheap when idle, precise when it counts** — `logger.opt(lazy=True)` defers expensive `meta` values until a record will really be emitted, `logger.opt(depth=N)` reports the right caller from inside a wrapper, `logquill.disable(__name__)` silences a library's own logs by default, and queued records are flushed automatically at interpreter exit — see [Lazy values, caller depth & disabling a library](#lazy-values-caller-depth--disabling-a-library) - **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 ## Install @@ -1410,6 +1411,79 @@ colors) when writing to a terminal; pass `--no-color` to disable that, or that isn't valid JSON, or isn't a JSON object, is skipped with a warning on stderr rather than aborting the whole tail. +### `logquill trace` — one run's span tree, from the command line + +`logquill trace --file logs.jsonl` reconstructs and prints one +agent run's span tree, annotated with each span's own duration and the +token/cost totals of everything nested under it — everything a hosted trace +UI shows you, from a plain JSONL file, no account or backend: + +```bash +logquill trace run-4f2a --file logs.jsonl +``` + +```text +└─ [INFO] run (812.5ms, 1540→412 tok, $0.0187) + ├─ [INFO] plan the work + ├─ [INFO] step (250.0ms, 1200→340 tok, $0.0123) + │ └─ [INFO] chat (1200→340 tok, $0.0123) + └─ [INFO] look it up (40.0ms) +``` + +It **streams** the file line by line — reconstructing one run out of a +multi-gigabyte log file costs memory proportional to that run, not the file +(tested at gigabyte scale in `benchmarks/test_trace_memory.py`). `--json` +prints the same tree as nested `{record, rollup, children}` objects instead, +for feeding into another tool. + +### `logquill serve` — a local, offline trace viewer + +`logquill serve --file logs.jsonl` runs a small web UI, entirely on the +stdlib (`http.server`) — no new dependency, no account, and nothing leaves +your machine: + +```bash +logquill serve --file logs.jsonl +# logquill serve: 12 run(s) found in logs.jsonl +# logquill serve: listening on http://127.0.0.1:52341/ — Ctrl+C to stop +``` + +Open the printed URL: a run list on the left (record count, token/cost +totals, an error badge), and a combined span-tree/waterfall view for +whichever run you click — indentation shows nesting, bar position and width +show timing. Search and the level filter both work inside a selected run +(client-side, instant) and, with nothing selected, across every run at once +(via `/api/search`, still streamed rather than loaded into memory). + +`--db logs.sqlite` reads from a `SQLiteTransport`-written database instead of +a JSONL file. One limitation, inherent to that transport's fixed table +schema: it doesn't store the `llm` block, so runs served from SQLite show no +token/cost annotations even if the original records had them — trace from +the JSONL file (or a transport that does keep `llm`) to see those. + +`logquill serve` computes the run list once at startup by streaming through +the source; it's a snapshot, not a live tail — restart it to pick up runs +logged after it started. + +### `logquill dev` — watch an agent run live + +`logquill dev logs.jsonl` follows a log file like `tail -f`, but redraws the +current run's span tree — colorized, screen cleared between redraws on a +terminal — every time a new record for it arrives, instead of printing flat +lines: + +```bash +logquill dev logs.jsonl +``` + +It tracks whichever run's records have arrived most recently by default, so +the view follows an agent from one run to the next without restarting; pass +`--run-id` to pin it to one run instead. `--backlog N` (default 5000) caps +how much of the file's *existing* content seeds the very first render, so +pointing it at a large pre-existing file doesn't stall before the first draw +— once running, though, a `dev` session keeps its own growing record list in +memory for as long as it runs, unlike `trace`'s bounded streaming. + ## API reference Every public class and function is documented with a docstring; the full diff --git a/benchmarks/test_trace_memory.py b/benchmarks/test_trace_memory.py new file mode 100644 index 0000000..93ef3cb --- /dev/null +++ b/benchmarks/test_trace_memory.py @@ -0,0 +1,133 @@ +"""The exit criterion for the local-first trace viewer: reconstructing one +run's span tree out of a multi-gigabyte log file costs memory proportional +to that run, not the file. Slow and memory-instrumented on purpose, 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 json +import tracemalloc +from pathlib import Path +from typing import Any, Iterator + +from logquill.trace_tree import build_trace + +#: Comfortably over 1 GB — the scale this test is meant to prove, without +#: depending on exactly hitting it. +TARGET_BYTES = 1_100_000_000 + +#: However large the file gets, reconstructing one small run out of it must +#: stay nowhere near proportional to the file — a few MB, not gigabytes. +MEMORY_BUDGET_BYTES = 100_000_000 + +NEEDLE_RUN_ID = "needle-run" + +_NOISE_META = {"user_id": 42, "route": "/checkout", "tags": ["a", "b", "c"], "note": "x" * 40} + + +def _record(*, run_id: str, message: str, **meta: Any) -> dict[str, Any]: + return { + "schema_version": "2.0", + "timestamp": "2026-01-01T00:00:00.000Z", + "level": "INFO", + "logger": "app.agent", + "message": message, + "meta": {"run_id": run_id, **meta}, + } + + +def _needle_run_records() -> list[dict[str, Any]]: + """A small, ordinary agent run — this is what the test must be able to + reconstruct out of the noise around it.""" + span_id, step_id = "a" * 16, "b" * 16 + return [ + _record( + run_id=NEEDLE_RUN_ID, + message="thought", + kind="thought", + parent_span_id=span_id, + ), + { + **_record(run_id=NEEDLE_RUN_ID, message="chat", kind="action", parent_span_id=step_id), + "llm": {"model": "m", "tokens_in": 10, "tokens_out": 5, "cost_usd": 0.01}, + }, + _record( + run_id=NEEDLE_RUN_ID, + message="step", + kind="span", + span_id=step_id, + parent_span_id=span_id, + duration_ms=5.0, + ), + _record( + run_id=NEEDLE_RUN_ID, + message="run", + kind="span", + span_id=span_id, + operation="invoke_agent", + agent_name="planner", + duration_ms=12.0, + ), + ] + + +def _write_large_jsonl(path: Path, *, target_bytes: int) -> int: + """Writes `target_bytes`+ of noise from many unrelated runs, with the + needle run's records inserted partway through — returns how many needle + records were written.""" + needle = [json.dumps(r, separators=(",", ":")) for r in _needle_run_records()] + noise_run_ids = [f"noise-{i}" for i in range(1000)] + written = 0 + needle_written = False + with path.open("w", encoding="utf-8") as f: + i = 0 + while written < target_bytes: + if not needle_written and written > target_bytes // 2: + for line in needle: + f.write(line + "\n") + written += len(line) + 1 + needle_written = True + line = json.dumps( + _record( + run_id=noise_run_ids[i % len(noise_run_ids)], message="noise", **_NOISE_META + ), + separators=(",", ":"), + ) + f.write(line + "\n") + written += len(line) + 1 + i += 1 + return len(needle) + + +def _read_jsonl(path: Path) -> Iterator[dict[str, Any]]: + with path.open("r", encoding="utf-8") as f: + for line in f: + yield json.loads(line) + + +def test_a_gigabyte_scale_log_file_is_traced_with_bounded_memory(tmp_path: Path) -> None: + path = tmp_path / "huge.jsonl" + needle_count = _write_large_jsonl(path, target_bytes=TARGET_BYTES) + file_size = path.stat().st_size + assert file_size >= TARGET_BYTES # the scale this test is actually about + + tracemalloc.start() + try: + baseline, _ = tracemalloc.get_traced_memory() + builder = build_trace(_read_jsonl(path), NEEDLE_RUN_ID) + _current, peak = tracemalloc.get_traced_memory() + finally: + tracemalloc.stop() + + assert builder.record_count == needle_count + (root,) = builder.roots + assert root.record["message"] == "run" + assert root.rollup().tokens_in == 10 + + peak_bytes = peak - baseline + assert peak_bytes < MEMORY_BUDGET_BYTES, ( + f"reconstructing one run used {peak_bytes:,} bytes against a " + f"{file_size:,}-byte file — memory grew with the file, not the run" + ) diff --git a/logquill/cli.py b/logquill/cli.py index d0ee650..179029e 100644 --- a/logquill/cli.py +++ b/logquill/cli.py @@ -8,10 +8,12 @@ import sys import time from pathlib import Path -from typing import IO, Any, Sequence +from typing import IO, Any, Callable, Sequence from logquill.formatters import format_text from logquill.levels import Level, parse_level +from logquill.trace_tree import build_trace, render_tree +from logquill.webui import JSONLSource, RecordSource, SQLiteSource, TraceViewerServer _COLORS = { Level.TRACE: "\x1b[90m", @@ -66,6 +68,64 @@ def build_parser() -> argparse.ArgumentParser: action="store_true", help="Disable ANSI colorization, even when writing to a terminal.", ) + + trace_parser = subparsers.add_parser( + "trace", + help="Print one agent run's span tree, reconstructed from a JSONL log file.", + ) + trace_parser.add_argument("run_id", help="The run_id (meta.run_id) to reconstruct.") + trace_parser.add_argument( + "--file", required=True, metavar="PATH", help="Path to a LogQuill JSONL log file." + ) + trace_parser.add_argument( + "--json", + action="store_true", + help="Print the tree as JSON (nested {record, rollup, children}) instead of text.", + ) + trace_parser.add_argument( + "--no-color", + action="store_true", + help="Disable ANSI colorization, even when writing to a terminal.", + ) + + serve_parser = subparsers.add_parser( + "serve", help="A local, offline web UI for browsing traced runs." + ) + source_group = serve_parser.add_mutually_exclusive_group(required=True) + source_group.add_argument("--file", metavar="PATH", help="A LogQuill JSONL log file.") + source_group.add_argument( + "--db", metavar="PATH", help="A SQLite database a SQLiteTransport wrote." + ) + serve_parser.add_argument( + "--table", default="logs", help="Table name, with --db (default: logs)." + ) + serve_parser.add_argument("--host", default="127.0.0.1", help="Bind host (default: 127.0.0.1).") + serve_parser.add_argument( + "--port", type=int, default=0, help="Bind port (default: 0, let the OS pick one)." + ) + + dev_parser = subparsers.add_parser( + "dev", + help="Follow a JSONL log file, live-rendering the current run's span tree as it grows.", + ) + dev_parser.add_argument("file", help="Path to a LogQuill JSONL log file.") + dev_parser.add_argument( + "--run-id", + default=None, + help="Track this run_id only; default is whichever run's records arrived most recently.", + ) + dev_parser.add_argument( + "--no-color", + action="store_true", + help="Disable ANSI colorization and screen-clearing, even on a terminal.", + ) + dev_parser.add_argument( + "--backlog", + type=int, + default=5000, + metavar="N", + help="At most this many existing lines seed the initial view (default: 5000).", + ) return parser @@ -215,6 +275,170 @@ def _run_tail( return 0 +def _iter_records(path: Path, *, warn_stream: IO[str]) -> Any: + """Yields every parseable JSON-object line in `path`, one at a time — + never loads the file into memory, so tracing one run out of a + multi-gigabyte file costs memory proportional to that run, not the file. + """ + with path.open("r", encoding="utf-8") as f: + for raw_line in f: + record = _parse_line(raw_line, warn_stream=warn_stream) + if record is not None: + yield record + + +def _run_trace(args: argparse.Namespace, *, out: IO[str], warn_stream: IO[str]) -> int: + path = Path(args.file) + if not path.exists(): + warn_stream.write(f"logquill trace: no such file: {args.file}\n") + return 1 + + builder = build_trace(_iter_records(path, warn_stream=warn_stream), args.run_id) + if builder.record_count == 0: + warn_stream.write( + f"logquill trace: no records with meta.run_id={args.run_id!r} found in " + f"{args.file} — check the run id, or that this file's records carry " + "run_id at all (see RunPlugin)\n" + ) + return 1 + + if args.json: + out.write(json.dumps([root.to_dict() for root in builder.roots], default=str) + "\n") + return 0 + + out.write(render_tree(builder.roots, colorize=_tree_colorizer(out, args.no_color)) + "\n") + return 0 + + +def _tree_colorizer(out: IO[str], no_color: bool) -> Callable[[str, str], str] | None: + """The `colorize` function `render_tree` wants — level-to-color, matching + `ConsoleTransport` — or `None` when colorizing doesn't make sense (piped + output, or explicitly disabled).""" + if no_color or not getattr(out, "isatty", lambda: False)(): + return None + + def colorize(text: str, level: str) -> str: + try: + color = _COLORS.get(parse_level(level)) + except (TypeError, ValueError): + color = None + return f"{color}{text}{_RESET}" if color else text + + return colorize + + +def _run_serve(args: argparse.Namespace, *, out: IO[str]) -> int: + path = Path(args.file or args.db) + if not path.exists(): + out.write(f"logquill serve: no such file: {path}\n") + return 1 + + # Deliberately not a ternary (`x if c else y`) despite what a linter + # suggests: older mypy (the pre-commit hook pins 1.11.2) infers a + # conditional expression between these two unrelated concrete classes as + # `object` rather than their common `RecordSource` Protocol, even with an + # explicit annotation — this if/else form types cleanly on every mypy + # version instead. + source: RecordSource + if args.db: # noqa: SIM108 + source = SQLiteSource(path, table=args.table) + else: + source = JSONLSource(path) + server = TraceViewerServer(source, host=args.host, port=args.port) + out.write(f"logquill serve: {len(server.runs)} run(s) found in {path}\n") + out.write(f"logquill serve: listening on {server.url} — Ctrl+C to stop\n") + with contextlib.suppress(KeyboardInterrupt): + server.serve_forever() + return 0 + + +_CLEAR_SCREEN = "\x1b[2J\x1b[H" + + +def _latest_run_id(records: list[dict[str, Any]]) -> str | None: + for record in reversed(records): + run_id = (record.get("meta") or {}).get("run_id") + if isinstance(run_id, str) and run_id: + return run_id + return None + + +def _run_dev( + args: argparse.Namespace, + *, + out: IO[str], + warn_stream: IO[str], + poll_interval: float = 0.5, + max_iterations: int | None = None, +) -> int: + """Backs `logquill dev`: follows `args.file`, redrawing the tracked run's + span tree — cleared and redrawn from scratch each time, like `top` — + whenever a new record for it arrives. Tracks `args.run_id` if given, + otherwise whichever run's records have arrived most recently, so the + view follows an agent from one run to the next without restarting. + + Only the last `args.backlog` lines of whatever the file already holds + seed the initial view — this isn't the bounded-memory streaming `trace` + promises for a huge historical file, since a live dev session instead + keeps growing its own in-memory record list for as long as it runs; + `--backlog` just keeps a large *pre-existing* file from being read in + full before the first redraw. + """ + path = Path(args.file) + if not path.exists(): + warn_stream.write(f"logquill dev: no such file: {args.file}\n") + return 1 + + colorize = _tree_colorizer(out, args.no_color) + records: list[dict[str, Any]] = [] + + def redraw() -> None: + run_id = args.run_id or _latest_run_id(records) + if run_id is None: + return + builder = build_trace(records, run_id) + if colorize is not None: + out.write(_CLEAR_SCREEN) + out.write(f"logquill dev — run {run_id}\n\n") + out.write(render_tree(builder.roots, colorize=colorize) + "\n") + out.flush() + + with path.open("r", encoding="utf-8") as f: + for raw_line in f: + record = _parse_line(raw_line, warn_stream=warn_stream) + if record is not None: + records.append(record) + offset = f.tell() + records = records[-args.backlog :] + redraw() + + iterations = 0 + with contextlib.suppress(KeyboardInterrupt): + while max_iterations is None or iterations < max_iterations: + iterations += 1 + try: + with path.open("r", encoding="utf-8") as f: + f.seek(offset) + new_lines = f.readlines() + offset = f.tell() + except FileNotFoundError: + time.sleep(poll_interval) + continue + + new_records = [ + record + for record in ( + _parse_line(raw_line, warn_stream=warn_stream) for raw_line in new_lines + ) + if record is not None + ] + if new_records: + records.extend(new_records) + redraw() + time.sleep(poll_interval) + return 0 + + def main(argv: Sequence[str] | None = None) -> int: """The `logquill` console-script entry point: parses `argv` (defaulting to `sys.argv`) and dispatches to the matching subcommand, returning the @@ -222,6 +446,15 @@ def main(argv: Sequence[str] | None = None) -> int: parser = build_parser() args = parser.parse_args(argv) + if args.command == "trace": + return _run_trace(args, out=sys.stdout, warn_stream=sys.stderr) + + if args.command == "serve": + return _run_serve(args, out=sys.stdout) + + if args.command == "dev": + return _run_dev(args, out=sys.stdout, warn_stream=sys.stderr) + if args.command != "tail": parser.error(f"Unknown command: {args.command}") diff --git a/logquill/static/trace_viewer.html b/logquill/static/trace_viewer.html new file mode 100644 index 0000000..844f466 --- /dev/null +++ b/logquill/static/trace_viewer.html @@ -0,0 +1,396 @@ + + + + + +LogQuill trace viewer + + + +
+

LogQuill

+ local trace viewer — nothing leaves this machine +
+
+ +
+
Pick a run on the left, or type a search above to look across every run.
+ + +
+
+
+ + + diff --git a/logquill/trace_tree.py b/logquill/trace_tree.py new file mode 100644 index 0000000..2426cd4 --- /dev/null +++ b/logquill/trace_tree.py @@ -0,0 +1,224 @@ +"""Reconstructing one agent run's span tree from its own records, and +rendering it as an indented, annotated tree — the shared engine behind +`logquill trace` and `logquill serve`. + +Streams: `TraceBuilder.feed()` takes one record at a time and only holds +onto records that belong to the run being built, so tracing one run out of +a huge log file costs memory proportional to that run, not the file (see +`logquill/cli.py`'s `trace` command, which reads a file line by line). +""" + +from __future__ import annotations + +from collections.abc import Callable, Iterable, Mapping +from dataclasses import dataclass, field +from typing import Any + + +@dataclass +class Rollup: + """Token/cost totals — one node's own, or its subtree's.""" + + tokens_in: int = 0 + tokens_out: int = 0 + cost_usd: float = 0.0 + + def __add__(self, other: Rollup) -> Rollup: + return Rollup( + self.tokens_in + other.tokens_in, + self.tokens_out + other.tokens_out, + self.cost_usd + other.cost_usd, + ) + + +@dataclass +class TraceNode: + """One record in a run's trace, plus everything nested directly inside + it — a mix of further spans and plain events (actions, observations, + thoughts, decisions, or an ordinary log line), in the order they were + written. That mix, not a spans-only tree, is what a waterfall view needs: + the true nesting order of everything that happened inside a span. + """ + + record: Mapping[str, Any] + children: list[TraceNode] = field(default_factory=list) + + @property + def _meta(self) -> Mapping[str, Any]: + meta = self.record.get("meta") + return meta if isinstance(meta, dict) else {} + + @property + def span_id(self) -> str | None: + """This node's own span id — set only on a record `Logger.span()` + wrote on exit, `None` for a plain event.""" + span_id = self._meta.get("span_id") + return span_id if isinstance(span_id, str) and span_id else None + + @property + def parent_span_id(self) -> str | None: + parent = self._meta.get("parent_span_id") + return parent if isinstance(parent, str) and parent else None + + @property + def is_span(self) -> bool: + return self._meta.get("kind") == "span" + + @property + def duration_ms(self) -> float | None: + value = self._meta.get("duration_ms") + return ( + float(value) + if isinstance(value, (int, float)) and not isinstance(value, bool) + else None + ) + + @property + def own_tokens_cost(self) -> Rollup: + llm = self.record.get("llm") + llm = llm if isinstance(llm, dict) else {} + tokens_in = llm.get("tokens_in") + tokens_out = llm.get("tokens_out") + cost = llm.get("cost_usd") + return Rollup( + tokens_in if isinstance(tokens_in, int) and not isinstance(tokens_in, bool) else 0, + tokens_out if isinstance(tokens_out, int) and not isinstance(tokens_out, bool) else 0, + float(cost) if isinstance(cost, (int, float)) and not isinstance(cost, bool) else 0.0, + ) + + def rollup(self) -> Rollup: + """This node's own tokens/cost plus everything nested under it.""" + total = self.own_tokens_cost + for child in self.children: + total = total + child.rollup() + return total + + def to_dict(self) -> dict[str, Any]: + """A plain, JSON-safe nested structure — what `logquill trace + --json` prints and what `logquill serve`'s `/api/trace` returns.""" + rollup = self.rollup() + return { + "record": dict(self.record), + "rollup": { + "tokens_in": rollup.tokens_in, + "tokens_out": rollup.tokens_out, + "cost_usd": rollup.cost_usd, + }, + "children": [child.to_dict() for child in self.children], + } + + +class TraceBuilder: + """Single-pass, streaming reconstruction of one `run_id`'s span tree. + + Feed records in file order via `feed()`; call `finalize()` once, after + the last one, to get the roots. Relies on how `Logger.span()` writes + records: everything logged inside a span appears in the file *before* + that span's own closing record, and a nested span closes before its + parent (ordinary `with`-block exit order) — so by the time a span's own + record is seen, everything nested inside it has already been seen too, + and this never needs to look ahead or buffer more than one run's worth + of records. + """ + + def __init__(self, run_id: str) -> None: + self.run_id = run_id + self._nodes: dict[str, TraceNode] = {} + #: Records/spans waiting for a parent span that hasn't closed yet, + #: keyed by the parent's `span_id`. + self._pending: dict[str, list[TraceNode]] = {} + self.roots: list[TraceNode] = [] + self.record_count = 0 + + def feed(self, record: Mapping[str, Any]) -> bool: + """Considers one record. Returns whether it belongs to this run (and + was kept) — records for any other run cost nothing beyond this call.""" + meta = record.get("meta") + meta = meta if isinstance(meta, dict) else {} + if meta.get("run_id") != self.run_id: + return False + + self.record_count += 1 + node = TraceNode(record) + if node.span_id is not None: + node.children = self._pending.pop(node.span_id, []) + self._nodes[node.span_id] = node + + parent_span_id = node.parent_span_id + if parent_span_id is None: + self.roots.append(node) + else: + parent = self._nodes.get(parent_span_id) + if parent is not None: + parent.children.append(node) + else: + self._pending.setdefault(parent_span_id, []).append(node) + return True + + def finalize(self) -> list[TraceNode]: + """Call once every record has been fed. Any node still waiting for a + parent that never closed (a still-open span in a crashed or + truncated run) is appended to the roots instead of being silently + dropped. Returns `self.roots`.""" + for orphans in self._pending.values(): + self.roots.extend(orphans) + self._pending.clear() + return self.roots + + +def build_trace(records: Iterable[Mapping[str, Any]], run_id: str) -> TraceBuilder: + """Convenience wrapper: feed every record in `records` and finalize.""" + builder = TraceBuilder(run_id) + for record in records: + builder.feed(record) + builder.finalize() + return builder + + +def _annotations(node: TraceNode) -> str: + parts = [] + if node.duration_ms is not None: + parts.append(f"{node.duration_ms:,.1f}ms") + rollup = node.rollup() + if rollup.tokens_in or rollup.tokens_out: + parts.append(f"{rollup.tokens_in}→{rollup.tokens_out} tok") + if rollup.cost_usd: + parts.append(f"${rollup.cost_usd:,.4f}") + return f" ({', '.join(parts)})" if parts else "" + + +def _label(node: TraceNode) -> str: + level = str(node.record.get("level", "?")) + message = str(node.record.get("message", "")) + return f"[{level}] {message}{_annotations(node)}" + + +def render_tree( + roots: list[TraceNode], + *, + colorize: Callable[[str, str], str] | None = None, +) -> str: + """Renders `roots` (from `TraceBuilder.finalize()`/`build_trace()`) as an + indented tree, one line per record, each annotated with its own + `duration_ms` and the token/cost totals of everything nested under it + (so an outer span shows the sum its children contributed). + + `colorize(text, level) -> text`, if given, wraps each line's label — + `logquill trace` passes the same level-to-color mapping `ConsoleTransport` + uses; left out (the default) for output going to a file or a non-TTY. + """ + lines: list[str] = [] + + def walk(node: TraceNode, prefix: str, is_last: bool) -> None: + connector = "└─ " if is_last else "├─ " + label = _label(node) + if colorize is not None: + label = colorize(label, str(node.record.get("level", ""))) + lines.append(f"{prefix}{connector}{label}") + child_prefix = prefix + (" " if is_last else "│ ") + for index, child in enumerate(node.children): + walk(child, child_prefix, index == len(node.children) - 1) + + for index, root in enumerate(roots): + walk(root, "", index == len(roots) - 1) + return "\n".join(lines) diff --git a/logquill/webui.py b/logquill/webui.py new file mode 100644 index 0000000..9633a6e --- /dev/null +++ b/logquill/webui.py @@ -0,0 +1,331 @@ +"""`logquill serve` — a local, offline web UI for browsing traced runs. + +Built entirely on the stdlib (`http.server`, `sqlite3`): no new dependency, +and nothing leaves the machine — the whole point of "local-first". Reads +records from a JSONL file, or (`db=True`) the `logs` table a `SQLiteTransport` +wrote. Run summaries are computed once, streaming through the source at +startup, and cached; `/api/trace` and `/api/search` stream through the +source again per request rather than holding it in memory, so a large file +costs request latency, not memory (see `logquill/trace_tree.py`). + +One documented gap: a source that's still being appended to (a live capture) +isn't picked up after startup — this is a snapshot of the file/database as +it was when the server started, not a live tail. Restart to pick up new runs. +""" + +from __future__ import annotations + +import json +import logging +import sqlite3 +from collections.abc import Iterable, Iterator, Mapping +from dataclasses import dataclass, field +from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer +from pathlib import Path +from typing import Any, Protocol +from urllib.parse import parse_qs, urlsplit + +from logquill.trace_tree import build_trace + +_logger = logging.getLogger("logquill") + +STATIC_DIR = Path(__file__).parent / "static" + +#: `/api/search` never returns more than this many records in one response — +#: a search box is for finding a handful of matches, not paginating a +#: multi-gigabyte file through the browser. +MAX_SEARCH_RESULTS = 500 + + +class RecordSource(Protocol): + """Where `logquill serve` reads records from. `read_all()` is called + once per request that needs them — implementations should stream, not + build a list, so a request costs memory proportional to what it keeps + (a run, a page of search results), not the whole source.""" + + def read_all(self) -> Iterator[dict[str, Any]]: + """Yields every record, in the order it was written.""" + ... + + +class JSONLSource: + """Reads records from a LogQuill JSONL file, one line at a time. A line + that isn't valid JSON, or isn't a JSON object, is skipped — the same + tolerance `logquill tail`/`trace` give a log file that isn't perfectly + clean.""" + + def __init__(self, path: str | Path) -> None: + self.path = Path(path) + + def read_all(self) -> Iterator[dict[str, Any]]: + with self.path.open("r", encoding="utf-8") as f: + for line in f: + line = line.strip() + if not line: + continue + try: + record = json.loads(line) + except json.JSONDecodeError: + continue + if isinstance(record, dict): + yield record + + +class SQLiteSource: + """Reads records from the fixed `logs` table schema every SQL transport + writes (see `BaseSQLTransport`) via a `SQLiteTransport`-created database. + + **Limitation, inherent to that schema**: only `meta`'s `run_id`/`span_id`/ + `parent_span_id`/`trace_id` are stored as queryable columns — an LLM + call's `llm` block (tokens, cost) is not a column `BaseSQLTransport` + writes, so runs traced from a SQLite source never show token/cost + annotations, regardless of what the original records carried. Trace from + the original JSONL file (or a transport that does keep `llm`) to see them. + """ + + def __init__(self, path: str | Path, *, table: str = "logs") -> None: + self.path = Path(path) + self.table = table + + def read_all(self) -> Iterator[dict[str, Any]]: + connection = sqlite3.connect(str(self.path)) + try: + cursor = connection.execute( + f"SELECT timestamp, level, logger, message, meta, run_id, span_id, " # noqa: S608 + f"parent_span_id, trace_id FROM {self.table} ORDER BY id" + ) + for row in cursor: + yield _row_to_record(row) + finally: + connection.close() + + +def _row_to_record(row: tuple[Any, ...]) -> dict[str, Any]: + timestamp, level, logger_name, message, meta_json, run_id, span_id, parent_span_id, trace_id = ( + row + ) + try: + meta = json.loads(meta_json) if meta_json else {} + except (json.JSONDecodeError, TypeError): + meta = {} + if not isinstance(meta, dict): + meta = {} + for key, value in ( + ("run_id", run_id), + ("span_id", span_id), + ("parent_span_id", parent_span_id), + ("trace_id", trace_id), + ): + if value is not None: + meta.setdefault(key, value) + return { + "schema_version": "1.0", + "timestamp": timestamp, + "level": level, + "logger": logger_name, + "message": message, + "meta": meta, + } + + +@dataclass +class RunSummary: + """One row of the run list.""" + + run_id: str + record_count: int = 0 + start: str | None = None + end: str | None = None + level_counts: dict[str, int] = field(default_factory=dict) + tokens_in: int = 0 + tokens_out: int = 0 + cost_usd: float = 0.0 + agent_name: str | None = None + + def to_dict(self) -> dict[str, Any]: + return { + "run_id": self.run_id, + "record_count": self.record_count, + "start": self.start, + "end": self.end, + "level_counts": self.level_counts, + "tokens_in": self.tokens_in, + "tokens_out": self.tokens_out, + "cost_usd": self.cost_usd, + "agent_name": self.agent_name, + "error_count": self.level_counts.get("ERROR", 0) + self.level_counts.get("FATAL", 0), + } + + +def summarize_runs(records: Iterable[Mapping[str, Any]]) -> list[RunSummary]: + """One pass over `records`, building one `RunSummary` per distinct + `meta.run_id` seen — a record with no `run_id` isn't part of any run and + is left out of the list (it still shows up inside `/api/search`). + Memory here is proportional to the number of distinct runs, not the + number of records. + """ + summaries: dict[str, RunSummary] = {} + for record in records: + meta = record.get("meta") + meta = meta if isinstance(meta, dict) else {} + run_id = meta.get("run_id") + if not isinstance(run_id, str) or not run_id: + continue + + summary = summaries.setdefault(run_id, RunSummary(run_id)) + summary.record_count += 1 + timestamp = record.get("timestamp") + if isinstance(timestamp, str): + if summary.start is None or timestamp < summary.start: + summary.start = timestamp + if summary.end is None or timestamp > summary.end: + summary.end = timestamp + level = record.get("level") + if isinstance(level, str): + summary.level_counts[level] = summary.level_counts.get(level, 0) + 1 + llm = record.get("llm") + if isinstance(llm, dict): + if isinstance(llm.get("tokens_in"), int): + summary.tokens_in += llm["tokens_in"] + if isinstance(llm.get("tokens_out"), int): + summary.tokens_out += llm["tokens_out"] + if isinstance(llm.get("cost_usd"), (int, float)): + summary.cost_usd += llm["cost_usd"] + if meta.get("operation") == "invoke_agent" and isinstance(meta.get("agent_name"), str): + summary.agent_name = meta["agent_name"] + + return sorted(summaries.values(), key=lambda s: s.start or "", reverse=True) + + +def _matches_search(record: Mapping[str, Any], query: str) -> bool: + if not query: + return True + haystack = f"{record.get('message', '')} {json.dumps(record.get('meta') or {}, default=str)}" + return query in haystack.lower() + + +def search_records( + records: Iterable[Mapping[str, Any]], + *, + query: str = "", + level: str | None = None, + run_id: str | None = None, + limit: int = MAX_SEARCH_RESULTS, +) -> list[dict[str, Any]]: + """Streams through `records`, keeping at most `limit` that match every + given filter — a substring of `query` in the message or `meta` + (case-insensitive), an exact `level`, and/or an exact `meta.run_id`.""" + query = query.lower() + results: list[dict[str, Any]] = [] + for record in records: + if level is not None and record.get("level") != level: + continue + if run_id is not None and (record.get("meta") or {}).get("run_id") != run_id: + continue + if not _matches_search(record, query): + continue + results.append(dict(record)) + if len(results) >= limit: + break + return results + + +class TraceViewerServer: + """The HTTP server behind `logquill serve`: one static page plus a small + JSON API, all reading from one `RecordSource`. + + Run summaries are computed once, at construction, by streaming through + the source — see the module docstring for what that means for a source + that's still growing. + """ + + def __init__(self, source: RecordSource, *, host: str = "127.0.0.1", port: int = 0) -> None: + self.source = source + self.runs = summarize_runs(source.read_all()) + self._host = host + handler = _make_handler(self) + self._httpd = ThreadingHTTPServer((host, port), handler) + + @property + def port(self) -> int: + """The port actually bound — resolved even when constructed with + `port=0` (let the OS pick one).""" + return int(self._httpd.server_address[1]) + + @property + def url(self) -> str: + return f"http://{self._host}:{self.port}/" + + def serve_forever(self) -> None: + """Blocks, serving requests until `shutdown()` is called (typically + from a signal handler or another thread).""" + self._httpd.serve_forever() + + def shutdown(self) -> None: + """Stops `serve_forever()` and releases the listening socket. Safe + to call from a different thread than the one running `serve_forever()`.""" + self._httpd.shutdown() + self._httpd.server_close() + + +def _make_handler(server: TraceViewerServer) -> type[BaseHTTPRequestHandler]: + class Handler(BaseHTTPRequestHandler): + def log_message(self, format: str, *args: Any) -> None: # noqa: A002 + _logger.debug("logquill serve: " + format, *args) + + def do_GET(self) -> None: # noqa: N802 + _route(self, server) + + return Handler + + +def _send_json(handler: BaseHTTPRequestHandler, payload: Any, *, status: int = 200) -> None: + body = json.dumps(payload, default=str).encode("utf-8") + handler.send_response(status) + handler.send_header("Content-Type", "application/json") + handler.send_header("Content-Length", str(len(body))) + handler.end_headers() + handler.wfile.write(body) + + +def _send_file(handler: BaseHTTPRequestHandler, path: Path, content_type: str) -> None: + try: + body = path.read_bytes() + except OSError: + handler.send_response(404) + handler.end_headers() + return + handler.send_response(200) + handler.send_header("Content-Type", content_type) + handler.send_header("Content-Length", str(len(body))) + handler.end_headers() + handler.wfile.write(body) + + +def _route(handler: BaseHTTPRequestHandler, server: TraceViewerServer) -> None: + parsed = urlsplit(handler.path) + query = {key: values[0] for key, values in parse_qs(parsed.query).items()} + + if parsed.path == "/": + _send_file(handler, STATIC_DIR / "trace_viewer.html", "text/html; charset=utf-8") + elif parsed.path == "/api/runs": + _send_json(handler, [run.to_dict() for run in server.runs]) + elif parsed.path == "/api/trace": + run_id = query.get("run_id") + if not run_id: + _send_json(handler, {"error": "run_id is required"}, status=400) + return + builder = build_trace(server.source.read_all(), run_id) + _send_json(handler, [node.to_dict() for node in builder.roots]) + elif parsed.path == "/api/search": + results = search_records( + server.source.read_all(), + query=query.get("q", ""), + level=query.get("level") or None, + run_id=query.get("run_id") or None, + limit=int(query.get("limit", MAX_SEARCH_RESULTS)), + ) + _send_json(handler, results) + else: + handler.send_response(404) + handler.end_headers() diff --git a/tests/test_cli_dev.py b/tests/test_cli_dev.py new file mode 100644 index 0000000..3f2f81b --- /dev/null +++ b/tests/test_cli_dev.py @@ -0,0 +1,185 @@ +from __future__ import annotations + +import io +from pathlib import Path + +from logquill import Logger, RunPlugin +from logquill.cli import _CLEAR_SCREEN, _run_dev, build_parser +from logquill.transports.file_transport import FileTransport + + +def _logger_for(path: Path, run_id: str) -> Logger: + return Logger("app.agent", transports=[FileTransport(path)], plugins=[RunPlugin(run_id=run_id)]) + + +class _FakeTTY(io.StringIO): + def isatty(self) -> bool: + return True + + +def _dev(file: str, *extra_args: str, **kwargs: object) -> tuple[str, str]: + args = build_parser().parse_args(["dev", file, *extra_args]) + out, warn = io.StringIO(), io.StringIO() + _run_dev(args, out=out, warn_stream=warn, poll_interval=0, **kwargs) # type: ignore[arg-type] + return out.getvalue(), warn.getvalue() + + +def test_dev_renders_the_most_recently_active_run_on_startup(tmp_path: Path) -> None: + path = tmp_path / "app.log" + logger = _logger_for(path, "run-1") + with logger.span("run"): + logger.thought("plan") + logger.close() + + output, warn = _dev(str(path), max_iterations=0) + + assert warn == "" + assert "run run-1" in output + assert "plan" in output + + +def test_a_missing_file_reports_a_helpful_error(tmp_path: Path) -> None: + args = build_parser().parse_args(["dev", str(tmp_path / "missing.log")]) + out, warn = io.StringIO(), io.StringIO() + + exit_code = _run_dev(args, out=out, warn_stream=warn) + + assert exit_code == 1 + assert "no such file" in warn.getvalue() + + +def test_dev_redraws_when_new_records_for_the_tracked_run_arrive(tmp_path: Path) -> None: + path = tmp_path / "app.log" + logger = _logger_for(path, "run-1") + logger.thought("first") + + args = build_parser().parse_args(["dev", str(path)]) + out = io.StringIO() + + def append_more(_: object = None) -> None: + logger.action("second", tool="search") + + # patch time.sleep so the poll loop's single iteration does real work + import logquill.cli as cli_module + + original_sleep = cli_module.time.sleep + cli_module.time.sleep = lambda _s: append_more() # type: ignore[assignment] + try: + _run_dev(args, out=out, warn_stream=io.StringIO(), poll_interval=0, max_iterations=2) + finally: + cli_module.time.sleep = original_sleep + + renders = out.getvalue().split("logquill dev") + assert len(renders) >= 2 # redrawn at least once beyond the initial render + assert "second" in out.getvalue() + + +def test_dev_follows_a_new_run_once_the_old_one_finishes(tmp_path: Path) -> None: + path = tmp_path / "app.log" + first = _logger_for(path, "run-1") + with first.span("run"): + pass + + args = build_parser().parse_args(["dev", str(path)]) + out = io.StringIO() + + import logquill.cli as cli_module + + second = _logger_for(path, "run-2") + + def switch(_: object = None) -> None: + second.thought("new run started") + + original_sleep = cli_module.time.sleep + cli_module.time.sleep = lambda _s: switch() # type: ignore[assignment] + try: + _run_dev(args, out=out, warn_stream=io.StringIO(), poll_interval=0, max_iterations=2) + finally: + cli_module.time.sleep = original_sleep + + assert "run run-2" in out.getvalue() + + +def test_an_explicit_run_id_is_never_overridden_by_a_newer_run(tmp_path: Path) -> None: + path = tmp_path / "app.log" + first = _logger_for(path, "run-1") + with first.span("run"): + pass + + args = build_parser().parse_args(["dev", str(path), "--run-id", "run-1"]) + out = io.StringIO() + + import logquill.cli as cli_module + + second = _logger_for(path, "run-2") + + def switch(_: object = None) -> None: + second.thought("a different run") + + original_sleep = cli_module.time.sleep + cli_module.time.sleep = lambda _s: switch() # type: ignore[assignment] + try: + _run_dev(args, out=out, warn_stream=io.StringIO(), poll_interval=0, max_iterations=2) + finally: + cli_module.time.sleep = original_sleep + + assert "run run-1" in out.getvalue() + assert "run-2" not in out.getvalue() + + +def test_backlog_caps_how_much_pre_existing_history_seeds_the_view(tmp_path: Path) -> None: + path = tmp_path / "app.log" + logger = _logger_for(path, "run-1") + for i in range(10): + logger.thought(f"step {i}") + + output, _warn = _dev(str(path), "--backlog", "3", max_iterations=0) + + assert "step 9" in output + assert "step 6" not in output # only the last 3 lines were kept + + +def test_no_run_at_all_yet_renders_nothing_but_does_not_crash(tmp_path: Path) -> None: + path = tmp_path / "app.log" + path.write_text( + '{"message": "no run_id here", "level": "INFO", "meta": {}}\n', encoding="utf-8" + ) + + output, warn = _dev(str(path), max_iterations=0) + + assert output == "" + assert warn == "" + + +def test_color_and_clear_screen_are_used_on_a_tty_and_suppressed_otherwise( + tmp_path: Path, +) -> None: + path = tmp_path / "app.log" + logger = _logger_for(path, "run-1") + with logger.span("run"): + pass + + args = build_parser().parse_args(["dev", str(path)]) + plain_out = io.StringIO() + _run_dev(args, out=plain_out, warn_stream=io.StringIO(), poll_interval=0, max_iterations=0) + assert _CLEAR_SCREEN not in plain_out.getvalue() + assert "\x1b[" not in plain_out.getvalue() + + tty_out = _FakeTTY() + _run_dev(args, out=tty_out, warn_stream=io.StringIO(), poll_interval=0, max_iterations=0) + assert _CLEAR_SCREEN in tty_out.getvalue() + assert "\x1b[32m" in tty_out.getvalue() # INFO green + + +def test_no_color_suppresses_color_and_clearing_even_on_a_tty(tmp_path: Path) -> None: + path = tmp_path / "app.log" + logger = _logger_for(path, "run-1") + with logger.span("run"): + pass + + args = build_parser().parse_args(["dev", str(path), "--no-color"]) + tty_out = _FakeTTY() + + _run_dev(args, out=tty_out, warn_stream=io.StringIO(), poll_interval=0, max_iterations=0) + + assert _CLEAR_SCREEN not in tty_out.getvalue() diff --git a/tests/test_cli_serve.py b/tests/test_cli_serve.py new file mode 100644 index 0000000..ac1abf6 --- /dev/null +++ b/tests/test_cli_serve.py @@ -0,0 +1,72 @@ +from __future__ import annotations + +import io +import threading +import time +import urllib.error +import urllib.request +from pathlib import Path + +import pytest + +from logquill import Logger +from logquill.cli import _run_serve, build_parser +from logquill.transports.file_transport import FileTransport + + +def _write_a_log(path: Path) -> None: + Logger("app", transports=[FileTransport(path)]).info("hello", run_id="run-1") + + +def test_file_and_db_are_mutually_exclusive_and_one_is_required() -> None: + with pytest.raises(SystemExit): + build_parser().parse_args(["serve"]) + with pytest.raises(SystemExit): + build_parser().parse_args(["serve", "--file", "a.log", "--db", "a.sqlite"]) + + +def test_serve_reports_a_helpful_error_for_a_missing_file(tmp_path: Path) -> None: + args = build_parser().parse_args(["serve", "--file", str(tmp_path / "missing.log")]) + out = io.StringIO() + + exit_code = _run_serve(args, out=out) + + assert exit_code == 1 + assert "no such file" in out.getvalue() + + +def test_serve_starts_a_real_server_and_prints_its_url(tmp_path: Path) -> None: + path = tmp_path / "app.log" + _write_a_log(path) + args = build_parser().parse_args(["serve", "--file", str(path), "--port", "0"]) + out = io.StringIO() + + # `_run_serve` blocks in `serve_forever()` until Ctrl+C; run it on a + # daemon thread so the test (and the process, if this thread never gets + # to stop cleanly) doesn't hang on it. + thread = threading.Thread(target=_run_serve, args=(args,), kwargs={"out": out}, daemon=True) + thread.start() + + for _ in range(100): + if "listening on" in out.getvalue(): + break + time.sleep(0.02) + assert "listening on" in out.getvalue() + assert "1 run(s) found" in out.getvalue() + + url = out.getvalue().split("listening on ")[1].split(" ")[0] + with urllib.request.urlopen(url, timeout=5) as response: + assert response.status == 200 + with urllib.request.urlopen(url + "api/runs", timeout=5) as response: + assert b"run-1" in response.read() + + +def test_serve_via_main_reports_the_error_for_a_bad_source( + tmp_path: Path, capsys: pytest.CaptureFixture[str] +) -> None: + from logquill.cli import main + + exit_code = main(["serve", "--db", str(tmp_path / "missing.sqlite")]) + + assert exit_code == 1 + assert "no such file" in capsys.readouterr().out diff --git a/tests/test_cli_trace.py b/tests/test_cli_trace.py new file mode 100644 index 0000000..557ceea --- /dev/null +++ b/tests/test_cli_trace.py @@ -0,0 +1,120 @@ +from __future__ import annotations + +import io +import json +from pathlib import Path + +import pytest + +from logquill import Logger, RunPlugin +from logquill.cli import _run_trace, build_parser, main +from logquill.transports.file_transport import FileTransport + + +def _write_a_run(path: Path, run_id: str = "run-1") -> None: + logger = Logger( + "app.agent", transports=[FileTransport(path)], plugins=[RunPlugin(run_id=run_id)] + ) + with logger.span("run", operation="invoke_agent", agent_name="planner"): + logger.thought("plan") + with logger.span("step"): + logger.llm_call("chat", model="m", tokens_in=10, tokens_out=5, cost_usd=0.02) + logger.action("call", tool="search") + logger.close() + + +def _trace(file: str, run_id: str, *extra_args: str) -> tuple[str, str, int]: + args = build_parser().parse_args(["trace", run_id, "--file", file, *extra_args]) + out, warn = io.StringIO(), io.StringIO() + exit_code = _run_trace(args, out=out, warn_stream=warn) + return out.getvalue(), warn.getvalue(), exit_code + + +def test_trace_prints_an_indented_annotated_tree(tmp_path: Path) -> None: + path = tmp_path / "app.log" + _write_a_run(path) + + output, warn, exit_code = _trace(str(path), "run-1") + + assert exit_code == 0 + assert warn == "" + assert "run (" in output + assert "10→5 tok" in output + assert "$0.0200" in output + assert "plan" in output and "call" in output + + +def test_trace_json_prints_a_nested_structure(tmp_path: Path) -> None: + path = tmp_path / "app.log" + _write_a_run(path) + + output, _warn, exit_code = _trace(str(path), "run-1", "--json") + + assert exit_code == 0 + (root,) = json.loads(output) + assert root["record"]["message"] == "run" + assert root["rollup"] == {"tokens_in": 10, "tokens_out": 5, "cost_usd": 0.02} + messages = [child["record"]["message"] for child in root["children"]] + assert messages == ["plan", "step", "call"] + + +def test_trace_ignores_records_from_other_runs_in_the_same_file(tmp_path: Path) -> None: + path = tmp_path / "app.log" + _write_a_run(path, run_id="run-1") + _write_a_run(path, run_id="run-2") + + output, _warn, exit_code = _trace(str(path), "run-2") + + assert exit_code == 0 + assert output.count("run (") == 1 + + +def test_trace_reports_a_helpful_error_for_an_unknown_run_id(tmp_path: Path) -> None: + path = tmp_path / "app.log" + _write_a_run(path) + + output, warn, exit_code = _trace(str(path), "no-such-run") + + assert exit_code == 1 + assert output == "" + assert "no-such-run" in warn + assert str(path) in warn + + +def test_trace_reports_a_helpful_error_for_a_missing_file(tmp_path: Path) -> None: + output, warn, exit_code = _trace(str(tmp_path / "missing.log"), "run-1") + + assert exit_code == 1 + assert "no such file" in warn + + +def test_trace_skips_malformed_lines_with_a_warning_but_still_traces(tmp_path: Path) -> None: + path = tmp_path / "app.log" + _write_a_run(path) + with path.open("a") as f: + f.write("not json at all\n") + f.write('["a", "json", "array", "not", "an", "object"]\n') + + output, warn, exit_code = _trace(str(path), "run-1") + + assert exit_code == 0 + assert "run (" in output + assert "malformed JSON" in warn + assert "non-object" in warn + + +def test_trace_via_main_entry_point(tmp_path: Path, capsys: pytest.CaptureFixture[str]) -> None: + path = tmp_path / "app.log" + _write_a_run(path) + + exit_code = main(["trace", "run-1", "--file", str(path)]) + + assert exit_code == 0 + assert "run (" in capsys.readouterr().out + + +def test_the_run_id_is_positional_and_file_is_required(tmp_path: Path) -> None: + with pytest.raises(SystemExit): + build_parser().parse_args(["trace"]) + with pytest.raises(SystemExit): + build_parser().parse_args(["trace", "run-1"]) diff --git a/tests/test_trace_tree.py b/tests/test_trace_tree.py new file mode 100644 index 0000000..b9220f4 --- /dev/null +++ b/tests/test_trace_tree.py @@ -0,0 +1,159 @@ +from __future__ import annotations + +from typing import Any + +from logquill import Logger, RunPlugin +from logquill.trace_tree import Rollup, TraceBuilder, build_trace, render_tree +from logquill.transports.transport import CollectingTransport + + +def _run(run_id: str = "run-1") -> list[dict[str, Any]]: + sink = CollectingTransport() + logger = Logger("app.agent", transports=[sink], plugins=[RunPlugin(run_id=run_id)]) + + with logger.span("run", operation="invoke_agent", agent_name="planner"): + logger.thought("plan the work") + with logger.span("plan_step"): + logger.llm_call( + "chat", + model="m1", + tokens_in=100, + tokens_out=20, + cost_usd=0.01, + finish_reason="stop", + ) + logger.action("look it up", tool="search") + logger.observation("found it", tool="search") + return sink.records + + +def test_a_full_run_reconstructs_into_one_root_matching_the_span_nesting() -> None: + records = _run() + + builder = build_trace(records, "run-1") + + assert builder.record_count == len(records) + (root,) = builder.roots + assert root.record["message"] == "run" + assert root.is_span + thought, step, action, observation = root.children + assert thought.record["message"] == "plan the work" + assert step.record["message"] == "plan_step" + assert action.record["message"] == "look it up" + assert observation.record["message"] == "found it" + (llm_call,) = step.children + assert llm_call.record["llm"]["model"] == "m1" + + +def test_records_from_other_runs_are_skipped_and_not_kept() -> None: + records = _run("run-1") + _run("run-2") + + builder = build_trace(records, "run-1") + + assert builder.record_count == len(_run("run-1")) + assert all(r.record["meta"]["run_id"] == "run-1" for r in _iter(builder.roots)) + + +def _iter(nodes): # noqa: ANN001, ANN202 + for node in nodes: + yield node + yield from _iter(node.children) + + +def test_feed_reports_whether_a_record_belonged_to_the_run() -> None: + records = _run("run-1") + _run("run-2") + builder = TraceBuilder("run-1") + + kept = [builder.feed(r) for r in records] + + assert kept.count(True) == len(_run("run-1")) + assert kept.count(False) == len(_run("run-2")) + + +def test_rollup_sums_tokens_and_cost_from_every_descendant() -> None: + records = _run() + builder = build_trace(records, "run-1") + + rollup = builder.roots[0].rollup() + + assert rollup == Rollup(tokens_in=100, tokens_out=20, cost_usd=0.01) + + +def test_a_span_with_no_llm_descendants_has_a_zero_rollup() -> None: + sink = CollectingTransport() + logger = Logger("app", transports=[sink], plugins=[RunPlugin(run_id="run-x")]) + with logger.span("work"): + logger.info("plain log line") + + builder = build_trace(sink.records, "run-x") + + assert builder.roots[0].rollup() == Rollup() + + +def test_records_with_no_span_at_all_are_still_roots_in_order() -> None: + sink = CollectingTransport() + logger = Logger("app", transports=[sink], plugins=[RunPlugin(run_id="run-x")]) + logger.thought("first") + logger.decision("second") + + builder = build_trace(sink.records, "run-x") + + assert [r.record["message"] for r in builder.roots] == ["first", "second"] + + +def test_an_orphaned_span_whose_parent_never_closed_still_surfaces_as_a_root() -> None: + # simulates a crashed/truncated run: the parent's own closing record is + # missing, but a record it logged is still on disk + sink = CollectingTransport() + logger = Logger("app", transports=[sink], plugins=[RunPlugin(run_id="run-x")]) + with logger.span("outer"): + logger.info("inside") + records = [r for r in sink.records if r["message"] != "outer"] # drop the parent's own record + + builder = build_trace(records, "run-x") + + assert [r.record["message"] for r in builder.roots] == ["inside"] + + +def test_two_records_sharing_a_span_id_do_not_crash_the_builder() -> None: + # adversarial/malformed input: two closing records claim the same span_id + sink = CollectingTransport() + logger = Logger("app", transports=[sink], plugins=[RunPlugin(run_id="run-x")]) + with logger.span("a", span_id="dupe" * 4): + pass + with logger.span("b", span_id="dupe" * 4): + pass + + builder = build_trace(sink.records, "run-x") # must not raise + + assert len(builder.roots) == 2 + + +def test_render_tree_shows_duration_tokens_and_cost() -> None: + builder = build_trace(_run(), "run-1") + + text = render_tree(builder.roots) + + assert "run (" in text + assert "100→20 tok" in text + assert "$0.0100" in text + assert "├─ " in text and "└─ " in text and "│ " in text + + +def test_render_tree_applies_the_given_colorize_function() -> None: + builder = build_trace(_run(), "run-1") + + text = render_tree(builder.roots, colorize=lambda label, level: f"<{level}>{label}") + + assert "[INFO] run" in text + + +def test_render_tree_of_an_empty_run_is_an_empty_string() -> None: + assert render_tree([]) == "" + + +def test_records_missing_meta_entirely_do_not_crash_feed() -> None: + builder = TraceBuilder("run-1") + + assert builder.feed({"message": "no meta key at all"}) is False + assert builder.feed({"meta": None}) is False diff --git a/tests/test_webui.py b/tests/test_webui.py new file mode 100644 index 0000000..3ad3d84 --- /dev/null +++ b/tests/test_webui.py @@ -0,0 +1,280 @@ +from __future__ import annotations + +import json +import sqlite3 +import threading +import time +import urllib.error +import urllib.request +from pathlib import Path +from typing import Any, Iterator + +import pytest + +from logquill import Logger, RunPlugin +from logquill.transports.file_transport import FileTransport +from logquill.transports.sql.sqlite_transport import SQLiteTransport +from logquill.webui import ( + STATIC_DIR, + JSONLSource, + SQLiteSource, + TraceViewerServer, + search_records, + summarize_runs, +) + + +def _write_a_run(path: Path, run_id: str = "run-1") -> None: + logger = Logger( + "app.agent", transports=[FileTransport(path)], plugins=[RunPlugin(run_id=run_id)] + ) + with logger.span("run", operation="invoke_agent", agent_name="planner"): + logger.thought("plan") + with logger.span("step"): + logger.llm_call("chat", model="m", tokens_in=10, tokens_out=5, cost_usd=0.02) + logger.action("call", tool="search") + logger.error("something broke") + logger.close() + + +# --- JSONLSource ----------------------------------------------------------- + + +def test_jsonl_source_reads_every_record_in_order(tmp_path: Path) -> None: + path = tmp_path / "app.log" + _write_a_run(path) + + records = list(JSONLSource(path).read_all()) + + # file order = write order: a span closes *after* what's nested inside it + assert [r["message"] for r in records] == [ + "plan", + "chat", + "step", + "call", + "run", + "something broke", + ] + + +def test_jsonl_source_skips_malformed_lines(tmp_path: Path) -> None: + path = tmp_path / "app.log" + path.write_text( + '{"message": "a", "meta": {}}\nnot json\n["not", "an", "object"]\n', encoding="utf-8" + ) + + records = list(JSONLSource(path).read_all()) + + assert [r["message"] for r in records] == ["a"] + + +# --- SQLiteSource ------------------------------------------------------------ + + +def test_sqlite_source_reconstructs_records_with_ids_from_columns(tmp_path: Path) -> None: + db_path = tmp_path / "logs.sqlite" + transport = SQLiteTransport(filename=str(db_path), ensure_schema=True, max_records=1) + logger = Logger("app", transports=[transport], plugins=[RunPlugin(run_id="run-1")]) + logger.info("hello", user_id=42) + logger.close() + + records = list(SQLiteSource(db_path).read_all()) + + assert len(records) == 1 + assert records[0]["message"] == "hello" + assert records[0]["meta"]["user_id"] == 42 + assert records[0]["meta"]["run_id"] == "run-1" + + +def test_sqlite_source_has_no_llm_block_even_if_meta_has_llm_looking_keys(tmp_path: Path) -> None: + db_path = tmp_path / "logs.sqlite" + transport = SQLiteTransport(filename=str(db_path), ensure_schema=True, max_records=1) + logger = Logger("app", transports=[transport], plugins=[RunPlugin(run_id="run-1")]) + logger.llm_call("chat", model="m", tokens_in=5) + logger.close() + + (record,) = list(SQLiteSource(db_path).read_all()) + + assert "llm" not in record # BaseSQLTransport's schema doesn't carry it + + +def test_sqlite_source_handles_a_null_meta_column(tmp_path: Path) -> None: + db_path = tmp_path / "logs.sqlite" + connection = sqlite3.connect(str(db_path)) + connection.execute( + "CREATE TABLE logs (id INTEGER PRIMARY KEY, timestamp TEXT, level TEXT, logger TEXT, " + "message TEXT, meta TEXT, run_id TEXT, span_id TEXT, parent_span_id TEXT, trace_id TEXT)" + ) + connection.execute( + "INSERT INTO logs (timestamp, level, logger, message, meta, run_id) VALUES (?,?,?,?,?,?)", + ("t", "INFO", "app", "hi", None, "run-1"), + ) + connection.commit() + connection.close() + + (record,) = list(SQLiteSource(db_path).read_all()) + + assert record["meta"] == {"run_id": "run-1"} + + +# --- summarize_runs / search_records ----------------------------------------- + + +def test_summarize_runs_aggregates_tokens_cost_levels_and_agent_name(tmp_path: Path) -> None: + path = tmp_path / "app.log" + _write_a_run(path) + + (summary,) = summarize_runs(JSONLSource(path).read_all()) + + assert summary.run_id == "run-1" + assert summary.record_count == 6 + assert summary.tokens_in == 10 + assert summary.tokens_out == 5 + assert summary.cost_usd == 0.02 + assert summary.agent_name == "planner" + assert summary.level_counts == {"INFO": 5, "ERROR": 1} + assert summary.to_dict()["error_count"] == 1 + + +def test_summarize_runs_skips_records_with_no_run_id() -> None: + records: list[dict[str, Any]] = [ + {"message": "loose", "meta": {}, "level": "INFO", "timestamp": "t"} + ] + + assert summarize_runs(records) == [] + + +def test_summarize_runs_orders_by_start_time_newest_first() -> None: + records = [ + { + "message": "a", + "level": "INFO", + "timestamp": "2026-01-01T00:00:00.000Z", + "meta": {"run_id": "old"}, + }, + { + "message": "b", + "level": "INFO", + "timestamp": "2026-01-02T00:00:00.000Z", + "meta": {"run_id": "new"}, + }, + ] + + runs = summarize_runs(records) + + assert [r.run_id for r in runs] == ["new", "old"] + + +def test_search_records_filters_by_query_level_and_run_id() -> None: + records = [ + {"message": "search failed", "level": "ERROR", "meta": {"run_id": "a"}}, + {"message": "search ok", "level": "INFO", "meta": {"run_id": "a"}}, + {"message": "search failed", "level": "ERROR", "meta": {"run_id": "b"}}, + {"message": "unrelated", "level": "ERROR", "meta": {"run_id": "a"}}, + ] + + assert len(search_records(records, query="search")) == 3 + assert len(search_records(records, query="search", level="ERROR")) == 2 + assert len(search_records(records, query="search", run_id="b")) == 1 + + +def test_search_records_matches_inside_meta_too() -> None: + records = [{"message": "x", "level": "INFO", "meta": {"user_id": "customer-42"}}] + + assert len(search_records(records, query="customer-42")) == 1 + assert len(search_records(records, query="nope")) == 0 + + +def test_search_records_is_capped_at_limit() -> None: + records = [{"message": f"m{i}", "level": "INFO", "meta": {}} for i in range(50)] + + assert len(search_records(records, limit=5)) == 5 + + +def test_search_records_with_no_filters_returns_everything_up_to_the_cap() -> None: + records = [{"message": "a", "level": "INFO", "meta": {}}] + + assert search_records(records) == records + + +# --- the real HTTP server ---------------------------------------------------- + + +@pytest.fixture() +def running_server(tmp_path: Path) -> Iterator[TraceViewerServer]: + path = tmp_path / "app.log" + _write_a_run(path) + server = TraceViewerServer(JSONLSource(path), port=0) + thread = threading.Thread(target=server.serve_forever, daemon=True) + thread.start() + try: + for _ in range(50): + try: + urllib.request.urlopen(server.url, timeout=1) + break + except (urllib.error.URLError, ConnectionError): + time.sleep(0.02) + yield server + finally: + server.shutdown() + thread.join(timeout=5) + + +def _get(server: TraceViewerServer, path: str) -> Any: + with urllib.request.urlopen(server.url.rstrip("/") + path, timeout=5) as response: + return response.status, json.loads(response.read()) + + +def test_root_serves_the_static_page(running_server: TraceViewerServer) -> None: + with urllib.request.urlopen(running_server.url, timeout=5) as response: + assert response.status == 200 + assert response.headers["Content-Type"].startswith("text/html") + body = response.read().decode("utf-8") + + assert "LogQuill trace viewer" in body + assert "/api/runs" in body # the page really does call the API + + +def test_api_runs_lists_the_run(running_server: TraceViewerServer) -> None: + status, runs = _get(running_server, "/api/runs") + + assert status == 200 + assert [r["run_id"] for r in runs] == ["run-1"] + + +def test_api_trace_returns_the_span_tree(running_server: TraceViewerServer) -> None: + status, tree = _get(running_server, "/api/trace?run_id=run-1") + + assert status == 200 + assert tree[0]["record"]["message"] == "run" + assert tree[0]["rollup"]["tokens_in"] == 10 + + +def test_api_trace_without_run_id_is_a_400(running_server: TraceViewerServer) -> None: + with pytest.raises(urllib.error.HTTPError) as exc_info: + urllib.request.urlopen(running_server.url + "api/trace", timeout=5) + + assert exc_info.value.code == 400 + + +def test_api_search_finds_the_error(running_server: TraceViewerServer) -> None: + status, hits = _get(running_server, "/api/search?q=broke") + + assert status == 200 + assert hits[0]["message"] == "something broke" + + +def test_an_unknown_path_is_a_404(running_server: TraceViewerServer) -> None: + with pytest.raises(urllib.error.HTTPError) as exc_info: + urllib.request.urlopen(running_server.url + "nope", timeout=5) + + assert exc_info.value.code == 404 + + +def test_port_zero_resolves_to_a_real_bound_port(running_server: TraceViewerServer) -> None: + assert running_server.port > 0 + assert str(running_server.port) in running_server.url + + +def test_the_static_page_ships_inside_the_package() -> None: + assert (STATIC_DIR / "trace_viewer.html").is_file()