From 271fdf9b5a641ba151715bf6a475895e0126e4d4 Mon Sep 17 00:00:00 2001 From: Yiorgis Gozadinos Date: Fri, 10 Jul 2026 13:23:17 +0300 Subject: [PATCH] Add Logfire debugging skills and worker-breaker event --- .claude/skills/debug-evals/SKILL.md | 151 +++++++++++++ .claude/skills/debug-ingestion/SKILL.md | 198 ++++++++++++++++++ .gitignore | 6 + CHANGELOG.md | 1 + docs/ingester.md | 8 +- evaluations/README.md | 7 + .../haiku/rag/ingester/workers/pool.py | 7 + tests/ingester/test_workers.py | 31 +++ 8 files changed, 408 insertions(+), 1 deletion(-) create mode 100644 .claude/skills/debug-evals/SKILL.md create mode 100644 .claude/skills/debug-ingestion/SKILL.md diff --git a/.claude/skills/debug-evals/SKILL.md b/.claude/skills/debug-evals/SKILL.md new file mode 100644 index 00000000..16f47c29 --- /dev/null +++ b/.claude/skills/debug-evals/SKILL.md @@ -0,0 +1,151 @@ +--- +name: debug-evals +description: Debug haiku.rag evaluation runs in Logfire. Use when asked to look at Logfire for an eval run, find failing or low-scoring eval cases, compare runs, check citation quality (cited_map) or judge pass rate (answer_equivalent), or explain why an eval case failed. Drives the Logfire MCP against the `evals` service. +--- + +# Debug eval runs in Logfire + +Eval runs (`evaluations/`) ship spans to Logfire under `service_name = 'evals'`. +This skill finds a run, surfaces its metrics and failures, and drills into a +single case. Read-only. + +## How to query + +1. Confirm the current schema with `mcp__logfire__schema_reference` (spans and + logs share the `records` table). +2. Run SQL with `mcp__logfire__arbitrary_query` (`query` + `age` in minutes, max + 30 days). The same SQL works pasted into Logfire's Explore UI. +3. Read span attributes as JSON: `attributes->>'key'`, nested as + `attributes->'a'->'b'->>'c'`. Cast when needed: `(...)::float`, `(...)::int`. +4. Hand back a clickable trace with `mcp__logfire__logfire_link(trace_id)`. + +## Vocabulary + +A run is one experiment span; its cases are direct children sharing its +`trace_id`. + +- Experiment span: `span_name = 'evaluate {name}'` (scope `pydantic-evals`). + - `attributes->>'name'` — run label (the `--name` arg, or `{dataset}_qa_evaluation` / `{dataset}_retrieval_evaluation`). + - `attributes->>'dataset_name'` — dataset. + - `(attributes->>'assertion_pass_rate')::float` — overall judge pass rate (QA runs). + - `attributes->'logfire.experiment.metadata'->'metadata'` — run config: `target` (`rag-skill`|`analysis-skill`), `qa_model`, `embedder_model`, `chunk_size`, `search_limit`, `rerank_model`, `judge_model`, `skill_model`, etc. + - `trace_id` — scopes the whole run. +- Case span: `span_name = 'case: {case_name}'` (scope `pydantic-evals`). + - `message` — `case: `. + - `attributes->'assertions'->'answer_equivalent'->>'value'` — `'true'`/`'false'` (LLM judge verdict). `->>'reason'` — why. + - `attributes->'scores'->'cited_map'->>'value'` — citation average precision (0..1). + - `attributes->'scores'->'number_match'->>'value'` — numeric-answer match (datasets that use it). + - `duration` — task time in seconds. +- Inside each case the skill under test emits agent spans (scope `pydantic-ai`): + `execute {task}`, `agent run`, `running tool`, `chat {model}`. + +The service is `evals` regardless of model, so filter on `service_name = 'evals'` +first. `otel_scope_name` separates the layers (`pydantic-evals` for run/case, +`pydantic-ai` for the agent). + +## Canned queries + +Recent runs (pick a `trace_id` to drill in): + +```sql +SELECT attributes->>'name' AS run, + attributes->>'dataset_name' AS dataset, + service_version, + (attributes->>'assertion_pass_rate')::float AS pass_rate, + start_timestamp, trace_id +FROM records +WHERE service_name='evals' AND span_name='evaluate {name}' +ORDER BY start_timestamp DESC +LIMIT 20; +``` + +Run summary (pass rate, mean citation score, mean task time): + +```sql +SELECT count(*) AS cases, + avg(CASE WHEN attributes->'assertions'->'answer_equivalent'->>'value'='true' + THEN 1.0 ELSE 0.0 END) AS pass_rate, + avg((attributes->'scores'->'cited_map'->>'value')::float) AS mean_cited_map, + avg(duration) AS mean_task_seconds +FROM records +WHERE service_name='evals' AND span_name='case: {case_name}' + AND trace_id=''; +``` + +Failing cases (judge said not equivalent), newest first, with the reason: + +```sql +SELECT message AS case_name, duration, + attributes->'assertions'->'answer_equivalent'->>'reason' AS reason +FROM records +WHERE service_name='evals' AND span_name='case: {case_name}' + AND trace_id='' + AND attributes->'assertions'->'answer_equivalent'->>'value'='false' +ORDER BY start_timestamp; +``` + +Low-citation cases (answer may be right but grounding is weak): + +```sql +SELECT message AS case_name, + (attributes->'scores'->'cited_map'->>'value')::float AS cited_map +FROM records +WHERE service_name='evals' AND span_name='case: {case_name}' + AND trace_id='' + AND (attributes->'scores'->'cited_map'->>'value')::float < 0.5 +ORDER BY cited_map; +``` + +Error / null cases in a run (an exception aborted the case): + +```sql +SELECT message, span_name, exception_type, exception_message +FROM records +WHERE service_name='evals' AND trace_id='' AND is_exception=true +ORDER BY start_timestamp +LIMIT 50; +``` + +Slowest cases (task time drives run cost): + +```sql +SELECT message AS case_name, duration +FROM records +WHERE service_name='evals' AND span_name='case: {case_name}' + AND trace_id='' +ORDER BY duration DESC +LIMIT 20; +``` + +Drill into one case's agent activity (all cases share the run `trace_id`, so +bound by the case's own time window): + +```sql +SELECT span_name, message, duration, is_exception +FROM records +WHERE service_name='evals' AND trace_id='' + AND otel_scope_name='pydantic-ai' + AND start_timestamp BETWEEN '' AND '' +ORDER BY start_timestamp; +``` + +## Workflow + +1. List recent runs, pick the one to inspect by `name` + `start_timestamp`, note + its `trace_id`. +2. Run the summary query for the headline numbers (pass rate, mean_cited_map, + mean_task_seconds — always report task time). +3. Pull failing and low-citation cases, read the judge `reason`. +4. To understand one case, take its start/end from the case query and run the + agent-activity query, then `logfire_link(trace_id)` so the user can expand + that case in the UI. + +## When a query returns nothing + +Span names or attributes may have changed. Probe: + +```sql +SELECT DISTINCT otel_scope_name, span_name +FROM records WHERE service_name='evals' +ORDER BY 1,2; +``` diff --git a/.claude/skills/debug-ingestion/SKILL.md b/.claude/skills/debug-ingestion/SKILL.md new file mode 100644 index 00000000..82d42a7d --- /dev/null +++ b/.claude/skills/debug-ingestion/SKILL.md @@ -0,0 +1,198 @@ +--- +name: debug-ingestion +description: Debug haiku.rag ingestion in Logfire. Use when asked to look at Logfire for ingestion, find failed or dead ingestion jobs, investigate retries or circuit-breaker events, trace a document through convert/chunk/embed/store, find which docling-serve instance served a request, spot slow conversions, or tell concurrent ingesters apart. Drives the Logfire MCP against the `haiku-ingester` service. +--- + +# Debug ingestion in Logfire + +The ingester ships spans to Logfire under `service_name = 'haiku-ingester'` (or a +custom `OTEL_SERVICE_NAME` if set per process). This skill finds failing jobs, +traces a document through the pipeline, and pinpoints the docling-serve instance +that served a request. Read-only. + +## How to query + +1. Confirm the current schema with `mcp__logfire__schema_reference` (spans and + logs share the `records` table). +2. Run SQL with `mcp__logfire__arbitrary_query` (`query` + `age` in minutes, max + 30 days). The same SQL works pasted into Logfire's Explore UI. +3. Read span attributes as JSON: `attributes->>'key'`, cast when needed + (`(attributes->>'attempt')::int`). +4. Hand back a clickable trace with `mcp__logfire__logfire_link(trace_id)`. +5. For recent exceptions tied to a file, `mcp__logfire__find_exceptions_in_file` + accepts `client/documents.py`, `ingester/workers/pool.py`, or + `ingester/pollers/base.py`. + +Adjust `service_name` if the operator set `OTEL_SERVICE_NAME` (e.g. per tenant). +Interactive `haiku-rag` ingests emit the same `document.*` spans under the CLI's +service name (or `unknown_service` for older runs), not `haiku-ingester`. + +## Vocabulary + +Span tree (all scope `haiku.rag`), each level nests under the one above and +shares its `trace_id`: + +- Poller: `ingester.poller.sweep` | `ingester.poller.dry_run` | + `ingester.poller.watch_event`. + - `attributes->>'source_id'`; watch adds `change`, `uri`. + - sweep sets `skipped`, `skip_reason` (`pending_work` / `circuit_open`), + `upsert`, `delete`, `unchanged`, `consecutive_failures`. A failed sweep + records the exception on the span (`exception_type` / `exception_message`). +- Job: `ingester.job` — `attributes->>'source_id'`, `->>'uri'`, `->>'op'` + (`UPSERT`/`DELETE`), `(->>'attempt')::int`. `is_exception=true` marks a job + that raised. +- Document pipeline: `document.fetch` (`bytes`, `content_hash`), + `document.convert`, `document.chunk` (`chunks_created`), `document.embed`, + `document.store` (`op` `create`/`update`, `document_id`). +- docling-serve: `docling_serve.request` — `attributes->>'name'` (operation), + `->>'url'` (instance), `(->>'attempt')::int`. A retry emits a new span with a + different `url`, so failover shows as sibling spans. + +Failures in Logfire are span-level: sweep and job exceptions sit on the span +(`is_exception`, `exception_type`, `exception_message`, `level >= 17`). A job that +raised `PermanentError` was dead-lettered on that attempt; `attempt >= 2` on an +`ingester.job` span is a retry. The worker circuit breaker opening emits a +dedicated event `span_name = 'ingester.worker breaker opened'` (attributes +`source_id`, `threshold`, `cooldown_s`), fired once per closed->open transition. +The ingester's per-job dead/reschedule narration stays on stderr and the queue +(`GET /jobs?status=dead` on the control-plane API), not Logfire. + +## Canned queries + +Job outcomes by source: + +```sql +SELECT attributes->>'source_id' AS source_id, + attributes->>'op' AS op, + is_exception, + count(*) AS n +FROM records +WHERE service_name='haiku-ingester' AND span_name='ingester.job' +GROUP BY 1,2,3 +ORDER BY n DESC; +``` + +Failed jobs with the exception and trace to drill in: + +```sql +SELECT attributes->>'uri' AS uri, + (attributes->>'attempt')::int AS attempt, + exception_type, exception_message, trace_id +FROM records +WHERE service_name='haiku-ingester' AND span_name='ingester.job' + AND is_exception=true +ORDER BY start_timestamp DESC +LIMIT 50; +``` + +Retried and dead-lettered jobs (a `PermanentError` was dead-lettered on that +attempt; `TransientError` at a high `attempt` was retried): + +```sql +SELECT attributes->>'uri' AS uri, + (attributes->>'attempt')::int AS attempt, + attributes->>'op' AS op, + exception_type, trace_id +FROM records +WHERE service_name='haiku-ingester' AND span_name='ingester.job' + AND (is_exception=true OR (attributes->>'attempt')::int > 1) +ORDER BY start_timestamp DESC +LIMIT 50; +``` + +Trace one document end-to-end (take `trace_id` from a job above): + +```sql +SELECT span_name, duration, is_exception, + attributes->>'url' AS docling_url, + attributes->>'op' AS op +FROM records +WHERE trace_id='' +ORDER BY start_timestamp; +``` + +Worker circuit-breaker trips (a source paused after consecutive transient +failures): + +```sql +SELECT start_timestamp, + attributes->>'source_id' AS source_id, + (attributes->>'cooldown_s')::float AS cooldown_s +FROM records +WHERE service_name='haiku-ingester' + AND span_name='ingester.worker breaker opened' +ORDER BY start_timestamp DESC +LIMIT 50; +``` + +docling-serve instance health and failover: + +```sql +SELECT attributes->>'url' AS instance, + (attributes->>'attempt')::int AS attempt, + is_exception, + count(*) AS n +FROM records +WHERE service_name='haiku-ingester' AND span_name='docling_serve.request' +GROUP BY 1,2,3 +ORDER BY n DESC; +``` + +Slowest pipeline stages: + +```sql +SELECT span_name, + attributes->>'uri' AS uri, + duration, trace_id +FROM records +WHERE service_name='haiku-ingester' + AND span_name IN ('document.convert','docling_serve.request','document.embed','document.chunk') +ORDER BY duration DESC +LIMIT 20; +``` + +Per-source sweep summary: + +```sql +SELECT attributes->>'source_id' AS source_id, + sum((attributes->>'upsert')::int) AS upserts, + sum((attributes->>'delete')::int) AS deletes, + sum((attributes->>'unchanged')::int) AS unchanged, + max(attributes->>'skip_reason') AS last_skip_reason +FROM records +WHERE service_name='haiku-ingester' AND span_name='ingester.poller.sweep' +GROUP BY 1; +``` + +Tell concurrent ingesters apart (set distinct `OTEL_SERVICE_NAME` per process; a +single host still separates by `process_pid` / `service_instance_id`): + +```sql +SELECT service_name, service_instance_id, process_pid, count(*) AS n +FROM records +WHERE span_name='ingester.job' +GROUP BY 1,2,3 +ORDER BY n DESC; +``` + +## Workflow + +1. Job outcomes by source shows where failures cluster. +2. Failed jobs surfaces the exception and each failing `trace_id`. +3. Trace one document end-to-end to see which stage failed and, for a docling + source, which `docling_url` served it (and whether it failed over). +4. Retried/dead-lettered jobs show what the queue kept struggling with; a + `ingester.worker breaker opened` event flags a source the pool paused. The + per-job dead/reschedule narration stays in the ingester console and the queue + API, not Logfire. +5. `logfire_link(trace_id)` for a document lets the user expand the full tree. + +## When a query returns nothing + +Span names or attributes may have changed. Probe: + +```sql +SELECT DISTINCT span_name +FROM records WHERE service_name='haiku-ingester' +ORDER BY 1; +``` diff --git a/.gitignore b/.gitignore index f83fce38..e1e4b608 100644 --- a/.gitignore +++ b/.gitignore @@ -35,3 +35,9 @@ site/ # Claude Code personal notes .claude.local.md + +# Publish shared Claude skills, keep local Claude settings private +!.claude/ +.claude/* +!.claude/skills/ +!.claude/skills/** diff --git a/CHANGELOG.md b/CHANGELOG.md index 309d9eff..d82a6087 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -5,6 +5,7 @@ - Logfire spans carry `service.version` and a per-process `service.name` (`haiku-ingester`, `haiku-rag`, `haiku-rag-app`); `OTEL_SERVICE_NAME` / `LOGFIRE_SERVICE_NAME` override the default. - Each docling-serve request emits a `docling_serve.request` span carrying the instance `url` and `attempt`. +- The ingester worker circuit breaker opening emits a `ingester.worker breaker opened` Logfire event with `source_id`, `threshold`, `cooldown_s`. ## [0.65.0] - 2026-07-09 diff --git a/docs/ingester.md b/docs/ingester.md index fdb7fa51..243f2fed 100644 --- a/docs/ingester.md +++ b/docs/ingester.md @@ -654,7 +654,13 @@ The span tree is `ingester.poller.sweep` -> `ingester.job` (tagged with `source_id` and `uri`) -> `document.convert` / `document.chunk`. When a source uses docling-serve, each request emits a `docling_serve.request` span carrying the instance `url` and `attempt`, so a failed conversion can be traced to the -exact instance that served it. +exact instance that served it. A worker circuit breaker opening emits an +`ingester.worker breaker opened` event with `source_id`, `threshold`, and +`cooldown_s`. + +The `debug-ingestion` skill in `.claude/skills/` turns these spans into +ready-made Logfire queries (failed jobs, docling-serve failover, per-source +sweeps, breaker trips) for use from Claude Code. ### Operating against the API diff --git a/evaluations/README.md b/evaluations/README.md index 78e59d42..37299dc3 100644 --- a/evaluations/README.md +++ b/evaluations/README.md @@ -55,6 +55,13 @@ the judge (defaults to `qa.model`, or `analysis.model` when set for the analysis-skill target). A citation retrieval metric (`cited_map`) is computed alongside QA accuracy from the URIs the skill registered via the `cite` tool. +### Debugging runs in Logfire + +With `LOGFIRE_TOKEN` set, runs ship spans under `service_name = 'evals'`. The +`debug-evals` skill in `.claude/skills/` turns these into ready-made Logfire +queries (recent runs, per-case pass rate and `cited_map`, failing and slowest +cases) for use from Claude Code. + ### Pre-built Databases Download pre-built evaluation databases from HuggingFace: diff --git a/haiku_rag_slim/haiku/rag/ingester/workers/pool.py b/haiku_rag_slim/haiku/rag/ingester/workers/pool.py index 501477bb..a7bbcf63 100644 --- a/haiku_rag_slim/haiku/rag/ingester/workers/pool.py +++ b/haiku_rag_slim/haiku/rag/ingester/workers/pool.py @@ -12,6 +12,7 @@ from haiku.rag.ingester.queue.models import Job, JobOp from haiku.rag.ingester.queue.repository import JobRepo, SyncStateRepo from haiku.rag.ingester.workers.pipeline import run_job from haiku.rag.ingester.workers.retry import RetryPolicy, compute_backoff +from haiku.rag.telemetry import logfire if TYPE_CHECKING: from collections.abc import Mapping @@ -352,6 +353,12 @@ class WorkerPool: was_closed = not breaker.is_open breaker.record_failure() if was_closed and breaker.is_open: + logfire.warn( + "ingester.worker breaker opened", + source_id=job.source_id, + threshold=_WORKER_BREAKER_THRESHOLD, + cooldown_s=_WORKER_BREAKER_COOLDOWN_S, + ) logger.warning( "Worker breaker opened for source %s after %d consecutive " "transient failures; pausing its claims for %.0fs", diff --git a/tests/ingester/test_workers.py b/tests/ingester/test_workers.py index c1602e40..bf6f6c28 100644 --- a/tests/ingester/test_workers.py +++ b/tests/ingester/test_workers.py @@ -934,6 +934,37 @@ async def test_breaker_opens_after_n_consecutive_transient_failures(client, jobs assert remaining_before >= 1 +@pytest.mark.asyncio +async def test_breaker_open_emits_logfire_event(client, jobs, sync, monkeypatch): + """The worker breaker opening emits exactly one structured Logfire event + (on the closed->open transition, not on every failure), tagged with the + source so it's queryable without scraping stderr logs.""" + from haiku.rag.ingester.workers import pool as pool_module + from haiku.rag.ingester.workers.pool import _WORKER_BREAKER_THRESHOLD + + events: list[dict] = [] + + def _capture(msg, /, **attrs): + events.append({"msg": msg, **attrs}) + + monkeypatch.setattr(pool_module.logfire, "warn", _capture) + + client.create_document_from_source.side_effect = TransientError("downstream down") + for i in range(_WORKER_BREAKER_THRESHOLD + 2): + await jobs.enqueue("src", f"u{i}", JobOp.UPSERT, max_attempts=5) + + pool = _pool( + client, jobs, sync, retry_policy=RetryPolicy(base_delay_s=60.0, jitter=0.0) + ) + # Drain past the threshold: the transition fires once, later failures don't. + for _ in range(_WORKER_BREAKER_THRESHOLD + 1): + await pool.drain_once() + assert pool.breaker_open is True + + assert len(events) == 1 + assert events[0]["source_id"] == "src" + + @pytest.mark.asyncio async def test_breaker_pauses_worker_loop_claims(client, jobs, sync): """Worker loop honours the breaker: an open source is excluded from