refactor(retrieval): the notes arms run on the one pipeline - auto_inject, its reuse and lesson slots, the write path by meaning, and rule_via_lesson are specs (milestone 456 step 4, #4906)
CI & Build / Python lint (push) Successful in 3s
CI & Build / Plugin hooks (push) Successful in 12s
CI & Build / TypeScript typecheck (push) Successful in 53s
CI & Build / integration (push) Successful in 1m2s
CI & Build / Python tests (push) Successful in 1m52s
CI & Build / Build & push image (push) Successful in 35s

retrieval_pipeline gains the notes half: NoteArm / NoteSlot / NoteMoment /
NoteIO / NoteResult and run_note_arm, which writes once the stages both
notes arms copied: search, withhold this response's own menu (#3739),
fresh/repeat split, the call row before any return (#3497, #3752), the
band, the reserved slots in their order (reuse evicts, lesson extends),
and the surfacing rows. The note renderer (_record_kind, _menu_name,
_menu_passage, the seen pointer, menu_entry) moves with it, and
run_via_lesson_arm takes rule_via_lesson.

Behaviour-preserving, with flags for today's differences: notes still log
BEFORE the band and rules after it (step 7's question). One deliberate
change: a notes arm now fails open like the rule arms, so a failing
search costs its lines and no longer the whole hook response.

The I/O is resolved from plugin_context at call time (_note_io), so the
existing patches keep working. The review re-run reads the AUTO_INJECT
spec instead of restating it, and its guard now compares the two live
searches. The registry declares the pipeline's notes fan-out sites.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
2026-10-05 20:02:37 -04:00
co-authored by Claude Opus 5.5
parent 4b1060ae8b
commit 2b8f41229d
10 changed files with 1172 additions and 763 deletions
+277
View File
@@ -0,0 +1,277 @@
"""The notes corpus on the one retrieval pipeline (milestone 456 step 4).
`run_note_arm` is driven directly here, with the search and both recorders
stubbed through `NoteIO`: the stage order, the invariants the rule arms
already pin (#3497's call row before any return, #3752's fresh-only cut,
#4101's ledger-renders-not-removes), the write path's withheld menu (#3739),
and the two reserved slots in their load-bearing order. The builders that
compose these arms keep their own tests; these pin the stages they share.
"""
from __future__ import annotations
import ast
from pathlib import Path
from unittest.mock import MagicMock
import pytest
from scribe.services import retrieval_pipeline as rp
from scribe.services.lessons import LESSON_NOTE_TYPE
from scribe.services.retrieval_registry import FAN_OUT_SITES, POINTS
from scribe.services.retrieval_surfaces import SURFACES
from tests.helpers import fake_note
def _io(route):
"""A NoteIO whose search answers by the kinds asked for, and records calls."""
calls: list[dict] = []
async def search(_uid, _q, **kw):
calls.append(kw)
return list(route(kw))
return rp.NoteIO(search=search, record_retrieval=MagicMock(),
record_surfaced=MagicMock()), calls
def _rows(mock, source):
return [c.kwargs for c in mock.call_args_list if c.kwargs["source"] == source]
def _note(nid, score=None, **kw):
note = fake_note(id=nid, title=f"record {nid}", user_id=1, **kw)
return (score, note) if score is not None else note
async def _run(arm, route, *, budget=3, **moment):
io, calls = _io(route)
result = await rp.run_note_arm(
arm, rp.NoteMoment(user_id=1, query="q", project_id=moment.pop("project_id", 2),
**moment),
floor=0.5, budget=budget, io=io,
)
return result, io, calls
# ── the call row (#3497, #3752) ───────────────────────────────────────────
@pytest.mark.asyncio
async def test_a_call_that_found_nothing_is_still_logged():
result, io, _ = await _run(rp.AUTO_INJECT, lambda kw: [])
assert result.menu == []
(row,) = _rows(io.record_retrieval, "auto_inject")
assert row["results"] == [] and row["suppressed"] == 0
io.record_surfaced.assert_not_called()
@pytest.mark.asyncio
async def test_repeats_and_named_records_are_counted_not_reported():
hits = [_note(11, 0.80), _note(22, 0.79), _note(33, 0.78)]
result, io, _ = await _run(
rp.AUTO_INJECT, lambda kw: [] if kw.get("note_type") else hits,
seen=frozenset({11}), named=frozenset({33}),
)
(row,) = _rows(io.record_retrieval, "auto_inject")
assert [int(n.id) for _s, n in row["results"]] == [22]
assert row["suppressed"] == 2
# The repeat stays on the menu (#4101); the named record does not — it is
# shown by its own block.
assert [int(n.id) for _s, n in result.menu] == [11, 22]
(surfaced,) = _rows(io.record_surfaced, "auto_inject")
assert surfaced["note_ids"] == [22]
def test_every_notes_stage_logs_before_it_can_return():
"""Structural, like the rule arms' guard: no `return` in a notes stage
precedes its call row. The via-lesson arm's one exception is the return
before any search runs — nothing was asked, so there is nothing to log."""
tree = ast.parse(Path("src/scribe/services/retrieval_pipeline.py").read_text())
fns = {n.name: n for n in ast.walk(tree) if isinstance(n, ast.AsyncFunctionDef)}
for name in ("run_note_arm", "_reserve_note_slot", "run_via_lesson_arm"):
fn = fns[name]
logged = [n.lineno for n in ast.walk(fn) if isinstance(n, ast.Call)
and getattr(n.func, "attr", None) == "record_retrieval"]
searched = [n.lineno for n in ast.walk(fn) if isinstance(n, ast.Call)
and getattr(n.func, "attr", None) == "search"]
assert logged and searched, f"{name} no longer searches and logs"
early = [n.lineno for n in ast.walk(fn) if isinstance(n, ast.Return)
and min(searched) < n.lineno < min(logged)]
assert not early, f"{name} can return at line {early[0]} before logging"
# ── the band, and the arm failing open ────────────────────────────────────
@pytest.mark.asyncio
async def test_the_band_narrows_the_menu_after_the_row_is_written():
"""Notes log BEFORE the band (#2085) — the row is the candidate set a
floor is tuned against, the surfacing rows are what the reader saw."""
hits = [_note(1, 0.90), _note(2, 0.70)]
result, io, _ = await _run(rp.WRITE_PATH, lambda kw: hits)
(row,) = _rows(io.record_retrieval, "write_path")
assert [int(n.id) for _s, n in row["results"]] == [1, 2]
assert [int(n.id) for _s, n in result.menu] == [1]
(surfaced,) = _rows(io.record_surfaced, "write_path_semantic")
assert surfaced["note_ids"] == [1]
@pytest.mark.asyncio
async def test_a_search_that_raises_costs_the_arm_not_the_response():
async def boom(*_a, **_kw):
raise RuntimeError("index unavailable")
io = rp.NoteIO(search=boom, record_retrieval=MagicMock(), record_surfaced=MagicMock())
result = await rp.run_note_arm(
rp.AUTO_INJECT, rp.NoteMoment(user_id=1, query="q", project_id=None),
floor=0.5, budget=3, io=io,
)
assert result.menu == [] and result.answered == []
@pytest.mark.asyncio
async def test_a_recorder_that_raises_costs_only_its_row():
hits = [_note(1, 0.90)]
io, _ = _io(lambda kw: [] if kw.get("note_type") else hits)
io.record_retrieval.side_effect = RuntimeError("telemetry down")
result = await rp.run_note_arm(
rp.AUTO_INJECT, rp.NoteMoment(user_id=1, query="q", project_id=None),
floor=0.5, budget=3, io=io,
)
assert [int(n.id) for _s, n in result.menu] == [1]
# ── the write path's withheld menu (#3739) ────────────────────────────────
@pytest.mark.asyncio
async def test_this_responses_menu_is_withheld_but_a_pulled_one_is_still_scored():
pulled, listed = 5, 6
answered = [_note(pulled, 0.91), _note(7, 0.88)]
result, io, calls = await _run(
rp.WRITE_PATH, lambda kw: answered, budget=2,
in_menu=frozenset({pulled, listed}), still_scored=frozenset({pulled}),
)
(call,) = calls
# The listed one never reaches the search; the pulled one does, and the
# limit covers it so the menu still gets its full budget.
assert call["exclude_ids"] == {listed}
assert call["limit"] == 3
assert [int(n.id) for _s, n in result.answered] == [pulled, 7]
assert [int(n.id) for _s, n in result.menu] == [7]
# Something was withheld after the search answered, so the score it
# reported may be one of ours: not measured on this call.
(row,) = _rows(io.record_retrieval, "write_path")
assert row["best_available"] is None and row["best_available_id"] is None
@pytest.mark.asyncio
async def test_the_prompt_menu_never_sends_the_ledger_into_its_search():
_result, _io_, calls = await _run(
rp.AUTO_INJECT, lambda kw: [], seen=frozenset({11}),
)
assert "exclude_ids" not in calls[0]
# ── the reserved slots, in their order ────────────────────────────────────
def _routed(main, reuse=(), lesson=()):
def route(kw):
kinds = kw.get("note_type") or ()
if LESSON_NOTE_TYPE in kinds:
return lesson
if kinds:
return reuse
return main
return route
@pytest.mark.asyncio
async def test_reuse_evicts_the_last_line_and_the_lesson_extends():
main = [_note(1, 0.70), _note(2, 0.69), _note(3, 0.68)]
reuse = [_note(9, 0.60, note_type="snippet")]
lesson = [_note(42, 0.58, note_type=LESSON_NOTE_TYPE)]
result, io, calls = await _run(
rp.AUTO_INJECT, _routed(main, reuse, lesson), budget=3,
)
# Reuse took the last of three; the lesson was added as a fourth.
assert [int(n.id) for _s, n in result.menu] == [1, 2, 9, 42]
assert result.slot_ids == {"reuse_slot": 9, "lesson_slot": 42}
# Reuse ran first: the lesson query was told about the snippet.
assert [c.get("note_type") for c in calls] == [
None, ("snippet", "process"), (LESSON_NOTE_TYPE,),
]
assert 9 in calls[2]["exclude_ids"]
# Each slot logs under its own name.
assert _rows(io.record_retrieval, "reuse_slot") and _rows(io.record_retrieval, "lesson_slot")
# The lesson books its own surfacing and never the arm's; the reuse line
# is counted under the arm, as it always was.
assert _rows(io.record_surfaced, "lesson_slot")[0]["note_ids"] == [42]
assert _rows(io.record_surfaced, "auto_inject")[0]["note_ids"] == [1, 2, 9]
@pytest.mark.asyncio
async def test_a_slot_stands_down_when_its_kind_won_on_score():
main = [_note(9, 0.80, note_type="snippet"), _note(42, 0.79, note_type=LESSON_NOTE_TYPE)]
_result, io, calls = await _run(rp.AUTO_INJECT, _routed(main))
assert len(calls) == 1
assert not _rows(io.record_retrieval, "reuse_slot")
assert not _rows(io.record_retrieval, "lesson_slot")
@pytest.mark.asyncio
async def test_a_slot_spends_itself_on_a_repeat_without_booking_it_again():
"""A record shown in an EARLIER call is exactly what a slot may spend
itself on (#4101) — it is rendered, but it is not a new surfacing."""
main = [_note(1, 0.70)]
lesson = [_note(42, 0.58, note_type=LESSON_NOTE_TYPE)]
result, io, _ = await _run(
rp.AUTO_INJECT, _routed(main, lesson=lesson), seen=frozenset({42}),
)
assert 42 in [int(n.id) for _s, n in result.menu]
assert not _rows(io.record_surfaced, "lesson_slot")
(row,) = _rows(io.record_retrieval, "lesson_slot")
assert row["results"] == [] and row["suppressed"] == 1
@pytest.mark.asyncio
async def test_the_write_path_reserves_no_slot():
main = [_note(1, 0.80)]
_result, _io_, calls = await _run(rp.WRITE_PATH, _routed(main, [_note(9, 0.7)]))
assert len(calls) == 1
# ── the specs are the declarations ────────────────────────────────────────
def test_every_notes_source_is_registered_and_every_arm_is_tunable():
for source in (*rp.NOTE_SOURCES, *rp.NOTE_SURFACED_SOURCES):
assert source in POINTS, f"{source} records telemetry but is not registered"
for arm in rp.NOTE_ARMS:
assert arm.source in SURFACES, f"{arm.source} has no floor or budget to tune"
# And the slots are measurable, never tunable: a budget of 1 is their feature.
for slot in rp.NOTE_SLOTS:
assert slot.source not in SURFACES
def test_the_pipelines_fan_out_sites_declare_every_notes_source():
logged = FAN_OUT_SITES[
"scribe/services/retrieval_pipeline.py::record_retrieval(source=source)"]
surfaced = FAN_OUT_SITES[
"scribe/services/retrieval_pipeline.py::record_surfaced(source=source)"]
assert set(rp.NOTE_SOURCES) <= set(logged)
assert set(rp.NOTE_SURFACED_SOURCES) == set(surfaced)
# Not the arm's: the slot whose line it is, and the lookups' own site.
assert "write_path_semantic" not in FAN_OUT_SITES[
"scribe/services/plugin_context.py::record_surfaced(source=arm)"]
def test_the_filters_say_which_kinds_and_whose_records():
assert rp.note_search_filters(rp.AUTO_INJECT) == {
"include_global_kinds": True, "scope": "browse",
}
assert rp.note_search_filters(rp.WRITE_PATH) == {
"note_type": ("snippet", "note", LESSON_NOTE_TYPE), "task_kind": "issue",
"include_global_kinds": True, "scope": "browse",
}
+37 -25
View File
@@ -5,12 +5,9 @@ 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
from unittest.mock import AsyncMock, MagicMock, patch
import pytest
import pytest_asyncio
@@ -54,31 +51,46 @@ async def test_an_empty_verdict_list_is_refused():
# ── 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 {}
_FILTERS = ("note_type", "task_kind", "include_global_kinds", "scope")
def test_the_re_run_searches_the_way_the_arm_does():
@pytest.mark.asyncio
async 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
a menu nobody was shown. Both searches are RUN and their keywords
compared, so a change to the arm's spec fails here until the review
follows it — and the arm has to have filters for the check to mean
anything (rule 167)."""
from scribe.services import retrieval_pipeline as rp
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
arm_kw: dict = {}
async def arm_search(_uid, _q, **kw):
arm_kw.update(kw)
return []
await rp.run_note_arm(
rp.AUTO_INJECT, rp.NoteMoment(user_id=1, query="q", project_id=None),
floor=0.5, budget=3,
io=rp.NoteIO(search=arm_search, record_retrieval=MagicMock(),
record_surfaced=MagicMock()),
)
rerun_kw: dict = {}
async def rerun_search(_uid, _q, **kw):
rerun_kw.update(kw)
return []
row = SimpleNamespace(query="q", threshold=0.5, limit_n=3, project_id=None,
created_at=CALL)
with patch("scribe.services.embeddings.semantic_search_notes", side_effect=rerun_search):
await review._rerun(1, row, 3)
arm = {k: arm_kw[k] for k in _FILTERS if k in arm_kw}
assert "scope" in arm and "include_global_kinds" in arm, (
f"the arm no longer passes its visibility filters: {arm_kw}"
)
assert {k: rerun_kw[k] for k in _FILTERS if k in rerun_kw} == arm
CALL = datetime(2026, 9, 20, 12, 0, tzinfo=timezone.utc)
+2 -2
View File
@@ -70,9 +70,9 @@ def test_every_surface_name_is_a_real_telemetry_source():
# The rule arms record through the one pipeline (milestone 456), where
# `source` is the spec's own field — so for them the spec IS the string
# `record_retrieval` receives, and the join key is checked against it.
from scribe.services.retrieval_pipeline import RULE_ARMS
from scribe.services.retrieval_pipeline import NOTE_ARMS, RULE_ARMS
via_pipeline = {arm.source for arm in RULE_ARMS}
via_pipeline = {arm.source for arm in (*RULE_ARMS, *NOTE_ARMS)}
missing = [
s.name for s in rs.SURFACES.values()
if f'source="{s.name}"' not in blob
+21 -9
View File
@@ -1727,16 +1727,28 @@ def test_every_hook_rule_search_says_which_project_it_is_for():
f"{path}:{call.lineno} searches every project's rules from a hook"
)
# The pipeline: exactly one search, in `_ranked`, whose keyword set is
# built as a dict literal — so the keys are read from that literal.
# The pipeline: exactly one RULE search, in `_ranked`, whose keyword set
# is built as a dict literal — so the keys are read from that literal.
# Since step 4 the module also searches the NOTES corpus (the two notes
# arms, their slots, and the via-lesson arm's lesson search), so the
# searches are counted per function: a rule search copied anywhere else
# is a new function in this map, and fails it.
pipeline = ast.parse(Path("src/scribe/services/retrieval_pipeline.py").read_text())
searches = [
n for n in ast.walk(pipeline)
if isinstance(n, ast.Call) and getattr(n.func, "attr", None) == "search"
]
assert len(searches) == 1, (
f"the pipeline has {len(searches)} rule searches, expected 1 — every "
f"arm is meant to reach the ranker through `_ranked`"
by_function = {}
for fn in ast.walk(pipeline):
if isinstance(fn, ast.AsyncFunctionDef):
n = sum(1 for c in ast.walk(fn) if isinstance(c, ast.Call)
and getattr(c.func, "attr", None) == "search")
if n:
by_function[fn.name] = n
assert by_function == {
"_ranked": 1, # the rule ranker, every rule arm
"_reserve_note_slot": 1, # notes: a reserved slot
"run_note_arm": 1, # notes: the arm itself
"run_via_lesson_arm": 1, # notes: lessons, then their linked rules
}, (
f"the pipeline's searches moved: {by_function} — every rule arm is "
f"meant to reach the ranker through `_ranked`"
)
ranked = next(n for n in ast.walk(pipeline)
if isinstance(n, ast.AsyncFunctionDef) and n.name == "_ranked")
+2 -1
View File
@@ -3,6 +3,7 @@ from unittest.mock import AsyncMock, MagicMock, patch
import pytest
from scribe.services import plugin_context as pc_module
from scribe.services import retrieval_surfaces as rs
from scribe.services.retrieval_pipeline import REUSE_SLOT
from scribe.services.lessons import LESSON_NOTE_TYPE
from tests.helpers import fake_note, writepath_cfg
@@ -340,7 +341,7 @@ _CFG = {"enabled": True, "threshold": 0.55, "top_k": 3}
def _asked_for_reuse(calls: list[dict]) -> bool:
"""Did the reuse slot issue its reserved query on this run?"""
return any(
tuple(c.get("note_type") or ()) == pc_module._REUSE_KINDS for c in calls
tuple(c.get("note_type") or ()) == REUSE_SLOT.kinds for c in calls
)
+2 -1
View File
@@ -621,9 +621,10 @@ def test_the_rule_arm_asks_for_a_set_and_lets_the_band_narrow_it():
sharper — are narrowed harder than the notes menu is.
"""
from scribe.services import plugin_context as pc
from scribe.services import retrieval_pipeline as rp
assert pc.RULEHINT_LIMIT > 1
assert 0 < pc._RULEHINT_BAND < pc._AUTOINJECT_BAND
assert 0 < pc._RULEHINT_BAND < rp._NOTE_BAND
# --- the minimum-substance floor on the semantic arm (#2223) ------------------