From 6598c7fa8536d97f6cc27dadb7c89124579ce4b9 Mon Sep 17 00:00:00 2001 From: Bryan Van Deusen Date: Sat, 3 Oct 2026 15:34:14 -0400 Subject: [PATCH] =?UTF-8?q?feat(retrieval):=20a=20review=20pass=20judges?= =?UTF-8?q?=20whether=20injected=20lines=20related=20=E2=80=94=20menus=5Ft?= =?UTF-8?q?o=5Freview,=20judge=5Fmenu=20and=20a=20judged=20readout=20(#477?= =?UTF-8?q?2)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit An open rate cannot say whether a menu line related: every line carries its matched passage (#4364), so "not opened" covers unrelated, enough as shown, and already in context. #4772 "Injected notes are never judged". - retrieval_judgments (0115): a reviewer verdict per line of a logged call, on_point / adjacent / unrelated, with its reason, rank, budget side and whether the agent opened it within the hour. - menus_to_review re-runs a random sample of unjudged auto_inject calls with the arm's own parameters, past its budget, passage on every line. judge_menu records verdicts, re-deriving rank from a fresh re-run. - retrieval_telemetry gains a judged block (by rank, within/beyond budget, on_point_unopened). surfaced_never_pulled stops blaming titles. - missed-retrieval guidance names the review before a budget move. Co-Authored-By: Claude Opus 5.5 --- alembic/versions/0115_retrieval_judgments.py | 56 +++ plugin/.claude-plugin/plugin.json | 2 +- .../skills/using-scribe/missed-retrieval.md | 9 + src/scribe/mcp/server.py | 7 + src/scribe/mcp/tools/__init__.py | 3 +- src/scribe/mcp/tools/retrieval_review.py | 74 ++++ src/scribe/mcp/tools/search.py | 20 +- src/scribe/models/__init__.py | 1 + src/scribe/models/retrieval_judgment.py | 82 ++++ src/scribe/services/backup.py | 4 + src/scribe/services/retrieval_review.py | 352 ++++++++++++++++++ src/scribe/services/retrieval_telemetry.py | 30 +- tests/test_retrieval_review.py | 250 +++++++++++++ 13 files changed, 879 insertions(+), 11 deletions(-) create mode 100644 alembic/versions/0115_retrieval_judgments.py create mode 100644 src/scribe/mcp/tools/retrieval_review.py create mode 100644 src/scribe/models/retrieval_judgment.py create mode 100644 src/scribe/services/retrieval_review.py create mode 100644 tests/test_retrieval_review.py diff --git a/alembic/versions/0115_retrieval_judgments.py b/alembic/versions/0115_retrieval_judgments.py new file mode 100644 index 00000000..6ed638bd --- /dev/null +++ b/alembic/versions/0115_retrieval_judgments.py @@ -0,0 +1,56 @@ +"""retrieval_judgments — a reviewer's verdict on each line of a logged menu +(#4772) + +Revision ID: 0115 +Revises: 0114 +Create Date: 2026-10-03 + +Whether an injected line RELATED, which the usage tables cannot say: since +every line carries its matched passage, "not opened" no longer implies "not +relevant". FK-free like the telemetry tables beside it; no CHECK on `verdict`, +which the service validates. +""" +import sqlalchemy as sa +from alembic import op + +revision = "0115" +down_revision = "0114" +branch_labels = None +depends_on = None + + +def upgrade() -> None: + op.create_table( + "retrieval_judgments", + sa.Column("id", sa.BigInteger(), primary_key=True), + sa.Column( + "created_at", sa.DateTime(timezone=True), nullable=False, + server_default=sa.text("now()"), + ), + sa.Column("user_id", sa.BigInteger(), nullable=True), + sa.Column("retrieval_log_id", sa.BigInteger(), nullable=False), + sa.Column("source", sa.Text(), nullable=False), + sa.Column("query", sa.Text(), nullable=True), + sa.Column("record_id", sa.BigInteger(), nullable=False), + sa.Column("rank", sa.Integer(), nullable=False), + sa.Column("score", sa.Float(), nullable=True), + sa.Column("within_budget", sa.Boolean(), nullable=False), + sa.Column("verdict", sa.Text(), nullable=False), + sa.Column("reason", sa.Text(), nullable=False), + sa.Column("opened_after", sa.Boolean(), nullable=True), + sa.UniqueConstraint( + "retrieval_log_id", "record_id", "user_id", + name="uq_retrieval_judgment_line", + ), + ) + op.create_index( + "ix_retrieval_judgments_source_created", "retrieval_judgments", + ["source", "created_at"], + ) + + +def downgrade() -> None: + op.drop_index( + "ix_retrieval_judgments_source_created", table_name="retrieval_judgments", + ) + op.drop_table("retrieval_judgments") diff --git a/plugin/.claude-plugin/plugin.json b/plugin/.claude-plugin/plugin.json index 7c38fe84..013d69dc 100644 --- a/plugin/.claude-plugin/plugin.json +++ b/plugin/.claude-plugin/plugin.json @@ -1,7 +1,7 @@ { "name": "scribe", "description": "Scribe for Claude Code: connects the scribe MCP server, adds the hooks that deliver live project state and relevant records at the right moment, ships the shared client-neutral Scribe skills (using-scribe, writing-plans, reporting-back, systematic-debugging, verification, brainstorming, reusing-code, shape-accounting), and syncs your saved Scribe Processes as skills (/scribe:sync).", - "version": "2026.10.03.0308", + "version": "2026.10.03.1934", "author": { "name": "Bryan Van Deusen" }, diff --git a/plugin/skills/using-scribe/missed-retrieval.md b/plugin/skills/using-scribe/missed-retrieval.md index d65750e8..c062d255 100644 --- a/plugin/skills/using-scribe/missed-retrieval.md +++ b/plugin/skills/using-scribe/missed-retrieval.md @@ -36,6 +36,15 @@ everything else that was sitting in the same band. `reason` is required and has to say what you read, because it is what lets the operator disagree with a number they did not choose. + **A budget is judged by what sits past it, and an open rate cannot judge + it.** Every menu line carries its matched passage, so a record left + unopened may have been unrelated, enough as shown, or already in context. + `menus_to_review` re-runs a sample of logged menus to a depth past the + budget; judge each line from what is shown with `judge_menu`, then read + `retrieval_telemetry`'s `judged` block. If the lines past the cut are + mostly `on_point`, the budget is costing hits; mostly `unrelated`, it is + doing its job. + Reaching for `tune_retrieval` before opening a single record is the wrong move, and it is the one that feels efficient. Worked example, measured on this install: a rule granting a routine push scored 0.6515 and ranked 5th for the diff --git a/src/scribe/mcp/server.py b/src/scribe/mcp/server.py index 8bdb3ba8..79ef5416 100644 --- a/src/scribe/mcp/server.py +++ b/src/scribe/mcp/server.py @@ -133,6 +133,10 @@ _READ_ONLY_TOOLS = frozenset({ # the prefixes the completeness test derives from, so nothing would have # prompted this decision. "notes_due_for_verification", + # The review pass's sample (#4772). Re-runs logged queries and writes no + # row of any kind — not even telemetry, since nothing is put in front of a + # working session. Spelled out for retrieval_telemetry's reason. + "menus_to_review", # Its rule twin and a rule's edit history (milestones 312 and 323). Both # pure reads, and both sat unlisted — so a read key was refused them — for # the same reason: no read prefix, back when the completeness test only @@ -198,6 +202,9 @@ _WRITE_TOOLS = frozenset({ # reads, and it appends the reason to the audit trail (#4102). "tune_retrieval", "migrate_retrieval_floor", + # A reviewer's verdicts on logged menu lines (#4772) — rows carrying free + # prose the agent authored, `rule_outcome`'s reason for being a write. + "judge_menu", # trash "restore", "purge_trash", }) diff --git a/src/scribe/mcp/tools/__init__.py b/src/scribe/mcp/tools/__init__.py index 59aeea15..1343d149 100644 --- a/src/scribe/mcp/tools/__init__.py +++ b/src/scribe/mcp/tools/__init__.py @@ -6,7 +6,7 @@ from `mcp.server.build_mcp_server`. """ from scribe.mcp.tools import ( design_systems, lessons, milestones, notes, processes, projects, recent, repos, - retrieval_tuning, + retrieval_review, retrieval_tuning, wide_net, rulebooks, search, shapes, snippets, systems, tags, tasks, trash, ) @@ -16,6 +16,7 @@ def register_all(mcp) -> None: """Register every tool module's tools on the given MCPServer instance.""" search.register(mcp) retrieval_tuning.register(mcp) + retrieval_review.register(mcp) wide_net.register(mcp) notes.register(mcp) tasks.register(mcp) diff --git a/src/scribe/mcp/tools/retrieval_review.py b/src/scribe/mcp/tools/retrieval_review.py new file mode 100644 index 00000000..92363473 --- /dev/null +++ b/src/scribe/mcp/tools/retrieval_review.py @@ -0,0 +1,74 @@ +"""Judging what a retrieval arm offered (#4772). + +`retrieval_telemetry` says what the reader did with each line — opened it or +not. These two say whether the line RELATED, which is a reading only someone +looking at the query beside the line can do. +""" +from scribe.mcp._context import current_user_id +from scribe.services import retrieval_review as review_svc + + +async def menus_to_review( + source: str = "auto_inject", n: int = 5, days: int = 14, depth: int = 6, +) -> dict: + """A sample of logged injection menus, re-run so you can judge each line. + + Reach for this before moving a budget or a floor, and whenever an open rate + is about to be read as a verdict on relevance. It is not one: every menu + line carries the passage that matched, so a record left unopened may have + been unrelated, enough as shown, or already in context, and the usage + counters cannot tell those apart. A judged sample can. + + Each menu is one logged call: the query the arm searched (enriched, as it + was asked), its `budget`, and `lines` — the query re-run with the arm's own + parameters to `depth`, each line with `rank`, `score`, `kind`, `name`, the + `passage` that matched, `within_budget` (was it inside the cut) and + `opened_after` (did your agent open it within an hour). Lines past the + budget are there on purpose: they are what a larger budget would add. + + `missing` names ids the call logged that the re-run no longer finds — the + corpus has moved since. Judge what is in front of you. + + JUDGE FROM THE LINE; DO NOT OPEN THE RECORD. Opening counts as a pull and + would credit the arm with an open it never earned. Then record verdicts + with `judge_menu`. `how_to_judge` in the response repeats the vocabulary. + + Args: + source: the arm to review. Only arms whose search can be re-run + exactly are offered; today that is `auto_inject`. + n: how many calls to sample, 1–20. Random among the calls you have not + judged, so a review does not describe one afternoon. + days: how far back to sample from. + depth: how many candidates to re-run per call, 1–10. Past the budget, + so the cut itself can be judged. + """ + return await review_svc.menus_to_review( + current_user_id(), source=source, n=n, days=days, depth=depth, + ) + + +async def judge_menu(log_id: int, verdicts: list[dict]) -> dict: + """Record your verdict on lines of one menu from `menus_to_review`. + + `verdicts` is a list of `{record_id, verdict, reason}`: + + "on_point" — someone in the query's situation should read this. + "adjacent" — the same area, but it would not change what they do. + "unrelated" — close in words only. + + `reason` is REQUIRED for every verdict: what in the query and the passage + decided it. A verdict nobody can re-read is a threshold, and this exists to + replace one. Rank, score and which side of the budget the line sat are + taken from a fresh re-run, not from you. Judging a line again replaces your + earlier verdict on it. + + The verdicts are read back as `retrieval_telemetry`'s `judged` block. + """ + return await review_svc.judge_menu( + current_user_id(), log_id=log_id, verdicts=verdicts, + ) + + +def register(mcp) -> None: + mcp.tool(name="menus_to_review")(menus_to_review) + mcp.tool(name="judge_menu")(judge_menu) diff --git a/src/scribe/mcp/tools/search.py b/src/scribe/mcp/tools/search.py index 129e3979..682d24ea 100644 --- a/src/scribe/mcp/tools/search.py +++ b/src/scribe/mcp/tools/search.py @@ -380,7 +380,7 @@ async def retrieval_telemetry( — so `near_miss_samples=5` and opening the ids it returns is the step that separates a real miss from a bar doing its job. - Four readouts, from the four tables built for them: + Five readouts, from the five tables built for them: `sources` — per retrieval surface (`auto_inject`, `write_path`, `mcp_search`, …), from `retrieval_logs`: `calls`, `zero_result_calls`, @@ -539,6 +539,19 @@ It is an UPPER BOUND per surface: a pull records the door it came a ratio of opens would read near zero on an arm that is working. `system_usage_failed: true` means its read broke and the zeros mean nothing. + `judged` — whether the lines RELATED (#4772), from `retrieval_judgments`: + per source, a reviewer's verdicts (`on_point` / `adjacent` / `unrelated`) + on a sample of logged menus, `within_budget` against `beyond_budget` and + `by_rank`. Every usage block above counts what the reader DID; a line + carries its matched passage, so "not opened" is no verdict on relevance, + and this is the block that holds one. `on_point_unopened` counts lines + that related and were left closed. Read it before moving a budget: if the + ranks past the cut are mostly `on_point`, the budget is costing hits; if + they are mostly `unrelated`, it is doing its job. Empty until someone + reviews — `menus_to_review` and `judge_menu` fill it. A sample is small by + nature, so the counts are printed, not rates. `judged_failed: true` means + the read broke. + EVERY COUNTER BLOCK CARRIES ITS OWN COVERAGE — `complete_from` and `covers_window`. `complete_from` is when the number became trustworthy: for one source, its first recorded row; for a section that sums several, @@ -608,8 +621,9 @@ It is an UPPER BOUND per surface: a pull records the door it came - `no_duration` — rows written without timings. A logging gap, not a slow arm, and it devalues every other number from that source. - `surfaced_never_pulled` — distinct records shown and never opened, per - corpus. Read their titles before touching a threshold: a record nobody - opens is usually one whose title does not say when it matters. + corpus. Not a verdict on them: a line carries its matched passage, so + an unopened record may have been unrelated, enough as shown, or already + in context. Judge a sample (`menus_to_review`) before acting on it. - `read_and_unacted` — distinct rules OPENED in the window that recorded no outcome, against the ones that did. The failure milestone 419 was opened on, and the worse sibling of `surfaced_never_pulled` above: a diff --git a/src/scribe/models/__init__.py b/src/scribe/models/__init__.py index ce9d7845..81d0d024 100644 --- a/src/scribe/models/__init__.py +++ b/src/scribe/models/__init__.py @@ -61,6 +61,7 @@ from scribe.models.retrieval_tuning import RetrievalTuningEvent # noqa: E402, F from scribe.models.note_usage import NoteUsageEvent # noqa: E402, F401 from scribe.models.rule_usage import RuleUsageEvent # noqa: E402, F401 from scribe.models.system_usage import SystemUsageEvent # noqa: E402, F401 +from scribe.models.retrieval_judgment import RetrievalJudgment # noqa: E402, F401 from scribe.models.project import Project # noqa: E402, F401 from scribe.models.milestone import Milestone # noqa: E402, F401 from scribe.models.task_log import TaskLog # noqa: E402, F401 diff --git a/src/scribe/models/retrieval_judgment.py b/src/scribe/models/retrieval_judgment.py new file mode 100644 index 00000000..cb630f74 --- /dev/null +++ b/src/scribe/models/retrieval_judgment.py @@ -0,0 +1,82 @@ +from sqlalchemy import BigInteger, Boolean, Float, Index, Integer, Text, UniqueConstraint +from sqlalchemy.orm import Mapped, mapped_column + +from scribe.models import Base +from scribe.models.base import CreatedAtMixin, iso + +ON_POINT = "on_point" +ADJACENT = "adjacent" +UNRELATED = "unrelated" +VERDICTS = (ON_POINT, ADJACENT, UNRELATED) + + +class RetrievalJudgment(Base, CreatedAtMixin): + """One reviewer's verdict on one candidate of one logged retrieval call + (#4772). + + The relevance half the usage tables cannot hold. `note_usage_events` says + whether a line was OPENED, and since every menu line carries its matched + passage (#4364) "not opened" covers an unrelated line, a line whose passage + was enough, and a line the reader already had. Only someone reading the + query beside the line can tell those apart, so this records that reading — + with its reason, because a verdict nobody can re-read is a threshold. + + Keyed to the `retrieval_logs` row it judges, by id and FK-free like every + telemetry table here. `query` is COPIED rather than joined: the verdict is + only re-readable beside the words it was judged against. + + `rank` and `score` are the candidate's place when the logged query was + re-run for review, which may run past the arm's budget on purpose — the + candidates just under the cut are the ones a budget change would add. + `within_budget` says which side of the cut the line was. + + `opened_after` is whether the reviewer's own agent pulled the record within + an hour of the call, read from `note_usage_events`. A correlation, not a + session join — the server has no session identity (see NoteUsageEvent) — + and None when it was not measured. + """ + + __tablename__ = "retrieval_judgments" + + id: Mapped[int] = mapped_column(BigInteger, primary_key=True) + user_id: Mapped[int | None] = mapped_column(BigInteger, nullable=True) + retrieval_log_id: Mapped[int] = mapped_column(BigInteger, nullable=False) + # The judged call's surface, copied so the readout groups without a join. + source: Mapped[str] = mapped_column(Text, nullable=False) + query: Mapped[str | None] = mapped_column(Text, nullable=True) + record_id: Mapped[int] = mapped_column(BigInteger, nullable=False) + # 1-based, in the review's re-run. + rank: Mapped[int] = mapped_column(Integer, nullable=False) + score: Mapped[float | None] = mapped_column(Float, nullable=True) + within_budget: Mapped[bool] = mapped_column(Boolean, nullable=False) + # One of VERDICTS. Plain Text, no CHECK, like the usage tables; the + # service refuses anything else. + verdict: Mapped[str] = mapped_column(Text, nullable=False) + reason: Mapped[str] = mapped_column(Text, nullable=False) + opened_after: Mapped[bool | None] = mapped_column(Boolean, nullable=True) + + __table_args__ = ( + # One verdict per reviewer per line; judging it again replaces it. + UniqueConstraint( + "retrieval_log_id", "record_id", "user_id", + name="uq_retrieval_judgment_line", + ), + Index("ix_retrieval_judgments_source_created", "source", "created_at"), + ) + + def to_dict(self) -> dict: + return { + "id": self.id, + "created_at": iso(self.created_at), + "user_id": self.user_id, + "retrieval_log_id": self.retrieval_log_id, + "source": self.source, + "query": self.query, + "record_id": self.record_id, + "rank": self.rank, + "score": self.score, + "within_budget": self.within_budget, + "verdict": self.verdict, + "reason": self.reason, + "opened_after": self.opened_after, + } diff --git a/src/scribe/services/backup.py b/src/scribe/services/backup.py index 8a90323c..2fd3282e 100644 --- a/src/scribe/services/backup.py +++ b/src/scribe/services/backup.py @@ -159,6 +159,10 @@ _NOT_INCLUDED = [ "app_logs", "notifications", "invitation_tokens", "password_reset_tokens", "user_profiles", "retrieval_logs", + # Verdicts on retrieval_logs rows (#4772), so they go where their subject + # goes: a restore that carried a judgment of a call it did not carry would + # leave a verdict on nothing. A sample, re-takeable by the review pass. + "retrieval_judgments", # Sensitive credentials, same reasoning as api_keys: a backup that carries # forge tokens is a token-exfiltration file. Users re-add connections # after a restore; the per-project pin (projects.forge_connection_id) is diff --git a/src/scribe/services/retrieval_review.py b/src/scribe/services/retrieval_review.py new file mode 100644 index 00000000..dbe5c490 --- /dev/null +++ b/src/scribe/services/retrieval_review.py @@ -0,0 +1,352 @@ +"""Judging what a retrieval arm offered — the review pass (#4772). + +The usage tables record what the reader DID with a line: opened it or not. +Since every menu line carries its matched passage (#4364), that action no +longer stands in for relevance — a line can be unrelated, related and enough +as shown, related and already in context, or related and ignored, and all four +are "not opened". Only someone reading the query beside the line can tell +them apart. This module is how that reading gets done and kept. + +THE SHAPE. `retrieval_logs` already stores each call's query (enriched, as the +arm searched it) and its ranked candidates. A review takes a sample of logged +calls the reviewer has not judged, RE-RUNS each query with the arm's own +parameters to a depth past its budget, and hands back every candidate with the +passage that matched. The reviewer records a verdict per line, with a reason. + +WHY RE-RUN rather than replay the logged ids. The log holds ids and scores, +not passages, and the reviewer needs the passage — it is what the reader saw. +And the depth past the budget is the point: the candidates just under the cut +are the ones a budget change would add, and the only way to judge them is to +look at them. The re-run is against today's corpus; `missing` names any +logged id the re-run no longer finds, so drift is visible rather than silent. + +WHY IT WRITES NO TELEMETRY. A review is not a retrieval surface — nothing is +put in front of a working session — so it records no `retrieval_logs` row and +no surfacing. A row here would be a call to the arm that the arm never made. + +ATTENDED, NEVER BULK. Every verdict needs a reason. A verdict without one is a +threshold with extra steps, which is the thing this exists to replace. +""" +import logging +from datetime import datetime, timedelta, timezone + +from sqlalchemy import and_, exists, func, select + +from scribe.models import async_session +from scribe.models.base import iso +from scribe.models.note_usage import PULLED, NoteUsageEvent +from scribe.models.retrieval_judgment import VERDICTS, RetrievalJudgment +from scribe.models.retrieval_log import RetrievalLog + +logger = logging.getLogger(__name__) + +# The arms a review can re-run FAITHFULLY. A re-run with different parameters +# from the arm's would judge a menu nobody was shown, so a surface joins this +# only with its own search spelled out in `_rerun`. +REVIEWABLE = ("auto_inject",) + +DEPTH_DEFAULT = 6 +DEPTH_MAX = 10 +SAMPLE_MAX = 20 +# How long after a call a pull still counts as "opened after". The server has +# no session identity (NoteUsageEvent), so this is a window, not a join. +OPENED_WINDOW = timedelta(hours=1) + +HOW_TO_JUDGE = ( + "Read each line as the reader did: the query, then the line's name, kind " + "and passage. Judge from those — do NOT open the record, which counts as a " + "pull and would credit the arm with an open it never earned. Verdicts: " + "`on_point` — someone in the query's situation should read this; " + "`adjacent` — same area, but it would not change what they do; " + "`unrelated` — a near neighbour in words only. Every verdict needs a " + "`reason` naming what in the query and the passage decided it. Judge the " + "lines past the budget too (`within_budget: false`): they are the " + "candidates a larger budget would add." +) + + +def _check_source(source: str) -> str: + if source not in REVIEWABLE: + raise ValueError(f"source must be one of {list(REVIEWABLE)}, got {source!r}") + return source + + +async def _rerun(user_id: int, row: RetrievalLog, depth: int) -> list[dict]: + """The logged call's query, searched again the way its arm searches. + + `auto_inject`'s parameters, from `build_autoinject_hint`: browse scope, + lessons reachable across projects, the arm's floor as logged. The reserved + slots (reuse, lesson) are left out — each logs its own source, and is + judged as that source when it becomes reviewable. + """ + # Imported here: plugin_context imports retrieval_telemetry, which reads + # this module's judged block. + from scribe.services.embeddings import document_title, semantic_search_notes + from scribe.services.plugin_context import _menu_name, _menu_passage, _record_kind + + rep: dict = {} + hits = await semantic_search_notes( + user_id, row.query or "", + limit=depth, + threshold=row.threshold if row.threshold is not None else 0.0, + project_id=row.project_id or None, + include_global_kinds=True, + scope="browse", + report=rep, + ) + chunks = rep.get("best_chunk") or {} + budget = int(row.limit_n or 0) + out = [] + for rank, (score, note) in enumerate(hits, start=1): + name = _menu_name(note.title, note.note_type, note.data, note.body) + out.append({ + "rank": rank, + "record_id": int(note.id), + "score": round(float(score), 4), + "within_budget": rank <= budget, + "kind": _record_kind(note), + "name": name, + "passage": _menu_passage( + document_title(note.title, note.note_type, note.data, note.body), + (chunks.get(int(note.id)) or {}).get("text"), name, + ) or None, + }) + return out + + +async def _opened_after(user_id: int, row: RetrievalLog, ids: list[int]) -> set[int]: + """Which of `ids` the user's AGENT pulled within OPENED_WINDOW of the call. + + Agent pulls only — the `mcp_` prefix, per NoteUsageEvent's convention — + because the question is whether the line moved the session, not whether + the operator later browsed to the record. + """ + if not ids or row.created_at is None: + return set() + async with async_session() as session: + found = await session.execute( + select(NoteUsageEvent.note_id).distinct().where( + NoteUsageEvent.user_id == user_id, + NoteUsageEvent.event == PULLED, + NoteUsageEvent.source.startswith("mcp_", autoescape=True), + NoteUsageEvent.note_id.in_(ids), + NoteUsageEvent.created_at >= row.created_at, + NoteUsageEvent.created_at <= row.created_at + OPENED_WINDOW, + ) + ) + return {int(i) for (i,) in found.all()} + + +def _eligible(user_id: int, source: str, since): + """Logged calls of `source` that offered something fresh and that this + reviewer has not judged a line of.""" + judged = exists().where(and_( + RetrievalJudgment.retrieval_log_id == RetrievalLog.id, + RetrievalJudgment.user_id == user_id, + )) + return ( + RetrievalLog.user_id == user_id, + RetrievalLog.source == source, + RetrievalLog.created_at >= since, + RetrievalLog.result_count > 0, + RetrievalLog.query.is_not(None), + ~judged, + ) + + +async def menus_to_review( + user_id: int, *, source: str = "auto_inject", n: int = 5, days: int = 14, + depth: int = DEPTH_DEFAULT, +) -> dict: + """A random sample of unjudged logged calls, each re-run for judging. + + Random, not newest: a review of the latest calls judges one afternoon's + work, and the readout would then describe that afternoon. + """ + _check_source(source) + n = max(1, min(int(n), SAMPLE_MAX)) + depth = max(1, min(int(depth), DEPTH_MAX)) + since = datetime.now(timezone.utc) - timedelta(days=max(1, int(days))) + where = _eligible(user_id, source, since) + async with async_session() as session: + rows = (await session.execute( + select(RetrievalLog).where(*where).order_by(func.random()).limit(n) + )).scalars().all() + remaining = (await session.execute( + select(func.count()).select_from(RetrievalLog).where(*where) + )).scalar_one() + + menus = [] + for row in rows: + lines = await _rerun(user_id, row, depth) + opened = await _opened_after(user_id, row, [ln["record_id"] for ln in lines]) + for ln in lines: + ln["opened_after"] = ln["record_id"] in opened + logged = [int(it["id"]) for it in (row.result_ids or [])] + rerun_ids = {ln["record_id"] for ln in lines} + menus.append({ + "log_id": int(row.id), + "created_at": iso(row.created_at), + "query": row.query, + "budget": row.limit_n, + "threshold": row.threshold, + "logged_ids": logged, + # Logged then, not found now: deleted, edited out of reach, or + # outranked by records written since. Judge what is here. + "missing": [i for i in logged if i not in rerun_ids], + "lines": lines, + }) + return { + "source": source, + "menus": menus, + "remaining": int(remaining), + "verdicts": list(VERDICTS), + "how_to_judge": HOW_TO_JUDGE, + } + + +async def judge_menu( + user_id: int, *, log_id: int, verdicts: list[dict], +) -> dict: + """Record a verdict, with its reason, on lines of one logged call. + + Rank, score and budget side are re-derived from a fresh re-run rather than + taken from the caller, so a verdict cannot be filed against a position the + record does not hold. Judging a line again replaces the earlier verdict. + """ + if not verdicts: + raise ValueError("pass at least one verdict") + parsed: list[tuple[int, str, str]] = [] + for v in verdicts: + rid = int(v.get("record_id") or 0) + verdict = str(v.get("verdict") or "").strip().lower() + reason = str(v.get("reason") or "").strip() + if verdict not in VERDICTS: + raise ValueError(f"verdict must be one of {list(VERDICTS)}, got {verdict!r}") + if not reason: + raise ValueError( + f"record {rid}: a verdict needs its reason — what in the query " + f"and the passage decided it. Without one it cannot be re-read." + ) + parsed.append((rid, verdict, reason)) + + async with async_session() as session: + row = await session.get(RetrievalLog, int(log_id)) + # Another user's call reads exactly like a missing one: its query is theirs. + if row is None or row.user_id != user_id: + raise ValueError(f"logged call {log_id} not found") + _check_source(row.source) + + lines = {ln["record_id"]: ln for ln in await _rerun(user_id, row, DEPTH_MAX)} + unknown = [rid for rid, _v, _r in parsed if rid not in lines] + if unknown: + raise ValueError( + f"records {unknown} are not in this call's re-run to depth " + f"{DEPTH_MAX} — judge the lines menus_to_review returned" + ) + opened = await _opened_after(user_id, row, [rid for rid, _v, _r in parsed]) + + async with async_session() as session: + existing = { + int(j.record_id): j + for j in (await session.execute( + select(RetrievalJudgment).where( + RetrievalJudgment.retrieval_log_id == row.id, + RetrievalJudgment.user_id == user_id, + RetrievalJudgment.record_id.in_([rid for rid, _v, _r in parsed]), + ) + )).scalars().all() + } + for rid, verdict, reason in parsed: + ln = lines[rid] + j = existing.get(rid) or RetrievalJudgment( + user_id=user_id, retrieval_log_id=int(row.id), record_id=rid, + ) + j.source = row.source + j.query = row.query + j.rank = ln["rank"] + j.score = ln["score"] + j.within_budget = ln["within_budget"] + j.verdict = verdict + j.reason = reason + j.opened_after = rid in opened + session.add(j) + await session.commit() + return { + "log_id": int(row.id), + "recorded": len(parsed), + "replaced": len(existing), + } + + +async def judged_block(user_id: int | None, since) -> dict: + """`retrieval_telemetry`'s judged block: verdicts by rank, per source. + + Counts, not rates, beside each other per rank — the sample is small by + nature and a percentage of seven reads as more than it is. `on_point_unopened` + is the number the open rate could never show: lines that related and were + left closed, which is the passage doing its job or the reader missing it. + + Raises on a failed read; the caller guards it. + """ + async with async_session() as session: + rows = (await session.execute( + select( + RetrievalJudgment.source, + RetrievalJudgment.rank, + RetrievalJudgment.within_budget, + RetrievalJudgment.verdict, + RetrievalJudgment.opened_after, + func.count().label("n"), + ) + .where( + RetrievalJudgment.user_id == user_id, + RetrievalJudgment.created_at >= since, + ) + .group_by( + RetrievalJudgment.source, RetrievalJudgment.rank, + RetrievalJudgment.within_budget, RetrievalJudgment.verdict, + RetrievalJudgment.opened_after, + ) + )).all() + calls = dict((await session.execute( + select( + RetrievalJudgment.source, + func.count(func.distinct(RetrievalJudgment.retrieval_log_id)), + ) + .where( + RetrievalJudgment.user_id == user_id, + RetrievalJudgment.created_at >= since, + ) + .group_by(RetrievalJudgment.source) + )).all()) + + def empty() -> dict: + return {v: 0 for v in VERDICTS} | {"opened": 0} + + out: dict = {} + for source, rank, within, verdict, opened, n in rows: + block = out.setdefault(source, { + "judged_calls": int(calls.get(source, 0)), + "judged_lines": 0, + "within_budget": empty(), + "beyond_budget": empty(), + "by_rank": {}, + "on_point_unopened": 0, + }) + n = int(n) + block["judged_lines"] += n + side = block["within_budget" if within else "beyond_budget"] + at = block["by_rank"].setdefault(int(rank), empty()) + for bucket in (side, at): + if verdict in VERDICTS: + bucket[verdict] += n + if opened: + bucket["opened"] += n + if verdict == VERDICTS[0] and not opened: + block["on_point_unopened"] += n + for block in out.values(): + block["by_rank"] = [ + {"rank": r, **c} for r, c in sorted(block["by_rank"].items()) + ] + return out diff --git a/src/scribe/services/retrieval_telemetry.py b/src/scribe/services/retrieval_telemetry.py index 2e1f8bed..60d269d3 100644 --- a/src/scribe/services/retrieval_telemetry.py +++ b/src/scribe/services/retrieval_telemetry.py @@ -40,6 +40,7 @@ from scribe.models.system_usage import SURFACED as SYSTEM_SURFACED from scribe.models.system_usage import SystemUsageEvent from scribe.services.rule_usage import is_ambient from scribe.models.retrieval_log import RetrievalLog +from scribe.services.retrieval_review import judged_block from scribe.services.retrieval_registry import ( POINTS, UNBIDDEN, get_point, is_registered, sources_expected_to_emit, ) @@ -745,9 +746,11 @@ def _compute_warnings(sources: dict, usage: dict, rule_usage: dict, # ── Surfaced and never pulled, for each corpus ─────────────────────── # # The one corpus-level check, and the only number here that judges the - # RECORDS rather than the arms. A record shown repeatedly and never opened - # is either badly titled or genuinely irrelevant, and both are actionable - # in a way "pull-through is 0.15" is not. + # RECORDS rather than the arms. It used to say a record shown repeatedly + # and never opened is "badly titled or genuinely irrelevant" — true when a + # line was a title, and not since it carries the passage (#4364): the + # record may have done its job unopened. So it now points at the judged + # sample rather than at a cause (#4772). # # Distinct records, not events: a note surfaced forty times and never # opened is one problem, not forty. @@ -764,9 +767,11 @@ def _compute_warnings(sources: dict, usage: dict, rule_usage: dict, out.append(_warn( "surfaced_never_pulled", f"{never} of {shown} distinct {label} were surfaced in this " - f"window and never opened. Read the titles before the " - f"threshold: a record nobody opens is usually one whose title " - f"does not say when it matters.", + f"window and never opened. That is not a verdict on them: a " + f"line carries its matched passage, so an unopened record may " + f"have been unrelated, enough as shown, or already in context. " + f"Before acting on it — a title, a floor, a budget — judge a " + f"sample with `menus_to_review` and read the `judged` block.", source=None, corpus=label, surfaced=int(shown), pulled=int(pulled or 0), never_pulled=never, )) @@ -995,6 +1000,7 @@ async def retrieval_summary( "usage": {}, "rule_usage": {}, "system_usage": {}, + "judged": {}, "read_failed": False, } @@ -1554,6 +1560,18 @@ async def retrieval_summary( out["rule_usage"] = rule_usage out["system_usage"] = await _system_usage(user_id, since) + # ── Judged lines (#4772) ───────────────────────────────────────────── + # + # The relevance readout the usage blocks cannot be: a reviewer's verdict + # on each line of a sample of logged menus, by rank. Its own guard, for + # `_system_usage`'s reason — a failure flags itself rather than printing + # an empty block, which reads as "nothing judged" (#2663). + try: + out["judged"] = await judged_block(user_id, since) + except Exception: + logger.warning("judged block read failed", exc_info=True) + out["judged"] = {"judged_failed": True} + # ── What is wrong (#3431) ──────────────────────────────────────────── # # Computed LAST, over the blocks above rather than over the database, so diff --git a/tests/test_retrieval_review.py b/tests/test_retrieval_review.py new file mode 100644 index 00000000..c952a21f --- /dev/null +++ b/tests/test_retrieval_review.py @@ -0,0 +1,250 @@ +"""The review pass: judging whether an injected line RELATED (#4772). + +Unit tests pin the refusals and the re-run's fidelity to the arm it replays. +Integration tests pin the parts a mock would make true by construction: which +logged calls are offered, that a verdict lands and replaces, that "opened +after" reads the usage table, and that the readout counts by rank. +""" +import ast +import inspect +import textwrap +from datetime import datetime, timedelta, timezone +from types import SimpleNamespace +from unittest.mock import AsyncMock, patch + +import pytest +import pytest_asyncio + +from scribe.models import async_session +from scribe.services import retrieval_review as review +from tests.helpers import ensure_user, fake_note + + +# ── refusals ───────────────────────────────────────────────────────────── + + +@pytest.mark.asyncio +async def test_only_an_arm_that_can_be_re_run_exactly_is_offered(): + with pytest.raises(ValueError, match="source must be one of"): + await review.menus_to_review(1, source="write_path") + + +@pytest.mark.asyncio +async def test_a_verdict_outside_the_vocabulary_is_refused(): + with pytest.raises(ValueError, match="verdict must be one of"): + await review.judge_menu(1, log_id=1, verdicts=[ + {"record_id": 5, "verdict": "relevant", "reason": "x"}, + ]) + + +@pytest.mark.asyncio +async def test_a_verdict_without_its_reason_is_refused(): + with pytest.raises(ValueError, match="needs its reason"): + await review.judge_menu(1, log_id=1, verdicts=[ + {"record_id": 5, "verdict": "unrelated", "reason": " "}, + ]) + + +@pytest.mark.asyncio +async def test_an_empty_verdict_list_is_refused(): + with pytest.raises(ValueError, match="at least one"): + await review.judge_menu(1, log_id=1, verdicts=[]) + + +# ── the re-run replays the arm ─────────────────────────────────────────── + + +def _search_kwargs(fn) -> dict: + """The constant keywords `fn` passes to semantic_search_notes.""" + tree = ast.parse(textwrap.dedent(inspect.getsource(fn))) + for node in ast.walk(tree): + if (isinstance(node, ast.Call) and getattr(node.func, "id", "") + == "semantic_search_notes"): + return { + kw.arg: kw.value.value for kw in node.keywords + if isinstance(kw.value, ast.Constant) + } + return {} + + +def test_the_re_run_searches_the_way_the_arm_does(): + """A re-run with different visibility or kinds from the arm's would judge + a menu nobody was shown. Read from both call sites, so a change to the + arm's search fails here until the review follows it.""" + from scribe.services.plugin_context import build_autoinject_hint + + arm = _search_kwargs(build_autoinject_hint) + rerun = _search_kwargs(review._rerun) + assert arm, "found no semantic_search_notes call in the arm" + for key in ("include_global_kinds", "scope"): + assert key in arm, f"the arm no longer passes {key} as a constant" + assert rerun.get(key) == arm[key], key + + +@pytest.mark.asyncio +async def test_the_re_run_ranks_from_one_and_marks_the_budget_cut(): + row = SimpleNamespace(query="q", threshold=0.55, limit_n=2, project_id=None) + hits = [ + (0.81, fake_note(id=11, title="first")), + (0.74, fake_note(id=12, title="second")), + (0.66, fake_note(id=13, title="third")), + ] + + async def search(user_id, query, **kw): + assert kw["limit"] == 3 and kw["threshold"] == 0.55 + assert kw["project_id"] is None + kw["report"]["best_chunk"] = {12: {"text": "second\nthe passage that matched"}} + return hits + + with patch("scribe.services.embeddings.semantic_search_notes", side_effect=search): + lines = await review._rerun(1, row, 3) + assert [(ln["rank"], ln["record_id"], ln["within_budget"]) for ln in lines] == [ + (1, 11, True), (2, 12, True), (3, 13, False), + ] + assert lines[1]["passage"] == "the passage that matched" + assert lines[0]["passage"] is None + + +# ── against Postgres ───────────────────────────────────────────────────── + + +@pytest_asyncio.fixture +async def logged_calls(): + """A reviewer's logged calls: two reviewable, one that offered nothing, + one of another arm — and another user's reviewable call.""" + from scribe.models.retrieval_log import RetrievalLog + + async with async_session() as s: + me = await ensure_user(s, "review_owner") + other = await ensure_user(s, "review_other") + now = datetime.now(timezone.utc) + + def log(uid, source="auto_inject", count=2, query="how is the cache invalidated"): + return RetrievalLog( + user_id=uid, source=source, query=query, threshold=0.55, + limit_n=3, result_count=count, created_at=now - timedelta(hours=2), + result_ids=[{"id": 501, "score": 0.8, "rank": 0}], + ) + + rows = { + "a": log(me.id), "b": log(me.id), + "empty": log(me.id, count=0), + "other_arm": log(me.id, source="write_path"), + "theirs": log(other.id), + } + s.add_all(rows.values()) + await s.commit() + ids = {k: int(v.id) for k, v in rows.items()} + ids |= {"me": me.id, "other": other.id} + yield ids + + from sqlalchemy import delete + from scribe.models.note_usage import NoteUsageEvent + from scribe.models.retrieval_judgment import RetrievalJudgment + async with async_session() as s: + logs = [ids[k] for k in ("a", "b", "empty", "other_arm", "theirs")] + await s.execute(delete(RetrievalJudgment).where( + RetrievalJudgment.retrieval_log_id.in_(logs))) + await s.execute(delete(RetrievalLog).where(RetrievalLog.id.in_(logs))) + await s.execute(delete(NoteUsageEvent).where( + NoteUsageEvent.user_id == ids["me"], NoteUsageEvent.note_id.in_([501, 502]))) + await s.commit() + + +def _lines(*specs): + """A stand-in re-run: (record_id, rank, within_budget).""" + return AsyncMock(return_value=[ + {"record_id": rid, "rank": rank, "score": 0.9 - rank / 10, + "within_budget": within, "kind": "note", "name": f"n{rid}", "passage": "p"} + for rid, rank, within in specs + ]) + + +@pytest.mark.integration +@pytest.mark.asyncio +@pytest.mark.usefixtures("_dispose_engine") +async def test_the_sample_offers_only_the_reviewers_unjudged_reviewable_calls(logged_calls): + ids = logged_calls + with patch.object(review, "_rerun", _lines((501, 1, True))): + out = await review.menus_to_review(ids["me"], n=20, days=1) + assert {m["log_id"] for m in out["menus"]} == {ids["a"], ids["b"]} + assert out["remaining"] == 2 + + await review.judge_menu(ids["me"], log_id=ids["a"], verdicts=[ + {"record_id": 501, "verdict": "on_point", "reason": "names the cache"}, + ]) + out = await review.menus_to_review(ids["me"], n=20, days=1) + assert {m["log_id"] for m in out["menus"]} == {ids["b"]} + + +@pytest.mark.integration +@pytest.mark.asyncio +@pytest.mark.usefixtures("_dispose_engine") +async def test_another_users_call_cannot_be_judged(logged_calls): + ids = logged_calls + with patch.object(review, "_rerun", _lines((501, 1, True))), \ + pytest.raises(ValueError, match="not found"): + await review.judge_menu(ids["me"], log_id=ids["theirs"], verdicts=[ + {"record_id": 501, "verdict": "unrelated", "reason": "r"}, + ]) + + +@pytest.mark.integration +@pytest.mark.asyncio +@pytest.mark.usefixtures("_dispose_engine") +async def test_a_record_the_re_run_does_not_hold_is_refused(logged_calls): + ids = logged_calls + with patch.object(review, "_rerun", _lines((501, 1, True))), \ + pytest.raises(ValueError, match=r"\[999\]"): + await review.judge_menu(ids["me"], log_id=ids["a"], verdicts=[ + {"record_id": 999, "verdict": "unrelated", "reason": "r"}, + ]) + + +@pytest.mark.integration +@pytest.mark.asyncio +@pytest.mark.usefixtures("_dispose_engine") +async def test_verdicts_land_replace_and_read_back_by_rank(logged_calls): + from scribe.models.note_usage import NoteUsageEvent + from scribe.services.retrieval_telemetry import retrieval_summary + + ids = logged_calls + # The reviewer's agent opened 501 an hour inside the call; nothing opened 502. + async with async_session() as s: + s.add(NoteUsageEvent( + user_id=ids["me"], note_id=501, event="pulled", source="mcp_get_note", + created_at=datetime.now(timezone.utc) - timedelta(hours=1, minutes=30), + )) + await s.commit() + + rerun = _lines((501, 1, True), (502, 4, False)) + with patch.object(review, "_rerun", rerun): + await review.judge_menu(ids["me"], log_id=ids["a"], verdicts=[ + {"record_id": 501, "verdict": "adjacent", "reason": "first look"}, + {"record_id": 502, "verdict": "on_point", "reason": "the exact fix"}, + ]) + again = await review.judge_menu(ids["me"], log_id=ids["a"], verdicts=[ + {"record_id": 501, "verdict": "on_point", "reason": "on reflection"}, + ]) + assert again == {"log_id": ids["a"], "recorded": 1, "replaced": 1} + + block = (await retrieval_summary(ids["me"], days=1))["judged"] + assert "judged_failed" not in block, "the readout did not execute" + ai = block["auto_inject"] + assert ai["judged_calls"] == 1 and ai["judged_lines"] == 2 + assert ai["within_budget"] == { + "on_point": 1, "adjacent": 0, "unrelated": 0, "opened": 1} + assert ai["beyond_budget"] == { + "on_point": 1, "adjacent": 0, "unrelated": 0, "opened": 0} + assert [r["rank"] for r in ai["by_rank"]] == [1, 4] + # Related and left closed: the number an open rate could never show. + assert ai["on_point_unopened"] == 1 + + +@pytest.mark.integration +@pytest.mark.asyncio +@pytest.mark.usefixtures("_dispose_engine") +async def test_a_fresh_install_reads_an_empty_judged_block_not_a_failure(): + from scribe.services.retrieval_telemetry import retrieval_summary + + assert (await retrieval_summary(990041, days=30))["judged"] == {}