Repair vector rows left behind by older databases at store startup. Keep the existing single-statement delete path from main and avoid reintroducing the redundant batched helper.
302 lines
12 KiB
Python
302 lines
12 KiB
Python
"""Document ingest / RAG store — hermetic (fake embeddings, temp DB)."""
|
|
import asyncio
|
|
import tempfile
|
|
from pathlib import Path
|
|
|
|
from synapse.memory.store import PersistentMemoryStore
|
|
|
|
|
|
def _store():
|
|
return PersistentMemoryStore(Path(tempfile.mkdtemp()) / "t.db")
|
|
|
|
|
|
def test_chunker_packs_and_splits():
|
|
s = _store()
|
|
one = s._chunk_text("short one.\n\nshort two.")
|
|
assert one == ["short one.\n\nshort two."] # both fit one chunk
|
|
many = s._chunk_text("a" * 700 + "\n\n" + "b" * 700)
|
|
assert len(many) == 2 # each paragraph near the size cap -> own chunk
|
|
|
|
|
|
async def _fake_embed(text):
|
|
kws = ["lego", "star", "wars", "gpu", "vega"]
|
|
v = [float(text.lower().count(k)) for k in kws]
|
|
return v if any(v) else None
|
|
|
|
|
|
def test_add_list_search_delete_roundtrip():
|
|
async def run():
|
|
s = _store()
|
|
r = await s.add_document(
|
|
"Guide",
|
|
"Beat the lego star wars boss with the force.\n\nUnrelated gpu vega notes.",
|
|
_fake_embed,
|
|
)
|
|
assert r["chunks"] >= 1
|
|
assert [d["title"] for d in s.list_documents()] == ["Guide"]
|
|
|
|
hits = await s.search_documents("lego star wars", _fake_embed, limit=2, min_score=0.1)
|
|
assert hits and "lego" in hits[0]["text"].lower()
|
|
# scores are sorted descending
|
|
assert all(hits[i]["score"] >= hits[i + 1]["score"] for i in range(len(hits) - 1))
|
|
|
|
# get_document returns ordered chunks for the viewer
|
|
chunks = s.get_document(r["doc_id"])
|
|
assert [c["chunk_idx"] for c in chunks] == list(range(len(chunks)))
|
|
|
|
assert s.delete_document(r["doc_id"]) is True
|
|
assert s.list_documents() == []
|
|
assert s.get_document(r["doc_id"]) == [] # gone -> no chunks
|
|
assert s.delete_document(r["doc_id"]) is False # already gone
|
|
|
|
asyncio.run(run())
|
|
|
|
|
|
def test_search_empty_query_returns_nothing():
|
|
s = _store()
|
|
assert asyncio.run(s.search_documents("", _fake_embed)) == []
|
|
|
|
|
|
def test_conversation_project_binding():
|
|
s = _store()
|
|
assert s.conversation_project("nope") is None # not created yet
|
|
s.create_conversation("c1", "projX")
|
|
assert s.conversation_project("c1") == "projX"
|
|
s.create_conversation("c1", "other") # idempotent: keeps projX
|
|
assert s.conversation_project("c1") == "projX"
|
|
s.create_conversation("c2")
|
|
assert s.conversation_project("c2") == "" # unscoped
|
|
|
|
|
|
def test_conversations_move_between_projects():
|
|
# The Projects page lists chats by project_id and moves them with a PATCH;
|
|
# if all_conversations() drops the column the list is silently empty.
|
|
s = _store()
|
|
s.create_conversation("c1", "projX")
|
|
s.create_conversation("c2")
|
|
assert {c.id: c.project_id for c in s.all_conversations()} == {"c1": "projX", "c2": ""}
|
|
|
|
s.set_conversation_project("c2", "projX")
|
|
s.set_conversation_project("c1", "") # removed from the project
|
|
assert {c.id: c.project_id for c in s.all_conversations()} == {"c1": "", "c2": "projX"}
|
|
|
|
|
|
def test_project_instructions_and_scoped_memory():
|
|
# The chat system prompt takes the project's instructions plus global facts
|
|
# and this project's facts only — another project's must never leak in.
|
|
from synapse.memory.store import MemoryItem
|
|
|
|
s = _store()
|
|
p = s.create_project("Roof rebuild")
|
|
assert s.project_instructions(p["id"]) == "" # default: no instructions
|
|
assert s.project_instructions("ghost") == "" # unknown project
|
|
assert s.set_project_instructions(p["id"], "answer as a roofer")
|
|
assert not s.set_project_instructions("ghost", "x")
|
|
assert s.project_instructions(p["id"]) == "answer as a roofer"
|
|
|
|
s.add(MemoryItem(id="g", text="lives in Ohio")) # global
|
|
s.add(MemoryItem(id="a", text="uses metal panels", project_id=p["id"])) # this project
|
|
s.add(MemoryItem(id="b", text="prefers Lua", project_id="other")) # elsewhere
|
|
in_scope = [m.id for m in s.all() if m.project_id in ("", p["id"])]
|
|
assert in_scope == ["g", "a"]
|
|
|
|
# Deleting a project keeps its chats and facts, unscoped.
|
|
s.create_conversation("c1", p["id"])
|
|
s.delete_project(p["id"])
|
|
assert s.conversation_project("c1") == ""
|
|
assert s.get("a").project_id == ""
|
|
|
|
|
|
def test_conversation_recall_uses_vec_and_matches_brute_force():
|
|
s = _store()
|
|
if not s.vec_enabled:
|
|
import pytest
|
|
pytest.skip("sqlite-vec not loadable on this host")
|
|
|
|
async def run():
|
|
s.create_conversation("c1")
|
|
s.add_message("c1", "user", "tell me about lego star wars")
|
|
s.add_message("c1", "assistant", "lego star wars is a fun game")
|
|
s.create_conversation("c2")
|
|
s.add_message("c2", "user", "gpu vega vram notes")
|
|
s.add_message("c2", "assistant", "vega has 4gb")
|
|
|
|
hits = await s.semantic_search_conversations("lego star wars", _fake_embed, limit=2, min_score=0.1)
|
|
assert hits and hits[0]["id"] == "c1"
|
|
|
|
conn = s._connect()
|
|
n = conn.execute("SELECT COUNT(*) FROM vec_messages").fetchone()[0]
|
|
conn.close()
|
|
assert n >= 2 # dual-write populated the message vec index
|
|
|
|
s.vec_enabled = False
|
|
bf = await s.semantic_search_conversations("lego star wars", _fake_embed, limit=2, min_score=0.1)
|
|
assert bf[0]["id"] == hits[0]["id"]
|
|
|
|
asyncio.run(run())
|
|
|
|
|
|
def test_startup_sweeps_pre_existing_orphan_vectors():
|
|
"""Databases written before delete_conversation cleaned up after itself are
|
|
repaired the next time the store opens them."""
|
|
import json
|
|
path = Path(tempfile.mkdtemp()) / "t.db"
|
|
s = PersistentMemoryStore(path)
|
|
s.create_conversation("c1")
|
|
mid = s.add_message("c1", "user", "lego star wars")
|
|
conn = s._connect()
|
|
conn.execute("INSERT INTO message_vectors (message_id, embedding) VALUES (?, ?)",
|
|
(mid, json.dumps([1.0, 0.0])))
|
|
conn.execute("DELETE FROM messages WHERE id = ?", (mid,)) # the old leaky delete
|
|
conn.commit()
|
|
conn.close()
|
|
|
|
reopened = PersistentMemoryStore(path)
|
|
conn = reopened._connect()
|
|
assert conn.execute("SELECT COUNT(*) FROM message_vectors").fetchone()[0] == 0
|
|
conn.close()
|
|
|
|
|
|
def test_projects_scope_documents_and_survive_delete():
|
|
s = _store()
|
|
|
|
async def run():
|
|
p = s.create_project("Star Wars")
|
|
assert [x["name"] for x in s.list_projects()] == ["Star Wars"]
|
|
|
|
await s.add_document("Lego", "lego star wars boss tips", _fake_embed, project_id=p["id"])
|
|
await s.add_document("GPU", "gpu vega notes", _fake_embed) # unscoped
|
|
|
|
# project sees only its own; unscoped/all sees both
|
|
assert [d["title"] for d in s.list_documents(p["id"])] == ["Lego"]
|
|
assert len(s.list_documents(None)) == 2
|
|
|
|
# scoped search only returns the project's docs
|
|
scoped = await s.search_documents("star wars", _fake_embed, min_score=0.1, project_id=p["id"])
|
|
assert scoped and all(h["title"] == "Lego" for h in scoped)
|
|
|
|
# deleting the project keeps the docs but unscopes them
|
|
assert s.delete_project(p["id"]) is True
|
|
assert s.list_projects() == []
|
|
assert len(s.list_documents(None)) == 2
|
|
assert len(s.list_documents(p["id"])) == 0 # nothing left in that project
|
|
|
|
asyncio.run(run())
|
|
|
|
|
|
def test_chunker_overlap_and_hard_split():
|
|
s = _store()
|
|
# a single oversized paragraph (no blank lines, as in PDF text) is split
|
|
big = "x" * 2000
|
|
parts = s._chunk_text(big, size=800, overlap=120)
|
|
# each chunk is one <=size unit, plus at most an overlap tail (+separator)
|
|
assert len(parts) >= 3 and all(len(p) <= 800 + 120 + 2 for p in parts)
|
|
# consecutive chunks share an overlap tail
|
|
two = s._chunk_text("A" * 700 + "\n\n" + "B" * 700, size=800, overlap=120)
|
|
assert len(two) == 2 and two[1].startswith("A" * 120)
|
|
|
|
|
|
def test_vec_index_used_and_matches_brute_force():
|
|
# On a host that can load sqlite-vec, the fast path must be exercised (not a
|
|
# silent fallback) and agree with brute force on the top hit.
|
|
s = _store()
|
|
if not s.vec_enabled:
|
|
import pytest
|
|
pytest.skip("sqlite-vec not loadable on this host")
|
|
|
|
async def run():
|
|
await s.add_document("Lego", "lego star wars boss fight tips", _fake_embed)
|
|
await s.add_document("GPU", "gpu vega vram notes", _fake_embed)
|
|
# vec table populated by the dual-write
|
|
conn = s._connect()
|
|
n = conn.execute("SELECT COUNT(*) FROM vec_documents").fetchone()[0]
|
|
conn.close()
|
|
assert n == 2
|
|
|
|
vec_hits = await s.search_documents("lego star wars", _fake_embed, limit=2, min_score=0.1)
|
|
assert vec_hits and vec_hits[0]["title"] == "Lego"
|
|
|
|
# force the brute-force path and compare the top title
|
|
s.vec_enabled = False
|
|
bf_hits = await s.search_documents("lego star wars", _fake_embed, limit=2, min_score=0.1)
|
|
assert bf_hits[0]["title"] == vec_hits[0]["title"]
|
|
|
|
asyncio.run(run())
|
|
|
|
|
|
def test_extract_text_by_type():
|
|
from synapse.main import _extract_text
|
|
# plain text / markdown -> UTF-8 decode
|
|
assert _extract_text("notes.md", b"# Title\n\nbody") == "# Title\n\nbody"
|
|
assert _extract_text("x.txt", "café".encode("utf-8")) == "café"
|
|
# a real (tiny) PDF built with pypdf -> text extracted back out
|
|
from pypdf import PdfWriter, PdfReader
|
|
import io
|
|
w = PdfWriter()
|
|
w.add_blank_page(width=200, height=200)
|
|
buf = io.BytesIO(); w.write(buf)
|
|
out = _extract_text("blank.pdf", buf.getvalue())
|
|
assert isinstance(out, str) # blank page -> "" or whitespace, never raises
|
|
assert PdfReader(io.BytesIO(buf.getvalue())).pages # sanity: it was a valid PDF
|
|
|
|
|
|
def test_delete_conversation_takes_its_embeddings_with_it(tmp_path, monkeypatch):
|
|
"""Stale vectors are inert — the search joins messages — but they still
|
|
occupy slots in the ANN over-fetch, so recall of the surviving
|
|
conversations quietly thins out as deleted ones pile up."""
|
|
import sqlite3
|
|
from synapse.memory.store import PersistentMemoryStore
|
|
|
|
db = tmp_path / "t.db"
|
|
s = PersistentMemoryStore(db)
|
|
s.create_conversation("keep", "")
|
|
s.create_conversation("drop", "")
|
|
kept = s.add_message("keep", "user", "hello")
|
|
doomed = s.add_message("drop", "user", "goodbye")
|
|
|
|
conn = sqlite3.connect(db)
|
|
for mid in (kept, doomed):
|
|
conn.execute(
|
|
"INSERT OR REPLACE INTO message_vectors (message_id, embedding) VALUES (?, ?)",
|
|
(mid, "[0.0, 1.0]"),
|
|
)
|
|
conn.commit()
|
|
|
|
s.delete_conversation("drop")
|
|
|
|
left = {r[0] for r in conn.execute("SELECT message_id FROM message_vectors")}
|
|
assert left == {kept}, left
|
|
|
|
|
|
def test_extraction_watermark_is_idempotent(tmp_path):
|
|
"""The curator reads a conversation when it goes idle, so the watermark is
|
|
what stops a restart (or a second sweep) from re-reading messages and
|
|
re-saving the facts it already saved."""
|
|
from synapse.memory.store import PersistentMemoryStore
|
|
|
|
s = PersistentMemoryStore(tmp_path / "t.db")
|
|
s.create_conversation("c", "")
|
|
s.add_message("c", "user", "i bought a bike")
|
|
last = s.add_message("c", "assistant", "nice")
|
|
|
|
pending, mark = s.pending_extraction("c")
|
|
assert [m["role"] for m in pending] == ["user", "assistant"]
|
|
assert mark == last
|
|
|
|
s.set_extracted_through("c", mark)
|
|
assert s.pending_extraction("c") == ([], 0) # nothing new -> no model call
|
|
|
|
s.add_message("c", "user", "a 2019 trek")
|
|
pending, _ = s.pending_extraction("c")
|
|
assert [m["content"] for m in pending] == ["a 2019 trek"] # only the unread tail
|
|
|
|
|
|
def test_idle_sweep_only_claims_quiet_conversations(tmp_path):
|
|
from synapse.memory.store import PersistentMemoryStore
|
|
|
|
s = PersistentMemoryStore(tmp_path / "t.db")
|
|
s.create_conversation("fresh", "")
|
|
s.add_message("fresh", "user", "still typing")
|
|
assert s.conversations_awaiting_extraction(3600) == [] # too recent to be "over"
|
|
assert s.conversations_awaiting_extraction(0) == ["fresh"]
|