diff --git a/src/scribe/services/dedup.py b/src/scribe/services/dedup.py index 107f61a..a86c9b6 100644 --- a/src/scribe/services/dedup.py +++ b/src/scribe/services/dedup.py @@ -69,6 +69,12 @@ _SEMANTIC_THRESHOLD = 0.90 # structural signals cannot see. _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 class DuplicateMatch: @@ -258,41 +264,45 @@ async def find_duplicate_note( # --- Signal 3: semantic similarity (only with a substantial body) --- if body and len(body.strip()) >= _MIN_BODY_FOR_SEMANTIC: - # Built by the SAME function the corpus was embedded with. This one is - # the copy that mattered most and was easiest to miss: it is a QUERY - # document, compared against embedded ones. Shaped differently from the - # corpus it searches, the gate degrades silently — it still returns - # neighbours, just less apt ones, and no signal says the query and the - # index stopped agreeing (found by the guard in test_embedding_text). - query = embeddings_svc.embedding_text(title, body) - # Scope the semantic check the same way as the title check: a record in - # project P compares only to P; a project-less (orphan) record compares - # only to other orphans (orphan_only), NOT across every project — without - # this, semantic_search_notes applies no project filter when project_id - # is None and would match an orphan note against any project's notes. - hits = await embeddings_svc.semantic_search_notes( - user_id, query, project_id=project_id, is_task=is_task, - orphan_only=(project_id is None), - limit=3, - threshold=(_SNIPPET_SEMANTIC_THRESHOLD - if note_type == SNIPPET_NOTE_TYPE else _SEMANTIC_THRESHOLD), - # Owner-only, deliberately: this gate BLOCKS a create and tells the - # caller to update the match instead. Matching someone else's record - # would refuse their write and point them at something they may not - # be able to edit. - scope="own", - # NOT demoted by supersession (#278). A superseded record is 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 it here would let - # the same note be recorded a second time, and the second copy would - # be the one nothing warns about. - demote_superseded=False, - ) - 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") + # Query with the SAME chunker the corpus was embedded with (#280). This + # was the copy that mattered most and was easiest to miss: these are + # QUERY documents, compared against embedded ones — shaped differently + # from the corpus, the gate degrades silently. Chunking also makes the + # gate see what the whole-document query diluted: a long candidate that + # duplicates an existing record IN ONE SECTION now matches on that + # section. Capped so one pathological paste can't turn a save into + # dozens of searches — a duplicate past the cap is the duplicate + # report's job, not the gate's. + for query in embeddings_svc.chunk_document(title, body)[:_GATE_MAX_CHUNKS]: + # Scope the semantic check the same way as the title check: a record + # in project P compares only to P; a project-less (orphan) record + # compares only to other orphans (orphan_only), NOT across every + # project — without this, semantic_search_notes applies no project + # filter when project_id is None and would match an orphan note + # against any project's notes. + hits = await embeddings_svc.semantic_search_notes( + user_id, query, project_id=project_id, is_task=is_task, + orphan_only=(project_id is None), + limit=3, + threshold=(_SNIPPET_SEMANTIC_THRESHOLD + if note_type == SNIPPET_NOTE_TYPE else _SEMANTIC_THRESHOLD), + # Owner-only, deliberately: this gate BLOCKS a create and tells + # the caller to update the match instead. Matching someone + # else's record would refuse their write and point them at + # something they may not be able to edit. + scope="own", + # NOT demoted by supersession (#278). A superseded record is + # 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 + # it here would let the same note be recorded a second time, and + # the second copy would be the one nothing warns about. + demote_superseded=False, + ) + 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 @@ -502,15 +512,22 @@ async def find_duplicate_records( left_note = aliased(Note, name="left_note") right_note = aliased(Note, name="right_note") 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]] = [] try: async with async_session() as session: stmt = ( - select(left.note_id, right.note_id, distance.label("distance")) + select(left.note_id, right.note_id, best.label("distance")) .select_from(left) # `<` 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(left_note, left_note.id == left.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. left_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)) ) rows = list((await session.execute(stmt)).all()) diff --git a/src/scribe/services/embeddings.py b/src/scribe/services/embeddings.py index 1a46c89..e5defcc 100644 --- a/src/scribe/services/embeddings.py +++ b/src/scribe/services/embeddings.py @@ -124,6 +124,14 @@ _SUPERSESSION_PENALTY = 0.05 # whose neighbours sit ~0.01-0.02 apart. _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( 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 # needs the true answer to be more than _SUPERSESSION_OVERFETCH # 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.where(distance <= max_distance) .order_by(distance.asc()) @@ -535,8 +545,19 @@ async def semantic_search_notes( logger.warning("Failed to query note embeddings", exc_info=True) return [] - # Recover similarity (1 - distance) and preserve the highest-first contract. - scored = [(1.0 - float(dist), note) for note, dist in rows] + # Collapse chunk rows to BEST-CHUNK-PER-NOTE (#280): rows arrive ordered by + # 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: return scored[:limit] return await _apply_supersession_penalty(scored, limit) diff --git a/src/scribe/services/notes.py b/src/scribe/services/notes.py index 4a05e75..7d46938 100644 --- a/src/scribe/services/notes.py +++ b/src/scribe/services/notes.py @@ -236,15 +236,25 @@ async def list_notes( if query_vec is not None: from scribe.models.embedding import NoteEmbedding from scribe.services.embeddings import INTERACTIVE_SEARCH_THRESHOLD - distance = NoteEmbedding.embedding.cosine_distance(query_vec) - sem_filter = distance <= (1.0 - INTERACTIVE_SEARCH_THRESHOLD) - query = query.join( - NoteEmbedding, NoteEmbedding.note_id == Note.id - ).where(sem_filter) - count_query = count_query.join( - NoteEmbedding, NoteEmbedding.note_id == Note.id - ).where(sem_filter) - semantic_order = distance.asc() + # Best-chunk-per-note as a correlated MIN, not a join (#280): + # a note stores one embedding row PER CHUNK, so the plain join + # this used to be would repeat a long note once per matching + # chunk — duplicated list rows and a total that counts chunks. + # This query is filter-heavy and paginated, never HNSW-bound, + # so the scalar subquery costs what the join did. + best_distance = ( + select( + 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: terms = _strip_type_nouns(q) for term in terms: diff --git a/tests/test_chunking.py b/tests/test_chunking.py index 1269862..15861cc 100644 --- a/tests/test_chunking.py +++ b/tests/test_chunking.py @@ -136,6 +136,36 @@ def test_a_monster_single_paragraph_is_hard_split_not_dropped(): 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) ------------------------- diff --git a/tests/test_integration_pgvector_search.py b/tests/test_integration_pgvector_search.py index efb1349..56547f6 100644 --- a/tests/test_integration_pgvector_search.py +++ b/tests/test_integration_pgvector_search.py @@ -68,7 +68,11 @@ async def seeded(): await s.flush() # query vector will be [1,0,0,...]; near ~ identical (sim≈1.0), # 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, 1, _vec(0.6, 0.8))) s.add(_emb(far.id, user.id, 0, _vec(0.0, 1.0))) await s.commit() 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 far_id not in ids 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] assert top_score == pytest.approx(1.0, abs=1e-3) diff --git a/tests/test_services_dedup.py b/tests/test_services_dedup.py index 5c333d7..9437174 100644 --- a/tests/test_services_dedup.py +++ b/tests/test_services_dedup.py @@ -68,6 +68,33 @@ async def test_semantic_match_when_body_substantial(): 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 async def test_semantic_match_of_other_note_type_is_ignored(): other = _fake_note(id=21, title="X", note_type="process")