Merge pull request 'A review pass judges whether injected lines related (#4772)' (#196) from dev into main
CI & Build / Python lint (push) Successful in 4s
CI & Build / Plugin hooks (push) Successful in 13s
CI & Build / TypeScript typecheck (push) Successful in 57s
CI & Build / integration (push) Successful in 1m19s
CI & Build / Python tests (push) Successful in 2m6s
CI & Build / Build & push image (push) Successful in 19s

This commit was merged in pull request #196.
This commit is contained in:
2026-10-03 18:17:49 -04:00
13 changed files with 879 additions and 11 deletions
@@ -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")
+1 -1
View File
@@ -1,7 +1,7 @@
{ {
"name": "scribe", "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).", "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": { "author": {
"name": "Bryan Van Deusen" "name": "Bryan Van Deusen"
}, },
@@ -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 `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. 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, 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 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 install: a rule granting a routine push scored 0.6515 and ranked 5th for the
+7
View File
@@ -133,6 +133,10 @@ _READ_ONLY_TOOLS = frozenset({
# the prefixes the completeness test derives from, so nothing would have # the prefixes the completeness test derives from, so nothing would have
# prompted this decision. # prompted this decision.
"notes_due_for_verification", "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 # 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 # 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 # 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). # reads, and it appends the reason to the audit trail (#4102).
"tune_retrieval", "tune_retrieval",
"migrate_retrieval_floor", "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 # trash
"restore", "purge_trash", "restore", "purge_trash",
}) })
+2 -1
View File
@@ -6,7 +6,7 @@ from `mcp.server.build_mcp_server`.
""" """
from scribe.mcp.tools import ( from scribe.mcp.tools import (
design_systems, lessons, milestones, notes, processes, projects, recent, repos, design_systems, lessons, milestones, notes, processes, projects, recent, repos,
retrieval_tuning, retrieval_review, retrieval_tuning,
wide_net, wide_net,
rulebooks, search, shapes, snippets, systems, tags, tasks, trash, 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.""" """Register every tool module's tools on the given MCPServer instance."""
search.register(mcp) search.register(mcp)
retrieval_tuning.register(mcp) retrieval_tuning.register(mcp)
retrieval_review.register(mcp)
wide_net.register(mcp) wide_net.register(mcp)
notes.register(mcp) notes.register(mcp)
tasks.register(mcp) tasks.register(mcp)
+74
View File
@@ -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)
+17 -3
View File
@@ -380,7 +380,7 @@ async def retrieval_telemetry(
— so `near_miss_samples=5` and opening the ids it returns is the step that — 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. 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`, `sources` — per retrieval surface (`auto_inject`, `write_path`,
`mcp_search`, …), from `retrieval_logs`: `calls`, `zero_result_calls`, `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. 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. `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 EVERY COUNTER BLOCK CARRIES ITS OWN COVERAGE — `complete_from` and
`covers_window`. `complete_from` is when the number became trustworthy: `covers_window`. `complete_from` is when the number became trustworthy:
for one source, its first recorded row; for a section that sums several, 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 - `no_duration` — rows written without timings. A logging gap, not a slow
arm, and it devalues every other number from that source. arm, and it devalues every other number from that source.
- `surfaced_never_pulled` — distinct records shown and never opened, per - `surfaced_never_pulled` — distinct records shown and never opened, per
corpus. Read their titles before touching a threshold: a record nobody corpus. Not a verdict on them: a line carries its matched passage, so
opens is usually one whose title does not say when it matters. 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 - `read_and_unacted` — distinct rules OPENED in the window that recorded
no outcome, against the ones that did. The failure milestone 419 was no outcome, against the ones that did. The failure milestone 419 was
opened on, and the worse sibling of `surfaced_never_pulled` above: a opened on, and the worse sibling of `surfaced_never_pulled` above: a
+1
View File
@@ -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.note_usage import NoteUsageEvent # noqa: E402, F401
from scribe.models.rule_usage import RuleUsageEvent # 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.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.project import Project # noqa: E402, F401
from scribe.models.milestone import Milestone # noqa: E402, F401 from scribe.models.milestone import Milestone # noqa: E402, F401
from scribe.models.task_log import TaskLog # noqa: E402, F401 from scribe.models.task_log import TaskLog # noqa: E402, F401
+82
View File
@@ -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,
}
+4
View File
@@ -159,6 +159,10 @@ _NOT_INCLUDED = [
"app_logs", "notifications", "app_logs", "notifications",
"invitation_tokens", "password_reset_tokens", "user_profiles", "invitation_tokens", "password_reset_tokens", "user_profiles",
"retrieval_logs", "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 # Sensitive credentials, same reasoning as api_keys: a backup that carries
# forge tokens is a token-exfiltration file. Users re-add connections # forge tokens is a token-exfiltration file. Users re-add connections
# after a restore; the per-project pin (projects.forge_connection_id) is # after a restore; the per-project pin (projects.forge_connection_id) is
+352
View File
@@ -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
+24 -6
View File
@@ -40,6 +40,7 @@ from scribe.models.system_usage import SURFACED as SYSTEM_SURFACED
from scribe.models.system_usage import SystemUsageEvent from scribe.models.system_usage import SystemUsageEvent
from scribe.services.rule_usage import is_ambient from scribe.services.rule_usage import is_ambient
from scribe.models.retrieval_log import RetrievalLog from scribe.models.retrieval_log import RetrievalLog
from scribe.services.retrieval_review import judged_block
from scribe.services.retrieval_registry import ( from scribe.services.retrieval_registry import (
POINTS, UNBIDDEN, get_point, is_registered, sources_expected_to_emit, 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 ─────────────────────── # ── Surfaced and never pulled, for each corpus ───────────────────────
# #
# The one corpus-level check, and the only number here that judges the # 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 # RECORDS rather than the arms. It used to say a record shown repeatedly
# is either badly titled or genuinely irrelevant, and both are actionable # and never opened is "badly titled or genuinely irrelevant" — true when a
# in a way "pull-through is 0.15" is not. # 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 # Distinct records, not events: a note surfaced forty times and never
# opened is one problem, not forty. # opened is one problem, not forty.
@@ -764,9 +767,11 @@ def _compute_warnings(sources: dict, usage: dict, rule_usage: dict,
out.append(_warn( out.append(_warn(
"surfaced_never_pulled", "surfaced_never_pulled",
f"{never} of {shown} distinct {label} were surfaced in this " f"{never} of {shown} distinct {label} were surfaced in this "
f"window and never opened. Read the titles before the " f"window and never opened. That is not a verdict on them: a "
f"threshold: a record nobody opens is usually one whose title " f"line carries its matched passage, so an unopened record may "
f"does not say when it matters.", 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, source=None, corpus=label,
surfaced=int(shown), pulled=int(pulled or 0), never_pulled=never, surfaced=int(shown), pulled=int(pulled or 0), never_pulled=never,
)) ))
@@ -995,6 +1000,7 @@ async def retrieval_summary(
"usage": {}, "usage": {},
"rule_usage": {}, "rule_usage": {},
"system_usage": {}, "system_usage": {},
"judged": {},
"read_failed": False, "read_failed": False,
} }
@@ -1554,6 +1560,18 @@ async def retrieval_summary(
out["rule_usage"] = rule_usage out["rule_usage"] = rule_usage
out["system_usage"] = await _system_usage(user_id, since) 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) ──────────────────────────────────────────── # ── What is wrong (#3431) ────────────────────────────────────────────
# #
# Computed LAST, over the blocks above rather than over the database, so # Computed LAST, over the blocks above rather than over the database, so
+250
View File
@@ -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"] == {}