feat(embeddings): best-chunk-per-note on every retrieval surface (#280 step 4)
CI & Build / Plugin hooks (push) Failing after 1s
CI & Build / Python lint (push) Failing after 3s
CI & Build / integration (push) Successful in 17s
CI & Build / TypeScript typecheck (push) Successful in 32s
CI & Build / Python tests (push) Successful in 46s
CI & Build / Build & push image (push) Skipped
CI & Build / Plugin hooks (push) Failing after 1s
CI & Build / Python lint (push) Failing after 3s
CI & Build / integration (push) Successful in 17s
CI & Build / TypeScript typecheck (push) Successful in 32s
CI & Build / Python tests (push) Successful in 46s
CI & Build / Build & push image (push) Skipped
A note's relevance is now its best chunk's similarity, everywhere: - semantic_search_notes keeps the indexed raw-distance top-k and over-fetches chunk rows (x4, composing with the x3 supersession over-fetch), then collapses to first-appearance-per-note — rows arrive distance-ordered, so first is best. Every ranked consumer (MCP/REST search, Browse, auto-inject, write-path, gate) inherits through the one function. - list_notes semantic q swaps its join for a correlated MIN-distance subquery — the join would have repeated a long note once per matching chunk and made total count chunks. - the duplicate report groups its self-join by note pair on MIN(distance): pair similarity = closest chunk pair, and the < join now also drops cross-chunk self-pairs that would flag every long note against itself. - the write gate queries once per chunk of the candidate (capped at 8), so a note duplicating an existing record in ONE SECTION is caught — the whole-document query diluted exactly the section that mattered. Integration test now seeds a two-chunk note and pins the collapse against real pgvector. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01UaYUaouG9jjhATyuxCKrQs
This commit is contained in:
@@ -69,6 +69,12 @@ _SEMANTIC_THRESHOLD = 0.90
|
|||||||
# structural signals cannot see.
|
# structural signals cannot see.
|
||||||
_SNIPPET_SEMANTIC_THRESHOLD = 0.96
|
_SNIPPET_SEMANTIC_THRESHOLD = 0.96
|
||||||
|
|
||||||
|
# The gate queries per CHUNK of the candidate (#280) — this caps how many
|
||||||
|
# searches one save may cost. Eight chunks ≈ five thousand words of candidate;
|
||||||
|
# a duplicate hiding past that is the duplicate report's job to find, not a
|
||||||
|
# reason to stall the write path.
|
||||||
|
_GATE_MAX_CHUNKS = 8
|
||||||
|
|
||||||
|
|
||||||
@dataclass
|
@dataclass
|
||||||
class DuplicateMatch:
|
class DuplicateMatch:
|
||||||
@@ -258,41 +264,45 @@ async def find_duplicate_note(
|
|||||||
|
|
||||||
# --- Signal 3: semantic similarity (only with a substantial body) ---
|
# --- Signal 3: semantic similarity (only with a substantial body) ---
|
||||||
if body and len(body.strip()) >= _MIN_BODY_FOR_SEMANTIC:
|
if body and len(body.strip()) >= _MIN_BODY_FOR_SEMANTIC:
|
||||||
# Built by the SAME function the corpus was embedded with. This one is
|
# Query with the SAME chunker the corpus was embedded with (#280). This
|
||||||
# the copy that mattered most and was easiest to miss: it is a QUERY
|
# was the copy that mattered most and was easiest to miss: these are
|
||||||
# document, compared against embedded ones. Shaped differently from the
|
# QUERY documents, compared against embedded ones — shaped differently
|
||||||
# corpus it searches, the gate degrades silently — it still returns
|
# from the corpus, the gate degrades silently. Chunking also makes the
|
||||||
# neighbours, just less apt ones, and no signal says the query and the
|
# gate see what the whole-document query diluted: a long candidate that
|
||||||
# index stopped agreeing (found by the guard in test_embedding_text).
|
# duplicates an existing record IN ONE SECTION now matches on that
|
||||||
query = embeddings_svc.embedding_text(title, body)
|
# section. Capped so one pathological paste can't turn a save into
|
||||||
# Scope the semantic check the same way as the title check: a record in
|
# dozens of searches — a duplicate past the cap is the duplicate
|
||||||
# project P compares only to P; a project-less (orphan) record compares
|
# report's job, not the gate's.
|
||||||
# only to other orphans (orphan_only), NOT across every project — without
|
for query in embeddings_svc.chunk_document(title, body)[:_GATE_MAX_CHUNKS]:
|
||||||
# this, semantic_search_notes applies no project filter when project_id
|
# Scope the semantic check the same way as the title check: a record
|
||||||
# is None and would match an orphan note against any project's notes.
|
# in project P compares only to P; a project-less (orphan) record
|
||||||
hits = await embeddings_svc.semantic_search_notes(
|
# compares only to other orphans (orphan_only), NOT across every
|
||||||
user_id, query, project_id=project_id, is_task=is_task,
|
# project — without this, semantic_search_notes applies no project
|
||||||
orphan_only=(project_id is None),
|
# filter when project_id is None and would match an orphan note
|
||||||
limit=3,
|
# against any project's notes.
|
||||||
threshold=(_SNIPPET_SEMANTIC_THRESHOLD
|
hits = await embeddings_svc.semantic_search_notes(
|
||||||
if note_type == SNIPPET_NOTE_TYPE else _SEMANTIC_THRESHOLD),
|
user_id, query, project_id=project_id, is_task=is_task,
|
||||||
# Owner-only, deliberately: this gate BLOCKS a create and tells the
|
orphan_only=(project_id is None),
|
||||||
# caller to update the match instead. Matching someone else's record
|
limit=3,
|
||||||
# would refuse their write and point them at something they may not
|
threshold=(_SNIPPET_SEMANTIC_THRESHOLD
|
||||||
# be able to edit.
|
if note_type == SNIPPET_NOTE_TYPE else _SEMANTIC_THRESHOLD),
|
||||||
scope="own",
|
# Owner-only, deliberately: this gate BLOCKS a create and tells
|
||||||
# NOT demoted by supersession (#278). A superseded record is still a
|
# the caller to update the match instead. Matching someone
|
||||||
# duplicate of what you are about to write — the claim is that it is
|
# else's record would refuse their write and point them at
|
||||||
# no longer CURRENT, not that it is gone. Demoting it here would let
|
# something they may not be able to edit.
|
||||||
# the same note be recorded a second time, and the second copy would
|
scope="own",
|
||||||
# be the one nothing warns about.
|
# NOT demoted by supersession (#278). A superseded record is
|
||||||
demote_superseded=False,
|
# still a duplicate of what you are about to write — the claim
|
||||||
)
|
# is that it is no longer CURRENT, not that it is gone. Demoting
|
||||||
for score, note in hits:
|
# it here would let the same note be recorded a second time, and
|
||||||
# semantic_search_notes doesn't filter note_type — enforce it here so
|
# the second copy would be the one nothing warns about.
|
||||||
# a note doesn't shadow a task of the same wording, etc.
|
demote_superseded=False,
|
||||||
if note.note_type == note_type:
|
)
|
||||||
return DuplicateMatch(note.id, note.title, round(score, 3), "semantic")
|
for score, note in hits:
|
||||||
|
# semantic_search_notes doesn't filter note_type — enforce it
|
||||||
|
# here so a note doesn't shadow a task of the same wording, etc.
|
||||||
|
if note.note_type == note_type:
|
||||||
|
return DuplicateMatch(note.id, note.title, round(score, 3), "semantic")
|
||||||
|
|
||||||
return None
|
return None
|
||||||
|
|
||||||
@@ -502,15 +512,22 @@ async def find_duplicate_records(
|
|||||||
left_note = aliased(Note, name="left_note")
|
left_note = aliased(Note, name="left_note")
|
||||||
right_note = aliased(Note, name="right_note")
|
right_note = aliased(Note, name="right_note")
|
||||||
distance = left.embedding.cosine_distance(right.embedding)
|
distance = left.embedding.cosine_distance(right.embedding)
|
||||||
|
# Chunk grain (#280): a note-pair's similarity is its closest CHUNK pair —
|
||||||
|
# two records duplicate each other where their most similar sections do,
|
||||||
|
# which is the honest definition when one section of a long note restates
|
||||||
|
# another record. GROUP BY collapses the chunk cross-product to one row
|
||||||
|
# per note pair.
|
||||||
|
best = func.min(distance)
|
||||||
|
|
||||||
pairs: list[tuple[int, int, float]] = []
|
pairs: list[tuple[int, int, float]] = []
|
||||||
try:
|
try:
|
||||||
async with async_session() as session:
|
async with async_session() as session:
|
||||||
stmt = (
|
stmt = (
|
||||||
select(left.note_id, right.note_id, distance.label("distance"))
|
select(left.note_id, right.note_id, best.label("distance"))
|
||||||
.select_from(left)
|
.select_from(left)
|
||||||
# `<` not `!=`: each unordered pair exactly once, and it drops
|
# `<` not `!=`: each unordered pair exactly once, and it drops
|
||||||
# the self-pair (distance 0) that would otherwise dominate.
|
# the self-pairs (including cross-chunk self-pairs, which would
|
||||||
|
# otherwise flag every multi-chunk note against itself).
|
||||||
.join(right, left.note_id < right.note_id)
|
.join(right, left.note_id < right.note_id)
|
||||||
.join(left_note, left_note.id == left.note_id)
|
.join(left_note, left_note.id == left.note_id)
|
||||||
.join(right_note, right_note.id == right.note_id)
|
.join(right_note, right_note.id == right.note_id)
|
||||||
@@ -523,9 +540,10 @@ async def find_duplicate_records(
|
|||||||
# report is bounded by what merge can actually act on.
|
# report is bounded by what merge can actually act on.
|
||||||
left_note.user_id == user_id,
|
left_note.user_id == user_id,
|
||||||
right_note.user_id == user_id,
|
right_note.user_id == user_id,
|
||||||
distance <= max_distance,
|
|
||||||
)
|
)
|
||||||
.order_by(distance.asc())
|
.group_by(left.note_id, right.note_id)
|
||||||
|
.having(best <= max_distance)
|
||||||
|
.order_by(best.asc())
|
||||||
.limit(max(1, limit))
|
.limit(max(1, limit))
|
||||||
)
|
)
|
||||||
rows = list((await session.execute(stmt)).all())
|
rows = list((await session.execute(stmt)).all())
|
||||||
|
|||||||
@@ -124,6 +124,14 @@ _SUPERSESSION_PENALTY = 0.05
|
|||||||
# whose neighbours sit ~0.01-0.02 apart.
|
# whose neighbours sit ~0.01-0.02 apart.
|
||||||
_SUPERSESSION_OVERFETCH = 3
|
_SUPERSESSION_OVERFETCH = 3
|
||||||
|
|
||||||
|
# Chunk rows fetched per requested result (#280). The HNSW top-k runs at CHUNK
|
||||||
|
# grain — several chunks of one strong note can occupy consecutive ranks, and
|
||||||
|
# each collapses into a single result. Four ranks of headroom per result keeps
|
||||||
|
# the top-k indexed while making it effectively impossible for collapsing to
|
||||||
|
# starve the result list: that would need every requested note to be shadowed
|
||||||
|
# by four chunks of notes ranked above it.
|
||||||
|
_CHUNK_OVERFETCH = 4
|
||||||
|
|
||||||
|
|
||||||
async def _apply_supersession_penalty(
|
async def _apply_supersession_penalty(
|
||||||
scored: list[tuple[float, "Note"]], limit: int
|
scored: list[tuple[float, "Note"]], limit: int
|
||||||
@@ -524,7 +532,9 @@ async def semantic_search_notes(
|
|||||||
# penalty far smaller than the window's score spread, that case
|
# penalty far smaller than the window's score spread, that case
|
||||||
# needs the true answer to be more than _SUPERSESSION_OVERFETCH
|
# needs the true answer to be more than _SUPERSESSION_OVERFETCH
|
||||||
# ranks down, which no observed query comes close to.
|
# ranks down, which no observed query comes close to.
|
||||||
fetch = limit * _SUPERSESSION_OVERFETCH if demote_superseded else limit
|
fetch = limit * _CHUNK_OVERFETCH * (
|
||||||
|
_SUPERSESSION_OVERFETCH if demote_superseded else 1
|
||||||
|
)
|
||||||
stmt = (
|
stmt = (
|
||||||
stmt.where(distance <= max_distance)
|
stmt.where(distance <= max_distance)
|
||||||
.order_by(distance.asc())
|
.order_by(distance.asc())
|
||||||
@@ -535,8 +545,19 @@ async def semantic_search_notes(
|
|||||||
logger.warning("Failed to query note embeddings", exc_info=True)
|
logger.warning("Failed to query note embeddings", exc_info=True)
|
||||||
return []
|
return []
|
||||||
|
|
||||||
# Recover similarity (1 - distance) and preserve the highest-first contract.
|
# Collapse chunk rows to BEST-CHUNK-PER-NOTE (#280): rows arrive ordered by
|
||||||
scored = [(1.0 - float(dist), note) for note, dist in rows]
|
# distance, so the first appearance of a note is its best chunk and later
|
||||||
|
# appearances are the same note matched less well. A note's relevance IS
|
||||||
|
# its best section's relevance — a query about one topic of a long record
|
||||||
|
# must find that record as strongly as if the topic were the whole record.
|
||||||
|
# Recover similarity (1 - distance); order stays highest-first.
|
||||||
|
scored: list[tuple[float, Note]] = []
|
||||||
|
seen: set[int] = set()
|
||||||
|
for note, dist in rows:
|
||||||
|
if int(note.id) in seen:
|
||||||
|
continue
|
||||||
|
seen.add(int(note.id))
|
||||||
|
scored.append((1.0 - float(dist), note))
|
||||||
if not demote_superseded:
|
if not demote_superseded:
|
||||||
return scored[:limit]
|
return scored[:limit]
|
||||||
return await _apply_supersession_penalty(scored, limit)
|
return await _apply_supersession_penalty(scored, limit)
|
||||||
|
|||||||
@@ -236,15 +236,25 @@ async def list_notes(
|
|||||||
if query_vec is not None:
|
if query_vec is not None:
|
||||||
from scribe.models.embedding import NoteEmbedding
|
from scribe.models.embedding import NoteEmbedding
|
||||||
from scribe.services.embeddings import INTERACTIVE_SEARCH_THRESHOLD
|
from scribe.services.embeddings import INTERACTIVE_SEARCH_THRESHOLD
|
||||||
distance = NoteEmbedding.embedding.cosine_distance(query_vec)
|
# Best-chunk-per-note as a correlated MIN, not a join (#280):
|
||||||
sem_filter = distance <= (1.0 - INTERACTIVE_SEARCH_THRESHOLD)
|
# a note stores one embedding row PER CHUNK, so the plain join
|
||||||
query = query.join(
|
# this used to be would repeat a long note once per matching
|
||||||
NoteEmbedding, NoteEmbedding.note_id == Note.id
|
# chunk — duplicated list rows and a total that counts chunks.
|
||||||
).where(sem_filter)
|
# This query is filter-heavy and paginated, never HNSW-bound,
|
||||||
count_query = count_query.join(
|
# so the scalar subquery costs what the join did.
|
||||||
NoteEmbedding, NoteEmbedding.note_id == Note.id
|
best_distance = (
|
||||||
).where(sem_filter)
|
select(
|
||||||
semantic_order = distance.asc()
|
func.min(
|
||||||
|
NoteEmbedding.embedding.cosine_distance(query_vec)
|
||||||
|
)
|
||||||
|
)
|
||||||
|
.where(NoteEmbedding.note_id == Note.id)
|
||||||
|
.scalar_subquery()
|
||||||
|
)
|
||||||
|
sem_filter = best_distance <= (1.0 - INTERACTIVE_SEARCH_THRESHOLD)
|
||||||
|
query = query.where(sem_filter)
|
||||||
|
count_query = count_query.where(sem_filter)
|
||||||
|
semantic_order = best_distance.asc()
|
||||||
else:
|
else:
|
||||||
terms = _strip_type_nouns(q)
|
terms = _strip_type_nouns(q)
|
||||||
for term in terms:
|
for term in terms:
|
||||||
|
|||||||
@@ -136,6 +136,36 @@ def test_a_monster_single_paragraph_is_hard_split_not_dropped():
|
|||||||
assert total_words == 2000
|
assert total_words == 2000
|
||||||
|
|
||||||
|
|
||||||
|
# --- the read path: best chunk wins (#280 step 4) ----------------------------
|
||||||
|
|
||||||
|
|
||||||
|
async def test_search_collapses_chunk_rows_to_best_chunk_per_note():
|
||||||
|
"""Rows arrive at CHUNK grain ordered by distance; a note appearing via
|
||||||
|
several chunks must come back ONCE, scored by its best chunk — otherwise a
|
||||||
|
long record fills the top-k with copies of itself."""
|
||||||
|
from unittest.mock import AsyncMock, MagicMock, patch
|
||||||
|
|
||||||
|
from scribe.services import embeddings as emb
|
||||||
|
|
||||||
|
note_a, note_b = MagicMock(id=1), MagicMock(id=2)
|
||||||
|
rows = [(note_a, 0.10), (note_b, 0.20), (note_a, 0.25), (note_a, 0.30)]
|
||||||
|
result = MagicMock()
|
||||||
|
result.all.return_value = rows
|
||||||
|
session, ctx = _session_ctx()
|
||||||
|
session.execute = AsyncMock(return_value=result)
|
||||||
|
|
||||||
|
with (
|
||||||
|
patch.object(emb, "async_session", return_value=ctx),
|
||||||
|
patch.object(emb, "get_embedding", AsyncMock(return_value=[0.0] * 384)),
|
||||||
|
):
|
||||||
|
out = await emb.semantic_search_notes(
|
||||||
|
1, "a query", limit=8, demote_superseded=False
|
||||||
|
)
|
||||||
|
|
||||||
|
assert [note.id for _s, note in out] == [1, 2]
|
||||||
|
assert out[0][0] == 1.0 - 0.10 # the BEST chunk's score, not a later one
|
||||||
|
|
||||||
|
|
||||||
# --- the write path: one row per chunk (#280 step 3) -------------------------
|
# --- the write path: one row per chunk (#280 step 3) -------------------------
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -68,7 +68,11 @@ async def seeded():
|
|||||||
await s.flush()
|
await s.flush()
|
||||||
# query vector will be [1,0,0,...]; near ~ identical (sim≈1.0),
|
# query vector will be [1,0,0,...]; near ~ identical (sim≈1.0),
|
||||||
# far is orthogonal (sim≈0.0 -> filtered by the default threshold).
|
# far is orthogonal (sim≈0.0 -> filtered by the default threshold).
|
||||||
|
# near gets a SECOND, weaker chunk (sim≈0.6) — the collapse to
|
||||||
|
# best-chunk-per-note (#280) is under test: near must come back once,
|
||||||
|
# at its best chunk's score, not twice.
|
||||||
s.add(_emb(near.id, user.id, 0, _vec(1.0)))
|
s.add(_emb(near.id, user.id, 0, _vec(1.0)))
|
||||||
|
s.add(_emb(near.id, user.id, 1, _vec(0.6, 0.8)))
|
||||||
s.add(_emb(far.id, user.id, 0, _vec(0.0, 1.0)))
|
s.add(_emb(far.id, user.id, 0, _vec(0.0, 1.0)))
|
||||||
await s.commit()
|
await s.commit()
|
||||||
ids = (user.id, near.id, far.id)
|
ids = (user.id, near.id, far.id)
|
||||||
@@ -96,6 +100,9 @@ async def test_semantic_search_ranks_and_thresholds_via_pgvector(seeded):
|
|||||||
assert near_id in ids
|
assert near_id in ids
|
||||||
assert far_id not in ids
|
assert far_id not in ids
|
||||||
assert ids[0] == near_id
|
assert ids[0] == near_id
|
||||||
|
# Chunk collapse (#280): near has TWO chunk rows above the floor (sim≈1.0
|
||||||
|
# and ≈0.6) and must appear exactly once, at its best chunk's score.
|
||||||
|
assert ids.count(near_id) == 1
|
||||||
top_score = results[0][0]
|
top_score = results[0][0]
|
||||||
assert top_score == pytest.approx(1.0, abs=1e-3)
|
assert top_score == pytest.approx(1.0, abs=1e-3)
|
||||||
|
|
||||||
|
|||||||
@@ -68,6 +68,33 @@ async def test_semantic_match_when_body_substantial():
|
|||||||
assert dup.similarity == 0.93
|
assert dup.similarity == 0.93
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_gate_catches_a_duplicate_hiding_in_a_later_chunk():
|
||||||
|
"""The capability #280 adds to the gate: a long candidate that duplicates
|
||||||
|
an existing record in ONE SECTION is caught, where the whole-document
|
||||||
|
query this replaces diluted exactly the section that mattered. The gate
|
||||||
|
queries once per chunk and any chunk's hit blocks."""
|
||||||
|
para = ("This section restates an existing decision in enough words to be "
|
||||||
|
"a real paragraph of content for the chunker to keep. ") * 4
|
||||||
|
body = "\n\n".join(f"## Topic {i}\n\n{para} (t{i})" for i in range(8))
|
||||||
|
|
||||||
|
from scribe.services.embeddings import chunk_document
|
||||||
|
n_chunks = len(chunk_document("Title", body))
|
||||||
|
assert n_chunks > 1, "test body must actually chunk"
|
||||||
|
|
||||||
|
hit = _fake_note(id=30, title="The existing decision", note_type="note")
|
||||||
|
# Every chunk misses except the LAST one the gate will ask about.
|
||||||
|
sem = AsyncMock(side_effect=[[] for _ in range(n_chunks - 1)] + [[(0.94, hit)]])
|
||||||
|
with patch("scribe.services.dedup.async_session",
|
||||||
|
return_value=_session_returning(None)), \
|
||||||
|
patch("scribe.services.dedup.embeddings_svc.semantic_search_notes", sem):
|
||||||
|
dup = await find_duplicate_note(
|
||||||
|
7, "Title", body=body, project_id=2, is_task=False, note_type="note",
|
||||||
|
)
|
||||||
|
assert dup is not None and dup.id == 30
|
||||||
|
assert sem.await_count == n_chunks
|
||||||
|
|
||||||
|
|
||||||
@pytest.mark.asyncio
|
@pytest.mark.asyncio
|
||||||
async def test_semantic_match_of_other_note_type_is_ignored():
|
async def test_semantic_match_of_other_note_type_is_ignored():
|
||||||
other = _fake_note(id=21, title="X", note_type="process")
|
other = _fake_note(id=21, title="X", note_type="process")
|
||||||
|
|||||||
Reference in New Issue
Block a user