Merge pull request 'The review drops records written after the call it re-runs (#4773)' (#197) from dev into main
CI & Build / Python lint (push) Successful in 5s
CI & Build / Plugin hooks (push) Successful in 18s
CI & Build / TypeScript typecheck (push) Successful in 55s
CI & Build / integration (push) Successful in 1m5s
CI & Build / Python tests (push) Successful in 1m59s
CI & Build / Build & push image (push) Successful in 23s

This commit was merged in pull request #197.
This commit is contained in:
2026-10-03 20:37:08 -04:00
3 changed files with 109 additions and 14 deletions
+7
View File
@@ -29,6 +29,13 @@ async def menus_to_review(
`missing` names ids the call logged that the re-run no longer finds — the `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. 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 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 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. with `judge_menu`. `how_to_judge` in the response repeats the vocabulary.
+65 -8
View File
@@ -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 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. 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 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 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. 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_DEFAULT = 6
DEPTH_MAX = 10 DEPTH_MAX = 10
SAMPLE_MAX = 20 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 # 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. # no session identity (NoteUsageEvent), so this is a window, not a join.
OPENED_WINDOW = timedelta(hours=1) OPENED_WINDOW = timedelta(hours=1)
@@ -61,7 +74,10 @@ HOW_TO_JUDGE = (
"`unrelated` — a near neighbour in words only. Every verdict needs a " "`unrelated` — a near neighbour in words only. Every verdict needs a "
"`reason` naming what in the query and the passage decided it. Judge the " "`reason` naming what in the query and the passage decided it. Judge the "
"lines past the budget too (`within_budget: false`): they are 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 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. """The logged call's query, searched again the way its arm searches.
`auto_inject`'s parameters, from `build_autoinject_hint`: browse scope, `auto_inject`'s parameters, from `build_autoinject_hint`: browse scope,
lessons reachable across projects, the arm's floor as logged. The reserved 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 slots (reuse, lesson) are left out — each logs its own source, and is
judged as that source when it becomes reviewable. 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 # Imported here: plugin_context imports retrieval_telemetry, which reads
# this module's judged block. # this module's judged block.
@@ -87,13 +112,17 @@ async def _rerun(user_id: int, row: RetrievalLog, depth: int) -> list[dict]:
rep: dict = {} rep: dict = {}
hits = await semantic_search_notes( hits = await semantic_search_notes(
user_id, row.query or "", user_id, row.query or "",
limit=depth, limit=depth + POSTDATED_SLACK,
threshold=row.threshold if row.threshold is not None else 0.0, threshold=row.threshold if row.threshold is not None else 0.0,
project_id=row.project_id or None, project_id=row.project_id or None,
include_global_kinds=True, include_global_kinds=True,
scope="browse", scope="browse",
report=rep, 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 {} chunks = rep.get("best_chunk") or {}
budget = int(row.limit_n or 0) budget = int(row.limit_n or 0)
out = [] out = []
@@ -104,6 +133,8 @@ async def _rerun(user_id: int, row: RetrievalLog, depth: int) -> list[dict]:
"record_id": int(note.id), "record_id": int(note.id),
"score": round(float(score), 4), "score": round(float(score), 4),
"within_budget": rank <= budget, "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), "kind": _record_kind(note),
"name": name, "name": name,
"passage": _menu_passage( "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()} 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): def _eligible(user_id: int, source: str, since):
"""Logged calls of `source` that offered something fresh and that this """Logged calls of `source` that offered something fresh and that this
reviewer has not judged a line of.""" reviewer has not judged a line of."""
@@ -178,12 +226,19 @@ async def menus_to_review(
menus = [] menus = []
for row in rows: for row in rows:
lines = await _rerun(user_id, row, depth) rep: dict = {}
opened = await _opened_after(user_id, row, [ln["record_id"] for ln in lines]) 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: for ln in lines:
ln["opened_after"] = ln["record_id"] in opened ln["opened_after"] = ln["record_id"] in opened
logged = [int(it["id"]) for it in (row.result_ids or [])] ln["changed_since_call"] = ln["changed_since_call"] or ln["record_id"] in relogged
rerun_ids = {ln["record_id"] for ln in lines} # 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({ menus.append({
"log_id": int(row.id), "log_id": int(row.id),
"created_at": iso(row.created_at), "created_at": iso(row.created_at),
@@ -192,8 +247,10 @@ async def menus_to_review(
"threshold": row.threshold, "threshold": row.threshold,
"logged_ids": logged, "logged_ids": logged,
# Logged then, not found now: deleted, edited out of reach, or # 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], "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, "lines": lines,
}) })
return { return {
+37 -6
View File
@@ -81,17 +81,22 @@ def test_the_re_run_searches_the_way_the_arm_does():
assert rerun.get(key) == arm[key], key 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 @pytest.mark.asyncio
async def test_the_re_run_ranks_from_one_and_marks_the_budget_cut(): 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 = [ hits = [
(0.81, fake_note(id=11, title="first")), (0.81, fake_note(id=11, title="first", created_at=BEFORE, updated_at=BEFORE)),
(0.74, fake_note(id=12, title="second")), (0.74, fake_note(id=12, title="second", created_at=BEFORE, updated_at=BEFORE)),
(0.66, fake_note(id=13, title="third")), (0.66, fake_note(id=13, title="third", created_at=BEFORE, updated_at=BEFORE)),
] ]
async def search(user_id, query, **kw): 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 assert kw["project_id"] is None
kw["report"]["best_chunk"] = {12: {"text": "second\nthe passage that matched"}} kw["report"]["best_chunk"] = {12: {"text": "second\nthe passage that matched"}}
return hits 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 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 ───────────────────────────────────────────────────── # ── against Postgres ─────────────────────────────────────────────────────
@@ -155,7 +185,8 @@ def _lines(*specs):
"""A stand-in re-run: (record_id, rank, within_budget).""" """A stand-in re-run: (record_id, rank, within_budget)."""
return AsyncMock(return_value=[ return AsyncMock(return_value=[
{"record_id": rid, "rank": rank, "score": 0.9 - rank / 10, {"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 for rid, rank, within in specs
]) ])