diff --git a/src/scribe/mcp/tools/retrieval_review.py b/src/scribe/mcp/tools/retrieval_review.py index 92363473..f07d7617 100644 --- a/src/scribe/mcp/tools/retrieval_review.py +++ b/src/scribe/mcp/tools/retrieval_review.py @@ -29,6 +29,13 @@ async def menus_to_review( `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. + The re-run is against today's corpus, so it guards what has been written + since. Records created after the call are left out and named in + `postdated`: the call could not have offered them, and the session that + made it often wrote them, quoting it. A line marked `changed_since_call` + was written to afterwards; if its passage is plainly the later text, leave + it unjudged. `logged` says the line was shown fresh at the time. + 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. diff --git a/src/scribe/services/retrieval_review.py b/src/scribe/services/retrieval_review.py index dbe5c490..29e498ac 100644 --- a/src/scribe/services/retrieval_review.py +++ b/src/scribe/services/retrieval_review.py @@ -20,6 +20,15 @@ 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. +WHAT TODAY'S CORPUS GETS WRONG. A record written AFTER the call could not have +been offered to it, yet it outranks everything: the session that made the call +goes on to write the decision, the fix, the work log — often quoting the very +message. Left in, it reads as the arm's best hit and inflates `on_point` +exactly where the budget question is decided. So the re-run drops records +created after the call (`postdated` names them), and marks a line +`changed_since_call` when the record or one of its work logs was written after +it — its passage may be text the reader never saw. + 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. @@ -48,6 +57,10 @@ REVIEWABLE = ("auto_inject",) DEPTH_DEFAULT = 6 DEPTH_MAX = 10 SAMPLE_MAX = 20 +# Extra candidates fetched so that dropping post-dated records still leaves +# `depth` lines. Records written after a call cluster at the top of its +# re-run, so the slack has to cover a session's worth of them. +POSTDATED_SLACK = 15 # 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) @@ -61,7 +74,10 @@ HOW_TO_JUDGE = ( "`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." + "candidates a larger budget would add. A line marked `changed_since_call` " + "was written to after the call: if its passage is plainly the later text " + "(dated after the call, or about what the call led to), leave it unjudged " + "— the reader never saw it." ) @@ -71,13 +87,22 @@ def _check_source(source: str) -> str: return source -async def _rerun(user_id: int, row: RetrievalLog, depth: int) -> list[dict]: +def _after(stamp, call) -> bool: + return stamp is not None and call is not None and stamp > call + + +async def _rerun( + user_id: int, row: RetrievalLog, depth: int, report: dict | None = None, +) -> 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. + + Records created after the call are dropped before ranking — the call could + not have been offered them — and named in `report["postdated"]`. """ # Imported here: plugin_context imports retrieval_telemetry, which reads # this module's judged block. @@ -87,13 +112,17 @@ async def _rerun(user_id: int, row: RetrievalLog, depth: int) -> list[dict]: rep: dict = {} hits = await semantic_search_notes( user_id, row.query or "", - limit=depth, + limit=depth + POSTDATED_SLACK, 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, ) + postdated = [int(n.id) for _s, n in hits if _after(n.created_at, row.created_at)] + if report is not None: + report["postdated"] = postdated + hits = [h for h in hits if int(h[1].id) not in postdated][:depth] chunks = rep.get("best_chunk") or {} budget = int(row.limit_n or 0) out = [] @@ -104,6 +133,8 @@ async def _rerun(user_id: int, row: RetrievalLog, depth: int) -> list[dict]: "record_id": int(note.id), "score": round(float(score), 4), "within_budget": rank <= budget, + # Work logs are folded in by menus_to_review, which has a session. + "changed_since_call": _after(note.updated_at, row.created_at), "kind": _record_kind(note), "name": name, "passage": _menu_passage( @@ -137,6 +168,23 @@ async def _opened_after(user_id: int, row: RetrievalLog, ids: list[int]) -> set[ return {int(i) for (i,) in found.all()} +async def _logged_after(row: RetrievalLog, ids: list[int]) -> set[int]: + """Which of `ids` gained a work log after the call. A log lands in its + task's embedded document without always touching the task's `updated_at` + (a closed task takes no claim stamp), so it is read here.""" + from scribe.models.task_log import TaskLog + + if not ids or row.created_at is None: + return set() + async with async_session() as session: + found = await session.execute( + select(TaskLog.task_id).distinct().where( + TaskLog.task_id.in_(ids), TaskLog.created_at > row.created_at, + ) + ) + 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.""" @@ -178,12 +226,19 @@ async def menus_to_review( 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]) + rep: dict = {} + lines = await _rerun(user_id, row, depth, report=rep) + ids = [ln["record_id"] for ln in lines] + opened = await _opened_after(user_id, row, ids) + relogged = await _logged_after(row, ids) + logged = [int(it["id"]) for it in (row.result_ids or [])] 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} + ln["changed_since_call"] = ln["changed_since_call"] or ln["record_id"] in relogged + # Shown fresh then. Inside the budget and not logged means a + # repeat the arm withheld, or a ranking the corpus has since moved. + ln["logged"] = ln["record_id"] in logged + rerun_ids = set(ids) menus.append({ "log_id": int(row.id), "created_at": iso(row.created_at), @@ -192,8 +247,10 @@ async def menus_to_review( "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. + # outranked since. Judge what is here. "missing": [i for i in logged if i not in rerun_ids], + # Written after the call, so never offered to it; left out. + "postdated": rep.get("postdated", []), "lines": lines, }) return { diff --git a/tests/test_retrieval_review.py b/tests/test_retrieval_review.py index c952a21f..8c6388f9 100644 --- a/tests/test_retrieval_review.py +++ b/tests/test_retrieval_review.py @@ -81,17 +81,22 @@ def test_the_re_run_searches_the_way_the_arm_does(): assert rerun.get(key) == arm[key], key +CALL = datetime(2026, 9, 20, 12, 0, tzinfo=timezone.utc) +BEFORE = CALL - timedelta(days=1) + + @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) + row = SimpleNamespace(query="q", threshold=0.55, limit_n=2, project_id=None, + created_at=CALL) 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")), + (0.81, fake_note(id=11, title="first", created_at=BEFORE, updated_at=BEFORE)), + (0.74, fake_note(id=12, title="second", created_at=BEFORE, updated_at=BEFORE)), + (0.66, fake_note(id=13, title="third", created_at=BEFORE, updated_at=BEFORE)), ] async def search(user_id, query, **kw): - assert kw["limit"] == 3 and kw["threshold"] == 0.55 + assert kw["limit"] == 3 + review.POSTDATED_SLACK and kw["threshold"] == 0.55 assert kw["project_id"] is None kw["report"]["best_chunk"] = {12: {"text": "second\nthe passage that matched"}} return hits @@ -105,6 +110,31 @@ async def test_the_re_run_ranks_from_one_and_marks_the_budget_cut(): assert lines[0]["passage"] is None +@pytest.mark.asyncio +async def test_a_record_written_after_the_call_is_dropped_before_ranking(): + """The session that made a call goes on to write about it, and that record + outranks everything in a re-run. The call could never have been offered + it, so it must not take a rank — or a verdict — from what was.""" + row = SimpleNamespace(query="q", threshold=0.55, limit_n=1, project_id=None, + created_at=CALL) + after = CALL + timedelta(hours=1) + hits = [ + (0.90, fake_note(id=21, title="the decision", created_at=after, updated_at=after)), + (0.80, fake_note(id=22, title="held then", created_at=BEFORE, updated_at=BEFORE)), + (0.70, fake_note(id=23, title="edited since", created_at=BEFORE, updated_at=after)), + ] + + async def search(user_id, query, **kw): + return hits + + rep: dict = {} + with patch("scribe.services.embeddings.semantic_search_notes", side_effect=search): + lines = await review._rerun(1, row, 2, report=rep) + assert rep["postdated"] == [21] + assert [(ln["rank"], ln["record_id"], ln["within_budget"], ln["changed_since_call"]) + for ln in lines] == [(1, 22, True, False), (2, 23, False, True)] + + # ── against Postgres ───────────────────────────────────────────────────── @@ -155,7 +185,8 @@ 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"} + "within_budget": within, "changed_since_call": False, + "kind": "note", "name": f"n{rid}", "passage": "p"} for rid, rank, within in specs ])