Merge pull request 'Merge dev: scope-then-rank searches (#4958, #4961), CI gate on integration, milestone 456 steps 4-6' (#205) from dev into main
CI & Build / Python lint (push) Successful in 2s
CI & Build / Plugin hooks (push) Successful in 14s
CI & Build / TypeScript typecheck (push) Successful in 52s
CI & Build / integration (push) Successful in 59s
CI & Build / Python tests (push) Successful in 1m52s
CI & Build / Build & push image (push) Successful in 18s

This commit was merged in pull request #205.
This commit is contained in:
2026-10-05 22:23:48 -04:00
21 changed files with 2149 additions and 1187 deletions
+7 -1
View File
@@ -328,7 +328,13 @@ jobs:
# wouldn't un-publish a bad hook anyway: the push already did that. A failed
# `plugin` job still turns the whole run red, which is the signal that
# matters.
needs: [typecheck, lint, test]
#
# `integration` IS in needs (rule 177, #4958). It was added after this
# gate and left out of it, so on 2026-10-05 main run 8212 published
# :latest with integration red. The image carries the migrations and the
# queries that lane runs against real Postgres; a red integration run is
# a red image.
needs: [typecheck, lint, test, integration]
# Build on dev, main, and v* tag pushes. dev → :dev, main → :latest,
# tag → :latest + :<version>; every build also gets an immutable :<sha>.
if: github.ref == 'refs/heads/dev' || github.ref == 'refs/heads/main' || startsWith(github.ref, 'refs/tags/v')
+4 -3
View File
@@ -7,9 +7,10 @@ from sqlalchemy.orm import Mapped, mapped_column
from scribe.models import Base
# bge-small-en-v1.5 produces 384-dim unit-normalized vectors. The column is a
# native pgvector `vector(384)` (see migration 0067) so similarity search runs
# as an indexed `ORDER BY embedding <=> :q LIMIT k` in Postgres rather than a
# full-table Python cosine scan.
# native pgvector `vector(384)` (see migration 0067) so similarity is ranked in
# Postgres by `<=>` rather than by a full-table Python cosine scan. The semantic
# searches rank the CALLER'S in-scope chunks exactly rather than through the
# HNSW index, which ranks before it filters (#4961; embeddings._rank_scoped).
EMBEDDING_DIM = 384
+124 -73
View File
@@ -856,6 +856,62 @@ def record_best_chunk(report: dict | None, chunks: dict[int, dict]) -> None:
report["best_chunk"] = chunks
# --- scope first, then rank (#4958, #4961) -----------------------------------
#
# Every semantic search here asks the same question of a shared table: the
# nearest chunks AMONG THE ONES THIS CALLER MAY SEE. Ordered straight off a
# `*_embeddings` table, the planner answers a different one. It walks the HNSW
# index, takes about `hnsw.ef_search` (40) nearest chunks from every owner and
# project, and only then applies the scope — so an in-scope record ranked past
# the 40th chunk overall is silently gone, and on a shared install other
# users' records are what fill those 40. The searches fail open, so nothing
# reports it: the result is just shorter.
#
# The fix is one shape, shared by every search so it cannot be fixed in one
# and missed in the next (which is how #4958 left three behind). The scope is
# a MATERIALIZED CTE with the distance computed inside it; the index cannot
# order a materialized CTE, so the ranking over it is exact.
#
# What exact costs is one distance per in-scope chunk. Rules, milestones and
# Systems are hundreds of chunks. Notes are the big corpus, but the hook arms
# that run on every prompt are project-scoped, so the pass is over one
# project's chunks rather than the whole table; `retrieval_logs.duration_ms`
# on those arms is where that claim is checked after the deploy. If it ever
# stops holding, pgvector 0.8's `hnsw.iterative_scan` keeps the index and is
# the next step — gated on the extension version, because an unknown `hnsw.*`
# setting errors and these searches would read that error as "nothing found".
def _scoped_chunks(embedding_model, record_key, distance):
"""A select over an embedding table's chunks, each with its distance — the
caller joins its record and adds its scope, and `_rank_scoped` orders it.
The columns are labelled so every search's CTE carries the same four.
"""
return select(
record_key.label("record_id"),
distance.label("distance"),
embedding_model.chunk_index.label("chunk_index"),
embedding_model.chunk_text.label("chunk_text"),
).select_from(embedding_model)
def _rank_scoped(record_model, scoped_chunks, *, name: str, limit: int):
"""(record, distance, chunk_index, chunk_text) rows, nearest first, drawn
ONLY from the in-scope chunks — exactly, never through the index.
The row shape is what each search's collapse already unpacks, so a search
moves onto this without its callers or its mocks noticing.
"""
scoped = scoped_chunks.cte(name).prefix_with("MATERIALIZED")
return (
select(record_model, scoped.c.distance, scoped.c.chunk_index, scoped.c.chunk_text)
.join(scoped, scoped.c.record_id == record_model.id)
.order_by(scoped.c.distance)
.limit(limit)
)
async def semantic_search_notes(
user_id: int,
query: str,
@@ -924,11 +980,14 @@ async def semantic_search_notes(
a caller that forgets is wrong in the safe direction.
Ranking and the top-k cut happen in Postgres via pgvector's cosine-distance
operator (`<=>`, exposed as ``Vector.cosine_distance``) backed by the HNSW
index from migration 0067 — so this is an indexed ``ORDER BY ... LIMIT k``
rather than a full-table scan. Cosine distance is ``1 - cosine_similarity``,
so a similarity floor of *threshold* is a distance ceiling of
``1 - threshold`` and similarity is recovered as ``1 - distance``.
operator (`<=>`, exposed as ``Vector.cosine_distance``), EXACTLY and over
the in-scope chunks only (#4961, see `_rank_scoped`). It used to be an
indexed ``ORDER BY ... LIMIT k`` straight off the HNSW index from
migration 0067, which ranks before it filters and so could drop an
in-scope record that sat behind ~40 nearer ones the caller cannot see.
Cosine distance is ``1 - cosine_similarity``, so a similarity floor of
*threshold* is a distance ceiling of ``1 - threshold`` and similarity is
recovered as ``1 - distance``.
`demote_superseded` applies the supersession penalty (#278): a record a
later note claims to have overtaken ranks below its equals. Callers asking
@@ -964,18 +1023,16 @@ async def semantic_search_notes(
# Scope on Note, not NoteEmbedding.user_id: the embedding row belongs
# to the note's owner, so filtering it would pin every scope to "own"
# and leave shared records unreachable by meaning.
#
# Every filter below is the SCOPE, built into the chunks the
# ranking draws from rather than applied to what it returns
# (#4961) — see `_rank_scoped`.
stmt = (
# chunk_index/chunk_text ride along so the collapse below can
# say WHICH passage matched. Without them the caller is left
# previewing the head of the body — a span this query has
# already determined is not why the record ranked (#4243).
select(
Note,
distance.label("distance"),
NoteEmbedding.chunk_index,
NoteEmbedding.chunk_text,
)
.select_from(NoteEmbedding)
_scoped_chunks(NoteEmbedding, NoteEmbedding.note_id, distance)
.join(Note, NoteEmbedding.note_id == Note.id)
.where(
notes_visibility_clause(user_id, scope),
@@ -1032,24 +1089,22 @@ async def semantic_search_notes(
# the results and the live record that should have replaced it was
# never fetched.
#
# Ordering stays on RAW distance so pgvector's HNSW index still
# serves it (migration 0067). Ordering by `distance + penalty`
# instead would be exact, and would turn an indexed top-k into a
# scan-and-sort of every embedded note.
#
# The cost of that trade, stated plainly: a live record outside the
# over-fetch window cannot be promoted into the results. With a
# penalty far smaller than the window's score spread, that case
# needs the true answer to be more than _SUPERSESSION_OVERFETCH
# ranks down, which no observed query comes close to.
# Ordering is on RAW distance; the penalty is applied after the
# cut, over the over-fetched window. The cost of that, stated
# plainly: a live record outside the window cannot be promoted
# into the results. With a penalty far smaller than the window's
# score spread, that case needs the true answer to be more than
# _SUPERSESSION_OVERFETCH ranks down, which no observed query
# comes close to.
fetch = limit * _CHUNK_OVERFETCH * (
_SUPERSESSION_OVERFETCH if demote_superseded else 1
)
# NO threshold predicate — see the note above this function. The
# bar is applied after the collapse, where the rejected scores can
# still be seen.
stmt = stmt.order_by(distance.asc()).limit(fetch)
rows = list((await session.execute(stmt)).all())
rows = list((await session.execute(
_rank_scoped(Note, stmt, name="scoped_note_chunks", limit=fetch)
)).all())
except Exception:
logger.warning("Failed to query note embeddings", exc_info=True)
return []
@@ -1550,29 +1605,31 @@ async def semantic_search_rules(
# function fails open.
home = await rule_home(user_id, project_id, everywhere=everywhere)
# SCOPE FIRST, THEN RANK — exactly (#4958; the shape is
# `_rank_scoped`'s). The home filter is built into the chunks the
# ranking draws from: applied to what an HNSW walk returned, it ran
# after ~40 nearest chunks from every owner and project had been
# taken, and an in-scope rule ranked past them was silently gone.
scoped = (
joined_to_homes(
_scoped_chunks(RuleEmbedding, RuleEmbedding.rule_id, distance)
.join(Rule, RuleEmbedding.rule_id == Rule.id)
)
.where(
Rule.deleted_at.is_(None),
# No threshold predicate — see the note above
# semantic_search_notes. Applied below, after the collapse.
home,
*( [Rule.kind == kind] if kind else [] ),
)
)
async with async_session() as session:
rows = (await session.execute(
joined_to_homes(
select(
Rule,
distance.label("distance"),
RuleEmbedding.chunk_index,
RuleEmbedding.chunk_text,
)
.select_from(RuleEmbedding)
.join(Rule, RuleEmbedding.rule_id == Rule.id)
)
.where(
Rule.deleted_at.is_(None),
# No threshold predicate — see the note above
# semantic_search_notes. Applied below, after the collapse.
home,
*( [Rule.kind == kind] if kind else [] ),
)
# Overfetch so collapsing chunks to their best row still fills
# the page — the same reason the note search overfetches.
.order_by(distance)
.limit(limit * _CHUNK_OVERFETCH)
_rank_scoped(Rule, scoped, name="scoped_rule_chunks",
limit=limit * _CHUNK_OVERFETCH)
)).all()
except Exception:
logger.warning("Rule semantic search failed", exc_info=True)
@@ -1762,23 +1819,20 @@ async def semantic_search_milestones(
scope = Milestone.project_id == project_id
else:
scope = Milestone.user_id == user_id
# Scope first, then rank (#4961): see `_rank_scoped`.
scoped = (
_scoped_chunks(MilestoneEmbedding, MilestoneEmbedding.milestone_id, distance)
.join(Milestone, MilestoneEmbedding.milestone_id == Milestone.id)
.where(
scope,
Milestone.deleted_at.is_(None),
*([Milestone.status == status] if status else []),
)
)
async with async_session() as session:
rows = (await session.execute(
select(
Milestone,
distance.label("distance"),
MilestoneEmbedding.chunk_index,
MilestoneEmbedding.chunk_text,
)
.select_from(MilestoneEmbedding)
.join(Milestone, MilestoneEmbedding.milestone_id == Milestone.id)
.where(
scope,
Milestone.deleted_at.is_(None),
*([Milestone.status == status] if status else []),
)
.order_by(distance)
.limit(limit * _CHUNK_OVERFETCH)
_rank_scoped(Milestone, scoped, name="scoped_milestone_chunks",
limit=limit * _CHUNK_OVERFETCH)
)).all()
except Exception:
logger.warning("Milestone semantic search failed", exc_info=True)
@@ -1937,23 +1991,20 @@ async def semantic_search_systems(
scope = System.project_id == project_id
else:
scope = System.user_id == user_id
# Scope first, then rank (#4961): see `_rank_scoped`.
scoped = (
_scoped_chunks(SystemEmbedding, SystemEmbedding.system_id, distance)
.join(System, SystemEmbedding.system_id == System.id)
.where(
scope,
System.deleted_at.is_(None),
System.status != "archived",
)
)
async with async_session() as session:
rows = (await session.execute(
select(
System,
distance.label("distance"),
SystemEmbedding.chunk_index,
SystemEmbedding.chunk_text,
)
.select_from(SystemEmbedding)
.join(System, SystemEmbedding.system_id == System.id)
.where(
scope,
System.deleted_at.is_(None),
System.status != "archived",
)
.order_by(distance)
.limit(limit * _CHUNK_OVERFETCH)
_rank_scoped(System, scoped, name="scoped_system_chunks",
limit=limit * _CHUNK_OVERFETCH)
)).all()
except Exception:
logger.warning("System semantic search failed", exc_info=True)
File diff suppressed because it is too large Load Diff
+912
View File
@@ -39,6 +39,8 @@ import time
from dataclasses import dataclass, field
from typing import Any, Callable
from scribe.services.lessons import LESSON_NOTE_TYPE, claim_line
logger = logging.getLogger(__name__)
# ── Shared rendering and ranking helpers (moved from plugin_context) ──────
@@ -329,6 +331,110 @@ def _rule_hint_line(
)
# ── What a spec declares about itself (milestone 456 step 5) ─────────────
#
# A ranked surface used to be described in three hand-kept lists besides its
# code: its tuning pair in `retrieval_surfaces.SURFACES`, its row in
# `retrieval_registry.POINTS`, and — for a rule surface — its membership in
# `rule_usage.RANKED_SOURCES`. Tests existed to make the four agree. Now the
# spec carries all of it, and the three lists are read off the specs, so an
# arm added here is tunable, measured and counted by construction.
#
# THE CLASSES LIVE HERE, NOT BESIDE THE LISTS, for the import graph: the
# registry, the tuning table and the usage counter all read the specs, so the
# specs cannot import any of them back.
@dataclass(frozen=True)
class Surface:
"""One push arm's tunable pair, plus enough prose to tune it responsibly.
`asks` / `over` / `fires` are not documentation for this file — they are
rendered by the tuning tool and the Settings UI. A floor cannot be moved
sensibly by anyone, model or human, who does not know what the query is, what
corpus it runs against, or how often it costs something. Those three facts
are exactly what separates these arms from each other, and they were
previously recoverable only by reading `plugin_context.py`.
"""
name: str
"""The telemetry `source` value, and the join key.
MUST equal the string this arm passes to `record_retrieval`. Everything
useful about tuning depends on that identity: the tool that moves a floor
and the table that says what the floor did have to be talking about the same
arm. A test asserts it rather than a comment asking nicely.
"""
floor_key: str
floor_default: float
budget_key: str
budget_default: int
asks: str
over: str
fires: str
measured_model: str = "BAAI/bge-small-en-v1.5"
measured_shape: int = 1
"""What the SHIPPED defaults above were measured against (#4104).
A floor is a distance in one embedding model's geometry, over documents cut
one particular way. Either can change, and when one does every number in
this table describes something that no longer exists.
TWO FIELDS, NEVER ONE FUSED STRING (rule 149). A mismatch has to be able to
say WHICH half moved: a new embedding model and a re-cut document shape
invalidate the same numbers for different reasons and call for different
responses. `"<model>@<n>"` could only report that something changed, which
is the answer nobody can act on. Same reason `calibration_stamp()` returns
a dict and the event table gives each half its own column.
Recorded per surface rather than once for the module because they need not
move together: a surface retuned after a model change carries the new stamp
while its untouched siblings still carry the old one, and telling those
apart is the whole job.
LITERALS, deliberately, rather than an import of the live values — a stamp
says what was true when the number was chosen, so one that tracked the
current model would always agree with it and could never report staleness.
"""
budget_falls_back_to: str = ""
"""A budget key to inherit when this surface has none of its own set.
Only `write_path` uses it, and only because it USED to share auto-inject's
`top_k` outright. Giving it a key without this would silently reset the
budget of every install that had tuned the shared one — a behaviour change
delivered as a default, which is the shape of regression nobody reports
because nothing looks broken.
"""
@dataclass(frozen=True)
class Declared:
"""What the telemetry readout must know about a ranked source and cannot
read off its rows — `retrieval_registry.Point`'s fields, for the specs
that produce one. Every ranked source is UNBIDDEN: nobody asked for it."""
what: str
"""One line, for an agent reading a warning that names this source."""
quiet_because: str = ""
"""Set when silence over an active window is correct, saying why (#2475)."""
fixed_query: bool = False
"""The arm always searches the same string, so its decline rate is 0% or
100% and `cannot_decline` says nothing about it."""
@dataclass(frozen=True)
class RankedSource:
"""A ranked source that is a stage of an arm rather than an arm: it
records under its own name, and is measured but never tuned."""
source: str
declared: Declared
# ── The specs ────────────────────────────────────────────────────────────
@@ -362,6 +468,12 @@ class RuleArm:
recorded as surfaced — the rest were ranked, and the call row counts them,
but nobody was shown them."""
tuning: Surface | None = None
"""Its floor and budget, and the prose a tuner reads. Every arm has one;
`SURFACES` is read off them."""
declared: Declared | None = None
# Today's differences, reproduced exactly (milestone 456 step 2). Whether the
# prompt arm should band, and whether the act arms should reserve a
@@ -369,14 +481,47 @@ class RuleArm:
PROMPT_RULE = RuleArm(
"prompt_rule", band=False, compact_tail=False, checkpoint=False,
preference_slot=True,
tuning=Surface(
name="prompt_rule",
floor_key="kb_promptrule_threshold",
floor_default=0.72,
budget_key="kb_promptrule_top_k",
budget_default=3,
asks="the operator's message, against rule triggers",
over="global rules plus the bound project's own",
fires="once per operator turn",
),
declared=Declared("rules that may govern what the operator just asked"),
)
PRE_TOOL_RULE = RuleArm(
"pre_tool_rule", band=True, compact_tail=True, checkpoint=True,
preference_slot=False,
tuning=Surface(
name="pre_tool_rule",
floor_key="kb_toolrule_threshold",
floor_default=0.68,
budget_key="kb_toolrule_top_k",
budget_default=5,
asks="the command about to run, against rule triggers",
over="global rules plus the bound project's own",
fires="before every Bash call — the busiest arm there is",
),
declared=Declared("rules that may govern a command about to run"),
)
WRITE_PATH_RULE = RuleArm(
"write_path_rule", band=True, compact_tail=True, checkpoint=True,
preference_slot=False,
tuning=Surface(
name="write_path_rule",
floor_key="kb_rulehint_threshold",
floor_default=0.72,
budget_key="kb_rulehint_top_k",
budget_default=5,
asks="the code being written, against rule triggers",
over="global rules plus the bound project's own",
fires="before every Write and Edit",
),
declared=Declared("rules that may govern the file being written"),
)
# The completion report's preferences (milestone 409 step 4): a FIXED query,
# preferences only, read by update_task as records rather than as lines. Its
@@ -384,6 +529,30 @@ WRITE_PATH_RULE = RuleArm(
REPORT_PREFERENCE = RuleArm(
"report_preference", band=False, compact_tail=False, checkpoint=False,
preference_slot=False, kind="preference",
tuning=Surface(
name="report_preference",
floor_key="kb_reportpref_threshold",
floor_default=0.72,
budget_key="kb_reportpref_top_k",
budget_default=3,
# THE ONE FIXED QUERY, and the reason this arm behaves unlike the rest.
# The others score something that varies per call; this one scores a
# constant string, so its top score for a given corpus is also a
# constant. A floor a hair above that constant is not a quiet arm, it
# is a dead one, and no amount of traffic will ever reveal it — which
# is precisely how this arm spent 69 calls declining the same record.
asks="a fixed question about how to lay out a completion report",
over="preferences",
fires="when a task finishes",
),
# `fixed_query`: COMPLETION_QUERY is a module constant, so this arm's top
# score is the same number on every call — measured at 0.791 across 45
# consecutive calls, with p10, p50, p90, min and max all identical. Five
# equal percentiles is the tell.
declared=Declared(
"the fixed question asked when a task finishes: how should this report read",
fixed_query=True,
),
)
# The backstop for every arm that ran earlier in the turn and missed: the
# finished reply against every rule's trigger. Its floor IS its stop bar
@@ -391,12 +560,35 @@ REPORT_PREFERENCE = RuleArm(
REPLY_RULE = RuleArm(
"reply_rule", band=False, compact_tail=False, checkpoint=True,
preference_slot=False, stop_only=True,
# THE REPLY BACKSTOP (milestone 458, folded in from 456 step 8). Its floor
# is a STOP bar, not a hint bar: at the end of a turn nothing can be shown
# beside the reply, so a hit either holds the reply for one read or says
# nothing. Hence a default at the checkpoint's level and a budget of one —
# the call row's results are then exactly the rule that would hold.
tuning=Surface(
name="reply_rule",
floor_key="kb_replyrule_threshold",
floor_default=0.80,
budget_key="kb_replyrule_top_k",
budget_default=1,
asks="the reply that ends a turn, against rule triggers",
over="global rules plus the bound project's own",
fires="once per turn, when the reply is finished",
),
declared=Declared(
"a rule that holds the finished reply for one read — the backstop "
"for whatever the earlier arms missed"
),
)
RULE_ARMS: tuple[RuleArm, ...] = (
WRITE_PATH_RULE, PRE_TOOL_RULE, PROMPT_RULE, REPORT_PREFERENCE, REPLY_RULE,
)
PREFERENCE_SLOT_SOURCE = "preference_slot"
PREFERENCE_SLOT = RankedSource(
PREFERENCE_SLOT_SOURCE,
Declared("the one line reserved for a preference at the prompt boundary"),
)
# Every source the shared stages below can record under — what the registry
# declares for this module's fan-out sites, since `source` reaches the
@@ -681,6 +873,726 @@ async def run_rule_arm(
return RuleResult()
# ── The notes corpus (milestone 456 step 4) ──────────────────────────────
#
# The same stages over the other corpus: search, withhold what this response
# already lists, split fresh from repeats, log the call before any early
# return, band, reserved slots, surfacing rows. Two arms take them — the
# prompt's menu (`auto_inject`) and the write path's match by meaning
# (`write_path`) — and they differed only in which kinds they ask for, whether
# a menu sits above them in the same response, and which slots they reserve.
#
# ONE DIFFERENCE FROM THE RULE ARMS IS KEPT, ON PURPOSE, FOR STEP 7. A notes
# arm logs the cut BEFORE its band; a rule arm logs after it. #2085 chose the
# first for notes (the log row is the candidate set a floor is tuned against,
# the surfacing rows are what the reader saw) and #3851 the second for rules.
# Both are recorded decisions, so this refactor reproduces both.
#
# What stays OUTSIDE: every lookup. The records named by number, the snippets
# recorded at a path, the rulings, the design system and the shape ledger
# have no score and no floor; the route builders compose them around what
# this returns.
# Margin gate: drop any hit more than this far below the top hit's score, so a
# single strong match doesn't drag in a wall of barely-passing neighbours.
# Twice the rule band, which is narrow because the rule corpus is flat (#3851).
_NOTE_BAND = 0.10
def _record_kind(note) -> str:
"""The kind marker for an injected menu line — and, for a task, its status.
The menu is drawn from every record that carries an embedding, so a snippet,
a stored process, an issue and a stray dev-log all arrive looking identical.
Recorded prior art only stands out if the line says what it is — and the kind
is also what tells the reader which tool opens it.
Task-ness wins over `note_type` because it's the more useful distinction at a
glance: "there's an open issue about this" beats "there's a note about this".
A TASK ALSO CARRIES ITS STATUS, because for that kind alone the line is
read as a claim about live work. A finished step and an open one rendered
identically is not a cosmetic gap: a done step was cited as a milestone's
open one on the strength of a line exactly like this, which carries an id,
a kind and a title and said nothing about where the work stood (#4154).
Only for tasks — a note or a snippet has no status to be wrong about.
"""
if note.is_task:
kind = "issue" if note.task_kind == "issue" else "task"
# No fallback for a missing status: `is_task` IS `status is not None`
# (models/note.py), so a branch for a task without one could never be
# taken, and a dead branch is a claim about the data that isn't true.
#
# Parenthesised rather than dot-joined: the write-path prior-art line
# joins its own fields with " · ", so a dotted status would read as
# another flag beside `seen` instead of as part of the kind.
return f"{kind} ({note.status})"
return note.note_type or "note"
# WHAT A MENU LINE CARRIES (#4364): the record's NAME, its kind and System,
# and the WHOLE passage that matched. Metadata plus the evidence, rather than a
# title asked to be both.
#
# The name, not the title. A snippet's or lesson's title is `name — when it
# applies` by construction (`embeddings.trigger_title`), because that join is
# what makes it rank on its situation. That is an EMBEDDING shape, and rendered
# as a menu line it ran to 1,500+ characters — the trigger paragraph spent
# again on every line, and again on every repeat. The trigger still arrives
# when it is what matched: it is in the chunk, and `_menu_passage` hands it
# over when the title was the whole match.
#
# The whole passage, not 200 characters of it. The search already chose the
# chunk that matched; the old cut kept its head and tail, and the head is the
# title every chunk is prefixed with — so the reader got the title twice and
# lost the middle, which is where the match was (lesson #4248). A chunk is at
# most ~1.4 KB (`embeddings._CHUNK_CHAR_BUDGET`), and it is shown once: a
# repeat is a one-line pointer (`_menu_seen_line`), not a second copy.
def _menu_name(title: str | None, note_type: str | None, data=None, body: str | None = "") -> str:
"""The record's name — its title without the trigger composed into it."""
title = (title or "(untitled)").replace("\n", " ").strip()
data = data if isinstance(data, dict) else {}
if note_type == "snippet":
from scribe.services.embeddings import TRIGGER_SEP
return (data.get("name") or title.partition(TRIGGER_SEP)[0]).strip() or title
if note_type == LESSON_NOTE_TYPE:
from types import SimpleNamespace
from scribe.services.embeddings import untrigger_title
from scribe.services.lessons import lesson_trigger
trigger = lesson_trigger(SimpleNamespace(data=data, body=body or ""))
# One line even when the stored name is a story (#4797).
return claim_line(untrigger_title(title, trigger).strip() or title)
return title
def _menu_passage(title: str | None, chunk_text: str | None, name: str = "") -> str:
"""The matched chunk on one line, without the title it was embedded under.
Every chunk is `title\nsection` (`embeddings.embedding_text`), and `title`
here must be the EMBEDDED one (`embeddings.document_title`) — for a snippet
or lesson that is `name — trigger`, not the stored name — so the prefix is
stripped exactly. A chunk that WAS only the title — a short
record, or the head chunk of one — matched on the title, and for a
trigger-keyed kind the part of it the name line no longer shows is the
trigger: that is returned, because it is precisely what matched.
One line, so the menu's blockquote survives it.
"""
title = (title or "").strip()
text = (chunk_text or "").strip()
if title and text.startswith(title):
text = text[len(title):]
text = " ".join(text.split())
if not text and name and title.startswith(name) and title != name:
text = " ".join(title[len(name):].lstrip(" —-").split())
return text
def _menu_label(kind: str, systems: list[str] | None) -> str:
"""`issue (done) · Plugin & hooks` — the kind, then where it belongs."""
return " · ".join([kind, *systems]) if systems else kind
def _menu_seen_line(note_id: int, kind: str, name: str) -> str:
"""A pointer to a record this session was already shown, not a copy of it."""
return f"> - #{note_id} [{kind} · seen] {name}"
def menu_entry(
note_id: int, *, kind: str, name: str, systems: list[str] | None = None,
seen: bool = False, stale: bool = False, score: float | None = None,
shared_by: str = "", under: str = "",
) -> list[str]:
"""One record on a notes menu: its line, and what sits under it.
The prompt menu's two blocks — records named by number and records ranked
by meaning — used to write this out twice, differing only in what they
pass: a ranked line carries its score, a named one does not, and each puts
its own text under the line (the matched passage, or the record's opening).
A repeat is a POINTER, not a copy (#4364). The record is in this session's
context already — the ledger is cleared at compaction, so "seen" stays true
— and re-rendering it spent its whole line again for nothing.
A superseded record is DEMOTED, not removed (#278), so one can still reach
a menu, and when it does the reader has to be told: an agent handed stale
material with nothing marking it acts on it with full confidence.
"""
if seen:
line = _menu_seen_line(note_id, kind, name)
return [line + (" — SUPERSEDED" if stale else "")]
line = f"> - #{note_id} [{_menu_label(kind, systems)}] \"{name}\""
if score is not None:
line += f" ({score:.2f})"
if stale:
line += " — SUPERSEDED, a later record covers this; check that first"
if shared_by:
line += f" — shared by {shared_by}"
return [line, f"> ↳ {under}"] if under else [line]
@dataclass(frozen=True)
class NoteSlot:
"""One line reserved for a kind the open ranking keeps losing."""
source: str
kinds: tuple[str, ...]
"""What the slot is FOR — asked for by the search and checked on the way
out, so a slot is never spent on a line indistinguishable from one that
earned its place on score."""
evicts: bool
"""Take the menu's last line when the menu is full, rather than adding one."""
include_global_kinds: bool = False
books_own: bool = False
"""The slot's line is recorded as surfaced under the slot's own source,
and so never again under the arm's."""
declared: Declared | None = None
# Order is load-bearing: reuse evicts the menu's weakest hit while the lesson
# slot extends, so running them the other way round would let a reserved
# lesson be the line reuse throws off — a slot another slot can silently undo
# is not a guarantee.
REUSE_SLOT = NoteSlot(
"reuse_slot", ("snippet", "process"), evicts=True,
declared=Declared("the one line reserved for a reusable snippet"),
)
LESSON_SLOT = NoteSlot(
"lesson_slot", (LESSON_NOTE_TYPE,), evicts=False,
include_global_kinds=True, books_own=True,
declared=Declared("the one line reserved for a lesson"),
)
@dataclass(frozen=True)
class NoteArm:
"""One ranked notes surface: what it searches, and which stages it takes."""
source: str
"""The telemetry `source` of its call row. MUST equal the SURFACES key."""
surfaced_as: str
"""The `source` its surfacing rows carry."""
note_type: tuple[str, ...] | None = None
task_kind: str | None = None
include_global_kinds: bool = True
"""Lessons are project-independent (#3730), so they join the candidate set
from wherever they were learned."""
scope: str = "browse"
"""Nobody asked for an injected line, so it takes the BROWSE scope: never
a record shared one-to-one with the operator."""
slots: tuple[NoteSlot, ...] = ()
withholds: bool = False
"""A menu sits above this arm in the same response, and what it lists is
kept out of the search (`exclude_ids`) — same-call duplication, which is a
different claim from the session ledger (#4101)."""
tuning: Surface | None = None
declared: Declared | None = None
AUTO_INJECT = NoteArm(
"auto_inject", surfaced_as="auto_inject", slots=(REUSE_SLOT, LESSON_SLOT),
tuning=Surface(
name="auto_inject",
floor_key="kb_autoinject_threshold",
floor_default=0.55,
budget_key="kb_autoinject_top_k",
budget_default=3,
asks="the operator's message, as they typed it",
over="notes, snippets, processes and issues",
fires="once per operator turn",
),
declared=Declared("the notes menu offered at the prompt boundary"),
)
# Snippets AND recorded experience (#2246): an issue saying "we tried this and
# it deadlocked" is prior art for the code about to be written. `task_kind`
# keeps the open to-do list out — a task resembles the code and answers
# nothing. Lessons too (milestone 385 step 5): the arm is kind-FILTERED, so a
# kind absent here is unreachable, not merely outranked (#3702). No reserved
# slot: this arm fires before every Write and Edit, and the field is already
# narrow enough that the 200:1 dilution a slot answers does not happen.
WRITE_PATH = NoteArm(
"write_path", surfaced_as="write_path_semantic",
note_type=("snippet", "note", LESSON_NOTE_TYPE), task_kind="issue",
withholds=True,
tuning=Surface(
name="write_path",
floor_key="kb_writepath_threshold",
floor_default=0.68,
budget_key="kb_writepath_top_k",
budget_default=3,
budget_falls_back_to="kb_autoinject_top_k",
asks="the code being written, rewritten as a concept query",
over="snippets and recorded issues",
fires="before every Write and Edit",
),
declared=Declared("prior art offered when a file is about to be written"),
)
NOTE_ARMS: tuple[NoteArm, ...] = (AUTO_INJECT, WRITE_PATH)
NOTE_SLOTS: tuple[NoteSlot, ...] = (REUSE_SLOT, LESSON_SLOT)
# What the notes stages record under, for the registry's fan-out sites.
NOTE_SOURCES: tuple[str, ...] = (
*(arm.source for arm in NOTE_ARMS), *(slot.source for slot in NOTE_SLOTS),
)
NOTE_SURFACED_SOURCES: tuple[str, ...] = (
*(arm.surfaced_as for arm in NOTE_ARMS),
*(slot.source for slot in NOTE_SLOTS if slot.books_own),
)
def note_search_filters(arm: NoteArm) -> dict:
"""The arm's constant search keywords — which kinds, whose records.
Read by `retrieval_review` too: a re-run that searched with different
visibility or kinds from the arm's would judge a menu nobody was shown.
"""
out: dict = {}
if arm.note_type:
out["note_type"] = arm.note_type
if arm.task_kind:
out["task_kind"] = arm.task_kind
if arm.include_global_kinds:
out["include_global_kinds"] = True
out["scope"] = arm.scope
return out
@dataclass(frozen=True)
class NoteIO:
"""The ranker and the two recorders a notes arm reports to."""
search: Callable[..., Any]
"""`semantic_search_notes`, or a stand-in with its signature."""
record_retrieval: Callable[..., Any]
record_surfaced: Callable[..., Any]
@dataclass(frozen=True)
class NoteMoment:
"""What a notes arm searches with, and what this response already holds."""
user_id: int
query: str
project_id: int | None
"""Searched and logged as `project_id or None`."""
seen: frozenset[int] = frozenset()
"""The session ledger: rendered as a pointer, never withheld (#4101)."""
named: frozenset[int] = frozenset()
"""Records this response shows by LOOKUP (named by number). Shown in their
own block and booked under `named_ref`, so this arm did not surface them."""
in_menu: frozenset[int] = frozenset()
"""Records listed earlier in this same response — withheld."""
still_scored: frozenset[int] = frozenset()
"""Of `in_menu`, those the search must still score (a pulled snippet's
resemblance to the payload is evidence for the shape ledger)."""
@dataclass
class NoteResult:
"""What a notes arm put on the menu, and what its search answered."""
menu: list = field(default_factory=list)
"""(score, note) best first: band, slots, and the named records removed."""
answered: list = field(default_factory=list)
"""Everything the search returned, before this response's own menu was
withheld — what `resembles` is read from."""
chunks: dict = field(default_factory=dict)
"""The passage each hit matched on, from the arm's OWN search — so a chunk
is only ever paired with the query that matched it. The slots' lines are
fetched by their own queries and so have none here."""
slot_ids: dict = field(default_factory=dict)
"""{slot source: the id it spent}."""
def _note_band(hits: list) -> list:
"""The top hit, plus every hit within `_NOTE_BAND` of it.
Computed over ALL hits, repeats included: the band measures distance from
the top SCORE, and letting the ledger move that cutoff would make "you
were shown this" change what counts as relevant (#3851's axis
independence).
"""
if not hits:
return []
top = hits[0][0]
return [(s, n) for s, n in hits if s >= top - _NOTE_BAND]
async def _reserve_note_slot(
io: NoteIO, arm: NoteArm, slot: NoteSlot, moment: NoteMoment, menu: list,
*, floor: float, budget: int, shown: set[int],
) -> tuple[list, int | None]:
"""Guarantee `slot.kinds` one line, if one clears the arm's own bar.
WHY A SLOT AT ALL (#2246, milestone 385 step 5). Ranking by raw cosine is
blind to what KIND of record answers what kind of ask, and the corpus
makes that fatal: Scribe's project records are about software work, so a
task about building a helper outranks the snippet that IS one. Snippets
were ~0.5% of the corpus when this was measured; no floor fixes 200:1. A
lesson crowded out is worse still — it exists only to be met at the moment
it applies, so the arm that surfaces it IS its delivery.
THE SLOT BUYS POSITION, NOT A LOWER BAR. It reserves at the menu's own
floor and is NOT held to the band (the top score is the very thing these
kinds lose to), so a weak record cannot buy the line.
ITS OWN SOURCE, from the first deploy. A guarantee has to be falsifiable,
and the hit a slot pushed out sits in the arm's row while the query that
pushed it out would otherwise be nowhere (#2463).
THE LEDGER IS NOT AN EXCLUSION HERE EITHER (#4101), but this call's own
menu is: a record already on it must not be shown twice, while one shown
in an EARLIER call is exactly what a slot may spend itself on — a snippet
relevant then and now is the reuse case, not a duplicate of it.
Returns the possibly-changed menu and the id the slot spent.
"""
if any(_record_kind(n) in slot.kinds for _s, n in menu):
return menu, None
t0 = time.perf_counter()
report: dict = {}
on_menu = {int(n.id) for _s, n in menu}
kwargs: dict = {
"limit": 1, "threshold": floor, "project_id": moment.project_id or None,
"exclude_ids": set(on_menu),
# KIND-FILTERED, so the slot can only be spent on what it is for.
"note_type": slot.kinds, "scope": arm.scope, "report": report,
}
if slot.include_global_kinds:
kwargs["include_global_kinds"] = True
found = await io.search(moment.user_id, moment.query, **kwargs)
fresh = [(s, n) for s, n in found if int(n.id) not in shown]
source = slot.source
try:
io.record_retrieval(
user_id=moment.user_id, source=source, query=moment.query,
threshold=floor, limit=1, project_id=moment.project_id or None,
is_task=None, results=fresh,
best_available=report.get("best_available_score"),
best_available_id=report.get("best_available_id"),
searched=bool(report.get("searched", True)),
suppressed=len(found) - len(fresh),
duration_ms=(time.perf_counter() - t0) * 1000.0,
)
except Exception: # noqa: BLE001 - observation never breaks the observed
_telemetry_failed(source)
# Verified, not trusted: the kind the query asked for, and not a line
# this menu already carries.
placed = [
(s, n) for s, n in found
if _record_kind(n) in slot.kinds and int(n.id) not in on_menu
][:1]
if not placed:
return menu, None
slot_id = int(placed[0][1].id)
# FRESH ONLY, matching the row above: a ledger repeat is rendered (#4101)
# but is not a new surfacing.
if slot.books_own and slot_id not in shown:
try:
io.record_surfaced(
user_id=moment.user_id, note_ids=[slot_id], source=source,
project_id=moment.project_id or None,
)
except Exception: # noqa: BLE001 - observation never breaks the observed
_telemetry_failed(source)
if not slot.evicts:
# IT EXTENDS, IT NEVER DISPLACES, siding with `preference_slot`: a
# displaced hit sits in the arm's row, and evicting it would make the
# two tables disagree about the same call (#3668). And a record that
# does not bind should not throw a better-scoring one off the menu.
return menu + placed, slot_id
# Take the LAST line, never the first: the strongest overall hit is still
# the best answer, and displacing it would trade one blindness for another.
if len(menu) >= budget:
return menu[:budget - 1] + placed, slot_id
return (menu + placed)[:budget], slot_id
async def run_note_arm(
arm: NoteArm, moment: NoteMoment, *, floor: float, budget: int, io: NoteIO,
) -> NoteResult:
"""Run one notes arm through every stage it takes, in the one order.
search → withhold this response's menu → fresh/repeat split → call row →
band → reserved slots → surfacing rows. Fails open to an empty result.
"""
try:
t0 = time.perf_counter()
report: dict = {}
scored = moment.in_menu & moment.still_scored
kwargs: dict = {
# The still-scored ids come back and are dropped below, so the
# limit has to cover them.
"limit": budget + len(scored), "threshold": floor,
"project_id": moment.project_id or None,
**note_search_filters(arm), "report": report,
}
if arm.withholds:
kwargs["exclude_ids"] = set(moment.in_menu - scored)
answered = await io.search(moment.user_id, moment.query, **kwargs)
hits = [(s, n) for s, n in answered if int(n.id) not in moment.in_menu]
# WHAT THIS ARM WITHHELD AFTER THE SEARCH ANSWERED (#3739). The score
# the search reports is measured BEFORE that drop, so a record already
# listed in this response could be logged as one the BAR turned away
# — live proof once read a "rejection" at 0.822 against a lowest
# acceptance of 0.6857. Whenever this removed anything, the honest
# best-available is null: "not measured on this call".
withheld = len(answered) - len(hits)
hits = hits[:budget]
# `results=fresh` and `suppressed` together (#3752): a rendered repeat
# is not a new surfacing, so it is COUNTED rather than reported, and an
# all-repeat zero reads apart from a bar nothing cleared. A record
# NAMED in this response is counted the same way — shown, but by the
# lookup and booked under `named_ref`.
fresh = [
(s, n) for s, n in hits
if int(n.id) not in moment.seen and int(n.id) not in moment.named
]
source = arm.source
try:
io.record_retrieval(
user_id=moment.user_id, source=source, query=moment.query,
threshold=floor, limit=budget,
project_id=moment.project_id or None, is_task=None,
results=fresh,
best_available=(
None if withheld else report.get("best_available_score")
),
# Withheld on the SAME condition: a surviving id beside a null
# score would name a record without saying what it scored.
best_available_id=(
None if withheld else report.get("best_available_id")
),
searched=bool(report.get("searched", True)),
suppressed=len(hits) - len(fresh),
duration_ms=(time.perf_counter() - t0) * 1000.0,
)
except Exception: # noqa: BLE001 - observation never breaks the observed
_telemetry_failed(source)
result = NoteResult(answered=answered, chunks=report.get("best_chunk") or {})
if not hits:
return result
menu = _note_band(hits)
# The slots are told about the named records as if already shown, so
# a slot that picks one does not book a surfacing the lookup booked.
shown = set(moment.seen | moment.named)
for slot in arm.slots:
menu, slot_id = await _reserve_note_slot(
io, arm, slot, moment, menu, floor=floor, budget=budget,
shown=shown,
)
if slot_id is not None:
result.slot_ids[slot.source] = slot_id
# A named record is shown once, in its own block, never again as a match.
menu = [(s, n) for s, n in menu if int(n.id) not in moment.named]
result.menu = menu
# What SURVIVED the band — the menu the reader saw — where the call
# row holds the candidate set a floor is tuned against (#2085). FRESH
# ONLY, the row's own cut (#4101, #3668). A slot that books its own
# line is not this arm's surfacing: counting it twice would leave the
# slot's two tables describing different numbers of one event.
own = {
sid for src, sid in result.slot_ids.items()
if any(slot.books_own and slot.source == src for slot in arm.slots)
}
ids = [
int(n.id) for _s, n in menu
if int(n.id) not in moment.seen and int(n.id) not in own
]
if ids:
source = arm.surfaced_as
try:
io.record_surfaced(
user_id=moment.user_id, note_ids=ids, source=source,
project_id=moment.project_id,
)
except Exception: # noqa: BLE001 - observation never breaks the observed
_telemetry_failed(source)
return result
except Exception: # noqa: BLE001 - a recall aid never breaks its act
logger.debug("%s arm failed", arm.source, exc_info=True)
return NoteResult()
# ── A rule reached through its lessons (milestone 440, #4633) ────────────
#
# A rule's own document is written in the rule's words, which are general by
# design; the situations that keep proving it are often closer to what a
# session is actually doing. A lesson JUDGED to be an instance of a rule (a
# CONFIRMED link) carries that situation, so a lesson matching the moment
# brings its rule along — in rule voice, naming the lesson that reached it.
# It searches the NOTES corpus and answers with rules, which is why it sits
# between the two halves of this module.
#
# ITS OWN SLOT, NOT A RULE SLOT, as a stated default. There was nothing to
# measure until links are confirmed (#4632), so the asymmetry decides it: a
# via-lesson line that took a rule slot could push out a rule the ranker
# matched DIRECTLY, a stronger claim displaced by a weaker one, while an extra
# line costs one line. Logged as its own source so it can be judged (#4636).
VIA_LESSON_SOURCE = "rule_via_lesson"
VIA_LESSON = RankedSource(
VIA_LESSON_SOURCE,
Declared(
"a rule reached through a lesson confirmed as an instance of it",
quiet_because="searches only once some lesson has a confirmed link to "
"a rule; an install where none has been judged is "
"correctly silent here",
),
)
VIA_LESSON_LIMIT = 1
# Lessons fetched before keeping only the linked ones. The search cannot be
# told "linked lessons only", so it overfetches and filters; on a corpus where
# most lessons are unlinked, a fetch of one would almost always be spent on a
# lesson that carries nothing.
_VIA_LESSON_OVERFETCH = 10
@dataclass(frozen=True)
class ViaLessonIO:
"""The lesson search, the link lookups, and the recorders."""
search: Callable[..., Any]
"""`semantic_search_notes`, or a stand-in with its signature."""
linked: Callable[..., Any]
"""`lesson_rules.confirmed_lessons`: the lessons with a confirmed link."""
rules_for: Callable[..., Any]
"""`lesson_rules.confirmed_rules_in_scope`: their rules, in scope."""
floor: Callable[..., Any]
"""The notes menu's own bar — a lesson too weak to be shown cannot carry a
rule in. Awaited only once a linked lesson exists."""
record_retrieval: Callable[..., Any]
record_rule_surfaced: Callable[..., Any]
async def run_via_lesson_arm(
moment: RuleMoment, *, skip: frozenset[int], io: ViaLessonIO,
) -> RuleResult:
"""Rule lines reached through a matching lesson's CONFIRMED links.
`skip` is every rule this response already names plus the session ledger:
suppression applies to the RULE, whichever lesson reached it. Fails open.
"""
try:
confirmed = await io.linked(moment.user_id)
if not confirmed:
return RuleResult()
bar = await io.floor()
t0 = time.perf_counter()
report: dict = {}
found = await io.search(
moment.user_id, moment.query, limit=_VIA_LESSON_OVERFETCH,
threshold=bar, project_id=moment.project_id,
note_type=(LESSON_NOTE_TYPE,), include_global_kinds=True,
scope="browse", report=report,
)
matched = [(s, n) for s, n in found if int(n.id) in confirmed]
by_lesson = await io.rules_for(
moment.user_id, [int(n.id) for _s, n in matched], moment.project_id,
)
candidates = [
(s, rule, n) for s, n in matched for rule in by_lesson.get(int(n.id), [])
]
chosen: list = []
taken: set[int] = set(skip)
for s, rule, n in candidates:
if rule.id in taken:
continue
taken.add(rule.id)
chosen.append((s, rule, n))
if len(chosen) >= VIA_LESSON_LIMIT:
break
# Logged whenever a search ran, results or not — the #3497 guard. The
# score is the LESSON's, because the lesson is what was matched; the
# best-available id is left off for the same reason, since the row's
# results are rules and an id beside them would read as a rule id.
try:
io.record_retrieval(
user_id=moment.user_id, source=VIA_LESSON_SOURCE,
query=moment.query, threshold=bar, limit=VIA_LESSON_LIMIT,
project_id=moment.project_id, is_task=None,
results=[(s, rule) for s, rule, _n in chosen],
duration_ms=(time.perf_counter() - t0) * 1000.0,
best_available=report.get("best_available_score"),
searched=bool(report.get("searched", True)),
suppressed=len({r.id for _s, r, _n in candidates}) - len(chosen),
)
except Exception: # noqa: BLE001 - observation never breaks the observed
_telemetry_failed(VIA_LESSON_SOURCE)
if not chosen:
return RuleResult()
rule_ids = [rule.id for _s, rule, _n in chosen]
try:
io.record_rule_surfaced(
user_id=moment.user_id, rule_ids=rule_ids, source=VIA_LESSON_SOURCE,
)
except Exception: # noqa: BLE001 - observation never breaks the observed
_telemetry_failed(VIA_LESSON_SOURCE)
lines = [
_rule_hint_line(rule, where=moment.where, seen=False,
held=rule.id in moment.held)
+ f" Reached through lesson #{n.id} "
+ f"“{_menu_name(n.title, n.note_type, n.data, n.body)}”, "
+ "a recorded instance of it."
for _s, rule, n in chosen
]
return RuleResult(
lines=lines, rule_ids=rule_ids, shown_rule_ids=list(rule_ids),
shown=[(s, rule) for s, rule, _n in chosen],
)
except Exception: # noqa: BLE001 - a recall aid never breaks its act
logger.debug("%s arm failed", VIA_LESSON_SOURCE, exc_info=True)
return RuleResult()
# ── The specs, read as lists (milestone 456 step 5) ──────────────────────
#
# What `retrieval_surfaces.SURFACES`, the ranked rows of
# `retrieval_registry.POINTS` and `rule_usage.RANKED_SOURCES` are read from.
# The order is the order the Settings page and the readouts list them in.
TUNED_ARMS: tuple = (*NOTE_ARMS, *RULE_ARMS)
"""Every arm with a floor and a budget: one per tunable surface."""
RANKED: tuple = (
*TUNED_ARMS, PREFERENCE_SLOT, *NOTE_SLOTS, VIA_LESSON,
)
"""Every source that RANKED what it showed — the arms, and the stages that
run a query of their own and record under their own name."""
RULE_RANKED_SOURCES: tuple[str, ...] = (
*(arm.source for arm in RULE_ARMS), PREFERENCE_SLOT_SOURCE, VIA_LESSON_SOURCE,
)
"""The ranked sources that surface RULES — a ranker chose each line, so a
pull can confirm or refute it (`rule_usage.RANKED_SOURCES`)."""
# ── The moment arm (milestone 458) ───────────────────────────────────────
#
# A LOOKUP beside the ranked arms, not one of them. A rule mounted on a moment
+31 -34
View File
@@ -31,7 +31,10 @@ reserved slots belong in it precisely because they are judgeable without being
tunable. One is a control panel; the other is an inventory. A test asserts
every tunable surface also appears here, so the two cannot drift apart.
ADDING AN ARM MEANS ADDING A ROW HERE. `tests/test_retrieval_registry.py`
ADDING A RANKED ARM MEANS WRITING ITS SPEC. Since milestone 456 step 5 the
ranked rows are read off `retrieval_pipeline.RANKED`, each spec declaring its
own (`Declared`). Any OTHER point — a lookup, an asked search, an ambient
carrier, a pull — still means adding a row here. `tests/test_retrieval_registry.py`
walks the call sites with the ast module and fails on a source it cannot find
below — deliberately not a grep, because two of the sources in this file
(`wide_net`, `preference_slot`) reach their recorder through a
@@ -42,7 +45,10 @@ from __future__ import annotations
from dataclasses import dataclass
from scribe.services.retrieval_pipeline import MOMENT_RULE_SOURCE, RULE_ARMS, RULE_SOURCES
from scribe.services.retrieval_pipeline import (
MOMENT_RULE_SOURCE, NOTE_SOURCES, NOTE_SURFACED_SOURCES, RANKED, RULE_ARMS,
RULE_SOURCES,
)
# ── How a point is reached ────────────────────────────────────────────────
#
@@ -130,18 +136,17 @@ def _p(source, kind, what, **kw) -> tuple[str, Point]:
POINTS: dict[str, Point] = dict([
# ── Unbidden push arms ───────────────────────────────────────────────
_p("auto_inject", UNBIDDEN, "the notes menu offered at the prompt boundary"),
_p("write_path", UNBIDDEN, "prior art offered when a file is about to be written"),
_p("write_path_rule", UNBIDDEN, "rules that may govern the file being written"),
_p("pre_tool_rule", UNBIDDEN, "rules that may govern a command about to run"),
_p("prompt_rule", UNBIDDEN, "rules that may govern what the operator just asked"),
_p("reply_rule", UNBIDDEN,
"a rule that holds the finished reply for one read — the backstop "
"for whatever the earlier arms missed"),
_p("preference_slot", UNBIDDEN,
"the one line reserved for a preference at the prompt boundary"),
_p("reuse_slot", UNBIDDEN, "the one line reserved for a reusable snippet"),
_p("lesson_slot", UNBIDDEN, "the one line reserved for a lesson"),
#
# The RANKED ones are read off the pipeline's specs (milestone 456 step
# 5): every arm, and every stage that runs a query of its own and records
# under its own name, declares its row where it is defined. What stays
# here is what no spec describes — the lookups below, and everything asked,
# ambient or pulled.
*(_p(spec.source, UNBIDDEN, spec.declared.what,
expects_traffic=not spec.declared.quiet_because,
quiet_because=spec.declared.quiet_because,
fixed_query=spec.declared.fixed_query)
for spec in RANKED),
# A LOOKUP, not a ranker (#4796): a record the operator named by number
# in the message, fetched by id. No score and no bar, so it writes no
# retrieval_logs row and no score-shaped warning can apply; its rows are
@@ -152,20 +157,6 @@ POINTS: dict[str, Point] = dict([
expects_traffic=False,
quiet_because="speaks only when a message names a record by its id; "
"a window in which none did is correctly silent here"),
_p("rule_via_lesson", UNBIDDEN,
"a rule reached through a lesson confirmed as an instance of it",
expects_traffic=False,
quiet_because="searches only once some lesson has a confirmed link to "
"a rule; an install where none has been judged is "
"correctly silent here"),
# `fixed_query`: COMPLETION_QUERY is a module constant in
# services/reply_preferences.py, so this arm's top score is the same number
# on every call — measured at 0.791 across 45 consecutive calls, with p10,
# p50, p90, min and max all identical. Five equal percentiles is the tell.
_p("report_preference", UNBIDDEN,
"the fixed question asked when a task finishes: how should this report read",
fixed_query=True),
# The moment arm (milestone 458). A LOOKUP, not a ranker: a rule mounted on
# a moment arrives when that moment happens, with no score, so it writes no
# retrieval_logs row and no score-shaped warning can apply. Its rows are in
@@ -262,8 +253,10 @@ POINTS: dict[str, Point] = dict([
# nothing, because the paths are relative to `src/` and so begin `scribe/` —
# every site sailed through a check that looked like it was running.
FAN_OUT_SITES: dict[str, tuple[str, ...]] = {
# The write path's two LOOKUPS. Its ranked arm (`write_path_semantic`)
# records through the pipeline's site below (milestone 456 step 4).
"scribe/services/plugin_context.py::record_surfaced(source=arm)": (
"write_path_place", "write_path_semantic", "write_path_sync",
"write_path_place", "write_path_sync",
),
"scribe/services/system_rulings.py::record_system_surfaced(source=source)": (
"rulings_pre_tool", "rulings_write_path",
@@ -272,12 +265,16 @@ FAN_OUT_SITES: dict[str, tuple[str, ...]] = {
"enter_project", "start_planning", "get_task", "get_project",
"get_milestone",
),
# The one retrieval pipeline (milestone 456): every rule arm's call row and
# surfacing rows are written by the same two sites, so `source` arrives as
# a value. The values are the pipeline's own specs, read rather than
# restated, so an arm added there is declared here by construction.
# The one retrieval pipeline (milestone 456): every ranked arm's call row
# and surfacing rows — both corpora — are written by these sites, so
# `source` arrives as a value. The values are the pipeline's own specs,
# read rather than restated, so an arm added there is declared here by
# construction.
"scribe/services/retrieval_pipeline.py::record_retrieval(source=source)": (
RULE_SOURCES
RULE_SOURCES + NOTE_SOURCES
),
"scribe/services/retrieval_pipeline.py::record_surfaced(source=source)": (
NOTE_SURFACED_SOURCES
),
"scribe/services/retrieval_pipeline.py::record_rule_surfaced(source=source)": (
tuple(arm.source for arm in RULE_ARMS)
+8 -7
View File
@@ -96,18 +96,20 @@ async def _rerun(
) -> 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.
`auto_inject`'s parameters, read from its spec (`retrieval_pipeline.
AUTO_INJECT`) rather than restated: which kinds, whose records, and 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.
from scribe.services import retrieval_pipeline as rp
from scribe.services.embeddings import document_title, semantic_search_notes
from scribe.services.plugin_context import _menu_name, _menu_passage, _record_kind
from scribe.services.retrieval_pipeline import _menu_name, _menu_passage, _record_kind
rep: dict = {}
hits = await semantic_search_notes(
@@ -115,8 +117,7 @@ async def _rerun(
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",
**rp.note_search_filters(rp.AUTO_INJECT),
report=rep,
)
postdated = [int(n.id) for _s, n in hits if _after(n.created_at, row.created_at)]
+9 -148
View File
@@ -64,9 +64,10 @@ pointed in opposite directions, and only opening the record could tell.
"""
from __future__ import annotations
from dataclasses import dataclass
from scribe.services.settings import bounded_float, get_setting
from scribe.services.retrieval_pipeline import ( # noqa: F401 - Surface re-exported
TUNED_ARMS, Surface,
)
# A budget nobody should be able to set past. Not a tuning value — a guard on
# the worst case, so a mistyped setting cannot turn a menu into a wall of text.
@@ -75,153 +76,13 @@ from scribe.services.settings import bounded_float, get_setting
MAX_BUDGET = 10
@dataclass(frozen=True)
class Surface:
"""One push arm's tunable pair, plus enough prose to tune it responsibly.
`asks` / `over` / `fires` are not documentation for this file — they are
rendered by the tuning tool and the Settings UI. A floor cannot be moved
sensibly by anyone, model or human, who does not know what the query is, what
corpus it runs against, or how often it costs something. Those three facts
are exactly what separates these arms from each other, and they were
previously recoverable only by reading `plugin_context.py`.
"""
name: str
"""The telemetry `source` value, and the join key.
MUST equal the string this arm passes to `record_retrieval`. Everything
useful about tuning depends on that identity: the tool that moves a floor
and the table that says what the floor did have to be talking about the same
arm. A test asserts it rather than a comment asking nicely.
"""
floor_key: str
floor_default: float
budget_key: str
budget_default: int
asks: str
over: str
fires: str
measured_model: str = "BAAI/bge-small-en-v1.5"
measured_shape: int = 1
"""What the SHIPPED defaults above were measured against (#4104).
A floor is a distance in one embedding model's geometry, over documents cut
one particular way. Either can change, and when one does every number in
this table describes something that no longer exists.
TWO FIELDS, NEVER ONE FUSED STRING (rule 149). A mismatch has to be able to
say WHICH half moved: a new embedding model and a re-cut document shape
invalidate the same numbers for different reasons and call for different
responses. `"<model>@<n>"` could only report that something changed, which
is the answer nobody can act on. Same reason `calibration_stamp()` returns
a dict and the event table gives each half its own column.
Recorded per surface rather than once for the module because they need not
move together: a surface retuned after a model change carries the new stamp
while its untouched siblings still carry the old one, and telling those
apart is the whole job.
LITERALS, deliberately, rather than an import of the live values — a stamp
says what was true when the number was chosen, so one that tracked the
current model would always agree with it and could never report staleness.
"""
budget_falls_back_to: str = ""
"""A budget key to inherit when this surface has none of its own set.
Only `write_path` uses it, and only because it USED to share auto-inject's
`top_k` outright. Giving it a key without this would silently reset the
budget of every install that had tuned the shared one — a behaviour change
delivered as a default, which is the shape of regression nobody reports
because nothing looks broken.
"""
# THE TABLE IS READ OFF THE SPECS (milestone 456 step 5). Each ranked arm in
# `retrieval_pipeline` carries its own `Surface` — the floor and budget pair,
# its settings keys and defaults, and the prose a tuner reads — so an arm added
# there is tunable by construction, and nothing here can drift from it. The
# keys and defaults are the ones this table held before the move, unchanged.
SURFACES: dict[str, Surface] = {
"auto_inject": Surface(
name="auto_inject",
floor_key="kb_autoinject_threshold",
floor_default=0.55,
budget_key="kb_autoinject_top_k",
budget_default=3,
asks="the operator's message, as they typed it",
over="notes, snippets, processes and issues",
fires="once per operator turn",
),
"write_path": Surface(
name="write_path",
floor_key="kb_writepath_threshold",
floor_default=0.68,
budget_key="kb_writepath_top_k",
budget_default=3,
budget_falls_back_to="kb_autoinject_top_k",
asks="the code being written, rewritten as a concept query",
over="snippets and recorded issues",
fires="before every Write and Edit",
),
"write_path_rule": Surface(
name="write_path_rule",
floor_key="kb_rulehint_threshold",
floor_default=0.72,
budget_key="kb_rulehint_top_k",
budget_default=5,
asks="the code being written, against rule triggers",
over="global rules plus the bound project's own",
fires="before every Write and Edit",
),
"pre_tool_rule": Surface(
name="pre_tool_rule",
floor_key="kb_toolrule_threshold",
floor_default=0.68,
budget_key="kb_toolrule_top_k",
budget_default=5,
asks="the command about to run, against rule triggers",
over="global rules plus the bound project's own",
fires="before every Bash call — the busiest arm there is",
),
"prompt_rule": Surface(
name="prompt_rule",
floor_key="kb_promptrule_threshold",
floor_default=0.72,
budget_key="kb_promptrule_top_k",
budget_default=3,
asks="the operator's message, against rule triggers",
over="global rules plus the bound project's own",
fires="once per operator turn",
),
"report_preference": Surface(
name="report_preference",
floor_key="kb_reportpref_threshold",
floor_default=0.72,
budget_key="kb_reportpref_top_k",
budget_default=3,
# THE ONE FIXED QUERY, and the reason this arm behaves unlike the rest.
# The others score something that varies per call; this one scores a
# constant string, so its top score for a given corpus is also a
# constant. A floor a hair above that constant is not a quiet arm, it is
# a dead one, and no amount of traffic will ever reveal it — which is
# precisely how this arm spent 69 calls declining the same record.
asks="a fixed question about how to lay out a completion report",
over="preferences",
fires="when a task finishes",
),
# THE REPLY BACKSTOP (milestone 458, folded in from 456 step 8). Its floor
# is a STOP bar, not a hint bar: at the end of a turn nothing can be shown
# beside the reply, so a hit either holds the reply for one read or says
# nothing. Hence a default at the checkpoint's level and a budget of one —
# the call row's results are then exactly the rule that would hold.
"reply_rule": Surface(
name="reply_rule",
floor_key="kb_replyrule_threshold",
floor_default=0.80,
budget_key="kb_replyrule_top_k",
budget_default=1,
asks="the reply that ends a turn, against rule triggers",
over="global rules plus the bound project's own",
fires="once per turn, when the reply is finished",
),
arm.tuning.name: arm.tuning for arm in TUNED_ARMS if arm.tuning is not None
}
# Reserved slots are deliberately absent. `preference_slot`, `reuse_slot` and
+9 -20
View File
@@ -89,6 +89,7 @@ from scribe.models.rule_usage import (
APPLIED, DEPARTED, OUTCOMES, PULLED, SURFACED, RuleUsageEvent,
)
from scribe.services.background import report_telemetry_failure, spawn
from scribe.services.retrieval_pipeline import MOMENT_RULE_SOURCE, RULE_RANKED_SOURCES
logger = logging.getLogger(__name__)
@@ -98,31 +99,19 @@ logger = logging.getLogger(__name__)
# Membership is the whole definition of the pull-through denominator: a ranked
# surfacing is a claim ("this rule may apply to what you are doing") that a pull
# can confirm or refute, while an ambient one is a delivery nobody decided on.
# Add a source here only when a ranker picked it.
# The ranked half is read off the retrieval pipeline's specs (milestone 456
# step 5): every rule arm — the reply backstop included, which records only
# the rule that HELD the reply — the preference slot, which is a ranker's
# choice twice over, and the rule reached through a confirmed lesson link,
# whose own source lets its pull-through read apart from a direct match
# (#4636). An arm added there is counted here by construction.
RANKED_SOURCES = (
"write_path_rule", "pre_tool_rule", "prompt_rule",
# The reply backstop records only the rule that HELD the reply — a claim
# put in front of the reader as plainly as any line, and one a pull
# confirms or refutes the same way.
"reply_rule",
*RULE_RANKED_SOURCES,
# A rule mounted on a moment (milestone 458). Not a ranker's pick, but
# not bulk either: somebody decided this rule applies at this moment, and
# that is exactly the claim a pull can confirm or refute. Its
# pull-through is the evidence for whether a mount earns its line.
"moment_rule",
# A reserved slot is a ranker's choice twice over — it ran a query AND
# decided a kind was worth guaranteeing a place. Left out, its line would
# be counted as bulk delivery and drop out of the denominator, so the one
# surface built because a record class kept losing would be the one whose
# hits nobody could confirm.
"preference_slot",
# The completion-report lookup on update_task (milestone 409 step 4). It
# runs its own query and shows only what cleared the bar — a ranker.
"report_preference",
# A rule reached through a lesson confirmed as an instance of it (#4633).
# Its own source so its pull-through reads apart from the rule's direct
# match — whether the lesson route earns its line is #4636's question.
"rule_via_lesson",
MOMENT_RULE_SOURCE,
)
+74 -2
View File
@@ -7,6 +7,7 @@ join doing the scoping: these seed real rules with hand-made vectors and stub
only the embedder, so every rule is an equally good match and the home alone
decides what comes back.
"""
import traceback
import uuid
from unittest.mock import AsyncMock, MagicMock, patch
@@ -74,10 +75,20 @@ async def homes():
async def _found(user_id: int, **scope) -> set[int]:
"""The rule ids a search finds — and a failure, not an empty set, when the
search never ran. `semantic_search_rules` fails open (an exception becomes
[]), so without `report["searched"]` an error and a scoping miss read the
same (#4958)."""
report: dict = {}
raised: list[str] = []
with patch("scribe.services.embeddings.get_embedding",
AsyncMock(return_value=QUERY_VEC)):
AsyncMock(return_value=QUERY_VEC)), \
patch("scribe.services.embeddings.logger") as log:
# Called inside the except block, so the traceback is still current.
log.warning.side_effect = lambda *_a, **_k: raised.append(traceback.format_exc())
hits = await semantic_search_rules(user_id, "anything", limit=10,
threshold=0.5, **scope)
threshold=0.5, report=report, **scope)
assert report.get("searched"), "the rule search did not run:\n" + "\n".join(raised)
return {rule.id for _score, rule in hits}
@@ -104,3 +115,64 @@ async def test_a_shared_project_brings_its_rules_to_a_collaborator(homes):
collaborator = homes["collaborator"]
assert await _found(collaborator, project_id=homes["a"]) == {homes["on_a"]}
assert await _found(collaborator, project_id=homes["b"]) == set()
async def test_an_in_scope_rule_is_found_however_many_closer_rules_are_out_of_scope():
"""#4958. Ordered through the HNSW index, the search took the ~40 nearest
chunks from every owner and filtered them AFTERWARDS — so a reader's own
rule ranked past the 40th chunk overall was never returned. Here another
user's 80 rules all sit nearer the query than the reader's one rule; the
reader's search must still find it."""
tag = uuid.uuid4().hex[:8]
async with async_session() as s:
reader = await ensure_user(s, f"rule_scope_reader_{tag}")
crowd = await ensure_user(s, f"rule_scope_crowd_{tag}")
await s.commit()
reader_id, crowd_id = reader.id, crowd.id
def near(i: int) -> list[float]:
# Cosine ~0.9996 to the query, each distinct so the graph is a graph.
vec = [1.0] + [0.0] * (EMBEDDING_DIM - 1)
vec[1 + i] = 0.02
return vec
mine_vec = [0.8, 0.6] + [0.0] * (EMBEDDING_DIM - 2) # cosine 0.8
books = []
try:
with patch("scribe.services.rulebooks._refresh_rule_embedding", MagicMock()):
book = await rulebooks_svc.create_rulebook(reader_id, "Reader's rules")
books.append((book.id, reader_id))
topic = await rulebooks_svc.create_topic(book.id, reader_id, "mine")
mine = await rulebooks_svc.create_rule(
topic.id, reader_id, "The reader's own rule", "Mine.", when_to_apply="always",
)
crowd_book = await rulebooks_svc.create_rulebook(crowd_id, "Someone else's rules")
books.append((crowd_book.id, crowd_id))
crowd_topic = await rulebooks_svc.create_topic(crowd_book.id, crowd_id, "theirs")
theirs = [
await rulebooks_svc.create_rule(
crowd_topic.id, crowd_id, f"Crowd rule {i}", "Not the reader's.",
when_to_apply="always",
)
for i in range(80)
]
async with async_session() as s:
s.add(RuleEmbedding(
rule_id=mine.id, chunk_index=0, embedding=mine_vec, chunk_text=mine.title,
chunker_version=CHUNKER_VERSION, embedding_model=EMBEDDING_MODEL,
))
for i, rule in enumerate(theirs):
s.add(RuleEmbedding(
rule_id=rule.id, chunk_index=0, embedding=near(i), chunk_text=rule.title,
chunker_version=CHUNKER_VERSION, embedding_model=EMBEDDING_MODEL,
))
await s.commit()
assert await _found(reader_id) == {mine.id}
# And the crowd still finds its own, not the reader's.
found = await _found(crowd_id)
assert mine.id not in found and len(found) == 10
finally:
# These vectors sit right beside the query every other test here uses.
for book_id, owner in books:
await rulebooks_svc.delete_rulebook(book_id, owner)
+189
View File
@@ -0,0 +1,189 @@
"""Real-Postgres crowd tests: an in-scope record is found however many nearer
records sit outside the scope (#4961; the rule search's is #4958's, in
test_integration_rule_scope).
Ordered straight off a `*_embeddings` table, the planner can walk the HNSW
index, which takes ~`hnsw.ef_search` (40) nearest chunks from every owner and
project and only then applies the scope. Each test here puts 80 out-of-scope
records nearer the query than the reader's one in-scope record, so the old
shape could not reach it when the index served the order. What a test cannot
force is the plan — the shape guard in test_search_scope_shape is what fails
on the old query whatever Postgres picks.
Rows are inserted directly with hand-made vectors and the embedder is stubbed,
so the scope alone decides what comes back.
"""
import traceback
import uuid
from unittest.mock import AsyncMock, patch
import pytest
from sqlalchemy import delete
from scribe.models import async_session
from scribe.models.embedding import (
EMBEDDING_DIM, MilestoneEmbedding, NoteEmbedding, SystemEmbedding,
)
from scribe.models.milestone import Milestone
from scribe.models.note import Note
from scribe.models.project import Project
from scribe.models.system import System
from scribe.models.user import User
from scribe.services import embeddings as emb
from scribe.services.embeddings import CHUNKER_VERSION, EMBEDDING_MODEL
from tests.helpers import ensure_user
pytestmark = [pytest.mark.integration, pytest.mark.usefixtures("_dispose_engine", "_no_embedding")]
CROWD = 80
QUERY_VEC = [1.0] + [0.0] * (EMBEDDING_DIM - 1)
# Cosine 0.8 to the query: comfortably above the bar, and behind every crowd row.
MINE_VEC = [0.8, 0.6] + [0.0] * (EMBEDDING_DIM - 2)
def _near(i: int) -> list[float]:
"""Cosine ~0.9996 to the query, each distinct so the graph is a graph."""
vec = [1.0] + [0.0] * (EMBEDDING_DIM - 1)
vec[1 + i] = 0.02
return vec
def _chunk(model, key: str, record_id: int, vec: list[float], **extra):
return model(**{key: record_id}, chunk_index=0, embedding=vec,
chunk_text=f"record {record_id}", chunker_version=CHUNKER_VERSION,
embedding_model=EMBEDDING_MODEL, **extra)
async def _found(search, user_id: int, **scope) -> set[int]:
"""The ids a search finds. Every search fails open, so a query that raised
reads as an empty result — the swallowed traceback is captured and named
in the failure instead (#4958)."""
raised: list[str] = []
with patch("scribe.services.embeddings.get_embedding",
AsyncMock(return_value=QUERY_VEC)), \
patch("scribe.services.embeddings.logger") as log:
log.warning.side_effect = lambda *_a, **_k: raised.append(traceback.format_exc())
hits = await search(user_id, "anything", limit=10, threshold=0.5, **scope)
assert not raised, "the search raised:\n" + "\n".join(raised)
return {int(record.id) for _score, record in hits}
async def _people(tag: str):
"""A reader and a crowd, each with a project of their own."""
async with async_session() as s:
reader = await ensure_user(s, f"scope_reader_{tag}")
crowd = await ensure_user(s, f"scope_crowd_{tag}")
mine = Project(user_id=reader.id, title="Reader's project")
other = Project(user_id=reader.id, title="Reader's other project")
theirs = Project(user_id=crowd.id, title="Crowd's project")
s.add_all([mine, other, theirs])
await s.commit()
return reader.id, crowd.id, mine.id, other.id, theirs.id
async def _cleanup(*user_ids: int) -> None:
# These vectors sit right beside the query every other vector test uses.
async with async_session() as s:
for model in (Note, Milestone, System):
await s.execute(delete(model).where(model.user_id.in_(user_ids)))
await s.execute(delete(Project).where(Project.user_id.in_(user_ids)))
await s.execute(delete(User).where(User.id.in_(user_ids)))
await s.commit()
@pytest.mark.asyncio
async def test_a_note_is_found_behind_another_users_nearer_notes():
reader, crowd, mine_p, _other, theirs_p = await _people(uuid.uuid4().hex[:8])
try:
async with async_session() as s:
mine = Note(user_id=reader, project_id=mine_p, title="mine", body="mine")
crowd_notes = [Note(user_id=crowd, project_id=theirs_p, title=f"crowd {i}",
body="not the reader's") for i in range(CROWD)]
s.add_all([mine, *crowd_notes])
await s.flush()
s.add(_chunk(NoteEmbedding, "note_id", mine.id, MINE_VEC, user_id=reader))
s.add_all(_chunk(NoteEmbedding, "note_id", n.id, _near(i), user_id=crowd)
for i, n in enumerate(crowd_notes))
await s.commit()
mine_id = mine.id
assert await _found(emb.semantic_search_notes, reader) == {mine_id}
found = await _found(emb.semantic_search_notes, crowd)
assert mine_id not in found and len(found) == 10
finally:
await _cleanup(reader, crowd)
@pytest.mark.asyncio
async def test_a_note_is_found_behind_the_readers_own_other_project():
"""The single-user case, and the one the hook arms live in: a project-scoped
search where the reader's OTHER projects hold the nearer chunks."""
reader, crowd, mine_p, other_p, _theirs = await _people(uuid.uuid4().hex[:8])
try:
async with async_session() as s:
mine = Note(user_id=reader, project_id=mine_p, title="mine", body="mine")
elsewhere = [Note(user_id=reader, project_id=other_p, title=f"elsewhere {i}",
body="another project") for i in range(CROWD)]
s.add_all([mine, *elsewhere])
await s.flush()
s.add(_chunk(NoteEmbedding, "note_id", mine.id, MINE_VEC, user_id=reader))
s.add_all(_chunk(NoteEmbedding, "note_id", n.id, _near(i), user_id=reader)
for i, n in enumerate(elsewhere))
await s.commit()
mine_id = mine.id
assert await _found(emb.semantic_search_notes, reader,
project_id=mine_p) == {mine_id}
finally:
await _cleanup(reader, crowd)
@pytest.mark.asyncio
async def test_a_milestone_is_found_behind_another_users_nearer_plans():
reader, crowd, mine_p, _other, theirs_p = await _people(uuid.uuid4().hex[:8])
try:
async with async_session() as s:
mine = Milestone(user_id=reader, project_id=mine_p, title="my plan")
plans = [Milestone(user_id=crowd, project_id=theirs_p, title=f"plan {i}")
for i in range(CROWD)]
s.add_all([mine, *plans])
await s.flush()
s.add(_chunk(MilestoneEmbedding, "milestone_id", mine.id, MINE_VEC))
s.add_all(_chunk(MilestoneEmbedding, "milestone_id", m.id, _near(i))
for i, m in enumerate(plans))
await s.commit()
mine_id = mine.id
assert await _found(emb.semantic_search_milestones, reader) == {mine_id}
assert await _found(emb.semantic_search_milestones, reader,
project_id=mine_p) == {mine_id}
found = await _found(emb.semantic_search_milestones, crowd)
assert mine_id not in found and len(found) == 10
finally:
await _cleanup(reader, crowd)
@pytest.mark.asyncio
async def test_a_system_is_found_behind_another_users_nearer_charters():
reader, crowd, mine_p, _other, theirs_p = await _people(uuid.uuid4().hex[:8])
try:
async with async_session() as s:
mine = System(user_id=reader, project_id=mine_p, name="mine",
description="the reader's area")
areas = [System(user_id=crowd, project_id=theirs_p, name=f"area {i}",
description="someone else's area") for i in range(CROWD)]
s.add_all([mine, *areas])
await s.flush()
s.add(_chunk(SystemEmbedding, "system_id", mine.id, MINE_VEC))
s.add_all(_chunk(SystemEmbedding, "system_id", a.id, _near(i))
for i, a in enumerate(areas))
await s.commit()
mine_id = mine.id
assert await _found(emb.semantic_search_systems, reader) == {mine_id}
assert await _found(emb.semantic_search_systems, reader,
project_id=mine_p) == {mine_id}
found = await _found(emb.semantic_search_systems, crowd)
assert mine_id not in found and len(found) == 10
finally:
await _cleanup(reader, crowd)
+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)
+71
View File
@@ -0,0 +1,71 @@
"""The specs are the registry (milestone 456 step 5).
Three lists used to be kept by hand beside the arms — the tuning table, the
ranked rows of the telemetry registry, and the rule-usage denominator — with
tests to make them agree. They are now read off the pipeline's specs. These
pin the derivation itself: every arm is tunable, every ranked source is
measured with the declaration its spec carries, and the rule-usage
denominator counts exactly the rule surfaces a ranker chose.
"""
from __future__ import annotations
from scribe.services import retrieval_pipeline as rp
from scribe.services.retrieval_registry import POINTS, UNBIDDEN
from scribe.services.retrieval_surfaces import SURFACES
from scribe.services.rule_usage import RANKED_SOURCES, is_ambient
def test_every_arm_is_one_tunable_surface_in_the_order_settings_lists_them():
assert [arm.source for arm in rp.TUNED_ARMS] == list(SURFACES)
for arm in rp.TUNED_ARMS:
assert arm.tuning is not None, f"{arm.source} has no floor or budget"
# The join key: the surface a tuner moves and the rows it is judged by
# must name the same arm.
assert arm.tuning.name == arm.source
assert SURFACES[arm.source] is arm.tuning
def test_every_ranked_source_is_measured_as_its_spec_declares():
sources = [spec.source for spec in rp.RANKED]
assert len(sources) == len(set(sources)), f"a source is declared twice: {sources}"
for spec in rp.RANKED:
d = spec.declared
assert d is not None and d.what.strip(), f"{spec.source} declares nothing"
point = POINTS[spec.source]
assert point.kind == UNBIDDEN
assert point.what == d.what
assert point.fixed_query == d.fixed_query
# A quiet source must say why, and only a quiet source may (#2475).
assert point.expects_traffic == (not d.quiet_because)
assert point.quiet_because == d.quiet_because
def test_the_slots_are_measured_but_never_tuned():
for slot in (rp.PREFERENCE_SLOT, *rp.NOTE_SLOTS):
assert slot.source in POINTS
assert slot.source not in SURFACES
def test_the_rule_denominator_is_every_rule_surface_a_ranker_chose():
rule_arms = {arm.source for arm in rp.RULE_ARMS}
assert rule_arms <= set(RANKED_SOURCES)
assert set(RANKED_SOURCES) == (
rule_arms | {rp.PREFERENCE_SLOT_SOURCE, rp.VIA_LESSON_SOURCE,
rp.MOMENT_RULE_SOURCE}
)
# And a notes arm is not a RULE surface: its lines are not rules, so a
# rule pull cannot confirm them.
for arm in rp.NOTE_ARMS:
assert is_ambient(arm.source)
def test_an_arm_added_without_tuning_would_be_caught():
"""Rule 167: the first test has to be able to bite. SURFACES skips an arm
with no tuning rather than raising, so an arm added without one would
silently be untunable — and the order-equality assertion is what notices,
as this replays with one such arm appended."""
bare = rp.RuleArm("x_rule", band=False, compact_tail=False,
checkpoint=False, preference_slot=False)
arms = (*rp.TUNED_ARMS, bare)
derived = {a.tuning.name: a.tuning for a in arms if a.tuning is not None}
assert [a.source for a in arms] != list(derived)
+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")
+49 -15
View File
@@ -18,6 +18,7 @@ import pytest
from scribe.services import plugin_context as pc
from scribe.services import rule_usage
from scribe.services.retrieval_registry import POINTS
from scribe.services.retrieval_pipeline import RuleMoment, RuleResult
# Bound before conftest's autouse stub replaces the module attribute.
_REAL = pc._rules_via_lessons
@@ -123,23 +124,56 @@ async def test_a_failure_brings_no_rule_and_raises_nothing():
@pytest.mark.asyncio
async def test_the_step_runs_only_when_the_arm_ran_and_never_leaks_its_key():
async def test_the_step_follows_the_arm_on_its_query_and_skips_what_it_named():
"""One composer for all three rule builders (milestone 456 step 6): the
via-lesson step searches with the query the arm searched with, skips the
session's ledger and every rule the arm's band named, and its lines and
fresh ids join the arm's — while `shown_rule_ids` stays the direct band."""
step = AsyncMock(return_value=(["Standing rule … Reached through lesson #41"], [7]))
direct = RuleResult(lines=["direct line"], rule_ids=[5], shown_rule_ids=[5, 6])
moment = RuleMoment(user_id=1, query="q", project_id=2, where="here",
exclude=frozenset({3, 6}))
with patch.object(pc, "_rules_via_lessons", step), \
patch.object(pc.rp, "run_rule_arm", AsyncMock(return_value=direct)):
out = await pc._rule_moment(pc.rp.PROMPT_RULE, moment, floor=0.5, budget=3)
assert step.await_args.args[1] == "q"
assert step.await_args.kwargs["skip"] == {3, 5, 6}
assert out.rule_ids == [5, 7]
assert out.lines == ["direct line", "Standing rule … Reached through lesson #41"]
assert out.shown_rule_ids == [5, 6]
@pytest.mark.asyncio
async def test_a_blank_prompt_runs_neither_the_arm_nor_the_step():
step = AsyncMock(return_value=([], []))
with patch.object(pc, "_rules_via_lessons", step):
idle = await pc._add_rules_via_lessons(
1, {"context": "", "rule_ids": []}, project_id=2,
exclude_rule_ids=[3], held_rule_ids=[], where="here",
)
ran = await pc._add_rules_via_lessons(
1, {"context": "direct line", "rule_ids": [5], "_via_query": "q"},
project_id=2, exclude_rule_ids=[3], held_rule_ids=[], where="here",
)
assert idle == {"context": "", "rule_ids": []}
assert step.await_count == 1
assert step.await_args.kwargs["skip"] == {3, 5}
assert "_via_query" not in ran
assert ran["rule_ids"] == [5, 7]
assert ran["context"].startswith("direct line\n")
out = await pc.build_prompt_rule_hint(1, " ")
assert out == {"context": "", "rule_ids": []}
step.assert_not_awaited()
def test_only_the_composer_runs_a_rule_arm_or_the_via_lesson_step():
"""Structural (rule 167): every rule builder goes through `_rule_moment`,
so a builder that runs the arm or the via-lesson step by hand again is the
copy this step removed coming back — and a count is what lets it fail."""
import ast
from pathlib import Path
tree = ast.parse(Path("src/scribe/services/plugin_context.py").read_text())
callers: dict[str, set[str]] = {}
for fn in ast.walk(tree):
if not isinstance(fn, (ast.FunctionDef, ast.AsyncFunctionDef)):
continue
for call in ast.walk(fn):
if isinstance(call, ast.Call):
name = getattr(call.func, "attr", None) or getattr(call.func, "id", None)
if name in ("run_rule_arm", "_rules_via_lessons", "_rule_moment"):
callers.setdefault(name, set()).add(fn.name)
assert callers["run_rule_arm"] == {"_rule_moment"}
assert callers["_rules_via_lessons"] == {"_rule_moment"}
assert callers["_rule_moment"] == {
"build_prompt_rule_hint", "build_tool_rule_hint", "build_write_path_hint",
}
def test_the_source_is_ranked_and_registered():
+94
View File
@@ -0,0 +1,94 @@
"""Every semantic search scopes first, then ranks (#4958, #4961).
Ordered straight off a `*_embeddings` table, the planner walks the HNSW index:
it takes ~`hnsw.ef_search` (40) nearest chunks from every owner and project
and filters them afterwards, so an in-scope record behind 40 nearer ones the
caller cannot see is silently dropped. The fix is one shape, `_rank_scoped`,
and these pin that every search goes through it — the rule search was fixed
alone first, and the other three kept the fault because nothing said they
were the same query.
The integration lane's crowd tests show the behaviour on a real index; these
are the guard that fails on the old shape whatever plan Postgres picks.
"""
from __future__ import annotations
import ast
from pathlib import Path
from sqlalchemy.dialects import postgresql
from scribe.models.embedding import NoteEmbedding
from scribe.models.note import Note
from scribe.services import embeddings as emb
from tests.helpers import compiled_sql
_SOURCE = Path("src/scribe/services/embeddings.py")
def _searches() -> dict[str, ast.AsyncFunctionDef]:
tree = ast.parse(_SOURCE.read_text())
return {
node.name: node for node in tree.body
if isinstance(node, ast.AsyncFunctionDef)
and node.name.startswith("semantic_search_")
}
def _calls(fn: ast.AST) -> list[str]:
return [
getattr(n.func, "attr", None) or getattr(n.func, "id", None)
for n in ast.walk(fn) if isinstance(n, ast.Call)
]
def test_every_semantic_search_ranks_through_the_scoped_shape():
searches = _searches()
# The four corpora with an embedding table: if one is renamed or a fifth
# is added, this has to be looked at rather than silently passing.
assert set(searches) == {
"semantic_search_notes", "semantic_search_rules",
"semantic_search_milestones", "semantic_search_systems",
}
for name, fn in searches.items():
calls = _calls(fn)
assert "_rank_scoped" in calls, f"{name} does not rank through _rank_scoped"
assert "_scoped_chunks" in calls, f"{name} does not build its scope as chunks"
# Its own ORDER BY is the old shape: a distance ordered on the
# embedding table, which the index serves before the scope applies.
assert "order_by" not in calls, f"{name} orders by distance itself"
def test_the_old_shape_would_be_caught():
"""Rule 167: replay the guard against the query every search used to run."""
old = ast.parse(
"async def semantic_search_x():\n"
" rows = await session.execute(select(Note, distance).select_from(E)"
".join(Note, E.note_id == Note.id).where(scope)"
".order_by(distance).limit(k))\n"
).body[0]
calls = _calls(old)
assert "_rank_scoped" not in calls and "order_by" in calls
def test_the_scope_is_materialized_and_the_order_is_over_it():
# Column against column, so there is no vector literal to inline: the
# shape is the subject here, not the query.
distance = NoteEmbedding.embedding.cosine_distance(NoteEmbedding.embedding)
scoped = (
emb._scoped_chunks(NoteEmbedding, NoteEmbedding.note_id, distance)
.join(Note, NoteEmbedding.note_id == Note.id)
.where(Note.user_id == 1)
)
sql = compiled_sql(
emb._rank_scoped(Note, scoped, name="scoped_x", limit=5),
dialect=postgresql.dialect(),
)
# MATERIALIZED is what keeps the planner from inlining the CTE and walking
# the index again; the ORDER BY must name the CTE's column, not the table's.
assert "scoped_x AS MATERIALIZED" in sql
assert "ORDER BY scoped_x.distance" in sql
assert "note_embeddings.embedding <=>" in sql.split("ORDER BY")[0]
# The scope sits inside the CTE, where the ranking draws from.
cte_body = sql.split("AS MATERIALIZED", 1)[1].split("SELECT notes", 1)[0]
assert "notes.user_id" in cte_body
+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
)
+4 -4
View File
@@ -22,6 +22,7 @@ import pytest
from scribe.services import plugin_context as pc
from scribe.services import system_rulings as sr
from scribe.services.retrieval_pipeline import RuleResult
# Bound at import, before conftest's autouse stub replaces the module attribute.
from scribe.services.system_rulings import rulings_for_paths as real_rulings_for_paths
from tests.helpers import writepath_cfg
@@ -198,8 +199,8 @@ RULING_LINE = "> Rulings for `src/a.py` — the operator's decisions about A (Sy
@pytest.mark.asyncio
async def test_the_tool_arm_leads_with_rulings_and_hands_back_what_it_showed():
rulings = AsyncMock(return_value={"lines": [RULING_LINE], "system_ids": [3]})
with patch.object(pc, "_tool_rule_hint", AsyncMock(return_value={
"context": "> a rule line", "rule_ids": [8], "checkpoint": {}})), \
with patch.object(pc, "_rule_moment", AsyncMock(return_value=RuleResult(
lines=["> a rule line"], rule_ids=[8]))), \
patch.object(pc, "get_writepath_config", AsyncMock(return_value=writepath_cfg())), \
patch.object(pc.system_rulings_svc, "rulings_for_paths", rulings):
out = await pc.build_tool_rule_hint(
@@ -215,8 +216,7 @@ async def test_the_tool_arm_leads_with_rulings_and_hands_back_what_it_showed():
@pytest.mark.asyncio
async def test_the_tool_arm_adds_no_key_when_no_ruling_was_shown():
with patch.object(pc, "_tool_rule_hint", AsyncMock(return_value={
"context": "", "rule_ids": [], "checkpoint": {}})), \
with patch.object(pc, "_rule_moment", AsyncMock(return_value=RuleResult())), \
patch.object(pc, "get_writepath_config", AsyncMock(return_value=writepath_cfg())):
out = await pc.build_tool_rule_hint(1, "Bash", "ls", project_id=2)
assert "ruling_system_ids" not in out
+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) ------------------