feat(retrieval): a review pass judges whether injected lines related — menus_to_review, judge_menu and a judged readout (#4772)
CI & Build / Python lint (push) Successful in 4s
CI & Build / Plugin hooks (push) Successful in 14s
CI & Build / TypeScript typecheck (push) Successful in 54s
CI & Build / integration (push) Successful in 1m1s
CI & Build / Python tests (push) Successful in 1m55s
CI & Build / Build & push image (push) Successful in 44s
CI & Build / Python lint (push) Successful in 4s
CI & Build / Plugin hooks (push) Successful in 14s
CI & Build / TypeScript typecheck (push) Successful in 54s
CI & Build / integration (push) Successful in 1m1s
CI & Build / Python tests (push) Successful in 1m55s
CI & Build / Build & push image (push) Successful in 44s
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 <noreply@anthropic.com>
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user