fix(retrieval): the review drops records written after the call it re-runs (#4773)
CI & Build / Python lint (push) Successful in 4s
CI & Build / Plugin hooks (push) Successful in 13s
CI & Build / TypeScript typecheck (push) Successful in 54s
CI & Build / integration (push) Successful in 1m8s
CI & Build / Python tests (push) Successful in 1m53s
CI & Build / Build & push image (push) Successful in 36s
CI & Build / Python lint (push) Successful in 4s
CI & Build / Plugin hooks (push) Successful in 13s
CI & Build / TypeScript typecheck (push) Successful in 54s
CI & Build / integration (push) Successful in 1m8s
CI & Build / Python tests (push) Successful in 1m53s
CI & Build / Build & push image (push) Successful in 36s
The first live sample ranked records the call could never have been offered: the session that made a call writes the decision, often quoting the message, and that record tops the re-run. Judged, it inflates on_point exactly where the budget is decided. - _rerun over-fetches by POSTDATED_SLACK, drops records created after the call before ranking, and names them in `postdated`. - each line carries `changed_since_call` (updated_at or a work log after the call) and `logged` (shown fresh then). - judge_menu shares the re-run, so a post-dated record cannot be judged. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
@@ -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.
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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
|
||||
])
|
||||
|
||||
|
||||
Reference in New Issue
Block a user