From 705ad030fe4ba1728b23823759521ca8815d23d0 Mon Sep 17 00:00:00 2001 From: Claude Agent Date: Mon, 6 Jul 2026 21:57:49 +0000 Subject: [PATCH] feat(embeddings): cache persistent de vectori in SQLite + warmup in fundal Vectorii corpusului k-NN persista in tabela embedding_cache (PK model+text_hash, blob float32 LE); la warmup se vectorizeaza doar textele lipsa din cache, deci restartul cu corpus neschimbat nu mai plateste ~1-2 min de embed (embed=0). - app/embedding_cache.py: serializare array('f'), load/save/purge chunk 500 cu BEGIN/COMMIT explicit (conexiuni autocommit), validare dimensiune la scriere si citire, orchestrare sync_corpus_vectors cu embed_fn injectat - index_corpus(vectors=): vectori precalculati cu validare aliniere; mismatch -> fallback embed complet - ensure_embeddings_corpus: warmup in thread la startup (block=True), calea de request ne-blocanta (acquire non-blocking pe lock; warmup in curs -> return imediat); purjare orfane + modele vechi doar dupa indexare reusita - log warmup: cache=N embed=M in Xs - 31 teste noi (cold/warm/incremental, model schimbat, concurenta, ranking exact, echivalenta float32); suita completa 1596 passed Co-Authored-By: Claude Fable 5 --- TODOS.md | 15 + app/embedding_cache.py | 210 ++++++++++++ app/embeddings.py | 77 ++++- app/main.py | 22 ++ app/mapping.py | 102 ++++-- app/schema.sql | 11 + docs/ROADMAP.md | 4 +- docs/prd/prd-5.21-embedding-cache-sqlite.md | 334 +++++++++++++++++++ tests/test_embedding_cache.py | 277 ++++++++++++++++ tests/test_embeddings.py | 73 ++++ tests/test_embeddings_warmup_cache.py | 347 ++++++++++++++++++++ 11 files changed, 1437 insertions(+), 35 deletions(-) create mode 100644 app/embedding_cache.py create mode 100644 docs/prd/prd-5.21-embedding-cache-sqlite.md create mode 100644 tests/test_embedding_cache.py create mode 100644 tests/test_embeddings_warmup_cache.py diff --git a/TODOS.md b/TODOS.md index 6806668..3fa7c57 100644 --- a/TODOS.md +++ b/TODOS.md @@ -2,6 +2,21 @@ Elemente deferate din review-uri. Negrupte de un PRD curent; de promovat cand devin prioritare. +## Din /autoplan PRD 5.21 (2026-07-06) + +- [ ] **CLI administrare cache embeddings (`python3 -m tools.embcache stats|clear|rebuild`)** — v1 foloseste + `DELETE FROM embedding_cache` manual (decizie Open Q1). De construit cand operarea manuala devine + frecventa sau cand apare al doilea operator. Effort: S. (CEO, low.) +- [ ] **Precompute vectori la build (artefact in imaginea Docker)** — modelul e pinned si corpusul seed e + comis, deci vectorii pot fi precalculati la build; ar face si PRIMUL start pe volum proaspat ieftin + (azi doar al doilea beneficiaza). Effort: M (pipeline build + fallback la runtime). (Outside voice CEO, medium.) +- [ ] **Graceful reload / zero-downtime restart** — procesul vechi serveste pana cel nou e cald; ar elimina + COMPLET fereastra post-restart, inclusiv incarcarea modelului (~10-20s) pe care cache-ul n-o poate evita. + Infra (Dokploy/compose), nu app. Effort: M. (Outside voice CEO, medium.) +- [ ] **Latenta per-query `suggest_nearest` (cosine pur Python)** — bucla Python peste 17k x 384 la fiecare + cautare; de MASURAT intai; daca >100ms, matrice numpy + dot product (efort mic). Designul curent e validat + pana la ~50k randuri de corpus (prag documentat in PRD 5.21 A12). (Eng F5, medium.) + ## Din hardening 80/20 (/autoplan, 2026-07-03) - [ ] **Traefik IP-allowlist pe /v1 + /metrics in fereastra de lansare** (T4) — cu zero clienti API, diff --git a/app/embedding_cache.py b/app/embedding_cache.py new file mode 100644 index 0000000..3cffee5 --- /dev/null +++ b/app/embedding_cache.py @@ -0,0 +1,210 @@ +"""Cache persistent de vectori embeddings in SQLite (tabela `embedding_cache`). + +Design: +- Cheie (model, text_hash): schimbarea modelului nu foloseste vectori vechi. +- Vectorii raman `array('f')` (float32) end-to-end -- fara conversie la list[float]. +- Scrieri/stergeri in tranzactii scurte, chunk-uite (conexiunile sunt autocommit, + vezi app/db.py -- BEGIN/COMMIT explicit per chunk, altfel fiecare INSERT/DELETE + e propria tranzactie). +- Degradare gratioasa: orice eroare SQLite -> log.warning, fara exceptie propagata + din load/save/purge (caller-ul, ensure_embeddings_corpus, ramane neschimbat). +""" +from __future__ import annotations + +import hashlib +import logging +import sqlite3 +import sys +from array import array +from typing import Callable, Iterable, Sequence + +log = logging.getLogger(__name__) + +EMB_DIM = 384 # paraphrase-multilingual-MiniLM-L12-v2 (app/embeddings.py::FASTEMBED_MODEL) +_CHUNK_SIZE = 500 # A2/A3: tranzactii scurte vs worker BEGIN IMMEDIATE pe submissions + + +def text_hash(text: str) -> str: + """SHA-256 hex al textului dat (apelantul normalizeaza inainte de hash).""" + return hashlib.sha256(str(text).encode("utf-8")).hexdigest() + + +def vector_to_blob(vector: Sequence[float]) -> bytes: + """Serializeaza un vector la float32 little-endian (independent de platforma).""" + arr = array("f", vector) + if sys.byteorder != "little": + arr = array("f", arr) + arr.byteswap() + return arr.tobytes() + + +def blob_to_vector(blob: bytes) -> array: + """Deserializeaza un blob float32 little-endian la `array('f')`.""" + arr = array("f") + arr.frombytes(blob) + if sys.byteorder != "little": + arr.byteswap() + return arr + + +def _chunks(seq: Sequence, size: int = _CHUNK_SIZE) -> Iterable[Sequence]: + for i in range(0, len(seq), size): + yield seq[i : i + size] + + +def load_cached_vectors(conn: sqlite3.Connection, model: str, hashes: Sequence[str]) -> dict[str, array]: + """Citeste vectorii existenti pentru `model` + `hashes`. Blob corupt (lungime + gresita) = tratat ca miss (nu apare in rezultat), cu log.warning. + + Degradare gratioasa: orice eroare SQLite -> dict gol/partial, fara exceptie. + """ + out: dict[str, array] = {} + if not hashes: + return out + expected_bytes = EMB_DIM * 4 + unique_hashes = list(dict.fromkeys(hashes)) + try: + for chunk in _chunks(unique_hashes): + placeholders = ",".join("?" for _ in chunk) + rows = conn.execute( + f"SELECT text_hash, vector FROM embedding_cache " + f"WHERE model = ? AND text_hash IN ({placeholders})", + (model, *chunk), + ).fetchall() + for row in rows: + blob = row["vector"] + if len(blob) != expected_bytes: + log.warning( + "embedding_cache: blob lungime gresita pentru hash=%s (asteptat %d, primit %d) -- tratat ca miss", + row["text_hash"], expected_bytes, len(blob), + ) + continue + out[row["text_hash"]] = blob_to_vector(blob) + except sqlite3.OperationalError as exc: + log.warning("embedding_cache: load_cached_vectors esuat (%s) -- fallback embed complet", exc) + return out + + +def save_vectors(conn: sqlite3.Connection, model: str, items: Sequence[tuple[str, Sequence[float]]]) -> None: + """INSERT OR REPLACE in chunk-uri de 500 randuri, BEGIN/COMMIT explicit per chunk + (conexiunile sunt autocommit -- conn.commit() singur e no-op). + + Valideaza `len(vector) == EMB_DIM` la scriere: vector gresit = respins (log.warning), + NU scris. Degradare gratioasa: eroare SQLite pe un chunk -> log.warning, chunk-urile + ramase continua (cache partial e idempotent, se completeaza la urmatorul warmup). + + Invarianta: `conn` trebuie sa fie in autocommit (fara tranzactie deschisa de + apelant) -- ROLLBACK-ul din except ar anula altfel tranzactia apelantului. + """ + valid = [] + for h, vec in items: + if len(vec) != EMB_DIM: + log.warning( + "embedding_cache: vector dimensiune gresita pentru hash=%s (asteptat %d, primit %d) -- respins", + h, EMB_DIM, len(vec), + ) + continue + valid.append((h, model, vector_to_blob(vec))) + + for chunk in _chunks(valid): + try: + conn.execute("BEGIN") + conn.executemany( + "INSERT OR REPLACE INTO embedding_cache (text_hash, model, vector) VALUES (?, ?, ?)", + chunk, + ) + conn.execute("COMMIT") + except sqlite3.OperationalError as exc: + try: + conn.execute("ROLLBACK") + except sqlite3.OperationalError: + pass + log.warning("embedding_cache: save_vectors esuat pe un chunk (%s) -- cache ramane partial", exc) + + +def purge_stale(conn: sqlite3.Connection, model: str, corpus_hashes: Iterable[str]) -> None: + """Sterge intrarile care nu mai apartin corpusului curent: orfane ale + modelului curent (text_hash absent din `corpus_hashes`) SI toate intrarile + modelelor VECHI (model != curent). Diff calculat in Python, DELETE chunk-uit + pe PK (model, text_hash) -- tranzactii scurte. + + Degradare gratioasa: eroare SQLite -> log.warning, orfanele raman pana la + urmatoarea trecere. + + Invarianta: `conn` trebuie sa fie in autocommit (fara tranzactie deschisa de + apelant) -- ROLLBACK-ul din except ar anula altfel tranzactia apelantului. + Apelantul trebuie sa cheme aceasta functie DOAR dupa o indexare reusita + (vezi `ensure_embeddings_corpus`) -- un esec de indexare nu trebuie sa goleasca + cache-ul de randuri inca valide. + """ + keep = set(corpus_hashes) + try: + rows = conn.execute("SELECT model, text_hash FROM embedding_cache").fetchall() + except sqlite3.OperationalError as exc: + log.warning("embedding_cache: purge_stale citire esuata (%s)", exc) + return + + to_delete = [ + (r["model"], r["text_hash"]) + for r in rows + if r["model"] != model or r["text_hash"] not in keep + ] + if not to_delete: + return + + for chunk in _chunks(to_delete): + try: + conn.execute("BEGIN") + conn.executemany( + "DELETE FROM embedding_cache WHERE model = ? AND text_hash = ?", + chunk, + ) + conn.execute("COMMIT") + except sqlite3.OperationalError as exc: + try: + conn.execute("ROLLBACK") + except sqlite3.OperationalError: + pass + log.warning("embedding_cache: purge_stale esuat pe un chunk (%s) -- orfanele raman", exc) + + +def sync_corpus_vectors( + conn: sqlite3.Connection, + model: str, + texts: Sequence[str], + embed_fn: Callable[[list[str]], Sequence[Sequence[float]]], +) -> list[array]: + """Orchestreaza hash -> load -> embed(doar miss-uri) -> save -> vectori aliniati. + + Returneaza o lista de `array('f')` aliniata pozitional cu `texts`. + + `embed_fn` primeste lista textelor lipsa din cache si intoarce vectorii lor + (aceeasi ordine). Daca `embed_fn` esueaza (arunca), exceptia se propaga; + apelantul (ensure_embeddings_corpus) are deja degradare gratioasa (except -> pass). + + Un save partial esuat (lock SQLite) NU opreste intoarcerea vectorilor din RAM + (A13c): vectorii noi raman in `cached` indiferent de rezultatul persistarii. + + NU purjeaza orfanele: purjarea e responsabilitatea apelantului, DUPA ce corpusul + a fost indexat cu succes (`index_corpus`) -- un esec de indexare nu trebuie sa + goleasca din greseala cache-ul de randuri inca valide. + """ + hashes = [text_hash(t) for t in texts] + cached = load_cached_vectors(conn, model, hashes) + + missing_positions = [i for i, h in enumerate(hashes) if h not in cached] + if missing_positions: + new_texts = [texts[i] for i in missing_positions] + new_vecs = embed_fn(new_texts) + if len(new_vecs) != len(new_texts): + raise ValueError( + f"embed_fn a intors {len(new_vecs)} vectori pentru {len(new_texts)} texte" + ) + to_save = [] + for pos, vec in zip(missing_positions, new_vecs): + arr = array("f", vec) + cached[hashes[pos]] = arr + to_save.append((hashes[pos], arr)) + save_vectors(conn, model, to_save) + + return [cached[h] for h in hashes] diff --git a/app/embeddings.py b/app/embeddings.py index 68d8fb7..17e4e0b 100644 --- a/app/embeddings.py +++ b/app/embeddings.py @@ -10,9 +10,9 @@ Design: - NU apelat din resolve_prestatii/load_mapping API public (nivel modul): - index_corpus(items) -> None - suggest_nearest(text, top_k) -> [{cod, is_nul, similaritate}] - is_available() -> bool + index_corpus(items, signature, vectors) -> None + suggest_nearest(text, top_k) -> [{cod, is_nul, similaritate}] + is_available() -> bool Clase (pentru teste / injectare backend): EmbeddingEngine(backend) -- motor testabil cu backend injectabil @@ -22,6 +22,7 @@ from __future__ import annotations import logging import math +import threading from typing import Protocol, runtime_checkable log = logging.getLogger(__name__) @@ -107,8 +108,19 @@ class EmbeddingEngine: doar cand semnatura nomenclatorului s-a schimbat (evita re-embed inutil).""" return self._corpus_sig - def index_corpus(self, items: list[dict], signature: str | None = None) -> None: - """Vectorizeaza corpus [{denumire, cod}] si il pastreaza in memorie. + def index_corpus( + self, + items: list[dict], + signature: str | None = None, + vectors: list | None = None, + ) -> None: + """Indexeaza corpus [{denumire, cod}] si il pastreaza in memorie. + + `vectors`: vectori precalculati, aliniati POZITIONAL cu `items` (ex. din + embedding_cache). `None` (default) = comportamentul existent, embed complet + prin backend. Daca `vectors` e furnizat dar lungimea nu corespunde cu `items` + sau contine `None`, se ignora (log.warning) si se cade pe embed complet (A8) -- + o dezaliniere silentioasa ar produce coduri sugerate GRESITE. Ignora silentios daca backend-ul lipseste, corpus-ul e gol sau apare orice exceptie la vectorizare (degradare gratioasa). @@ -120,16 +132,35 @@ class EmbeddingEngine: if not items or not self.is_available(): return + if vectors is not None and (len(vectors) != len(items) or any(v is None for v in vectors)): + log.warning( + "embeddings: index_corpus vectors (%d) nealiniat cu items (%d) sau contine None -- fallback embed complet", + len(vectors), len(items), + ) + vectors = None + try: - texts = [str(item["denumire"]) for item in items] - vecs = self._backend.embed(texts) - self._corpus_vecs = vecs + if vectors is not None: + self._corpus_vecs = list(vectors) + else: + texts = [str(item["denumire"]) for item in items] + self._corpus_vecs = self._backend.embed(texts) self._corpus_items = list(items) self._corpus_sig = signature except Exception as exc: log.warning("embeddings: index_corpus esuat: %s", exc) # corpus ramane gol -- suggest_nearest va returna [] + def embed(self, texts: list[str]) -> list[list[float]]: + """Vectorizeaza texte brute prin backend (folosit la miss-uri de cache). + + Arunca daca backend-ul lipseste sau embed() esueaza -- apelantul + (sync_corpus_vectors) propaga eroarea, fara degradare gratioasa aici. + """ + if not self.is_available(): + raise RuntimeError("embeddings: backend indisponibil") + return self._backend.embed(texts) + def suggest_nearest( self, denumire: str, @@ -168,6 +199,7 @@ class EmbeddingEngine: # --------------------------------------------------------------------------- # _engine: EmbeddingEngine | None = None +_engine_lock = threading.Lock() def _load_engine() -> EmbeddingEngine: @@ -194,13 +226,24 @@ def _load_engine() -> EmbeddingEngine: def _get_engine() -> EmbeddingEngine: - """Returneaza engine-ul global (lazy-init).""" + """Returneaza engine-ul global (lazy-init, thread-safe). + + Lock-ul previne incarcarea dubla a modelului cand warmup-ul de la startup + si un request concurent ajung aici simultan. + """ global _engine if _engine is None: - _engine = _load_engine() + with _engine_lock: + if _engine is None: + _engine = _load_engine() return _engine +def is_loaded() -> bool: + """True daca engine-ul global a fost deja construit. NU forteaza incarcarea.""" + return _engine is not None + + # --------------------------------------------------------------------------- # # API public la nivel de modul (wiring L14-S6) # # --------------------------------------------------------------------------- # @@ -233,12 +276,22 @@ def corpus_signature() -> str | None: return _engine.corpus_signature() -def index_corpus(items: list[dict], signature: str | None = None) -> None: +def embed_texts(texts: list[str]) -> list[list[float]]: + """Vectorizeaza texte brute prin motorul global (folosit la miss-uri de cache). + + Arunca daca engine-ul e indisponibil -- apelantul (embedding_cache.sync_corpus_vectors) + propaga eroarea catre ensure_embeddings_corpus (degradare gratioasa acolo). + """ + return _get_engine().embed(texts) + + +def index_corpus(items: list[dict], signature: str | None = None, vectors: list | None = None) -> None: """Vectorizeaza corpus [{denumire, cod}] in motorul global. + `vectors`: vezi EmbeddingEngine.index_corpus (precalculati, aliniati cu `items`). Silentios pe eroare (degradare gratioasa). """ - _get_engine().index_corpus(items, signature=signature) + _get_engine().index_corpus(items, signature=signature, vectors=vectors) def suggest_nearest(denumire: str, top_k: int = 3) -> list[dict]: diff --git a/app/main.py b/app/main.py index 3a2f884..cbbceaf 100644 --- a/app/main.py +++ b/app/main.py @@ -7,6 +7,7 @@ un worker mort nu trebuie sa lase containerul "sanatos". from __future__ import annotations import secrets +import threading from contextlib import asynccontextmanager from datetime import datetime, timezone from pathlib import Path @@ -39,6 +40,25 @@ from .web.csrf import CsrfError from .web.session import AdminRequired, LoginRequired +def _warmup_embeddings() -> None: + """Incarca modelul de embeddings si indexeaza corpusul SILVER, in fundal. + + Ruleaza intr-un thread daemon la startup: incarcarea modelului (~230MB) plus + vectorizarea corpusului dureaza zeci de secunde si NU are voie sa blocheze + primul request pe /mapari. Pana termina, sugestiile embeddings lipsesc + (degradare gratioasa); GOLD/SILVER/fuzzy functioneaza normal. + """ + from .mapping import ensure_embeddings_corpus + try: + conn = get_connection() + try: + ensure_embeddings_corpus(conn, block=True) + finally: + conn.close() + except Exception: + pass # best-effort: esecul warmup-ului nu opreste API-ul + + @asynccontextmanager async def lifespan(app: FastAPI): install_log_redaction() @@ -49,6 +69,8 @@ async def lifespan(app: FastAPI): # cheia API si secretul de sesiune, in loc de o instanta descoperita post-deploy. validate_prod_invariants(get_settings()) init_db() + if get_settings().embeddings_enabled: + threading.Thread(target=_warmup_embeddings, name="emb-warmup", daemon=True).start() yield diff --git a/app/mapping.py b/app/mapping.py index 98a5e6b..2f4ffd7 100644 --- a/app/mapping.py +++ b/app/mapping.py @@ -16,7 +16,10 @@ from __future__ import annotations import hashlib import json +import logging import re +import threading +import time import unicodedata from typing import Any @@ -27,6 +30,8 @@ from .accounts import held_for_account from .nomenclator_seed import FALLBACK_NOMENCLATOR from .validation import validate_prezentare +log = logging.getLogger(__name__) + # Cont implicit cat timp auth API-key (CORE) nu e implementat: ingestiile vin cu # account_id NULL si le atribuim contului seed-at in schema (id=1). DEFAULT_ACCOUNT_ID = 1 @@ -631,6 +636,13 @@ def delete_text_rule(conn, account_id: int | None, pattern: str) -> None: # irelevante cand corpus-ul e mic sau neindexat corect). EMB_MIN_SIMILARITATE = 0.5 +# Protejeaza secventa hash->load->embed->save->purge->index (embedding_cache) de +# executie concurenta intre warmup-ul de fundal (block=True) si calea de request +# (block=False, dupa ce modelul e deja incarcat) -- altfel purjarea uneia ar sterge +# randuri tocmai scrise de cealalta. Calea de request obtine lock-ul neblocant +# (nu asteapta warmup-ul in curs); doar warmup-ul asteapta normal. +_embeddings_lock = threading.Lock() + def _corpus_signature_silver(rows: list) -> str: """Semnatura stabila a corpusului SILVER (mapping_suggestions) pentru cache. @@ -646,7 +658,7 @@ def _corpus_signature_silver(rows: list) -> str: return hashlib.sha256(blob.encode("utf-8")).hexdigest() -def ensure_embeddings_corpus(conn, nomenclator: list[dict] | None = None) -> None: +def ensure_embeddings_corpus(conn, nomenclator: list[dict] | None = None, *, block: bool = False) -> None: """Construieste/actualizeaza corpusul embeddings din corpusul ETICHETAT. Sursa corpusului = `mapping_suggestions` (SILVER): exemple reale etichetate @@ -658,9 +670,24 @@ def ensure_embeddings_corpus(conn, nomenclator: list[dict] | None = None) -> Non Gated pe `AUTOPASS_EMBEDDINGS_ENABLED` (default ON; OFF in teste): cand e dezactivat, e un no-op total -> /mapari instant + suita de teste rapida. - Cand e activat: indexeaza corpusul o singura data (lazy-load modelul ~230MB la - prima chemare), re-indexeaza doar cand semnatura corpusului SILVER s-a schimbat. - Itemii NUL (is_nul=1, cod NULL) raman in corpus: un vecin NUL e semnal de supresie. + Cand e activat: indexeaza corpusul o singura data, re-indexeaza doar cand + semnatura corpusului SILVER s-a schimbat. Itemii NUL (is_nul=1, cod NULL) raman + in corpus: un vecin NUL e semnal de supresie. + + `block=False` (default, calea de request): daca modelul NU e inca incarcat, + return imediat — incarcarea modelului (~230MB, zeci de secunde) NU are voie sa + blocheze un request HTTP; o face warmup-ul de la startup (block=True, in thread). + Odata modelul incarcat insa, warmup-ul mai poate fi INCA in curs de vectorizare + a corpusului (~1-2 min): calea de request NU asteapta dupa lock in acest caz — + incearca sa il obtina neblocant, iar daca e ocupat, iese imediat (degradare + gratioasa, sugestii lipsa pana termina warmup-ul). block=True (warmup) asteapta + normal dupa lock. + + Cache persistent (embedding_cache): hash-ul se calculeaza pe lista FILTRATA de + `denumire` (EXACT ce intra in `index_corpus`), citeste vectorii existenti pentru + modelul curent, vectorizeaza doar miss-urile si salveaza-i. Purjarea orfanelor + ruleaza DUPA indexare, doar daca indexarea a reusit efectiv (semnatura noua + confirmata) — un esec de indexare nu trebuie sa goleasca cache-ul. Degradare gratioasa: orice eroare lasa corpusul gol -> enrich cade pe restul. """ from .config import get_settings @@ -668,24 +695,55 @@ def ensure_embeddings_corpus(conn, nomenclator: list[dict] | None = None) -> Non return try: from . import embeddings as _emb - rows = conn.execute( - "SELECT denumire_normalizata, cod_prestatie, is_nul FROM mapping_suggestions" - ).fetchall() - if not rows: - return - sig = _corpus_signature_silver(rows) - if _emb.corpus_signature() == sig and _emb.has_corpus(): - return # deja indexat pe acelasi corpus SILVER -> nimic de facut - items = [ - { - "denumire": str(r["denumire_normalizata"]), - "cod": (str(r["cod_prestatie"]) if r["cod_prestatie"] is not None else None), - "is_nul": bool(r["is_nul"]), - } - for r in rows - if r["denumire_normalizata"] - ] - _emb.index_corpus(items, signature=sig) + if not block and not _emb.is_loaded(): + return # warmup-ul din fundal nu a terminat inca; nu bloca request-ul + + if not _embeddings_lock.acquire(blocking=block): + return # warmup in curs; calea de request nu asteapta (nu bloca request-ul) + try: + rows = conn.execute( + "SELECT denumire_normalizata, cod_prestatie, is_nul FROM mapping_suggestions" + ).fetchall() + if not rows: + return + sig = _corpus_signature_silver(rows) + if _emb.corpus_signature() == sig and _emb.has_corpus(): + return # deja indexat pe acelasi corpus SILVER -> nimic de facut + items = [ + { + "denumire": str(r["denumire_normalizata"]), + "cod": (str(r["cod_prestatie"]) if r["cod_prestatie"] is not None else None), + "is_nul": bool(r["is_nul"]), + } + for r in rows + if r["denumire_normalizata"] + ] + if not items: + return + + from . import embedding_cache as _cache + + texts = [item["denumire"] for item in items] + hashes = [_cache.text_hash(t) for t in texts] + miss_count = 0 + + def _embed_fn(missing_texts: list[str]) -> list: + nonlocal miss_count + miss_count += len(missing_texts) + return _emb.embed_texts(missing_texts) + + t0 = time.monotonic() + vectors = _cache.sync_corpus_vectors(conn, _emb.FASTEMBED_MODEL, texts, _embed_fn) + _emb.index_corpus(items, signature=sig, vectors=vectors) + if _emb.corpus_signature() == sig and _emb.has_corpus(): + _cache.purge_stale(conn, _emb.FASTEMBED_MODEL, set(hashes)) + eticheta = "warmup ok" if block else "corpus reindexat" + log.info( + "embeddings: %s cache=%d embed=%d in %.1fs", + eticheta, len(texts) - miss_count, miss_count, time.monotonic() - t0, + ) + finally: + _embeddings_lock.release() except Exception: pass # degradare gratioasa: esecul indexarii nu blocheaza editorul diff --git a/app/schema.sql b/app/schema.sql index 50acc1b..2b60113 100644 --- a/app/schema.sql +++ b/app/schema.sql @@ -281,6 +281,17 @@ CREATE TABLE IF NOT EXISTS shared_mappings ( updated_at TEXT NOT NULL DEFAULT (datetime('now')) ); +-- Cache persistent vectori embeddings corpus (mapping_suggestions). Cheia (model, +-- text_hash) leaga cache-ul de numele modelului: schimbare model = randuri vechi inerte, +-- fara sa fie folosite (curatate ulterior de purge_stale). +CREATE TABLE IF NOT EXISTS embedding_cache ( + text_hash TEXT NOT NULL, + model TEXT NOT NULL, + vector BLOB NOT NULL, + created_at TEXT NOT NULL DEFAULT (datetime('now')), + PRIMARY KEY (model, text_hash) +); + -- Heartbeat worker (un singur rand, id=1). /healthz citeste de aici. CREATE TABLE IF NOT EXISTS worker_heartbeat ( id INTEGER PRIMARY KEY CHECK (id = 1), diff --git a/docs/ROADMAP.md b/docs/ROADMAP.md index a13e42b..5642c50 100644 --- a/docs/ROADMAP.md +++ b/docs/ROADMAP.md @@ -48,7 +48,9 @@ Reguli de contract (detalii in `docs/api-rar-contract.md`): `FINALIZATA` e termi > PRD-uri (`docs/prd/prd-X.Y-*.md`), linkate in coloana Detalii. La fiecare livrabila terminata: > schimba statusul + data + linkul PRD si actualizeaza "Ultima actualizare". -**Ultima actualizare**: 2026-07-06 — **HARDENING 80/20 LIVRAT + PUSH** (`main` la `63b6cbc`; plan /autoplan `main-hardening-8020-plan-20260703.md`, audit securitate 2026-07-03, inainte de expunerea publica autopass.romfast.ro). **P0**: compose fail-safe — `:?` obligatoriu pe `AUTOPASS_REQUIRE_API_KEY`/`AUTOPASS_WORKER_SEND_ENABLED`/`AUTOPASS_RAR_ENV` (split-brain api-prod/worker-test eliminat; post-5.20 variabila e doar ancora de fallback, nu tinta trimiterilor), `FORWARDED_ALLOW_IPS="*"` ca rate-limit-ul sa vada IP-ul real dupa Traefik (invariante in comentariu: api fara `ports:` pe host, Traefik fara `forwardedHeaders.insecure`), `AUTOPASS_SESSION_SECRET` obligatoriu + `SESSION_HTTPS_ONLY` default true; invarianta E1 fail-fast la boot (`validate_prod_invariants`): prod fara cheie API sau session secret -> RuntimeError cu mesaj actionabil. **P0-4 backup criptat SQLite**: `tools/backup_db.sh` (snapshot online stdlib + gpg AES256 + retentie + rclone optional), `tools/restore_check.sh` (integrity_check pe restore), `docs/backup.md`; trigger dur: configurat INAINTE de prima declaratie reala prod. **P1**: headers de securitate in-app (`SecurityHeadersMiddleware`, HSTS doar pe https) + body-cap global 10MB (`BodyCapMiddleware` ASGI pur, 413 inainte de parserul multipart/JSON; verificarea per-endpoint ramane strat 2) + imagine Docker non-root (uid 10001, `HOME=/home/app` pt fastembed, chown /data INAINTE de VOLUME, loguri pe `/data/logs`, port 8010 aliniat EXPOSE/CMD; migrare one-off volum existent: `docker compose run --rm --user root api chown -R app:app /data`). **P2**: print-uri cu email in stdout eliminate (signup + notify degradat din `app/email.py` -> `log_event` fara PII) + fix crestere monotona chei `ratelimit._hits`. Executie multi-agent (4 agenti implementare pe fisiere disjuncte + review dedicat: 0 constatari blocante, 3 nit-uri acceptate). Suita completa: 1557 passed, 1 skipped (3 rulari independente). **Post-deploy HOUR-1 (Dokploy)**: seteaza env-urile obligatorii INAINTE de redeploy (`:?` pica pornirea), creeaza cheia API + tier pro/trial pe contul propriu (altfel 403 PLAN_FARA_API), ruleaza migrarea chown, seteaza passphrase-ul de backup. +**Ultima actualizare**: 2026-07-06 — **5.21 CACHE PERSISTENT EMBEDDINGS IN SQLITE — LIVRAT** (PRD: [prd-5.21](prd/prd-5.21-embedding-cache-sqlite.md); executie echipa de agenti Sonnet: cache-core T1+T2, wiring T3+T5, tests T4+T6, reviewer adversarial). Vectorii corpusului k-NN (17.181 exemple SILVER) persista in tabela noua `embedding_cache` (PK `(model, text_hash)`, blob float32 LE 384, `array('f')` end-to-end); la warmup se vectorizeaza DOAR miss-urile (hash sha256 pe `denumire_normalizata` din lista FILTRATA, identica cu ce intra in `index_corpus`), restarturile cu corpus neschimbat = `embed=0`. Modul nou `app/embedding_cache.py` (fara dependinte noi): `sync_corpus_vectors(conn, model, texts, embed_fn)` hash->load->embed(miss)->save; scrieri/stergeri in chunk-uri de 500 cu `BEGIN`/`COMMIT` EXPLICIT (conexiunile sunt autocommit — `conn.commit()` e no-op); validare dimensiune la scriere SI citire (blob corupt = miss re-vectorizat); purjare orfane + modele vechi DOAR dupa indexare reusita; toate erorile SQLite -> `log.warning` + fallback embed complet (degradare gratioasa, nimic propagat). `EmbeddingEngine.index_corpus(vectors=)` valideaza alinierea `len(vectors)==len(items)` + absenta None (mismatch -> fallback embed complet). Secventa protejata de `_embeddings_lock` (mapping.py); pe calea de request (`block=False`) acquire NE-blocant — daca warmup-ul tine lock-ul, return imediat (fix MAJOR din review: altfel primul request in fereastra de cold start ingheta ~1-2 min). Log observabilitate: `embeddings: warmup ok cache=N embed=M in Xs` (eticheta `corpus reindexat` pe calea de request). Semnatura corpusului (`_corpus_signature_silver`) = mecanism de decizie NESCHIMBAT; worker-ul nu atinge cache-ul. Review adversarial: 1 MAJOR + 2 MINOR + 1 NIT, toate reparate. Suita completa: **1596 passed, 1 skipped** (+31 teste noi: 15/15 cai din diagrama de acoperire, toate cu `AUTOPASS_EMBEDDINGS_ENABLED=1` explicit, backend mock). Crestere DB ~26-30MB acceptata; deferate in TODOS.md: CLI embcache, precompute la build, graceful reload, numpy la >50k randuri. **De verificat pe instanta reala (US-006)**: doua restarturi `./start.sh`, al doilea cu `embed=0` si sugestia "similaritate" in <30s. + +> 2026-07-06 — **HARDENING 80/20 LIVRAT + PUSH** (`main` la `63b6cbc`; plan /autoplan `main-hardening-8020-plan-20260703.md`, audit securitate 2026-07-03, inainte de expunerea publica autopass.romfast.ro). **P0**: compose fail-safe — `:?` obligatoriu pe `AUTOPASS_REQUIRE_API_KEY`/`AUTOPASS_WORKER_SEND_ENABLED`/`AUTOPASS_RAR_ENV` (split-brain api-prod/worker-test eliminat; post-5.20 variabila e doar ancora de fallback, nu tinta trimiterilor), `FORWARDED_ALLOW_IPS="*"` ca rate-limit-ul sa vada IP-ul real dupa Traefik (invariante in comentariu: api fara `ports:` pe host, Traefik fara `forwardedHeaders.insecure`), `AUTOPASS_SESSION_SECRET` obligatoriu + `SESSION_HTTPS_ONLY` default true; invarianta E1 fail-fast la boot (`validate_prod_invariants`): prod fara cheie API sau session secret -> RuntimeError cu mesaj actionabil. **P0-4 backup criptat SQLite**: `tools/backup_db.sh` (snapshot online stdlib + gpg AES256 + retentie + rclone optional), `tools/restore_check.sh` (integrity_check pe restore), `docs/backup.md`; trigger dur: configurat INAINTE de prima declaratie reala prod. **P1**: headers de securitate in-app (`SecurityHeadersMiddleware`, HSTS doar pe https) + body-cap global 10MB (`BodyCapMiddleware` ASGI pur, 413 inainte de parserul multipart/JSON; verificarea per-endpoint ramane strat 2) + imagine Docker non-root (uid 10001, `HOME=/home/app` pt fastembed, chown /data INAINTE de VOLUME, loguri pe `/data/logs`, port 8010 aliniat EXPOSE/CMD; migrare one-off volum existent: `docker compose run --rm --user root api chown -R app:app /data`). **P2**: print-uri cu email in stdout eliminate (signup + notify degradat din `app/email.py` -> `log_event` fara PII) + fix crestere monotona chei `ratelimit._hits`. Executie multi-agent (4 agenti implementare pe fisiere disjuncte + review dedicat: 0 constatari blocante, 3 nit-uri acceptate). Suita completa: 1557 passed, 1 skipped (3 rulari independente). **Post-deploy HOUR-1 (Dokploy)**: seteaza env-urile obligatorii INAINTE de redeploy (`:?` pica pornirea), creeaza cheia API + tier pro/trial pe contul propriu (altfel 403 PLAN_FARA_API), ruleaza migrarea chown, seteaza passphrase-ul de backup. > 2026-06-18 — 3.4 LIVRAT (interfata web ergonomica: tab-uri + wizard + microcopy). US-001 modul pur `app/web/labels.py` (stari tehnice→text uman + clasa CSS; test parametrizat din CHECK-ul `schema.sql` iese rosu la stare nemapata). US-002 bara status `/_fragments/status` + `_status.html` (etichete umane, defalcare blocate pe motiv, poll 15s, scoped pe cont). US-003 shell 6 tab-uri (Acasa·Import·Coada·Mapari·Cont·Nomenclator) cu deep-link `?tab=`, panou activ randat server-side, fragmente inactive lazy pe click, ARIA real (tablist/tab/tabpanel + aria-selected + navigare cu sageti). US-004 stepper import 4 pasi (PUR vizual, `hx-target="#import-section"` + csrf pastrate). US-005 Acasa onboarding checklist auto-bifat (are_creds/are_trimiteri) + colaps cand totul gata + empty states prietenoase Coada/Mapari. VERIFY lead-driven (TestClient ACs + 434 pytest pass; E2E browser/RAR LIVE neprobat in sesiune — recomandata probare manuala `--send`). Fix izolare teste (reset `ratelimit._hits` in fixturi, 429 la rulare subset). `/code-review` high: regasit avertisment "cont in asteptare de activare" (regresie din scoaterea `/_fragments/banner`) re-introdus in bara status + culori hardcodate→variabile paleta. 434 teste pass. Backend trimitere neatins. PRD: [prd-3.4](prd/prd-3.4-ux-dashboard-web.md). Urmeaza Etapa 4 (4.1 mapare AI/MCP). diff --git a/docs/prd/prd-5.21-embedding-cache-sqlite.md b/docs/prd/prd-5.21-embedding-cache-sqlite.md new file mode 100644 index 0000000..1c1e85d --- /dev/null +++ b/docs/prd/prd-5.21-embedding-cache-sqlite.md @@ -0,0 +1,334 @@ +# PRD: Cache persistent de vectori embeddings in SQLite + +## 1. Introducere + +Vectorii corpusului k-NN (17.181 exemple etichetate din `mapping_suggestions`) traiesc doar in RAM si se recalculeaza integral la fiecare pornire a API-ului (~1-2 minute de CPU), desi textele nu s-au schimbat. Warmup-ul in fundal (deja implementat) tine pagina Mapari rapida, dar lasa o fereastra de 1-2 minute dupa restart in care sugestiile embeddings lipsesc. Acest feature persista vectorii intr-o tabela SQLite si vectorizeaza incremental doar textele noi/modificate, reducand fereastra la doar incarcarea modelului (~10-20s). + +## 2. Obiective + +### Obiectiv Principal +- Elimina re-vectorizarea integrala a corpusului la fiecare restart: vectorii se calculeaza o singura data si se refolosesc. + +### Obiective Secundare +- Fereastra fara sugestii embeddings dupa restart scade de la ~1-2 minute la ~10-20s (doar incarcarea modelului ONNX). +- Cresterea corpusului SILVER nu mai scumpeste restarturile (cost proportional doar cu randurile NOI). +- Cache reconstruibil oricand (`DELETE FROM embedding_cache` + restart), fara alte efecte. + +### Metrici de Succes +- La corpus neschimbat: zero apeluri `embed()` pe corpus la startup; indexarea din cache < 5s. +- La N randuri noi in corpus: exact N texte vectorizate la urmatorul warmup. +- Sugestiile embeddings (sursa "similaritate") functioneaza identic inainte/dupa (acelasi cod sugerat pentru acelasi query). + +## 3. User Stories + +### US-001: Tabela embedding_cache in schema +**Ca** sistem +**Vreau** o tabela dedicata pentru vectorii corpusului +**Pentru ca** vectorii sa supravietuiasca restartului fara sa ating schema `mapping_suggestions`. + +**Acceptance Criteria:** +- [ ] `app/schema.sql` contine tabela `embedding_cache` cu: `text_hash TEXT NOT NULL` (SHA-256 al textului normalizat), `model TEXT NOT NULL` (numele modelului fastembed), `vector BLOB NOT NULL` (float32 little-endian, 384 valori), `created_at TEXT NOT NULL DEFAULT (datetime('now'))`, `PRIMARY KEY (model, text_hash)`. +- [ ] `init_db()` creeaza tabela idempotent (CREATE TABLE IF NOT EXISTS), pe baza de date existenta si pe una noua. +- [ ] `python3 -m pytest -q` trece. + +### US-002: Serializare + acces cache (functii pure) +**Ca** developer +**Vreau** functii de citire/scriere vectori in cache +**Pentru ca** logica de (de)serializare sa fie testabila izolat. + +**Acceptance Criteria:** +- [ ] Functii in `app/embeddings.py` (sau modul nou `app/embedding_cache.py`): `vector_to_blob(list[float]) -> bytes` si `blob_to_vector(bytes) -> list[float]` (float32 LE, `array('f')` din stdlib, fara dependinta noua). +- [ ] `load_cached_vectors(conn, model, hashes) -> dict[hash, vector]` — un singur SELECT (chunk-uit la limita de parametri SQLite), intoarce doar hash-urile gasite. +- [ ] `save_vectors(conn, model, items: list[(hash, vector)])` — INSERT OR REPLACE batch, commit la final. +- [ ] Round-trip exact: `blob_to_vector(vector_to_blob(v)) == pytest.approx(v)` pentru vectori de 384 float. +- [ ] `python3 -m pytest -q` trece. + +### US-003: ensure_embeddings_corpus foloseste cache-ul +**Ca** operator +**Vreau** ca warmup-ul sa vectorizeze doar textele lipsa din cache +**Pentru ca** restartul sa nu mai coste 1-2 minute de CPU. + +**Acceptance Criteria:** +- [ ] `ensure_embeddings_corpus(conn, block=True)`: calculeaza hash per `denumire_normalizata`, citeste vectorii existenti din `embedding_cache` pentru modelul curent, apeleaza `embed()` DOAR pe textele fara vector in cache, salveaza vectorii noi in cache, apoi indexeaza corpusul complet (cache + noi) in `EmbeddingEngine`. +- [ ] `EmbeddingEngine.index_corpus` accepta vectori precalculati (parametru nou `vectors` sau metoda `index_corpus_precomputed`) fara sa apeleze backend-ul pentru ei; comportamentul existent (fara vectori -> embed tot) ramane pentru compatibilitate. +- [ ] Semantica semnaturii corpusului ramane neschimbata: acelasi corpus SILVER -> nu se reindexeaza; corpus schimbat -> reindexare (dar embed doar pe diferenta). +- [ ] Calea de request (block=False) ramane non-blocanta, neatinsa. +- [ ] Degradare gratioasa pastrata: orice eroare de cache (tabela lipsa, blob corupt) -> fallback pe embed complet, fara exceptie propagata. +- [ ] Blob corupt (lungime gresita) e ignorat si re-vectorizat, nu crapa indexarea. +- [ ] `python3 -m pytest -q` trece. + +### US-004: Invalidare pe model + curatare orfane +**Ca** sistem +**Vreau** cache-ul legat de numele modelului si curatat de intrari moarte +**Pentru ca** schimbarea modelului sa nu serveasca vectori incompatibili, iar cache-ul sa nu creasca nelimitat. + +**Acceptance Criteria:** +- [ ] Cheia de cache include `model` (= `FASTEMBED_MODEL`): la schimbarea constantei, vectorii vechi NU sunt folositi (se re-vectorizeaza sub noul model, intrarile vechi raman inerte). +- [ ] La finalul unei indexari reusite, intrarile `embedding_cache` ale modelului curent cu `text_hash` care nu mai exista in corpus sunt sterse (DELETE orfane). +- [ ] Test: schimbarea modelului (mock cu nume diferit) -> zero hit-uri din cache; stergerea unui rand din corpus -> intrarea orfana dispare dupa reindexare. +- [ ] `python3 -m pytest -q` trece. + +### US-005: Teste de flux cold/warm start +**Ca** developer +**Vreau** teste care fixeaza comportamentul incremental +**Pentru ca** regresiile de performanta la startup sa fie prinse de suita. + +**Acceptance Criteria:** +- [ ] Test cold start: cache gol + corpus de N randuri -> backend mock primeste exact N texte, cache-ul contine N intrari dupa. +- [ ] Test warm start: al doilea proces (engine nou, acelasi conn) -> backend mock primeste ZERO texte de corpus, `has_corpus()` True, `suggest_nearest` functioneaza cu vectorii din cache. +- [ ] Test incremental: se adauga 1 rand in `mapping_suggestions` -> backend mock primeste exact 1 text la reindexare. +- [ ] Testele folosesc backend mock (fara fastembed real), ruleaza in suita implicita (fara marker `live`). +- [ ] `python3 -m pytest -q` trece. + +### US-006: Verificare end-to-end pe instanta reala +**Ca** operator +**Vreau** confirmarea pe serverul real ca restartul e ieftin +**Pentru ca** metrica de succes sa fie masurata, nu presupusa. + +**Acceptance Criteria:** +- [ ] Log la finalul warmup-ului: numar vectori din cache / numar vectori calculati / durata totala (ex. `embeddings: warmup ok cache=17181 embed=0 in 12.3s`). +- [ ] Verify in browser: dupa `./start.sh stop && ./start.sh test both --send`, primul restart populeaza cache-ul; la al DOILEA restart, log-ul arata `embed=0` si sugestia "similaritate" apare in preview-ul regulilor text in < 30s de la pornire. +- [ ] Dimensiunea bazei creste cu ~26-30MB (17k x 384 float32) — documentat in log/PRD, acceptat. + +## 4. Cerinte Functionale + +1. [REQ-001] Sistemul trebuie sa persiste vectorii corpusului SILVER in tabela `embedding_cache`, cheie `(model, text_hash)`. +2. [REQ-002] La warmup, sistemul trebuie sa vectorizeze DOAR textele fara intrare in cache pentru modelul curent. +3. [REQ-003] Cand corpusul SILVER se schimba (semnatura diferita), sistemul reindexeaza folosind cache-ul si vectorizeaza doar diferenta. +4. [REQ-004] Cand modelul se schimba, cache-ul vechi nu este folosit; vectorii se recalculeaza integral sub noua cheie de model. +5. [REQ-005] Orice eroare legata de cache degradeaza la comportamentul actual (embed complet), fara a bloca API-ul sau a arunca in request. +6. [REQ-006] Intrarile orfane (texte disparute din corpus) se curata la reindexare. +7. [REQ-007] Vectorii query (textul operatiei/pattern-ului cautat) NU se cacheaza — se calculeaza la cerere (milisecunde). + +## 5. Non-Goals (Ce NU facem) + +- NU introducem vector store extern (sqlite-vec, Qdrant, FAISS) — cosine in Python peste 17k vectori ramane suficient. +- NU cacheam embeddings pentru query-uri (texte noi cautate de utilizatori). +- NU modificam pipeline-ul de etichetare (`tools/mapare-llm/`) sau formatul seedului `app/data/operatii-etichetate.json`. +- NU vectorizam GOLD (`shared_mappings`) sau nomenclatorul — sursa corpusului ramane exclusiv `mapping_suggestions`. +- NU schimbam modelul de embeddings, pragul `EMB_MIN_SIMILARITATE` sau precedenta GOLD > SILVER > embeddings. +- NU eliminam warmup-ul din fundal — ramane, doar devine ieftin. + +## 6. Consideratii Tehnice + +### Stack/Tehnologii +- SQLite existent (acelasi fisier `data/autopass.db`, WAL), stdlib `array`/`hashlib` — zero dependinte noi. +- fastembed/ONNX ramane sursa vectorilor la miss. + +### Patterns de Urmat +- Degradare gratioasa ca in `ensure_embeddings_corpus`/`enrich_suggestions` (except -> pass, corpus gol e acceptabil). +- Schema idempotenta in `app/schema.sql` + `init_db()` (CREATE IF NOT EXISTS), ca restul tabelelor. +- Semnatura corpusului (`_corpus_signature_silver`) ramane mecanismul de decizie "reindexez sau nu". +- Backend injectabil in `EmbeddingEngine` pentru teste (mock, fara model real). + +### Dependente +- Warmup-ul din fundal existent (`app/main.py::_warmup_embeddings`, `ensure_embeddings_corpus(block=True)`). +- `mapping_suggestions` populat la `init_db` din seedul comis (`app/operatii_seed.py`). + +### Riscuri Tehnice +- Scriere concurenta API/worker pe SQLite: scrierile de cache se fac in tranzactii scurte, batch, doar din thread-ul de warmup al API-ului (worker-ul NU incarca modelul — invariant existent). +- Blob corupt/lungime gresita: tratat ca miss (re-embed), nu ca eroare. +- Crestere DB ~26-30MB: acceptata; curatarea orfanelor previne cresterea nelimitata. +- Hash-ul textului trebuie calculat pe `denumire_normalizata` EXACT cum intra in `index_corpus` — altfel miss permanent si cache inutil (test dedicat in US-005). + +## 7. Consideratii UI/UX + +Fara schimbari de UI. Efect indirect: sugestia "similaritate" (editor mapari + preview reguli text) devine disponibila la ~10-20s dupa restart in loc de 1-2 minute. + +## 8. Success Metrics + +- Apeluri `embed()` pe corpus la restart cu corpus neschimbat: 0 (log `embed=0`). +- Durata warmup la warm start: < 30s total (dominata de incarcarea modelului), fata de ~1-2 min acum. +- Zero regresii in suita: `python3 -m pytest -q` verde. + +## 9. Open Questions + +- [x] Vrem o comanda CLI de administrare (`python3 -m tools.embcache rebuild|clear|stats`)? DECIS (/autoplan): manual (`DELETE FROM embedding_cache`) suficient in v1; CLI deferat la TODOS.md. +- [x] Purjarea orfanelor sterge si intrarile modelelor VECHI? DECIS (/autoplan): DA, la aceeasi trecere — simplu si sigur. + +## 10. Amendamente /autoplan (2026-07-06) + +Acceptate in scope (auto-decizii, vezi Decision Audit Trail): + +1. **A1** Codul de cache traieste in modul NOU `app/embedding_cache.py` (functiile cu `conn`); `app/embeddings.py` ramane fara acces DB. `EmbeddingEngine.index_corpus` primeste parametru nou `vectors: list[list[float]] | None`, aliniat pozitional cu `items`; `None` = comportamentul actual (embed tot). +2. **A2** `save_vectors` scrie in tranzactii chunk-uite de **500 randuri/commit** — AMENDEAZA AC-ul US-002 "commit la final". Motiv: un commit unic de ~26MB tine write-lock in timp ce worker-ul face `BEGIN IMMEDIATE` pe `submissions`. +3. **A3** Purjarea orfanelor: diff in Python (set hash-uri corpus vs set hash-uri cache), apoi `DELETE` pe chunk-uri de 500 pe PK. Include intrarile modelelor vechi (`model != curent`). Motiv: tranzactii scurte; (limita de parametri SQLite e 32766 pe 3.32+, nu 999 — chunking-ul e pentru durata lock-ului). +4. **A4** Orice eroare a cailor de cache (SELECT/INSERT/blob corupt) se logheaza cu `log.warning` INAINTE de fallback-ul pe embed complet. Degradarea existenta (`except -> pass` pe indexare) ramane neschimbata — scoping strict pe codul nou. +5. **A5** Metrica "acelasi cod sugerat" se verifica cu toleranta float32 (round-trip float64->float32 poate muta scoruri la a 7-a zecimala): test de echivalenta a RANKING-ULUI, nu egalitate stricta de scor. +6. **A6** Teste suplimentare la US-005: (a) chunking >500 randuri (scriere multi-tranzactie completa), (b) idempotenta la cache partial (populare intrerupta -> urmatorul warmup completeaza, fara duplicate, `INSERT OR REPLACE`). +7. **A7** Nota bounded acceptata: purjarea ruleaza doar la finalul unei indexari reusite; cu embeddings dezactivat sau warmup esuat permanent, cache-ul mort ramane pana la reset manual (`DELETE FROM embedding_cache`). + +Adaugate de review-ul Eng (subagent independent, toate verificate in cod): + +8. **A8** `index_corpus(items, signature, vectors)` VALIDEAZA `len(vectors) == len(items)` si absenta `None`-urilor; mismatch => `log.warning` + fallback pe embed complet. Test obligatoriu de ranking EXACT la warm start: >=2 query-uri distincte care aserteaza CODUL vecinului cel mai apropiat (nu doar non-empty) — singurul test care prinde o permutare de aliniere (misalignment = coduri gresite sugerate, silentios). +9. **A9** Secventa hash->load->embed->save->purge->index e protejata de un lock de modul (`threading.Lock`, pattern `_engine_lock` existent). CORECTEAZA riscul din sectiunea 6: "scrierile doar din thread-ul de warmup" e FALS — calea de request (`block=False`) trece de gate-ul `is_loaded()` (mapping.py:675) dupa incarcarea modelului si ar scrie/purja concurent cu warmup-ul (purjare falsa de randuri noi). +10. **A10** Conexiunile sunt autocommit (`isolation_level=None`, db.py:16): `save_vectors`/`purge_stale` emit EXPLICIT `BEGIN`/`COMMIT` per chunk — `conn.commit()` singur e no-op si ar produce silentios o tranzactie per rand (17k commit-uri), anuland exact motivul A2. +11. **A11** Numele modelului = PARAMETRU explicit pe tot lantul (inclusiv `ensure_embeddings_corpus`); nu se hardcodeaza `FASTEMBED_MODEL` in adancime — altfel testele cu mock scriu sub cheia modelului real si US-004 nu are unde injecta numele. +12. **A12** Vectorii din cache raman `array('f')` end-to-end (fara conversie la `list[float]`) — RAM ~8x mai mic; prag documentat: designul curent (cosine pur Python) e validat pana la ~50k randuri de corpus; peste, vezi TODOS (numpy). +13. **A13** Teste suplimentare obligatorii: (a) purjarea NU ruleaza cand indexarea a esuat (altfel un esec golește cache-ul permanent); (b) `len(vectors) != len(items)` => fallback gratios; (c) save partial esuat => indexarea continua din vectorii din RAM; (d) TOATE testele noi seteaza explicit `AUTOPASS_EMBEDDINGS_ENABLED=1` — gate-ul din mapping.py:671 face altfel testele vacuos-verzi. +14. **A14** Validare `len(vector) == 384` si la SCRIERE (nu doar la citire) — un backend defect nu umple DB-ul cu blob-uri arbitrare. +15. **A15** Orchestrarea hash->load->embed->save->purge se extrage intr-o functie in `embedding_cache.py` care primeste `embed_fn` — `ensure_embeddings_corpus` ramane citibila (~30 linii, nu ~80 cu 6 responsabilitati). + +### Diagrama arhitectura + +``` + API startup (lifespan) request path (neatins) + | | + _warmup_embeddings (thread) enrich_suggestions + | | + ensure_embeddings_corpus(conn, block=True) suggest_nearest (query embed la cerere) + | + |-- SELECT mapping_suggestions ----------- corpus SILVER (sursa) + |-- _corpus_signature_silver ------------- decizia reindexarii (NESCHIMBATA) + |-- [NOU] embedding_cache.load_cached_vectors(conn, model, hashes) + | (SELECT chunk-uit; blob corupt => miss) + |-- embed() DOAR pe miss-uri (backend fastembed/mock) + |-- [NOU] embedding_cache.save_vectors(conn, model, items_noi) + | (INSERT OR REPLACE, commit la 500) + |-- [NOU] purge_stale(conn, model, hash-uri_corpus) + | (orfane + modele vechi, DELETE chunk-uit) + `-- EmbeddingEngine.index_corpus(items, signature, vectors=toate) + + SQLite (WAL, busy_timeout 15s): embedding_cache (model, text_hash) PK + Worker: NU atinge embedding_cache (nu incarca modelul — invariant existent) +``` + +### Data flow cu shadow paths + +``` + corpus rows ─▶ filtrare denumire goala ─▶ hash sha256(text) ─▶ cache lookup ─▶ embed(miss) ─▶ save ─▶ index + │ │ │ │ │ │ + [gol? => return] [toate goale? => [hash pe lista [tabela lipsa/ [backend [OperationalError + cache neatins] corpus gol, no-op] FILTRATA, nu pe DB locked => esueaza => la commit => + rows brute!] warning + warning + warning + fallback; + embed complet] corpus gol] cache partial OK] +``` + +### Error & Rescue Registry + +| Codepath | Ce poate esua | Exceptie | Rescued? | Actiune | User vede | +|---|---|---|---|---|---| +| load_cached_vectors | DB locked/busy | sqlite3.OperationalError | DA | log.warning + fallback embed complet | nimic (warmup mai lung) | +| load_cached_vectors | blob lungime gresita | (validare explicita) | DA | tratat ca miss, re-embed, log count | nimic | +| save_vectors | lock la commit chunk | sqlite3.OperationalError | DA | log.warning; cache partial ramane valid (idempotent) | nimic | +| purge_stale | lock la DELETE | sqlite3.OperationalError | DA | log.warning; orfanele raman pana la urmatoarea trecere | nimic | +| hash drift (normalizare schimbata) | miss permanent | — (nu e exceptie) | DA (observabil) | log warmup arata embed=N cand se astepta 0 (US-006) | nimic | +| ensure_embeddings_corpus | orice alta exceptie | Exception (catch-all EXISTENT) | DA | pass (degradare documentata REQ-005) | sugestii embeddings lipsesc | + +### Failure Modes Registry + +| Codepath | Failure mode | Rescued? | Test? | User vede? | Logged? | +|---|---|---|---|---|---| +| cold start | cache gol | DA (embed tot) | US-005 | fereastra 1-2 min (o data) | DA (US-006) | +| warm start | cache complet | — | US-005 | fereastra ~10-20s | DA | +| incremental | N texte noi | DA | US-005 | — | DA | +| blob corupt | re-embed silentios | DA | US-003 AC | nimic | DA (A4) | +| populare intrerupta | cache partial | DA (idempotent) | A6b | nimic | DA | +| model schimbat | zero hit, re-embed sub cheie noua | DA | US-004 | fereastra 1-2 min (o data) | DA | +| warmup esuat permanent | cache mort ne-purjat | partial (A7, manual) | — | nimic (bounded, ~30MB) | DA | + +**Zero CRITICAL GAP** (niciun rand cu Rescued=N + Test=N + Silent). + +### NOT in scope (deferate, cu motiv) + +- CLI `tools/embcache` — manual suficient in v1 (Open Q1). +- Precompute vectori la build (artefact Docker) — ar face si primul start ieftin; separat. +- Graceful reload / zero-downtime — ar elimina complet fereastra post-restart; infra, nu app. +- Optimizare latenta per-query `suggest_nearest` (numpy) — de masurat intai; PRD-ul declara cosine Python suficient (non-goal). +- Cache query-uri (REQ-007), vector store extern, panou admin stats — non-goals/skip. + +### What already exists (reuse) + +- `_corpus_signature_silver` (mapping.py:635) — decizia reindexarii, pastrata identic. +- `ensure_embeddings_corpus` (mapping.py:649) — singurul punct de interceptie pe flux. +- `EmbeddingEngine(backend=...)` (embeddings.py:84) — mock-ul de test existent acopera US-005. +- Schema idempotenta `schema.sql` + `init_db()` — pattern-ul tabelei noi. +- `PRAGMA busy_timeout=15000` + WAL (db.py:18-20) — deja setate, sustin scrierile concurente. + +### Dream state delta + +Planul duce sistemul de la "restart plateste intreg corpusul" la "restart plateste doar delta"; cheia per-text face re-etichetarile (cod/is_nul) GRATUITE la restart — exact traiectoria 12 luni (corpus SILVER in crestere organica). Ramane in afara idealului: fereastra de incarcare a modelului (~10-20s) si costul per-query al cosine-ului pur Python (ambele in TODOS). + +### Diagrama de acoperire teste (Eng S3) + +``` +CODE PATHS (toate NOI — teste planificate, nu existente) ACOPERIT DE +[+] app/embedding_cache.py + ├── vector_to_blob/blob_to_vector round-trip US-002 AC (+A5 toleranta) + ├── blob corupt (lungime gresita) => miss US-003 AC + ├── len(vector)!=384 la SCRIERE => respins A14 + ├── save_vectors chunk 500 + BEGIN/COMMIT explicit A6a + A10 + ├── save partial esuat => indexare continua din RAM A13c + └── purge orfane + modele vechi; NU ruleaza la index esuat US-004 + A13a +[+] app/embeddings.py :: index_corpus(vectors=) + ├── len(vectors)!=len(items) => fallback embed complet A8/A13b + └── ranking EXACT la warm start (>=2 query-uri, cod asertat) A8 +[+] app/mapping.py :: ensure_embeddings_corpus + ├── cold start: exact N texte la backend mock US-005 + ├── warm start: ZERO texte, has_corpus True, suggest ok US-005 + ├── incremental: +1 rand => exact 1 embed US-005 + ├── model schimbat => zero hit-uri cache US-004 + ├── hash pe lista FILTRATA (denumire goala exclusa) risc §6, test dedicat + └── concurenta warmup/request sub lock A9 (test cu 2 threaduri) +TOATE testele: AUTOPASS_EMBEDDINGS_ENABLED=1 explicit A13d (anti-vacuos) + +COVERAGE PLANIFICAT: 15/15 cai identificate au test specificat | GAPS ramase: 0 +E2E: verify pe instanta reala (US-006) — operational, in afara suitei +``` + +### Paralelizare worktree (Eng) + +Implementare secventiala, fara oportunitate de paralelizare: lantul schema.sql → embedding_cache.py → embeddings.py (vectors=) → mapping.py (wiring) → teste e strict dependent, toate in acelasi modul de suggestii. Un singur lane. + +## Implementation Tasks +Sintetizate din finding-urile review-ului. P1 blocheaza ship-ul. + +- [ ] **T1 (P1, human: ~3h / CC: ~20min)** — embedding_cache — Modul nou: serializare array('f') end-to-end, load/save chunk 500 cu BEGIN/COMMIT explicit, validare dim la scriere/citire, purge diff-Python (orfane + modele vechi), log.warning pe erori — Surfaced by: US-001/002/004 + A2/A3/A4/A10/A12/A14 — Files: app/schema.sql, app/embedding_cache.py — Verify: pytest teste noi serializare/chunk/purge +- [ ] **T2 (P1, human: ~2h / CC: ~15min)** — embeddings — `index_corpus(vectors=)` cu validare aliniere + fallback; model ca parametru pe lant — Surfaced by: A8/A11 — Files: app/embeddings.py — Verify: test ranking exact + test mismatch +- [ ] **T3 (P1, human: ~2h / CC: ~15min)** — mapping — Wiring `ensure_embeddings_corpus`: orchestrare extrasa cu embed_fn (A15), lock de modul (A9), hash pe lista filtrata — Surfaced by: US-003 + A9/A15 — Files: app/mapping.py, app/embedding_cache.py — Verify: teste cold/warm/incremental + test concurenta +- [ ] **T4 (P1, human: ~2h / CC: ~15min)** — teste — Suita completa din diagrama (15 cai), toate cu AUTOPASS_EMBEDDINGS_ENABLED=1 — Surfaced by: US-005 + A6/A13 — Files: tests/ — Verify: python3 -m pytest -q +- [ ] **T5 (P2, human: ~30min / CC: ~5min)** — observabilitate — Log warmup `cache=N embed=M in Xs` — Surfaced by: US-006 — Files: app/mapping.py sau app/main.py — Verify: manual pe ./start.sh, doua restarturi +- [ ] **T6 (P2, human: ~15min / CC: ~2min)** — teste — Test echivalenta ranking cu toleranta float32 — Surfaced by: A5 — Files: tests/ — Verify: pytest + +### Decision Audit Trail + + + +| # | Phase | Decision | Classification | Principle | Rationale | Rejected | +|---|-------|----------|----------------|-----------|-----------|----------| +| 1 | CEO 0C-bis | Design per-text (model, text_hash), nu blob monolit pe semnatura | Mechanical | P1 | Doar per-text da embed=0 la relabel (workflow activ: re-etichetare Haiku); semnatura include cod/is_nul, textul nu se schimba | blob monolit (re-embed tot la orice relabel) | +| 2 | CEO 0E | Modul nou app/embedding_cache.py; embeddings.py ramane fara DB | Mechanical | P5 | Separare existenta curata (engine pur / mapping face DB) | functii in embeddings.py | +| 3 | CEO 0D | Open Q2: purjare include modelele vechi | Mechanical | P2/P3 | ~5 linii, previne crestere moarta; propunerea PRD | lasare inerta | +| 4 | CEO 0D | CLI embcache deferat la TODOS | Mechanical | P3/P6 | PRD propune manual in v1 | build acum | +| 5 | CEO S2 | log.warning pe erorile cailor de cache (A4) | Mechanical | P1 (zero silent failures) | fallback-ul ramane, dar devine vizibil | pass mut | +| 6 | CEO S6/S5 | Metrica cu toleranta float32 + test echivalenta ranking (A5) | Mechanical | P1 | egalitate stricta ar esua legitim pe round-trip float32 | egalitate stricta | +| 7 | CEO S7/S1 | Commit chunk-uit 500 (A2) + purjare diff-Python chunk-uita (A3) | Mechanical | P5 | tranzactii scurte vs worker BEGIN IMMEDIATE | commit unic ~26MB | +| 8 | CEO S6 | Teste chunking + idempotenta cache partial (A6) | Mechanical | P1 | fixeaza comportamentul multi-tranzactie | doar testele US-005 | +| 9 | CEO 0D | Metrici Prometheus warmup | TASTE | — | log-ul US-006 acopera nevoia; counterul e nice-to-have | — (la gate) | +| 10 | CEO OV | Sidecar DB separat (voce externa) vs acelasi autopass.db | TASTE | — | user a confirmat P3 la D1; vocea externa ridica bloat backup + no-shrink fara VACUUM | — (la gate) | +| 11 | CEO 0A | Landscape check din cunostinte in-distributie (fara WebSearch) | Mechanical | P3 | pattern Layer 1 arhicunoscut (cache pe (model, hash(text))) | cautare web | +| 12 | CEO S4 | Dublu-index teoretic warmup/request: no action | Mechanical | P6 | SUPERSEDED de #14 — evaluarea era corecta doar pentru planul FARA cache; cu scrieri+DELETE, riscul devine real | lock nou | +| 13 | Eng S1 | A8: validare aliniere vectors/items + test ranking exact | Mechanical | P1 | misalignment = coduri gresite silentios; zip nu detecteaza nimic (embeddings.py:158) | liste paralele nevalidate | +| 14 | Eng S1 | A9: lock de modul pe secventa de cache (supersedes #12) | Mechanical | P1 | calea de request scrie si ea dupa is_loaded() (mapping.py:675); purjare concurenta = stergere falsa de randuri noi | write-uri doar pe block=True (re-embed repetat pe request path) | +| 15 | Eng S2 | A10: BEGIN/COMMIT explicit per chunk (autocommit db.py:16) | Mechanical | P5 | conn.commit() e no-op in autocommit; chunking-ul A2 nu s-ar intampla de fapt | commit() naiv | +| 16 | Eng S2 | A11: model ca parametru explicit pe tot lantul | Mechanical | P5 | testabilitate US-004; mock-urile nu trebuie sa scrie sub cheia modelului real | FASTEMBED_MODEL hardcodat | +| 17 | Eng S4 | A12: array('f') end-to-end + prag documentat ~50k | Mechanical | P1 | RAM ~8x mai mic gratis; blob_to_vector trece oricum prin array('f') | list[float] boxed (~2GB la 170k) | +| 18 | Eng S3 | A13: 4 teste suplimentare (purge-dupa-esec, mismatch, save partial, env flag explicit) | Mechanical | P1 | teste vacuos-verzi cu embeddings dezactivat = mai rau decat lipsa lor | doar US-005+A6 | +| 19 | Eng S3 | A14: validare dimensiune vector la scriere | Mechanical | P1 | o linie; previne umplerea DB cu blob-uri arbitrare | validare doar la citire | +| 20 | Eng S2 | A15: orchestrare extrasa cu embed_fn injectat | Mechanical | P5 | ensure_embeddings_corpus ramane citibila | functie de 80 linii cu 6 responsabilitati | + +## GSTACK REVIEW REPORT + +| Review | Trigger | Why | Runs | Status | Findings | +|--------|---------|-----|------|--------|----------| +| CEO Review | `/plan-ceo-review` | Scope & strategy | 1 | issues_open (via /autoplan) | 6 propuneri, 2 acceptate, 2 deferate, 2 taste la gate | +| Codex Review | `/codex review` | Independent 2nd opinion | 1 | issues_found (subagent-only) | Codex indisponibil (usage limit); voce externa = subagent Claude, 6 finding-uri CEO + 8 Eng | +| Eng Review | `/plan-eng-review` | Architecture & tests (required) | 1 | clean (via /autoplan) | 8 issues, 0 critical gaps — toate absorbite ca A8-A15 | +| Design Review | `/plan-design-review` | UI/UX gaps | 0 | SKIPPED | fara scope UI (grep 0 potriviri) | +| DX Review | `/plan-devex-review` | Developer experience gaps | 0 | SKIPPED | fara suprafata developer-facing noua | + +**CROSS-MODEL:** indisponibil — Codex pe usage limit pana la 18 iul; toate vocile externe au rulat ca subagenti Claude independenti (context proaspat, fara istoricul review-ului). + +**VERDICT:** CEO + ENG CLEARED — APPROVED la gate-ul /autoplan (2026-07-06): taste #9 = SKIP metrici Prometheus (log-ul US-006 suficient); taste #10 = cache ramane in autopass.db (sidecar respins; caveat VACUUM documentat). Ready to implement. + +NO UNRESOLVED DECISIONS diff --git a/tests/test_embedding_cache.py b/tests/test_embedding_cache.py new file mode 100644 index 0000000..4cb1a7a --- /dev/null +++ b/tests/test_embedding_cache.py @@ -0,0 +1,277 @@ +"""Teste pentru app/embedding_cache.py -- cache persistent de vectori in SQLite. + +Acopera: round-trip serializare, blob corupt = miss, validare dimensiune la scriere, +chunking >500 randuri, save partial esuat, purge orfane + modele vechi, purge NU +ruleaza cand embed_fn esueaza. +""" +from __future__ import annotations + +import os +import sqlite3 +import tempfile +from array import array + +import pytest + +from app.embedding_cache import ( + EMB_DIM, + blob_to_vector, + load_cached_vectors, + purge_stale, + save_vectors, + sync_corpus_vectors, + text_hash, + vector_to_blob, +) + + +@pytest.fixture() +def conn(monkeypatch): + tmp = tempfile.mkdtemp() + monkeypatch.setenv("AUTOPASS_DB_PATH", os.path.join(tmp, "embcache.db")) + monkeypatch.setenv("AUTOPASS_WEB_AUTH_REQUIRED", "false") + monkeypatch.setenv("AUTOPASS_EMBEDDINGS_ENABLED", "true") # A13d: anti-vacuos + from app.config import get_settings + get_settings.cache_clear() + from app.db import init_db, get_connection + init_db() + c = get_connection() + yield c + c.close() + get_settings.cache_clear() + + +def _vec(seed: float = 1.0, dim: int = EMB_DIM) -> list[float]: + return [seed + i * 0.001 for i in range(dim)] + + +# --------------------------------------------------------------------------- # +# Serializare # +# --------------------------------------------------------------------------- # + +def test_roundtrip_vector_to_blob_blob_to_vector(): + v = _vec(3.5) + blob = vector_to_blob(v) + assert isinstance(blob, bytes) + assert len(blob) == EMB_DIM * 4 + out = blob_to_vector(blob) + assert isinstance(out, array) + assert out == pytest.approx(v, rel=1e-6) + + +def test_text_hash_deterministic_and_distinct(): + assert text_hash("SCHIMB ULEI") == text_hash("SCHIMB ULEI") + assert text_hash("SCHIMB ULEI") != text_hash("SCHIMB FILTRU") + + +# --------------------------------------------------------------------------- # +# load_cached_vectors # +# --------------------------------------------------------------------------- # + +def test_load_cached_vectors_roundtrip(conn): + h = text_hash("SCHIMB ULEI") + save_vectors(conn, "model-a", [(h, _vec(1.0))]) + out = load_cached_vectors(conn, "model-a", [h]) + assert h in out + assert out[h] == pytest.approx(_vec(1.0), rel=1e-6) + + +def test_load_cached_vectors_miss_on_missing_hash(conn): + out = load_cached_vectors(conn, "model-a", [text_hash("NECUNOSCUT")]) + assert out == {} + + +def test_load_cached_vectors_blob_corupt_e_miss(conn): + h = text_hash("SCHIMB ULEI") + conn.execute("BEGIN") + conn.execute( + "INSERT INTO embedding_cache (text_hash, model, vector) VALUES (?, ?, ?)", + (h, "model-a", b"\x00\x01\x02"), # lungime gresita + ) + conn.execute("COMMIT") + out = load_cached_vectors(conn, "model-a", [h]) + assert h not in out + + +def test_load_cached_vectors_scoped_pe_model(conn): + h = text_hash("SCHIMB ULEI") + save_vectors(conn, "model-a", [(h, _vec(1.0))]) + out = load_cached_vectors(conn, "model-b", [h]) + assert out == {} + + +# --------------------------------------------------------------------------- # +# save_vectors: validare dimensiune + chunking # +# --------------------------------------------------------------------------- # + +def test_save_vectors_respinge_dimensiune_gresita(conn): + h = text_hash("SCHIMB ULEI") + save_vectors(conn, "model-a", [(h, [1.0, 2.0, 3.0])]) # nu e EMB_DIM + out = load_cached_vectors(conn, "model-a", [h]) + assert h not in out + + +def test_save_vectors_chunking_peste_500_randuri(conn): + items = [(text_hash(f"text-{i}"), _vec(float(i))) for i in range(1200)] + save_vectors(conn, "model-a", items) + hashes = [h for h, _ in items] + out = load_cached_vectors(conn, "model-a", hashes) + assert len(out) == 1200 + for h, vec in items: + assert out[h] == pytest.approx(vec, rel=1e-6) + + +def test_save_vectors_insert_or_replace_idempotent(conn): + h = text_hash("SCHIMB ULEI") + save_vectors(conn, "model-a", [(h, _vec(1.0))]) + save_vectors(conn, "model-a", [(h, _vec(2.0))]) # populare intrerupta -> completare + row = conn.execute( + "SELECT COUNT(*) AS n FROM embedding_cache WHERE model=? AND text_hash=?", + ("model-a", h), + ).fetchone() + assert row["n"] == 1 + out = load_cached_vectors(conn, "model-a", [h]) + assert out[h] == pytest.approx(_vec(2.0), rel=1e-6) + + +# --------------------------------------------------------------------------- # +# purge_stale # +# --------------------------------------------------------------------------- # + +def test_purge_stale_sterge_orfane_model_curent(conn): + h1, h2 = text_hash("A"), text_hash("B") + save_vectors(conn, "model-a", [(h1, _vec(1.0)), (h2, _vec(2.0))]) + purge_stale(conn, "model-a", corpus_hashes=[h1]) # h2 nu mai e in corpus + out = load_cached_vectors(conn, "model-a", [h1, h2]) + assert h1 in out + assert h2 not in out + + +def test_purge_stale_sterge_modele_vechi(conn): + h = text_hash("A") + save_vectors(conn, "model-old", [(h, _vec(1.0))]) + save_vectors(conn, "model-a", [(h, _vec(2.0))]) + purge_stale(conn, "model-a", corpus_hashes=[h]) + assert load_cached_vectors(conn, "model-old", [h]) == {} + assert h in load_cached_vectors(conn, "model-a", [h]) + + +# --------------------------------------------------------------------------- # +# sync_corpus_vectors # +# --------------------------------------------------------------------------- # + +def test_sync_corpus_vectors_cold_start(conn): + calls = [] + + def embed_fn(texts): + calls.append(list(texts)) + return [_vec(float(i)) for i in range(len(texts))] + + texts = ["A", "B", "C"] + vecs = sync_corpus_vectors(conn, "model-a", texts, embed_fn) + assert len(vecs) == 3 + assert calls == [texts] # toate 3 lipsesc din cache -> toate trimise la embed + hashes = [text_hash(t) for t in texts] + cached = load_cached_vectors(conn, "model-a", hashes) + assert len(cached) == 3 + + +def test_sync_corpus_vectors_warm_start_zero_embed_calls(conn): + calls = [] + + def embed_fn(texts): + calls.append(list(texts)) + return [_vec(float(i)) for i in range(len(texts))] + + texts = ["A", "B", "C"] + sync_corpus_vectors(conn, "model-a", texts, embed_fn) + calls.clear() + vecs2 = sync_corpus_vectors(conn, "model-a", texts, embed_fn) + assert calls == [] # nimic nou de vectorizat + assert len(vecs2) == 3 + + +def test_sync_corpus_vectors_incremental_un_text_nou(conn): + calls = [] + + def embed_fn(texts): + calls.append(list(texts)) + return [_vec(float(i)) for i in range(len(texts))] + + sync_corpus_vectors(conn, "model-a", ["A", "B"], embed_fn) + calls.clear() + sync_corpus_vectors(conn, "model-a", ["A", "B", "C"], embed_fn) + assert calls == [["C"]] + + +def test_sync_corpus_vectors_nu_purjeaza_singur_apelantul_decide(conn): + """sync_corpus_vectors NU mai purjeaza -- e responsabilitatea apelantului, DUPA + o indexare reusita (vezi ensure_embeddings_corpus). purge_stale ramane apelabil + separat, explicit, de catre apelant.""" + def embed_fn(texts): + return [_vec(float(i)) for i in range(len(texts))] + + sync_corpus_vectors(conn, "model-a", ["A", "B"], embed_fn) + sync_corpus_vectors(conn, "model-a", ["A"], embed_fn) # B disparut din corpus solicitat + out = load_cached_vectors(conn, "model-a", [text_hash("A"), text_hash("B")]) + assert text_hash("A") in out + assert text_hash("B") in out # nepurjat automat -- sync_corpus_vectors nu mai face asta + + purge_stale(conn, "model-a", {text_hash("A")}) # apelantul purjeaza dupa indexare reusita + out2 = load_cached_vectors(conn, "model-a", [text_hash("A"), text_hash("B")]) + assert text_hash("A") in out2 + assert text_hash("B") not in out2 + + +def test_sync_corpus_vectors_nu_purjeaza_cand_embed_fn_esueaza(conn): + def embed_fn_ok(texts): + return [_vec(float(i)) for i in range(len(texts))] + + sync_corpus_vectors(conn, "model-a", ["A", "B"], embed_fn_ok) + + def embed_fn_broken(texts): + raise RuntimeError("model indisponibil") + + with pytest.raises(RuntimeError): + sync_corpus_vectors(conn, "model-a", ["A", "C"], embed_fn_broken) + + # "B" nu a fost purjat -- sync_corpus_vectors nu purjeaza niciodata singur. + out = load_cached_vectors(conn, "model-a", [text_hash("A"), text_hash("B")]) + assert text_hash("A") in out + assert text_hash("B") in out + + +class _LockingConn: + """Wrapper peste o conexiune reala: simuleaza `database is locked` la BEGIN + (exercita try/except-ul din save_vectors, nu il ocoleste).""" + + def __init__(self, real): + self._real = real + + def execute(self, sql, *a, **kw): + if sql.strip() == "BEGIN": + raise sqlite3.OperationalError("database is locked") + return self._real.execute(sql, *a, **kw) + + def executemany(self, *a, **kw): + return self._real.executemany(*a, **kw) + + def __getattr__(self, name): + return getattr(self._real, name) + + +def test_sync_corpus_vectors_save_partial_esuat_continua_din_ram(conn): + """Daca save_vectors esueaza (ex. DB locked la BEGIN), vectorii noi tot se + intorc din RAM -- indexarea continua, doar persistarea in cache rateaza.""" + locking = _LockingConn(conn) + + def embed_fn(texts): + return [_vec(float(i)) for i in range(len(texts))] + + vecs = sync_corpus_vectors(locking, "model-a", ["A", "B"], embed_fn) + assert len(vecs) == 2 + assert vecs[0] == pytest.approx(_vec(0.0), rel=1e-6) + + # Nimic nu a fost persistat (BEGIN a esuat la fiecare chunk). + out = load_cached_vectors(conn, "model-a", [text_hash("A"), text_hash("B")]) + assert out == {} diff --git a/tests/test_embeddings.py b/tests/test_embeddings.py index 4ab2978..023993d 100644 --- a/tests/test_embeddings.py +++ b/tests/test_embeddings.py @@ -166,6 +166,79 @@ def test_index_corpus_no_exception_on_backend_error(): assert engine.suggest_nearest("CEVA") == [] +# --------------------------------------------------------------------------- # +# index_corpus(vectors=) -- precalculati, aliniati cu items (A8) # +# --------------------------------------------------------------------------- # + +def test_index_corpus_vectors_precalculati_nu_apeleaza_backend(): + """Cand `vectors` e furnizat, backend-ul NU e apelat pentru corpus (embed=0).""" + + class NoCallBackend: + def embed(self, texts): + raise AssertionError("backend.embed() NU trebuie apelat cand vectors e furnizat") + + corpus = [ + {"denumire": "SCHIMB ULEI", "cod": "OE-3"}, + {"denumire": "REPARATIE MOTOR", "cod": "OE-1"}, + ] + vectors = [_vec("SCHIMB ULEI"), _vec("REPARATIE MOTOR")] + engine = EmbeddingEngine(backend=NoCallBackend()) + engine.index_corpus(corpus, vectors=vectors) + assert engine.has_corpus() + + +def test_index_corpus_vectors_ranking_exact_warm_start(): + """Vectori precalculati din 'cache' produc EXACT acelasi cod ca embed direct + (nu doar non-empty) -- prinde o eventuala dezaliniere intre items si vectors.""" + corpus = [ + {"denumire": "SCHIMB ULEI MOTOR", "cod": "OE-3"}, + {"denumire": "REPARATIE CUTIE VITEZE", "cod": "OE-1"}, + {"denumire": "VERIFICARE DIRECTIE VOLAN", "cod": "OE-4"}, + {"denumire": "INLOCUIT PLACUTE FRANA", "cod": "OE-2"}, + ] + vectors = [_vec(item["denumire"]) for item in corpus] + engine = EmbeddingEngine(backend=MockBackend()) + engine.index_corpus(corpus, vectors=vectors) + + assert engine.suggest_nearest("SCHIMB ULEI MOTOR", top_k=1)[0]["cod"] == "OE-3" + assert engine.suggest_nearest("VERIFICARE DIRECTIE VOLAN", top_k=1)[0]["cod"] == "OE-4" + + +def test_index_corpus_vectors_mismatch_lungime_fallback_embed_complet(): + """len(vectors) != len(items) -> fallback pe embed complet (backend chemat), fara exceptie.""" + corpus = [ + {"denumire": "SCHIMB ULEI", "cod": "OE-3"}, + {"denumire": "REPARATIE MOTOR", "cod": "OE-1"}, + ] + engine = EmbeddingEngine(backend=MockBackend()) + engine.index_corpus(corpus, vectors=[_vec("SCHIMB ULEI")]) # un singur vector pentru 2 itemi + + assert engine.has_corpus() + results = engine.suggest_nearest("SCHIMB ULEI", top_k=1) + assert results and results[0]["cod"] == "OE-3" + + +def test_index_corpus_vectors_contine_none_fallback_embed_complet(): + """Un `None` in `vectors` -> fallback pe embed complet, fara exceptie.""" + corpus = [ + {"denumire": "SCHIMB ULEI", "cod": "OE-3"}, + {"denumire": "REPARATIE MOTOR", "cod": "OE-1"}, + ] + engine = EmbeddingEngine(backend=MockBackend()) + engine.index_corpus(corpus, vectors=[_vec("SCHIMB ULEI"), None]) + + assert engine.has_corpus() + assert engine.suggest_nearest("SCHIMB ULEI", top_k=1)[0]["cod"] == "OE-3" + + +def test_index_corpus_vectors_none_e_comportamentul_existent(): + """`vectors=None` (default) -> embed complet prin backend, ca inainte.""" + corpus = [{"denumire": "SCHIMB ULEI", "cod": "OE-3"}] + engine = EmbeddingEngine(backend=MockBackend()) + engine.index_corpus(corpus, vectors=None) + assert engine.suggest_nearest("SCHIMB ULEI", top_k=1)[0]["cod"] == "OE-3" + + # --------------------------------------------------------------------------- # # API la nivel de modul (singleton global) # # --------------------------------------------------------------------------- # diff --git a/tests/test_embeddings_warmup_cache.py b/tests/test_embeddings_warmup_cache.py new file mode 100644 index 0000000..bd4bfac --- /dev/null +++ b/tests/test_embeddings_warmup_cache.py @@ -0,0 +1,347 @@ +"""Teste de flux warmup + cache persistent la nivelul `ensure_embeddings_corpus` +(app/mapping.py): cold/warm start, incremental, model schimbat, hash pe lista +filtrata, concurenta warmup/request sub lock (A9), echivalenta ranking cu +toleranta float32 (A5). + +Backend mock determinist (fara fastembed real). Toate testele seteaza explicit +AUTOPASS_EMBEDDINGS_ENABLED=1 (A13d). +""" +from __future__ import annotations + +import hashlib +import os +import tempfile +import threading +import time + +import pytest + +from app.embedding_cache import EMB_DIM + + +def _det_vector(text: str, dim: int = EMB_DIM) -> list[float]: + """Vector determinist (384-dim) derivat din hash-ul textului. Suficient pentru + ranking cosine in teste -- nu evaluam calitatea semantica, doar alinierea/cache-ul.""" + digest = hashlib.sha256(text.encode("utf-8")).digest() + return [((digest[i % len(digest)] + i) % 256) / 255.0 for i in range(dim)] + + +class CountingMockBackend: + """Backend determinist care numara textele primite la fiecare apel embed().""" + + def __init__(self): + self.calls: list[list[str]] = [] + + def embed(self, texts): + self.calls.append(list(texts)) + return [_det_vector(t) for t in texts] + + @property + def total_texts(self) -> int: + return sum(len(c) for c in self.calls) + + +@pytest.fixture() +def env(monkeypatch): + tmp = tempfile.mkdtemp() + monkeypatch.setenv("AUTOPASS_DB_PATH", os.path.join(tmp, "warmup.db")) + monkeypatch.setenv("AUTOPASS_WEB_AUTH_REQUIRED", "false") + monkeypatch.setenv("AUTOPASS_EMBEDDINGS_ENABLED", "true") # A13d: anti-vacuos + from app.config import get_settings + get_settings.cache_clear() + from app.db import init_db + init_db() + yield monkeypatch + get_settings.cache_clear() + + +@pytest.fixture() +def conn(env): + from app.db import get_connection + c = get_connection() + yield c + c.close() + + +def _inject_engine(backend): + import app.embeddings as emb + from app.embeddings import EmbeddingEngine + emb._engine = EmbeddingEngine(backend=backend) + return emb + + +def _seed_silver(conn, rows): + """rows = [(denumire_normalizata, cod, is_nul)].""" + conn.executemany( + "INSERT OR IGNORE INTO mapping_suggestions " + "(denumire_normalizata, cod_prestatie, is_nul, source, confidence) VALUES (?, ?, ?, 'llm_seed', 0.7)", + rows, + ) + conn.commit() + + +# --------------------------------------------------------------------------- # +# Cold / warm / incremental start (US-005) # +# --------------------------------------------------------------------------- # + +def test_cold_start_trimite_exact_n_texte_si_populeaza_cache(conn): + backend = CountingMockBackend() + emb = _inject_engine(backend) + denumiri = ["SCHIMB ULEI MOTOR", "INLOCUIT PLACUTE FRANA", "VERIFICARE DIRECTIE"] + _seed_silver(conn, [(d, "OE-1", 0) for d in denumiri]) + + from app.mapping import ensure_embeddings_corpus + ensure_embeddings_corpus(conn) + assert backend.total_texts == 3 + + from app.embedding_cache import load_cached_vectors, text_hash + cached = load_cached_vectors(conn, emb.FASTEMBED_MODEL, [text_hash(d) for d in denumiri]) + assert len(cached) == 3 + + +def test_warm_start_al_doilea_proces_zero_texte_embed(conn): + denumiri = ["SCHIMB ULEI MOTOR", "INLOCUIT PLACUTE FRANA"] + _seed_silver(conn, [(d, "OE-1", 0) for d in denumiri]) + from app.mapping import ensure_embeddings_corpus + + backend1 = CountingMockBackend() + _inject_engine(backend1) + ensure_embeddings_corpus(conn) + assert backend1.total_texts == 2 + + # Simuleaza un al doilea proces: engine NOU, acelasi conn (cache-ul e in DB, nu in RAM). + backend2 = CountingMockBackend() + emb2 = _inject_engine(backend2) + ensure_embeddings_corpus(conn) + assert backend2.total_texts == 0 + assert emb2.has_corpus() + + res = emb2.suggest_nearest("SCHIMB ULEI MOTOR", top_k=1) + assert res and res[0]["cod"] == "OE-1" + + +def test_incremental_un_rand_nou_trimite_exact_un_text(conn): + from app.mapping import ensure_embeddings_corpus + _seed_silver(conn, [("SCHIMB ULEI MOTOR", "OE-3", 0)]) + _inject_engine(CountingMockBackend()) + ensure_embeddings_corpus(conn) + + _seed_silver(conn, [("INLOCUIT BATERIE", "OE-1", 0)]) + backend2 = CountingMockBackend() + _inject_engine(backend2) + ensure_embeddings_corpus(conn) + assert backend2.total_texts == 1 + assert backend2.calls == [["INLOCUIT BATERIE"]] + + +# --------------------------------------------------------------------------- # +# Model schimbat (US-004) # +# --------------------------------------------------------------------------- # + +def test_model_schimbat_zero_hituri_cache_reindexare_completa_si_purjare(conn, monkeypatch): + import app.embeddings as emb_module + from app.mapping import ensure_embeddings_corpus + + denumiri = ["SCHIMB ULEI MOTOR", "INLOCUIT PLACUTE FRANA"] + _seed_silver(conn, [(d, "OE-1", 0) for d in denumiri]) + model_vechi = emb_module.FASTEMBED_MODEL + + backend_old = CountingMockBackend() + _inject_engine(backend_old) + ensure_embeddings_corpus(conn) + assert backend_old.total_texts == 2 + + monkeypatch.setattr(emb_module, "FASTEMBED_MODEL", "model-nou-v2") + backend_new = CountingMockBackend() + _inject_engine(backend_new) + ensure_embeddings_corpus(conn) + assert backend_new.total_texts == 2 # zero hit-uri sub noul model -> re-vectorizare integrala + + from app.embedding_cache import load_cached_vectors, text_hash + hashes = [text_hash(d) for d in denumiri] + # Intrarile modelului vechi au fost purjate la reindexarea reusita sub noul model (A3). + assert load_cached_vectors(conn, model_vechi, hashes) == {} + assert len(load_cached_vectors(conn, "model-nou-v2", hashes)) == 2 + + +# --------------------------------------------------------------------------- # +# Hash pe lista FILTRATA (denumire goala exclusa din corpus) # +# --------------------------------------------------------------------------- # + +def test_denumire_goala_nu_strica_alinierea_si_nu_produce_miss_permanent(conn): + from app.mapping import ensure_embeddings_corpus + _seed_silver(conn, [ + ("", "OE-9", 0), # denumire_normalizata goala -- exclusa din corpus la filtrare + ("SCHIMB ULEI MOTOR", "OE-3", 0), + ("INLOCUIT PLACUTE FRANA", "OE-1", 0), + ]) + backend1 = CountingMockBackend() + _inject_engine(backend1) + ensure_embeddings_corpus(conn) + assert backend1.total_texts == 2 # doar cele 2 randuri cu denumire nevida + + backend2 = CountingMockBackend() + emb2 = _inject_engine(backend2) + ensure_embeddings_corpus(conn) + assert backend2.total_texts == 0 # warm: fara miss permanent din cauza filtrarii + assert emb2.suggest_nearest("SCHIMB ULEI MOTOR", top_k=1)[0]["cod"] == "OE-3" + + +# --------------------------------------------------------------------------- # +# Concurenta warmup/request sub lock de modul (A9) # +# --------------------------------------------------------------------------- # + +def test_concurenta_doua_threaduri_fara_dubla_vectorizare(conn): + class ReentrancyDetectingBackend: + """Detecteaza executie concurenta reala in embed(): daca lock-ul de modul + NU serializeaza secventa hash->embed->save->purge, `max_active` ar depasi 1.""" + + def __init__(self): + self.active = 0 + self.max_active = 0 + self.total_calls = 0 + self._guard = threading.Lock() + + def embed(self, texts): + with self._guard: + self.active += 1 + self.max_active = max(self.max_active, self.active) + self.total_calls += 1 + time.sleep(0.05) # largeste deliberat fereastra de suprapunere + with self._guard: + self.active -= 1 + return [_det_vector(t) for t in texts] + + backend = ReentrancyDetectingBackend() + emb = _inject_engine(backend) + denumiri = ["SCHIMB ULEI MOTOR", "INLOCUIT PLACUTE FRANA"] + _seed_silver(conn, [(d, "OE-1", 0) for d in denumiri]) + + from app.db import get_connection + from app.mapping import ensure_embeddings_corpus + + errors: list[Exception] = [] + + def _run(): + try: + c = get_connection() + try: + ensure_embeddings_corpus(c, block=True) + finally: + c.close() + except Exception as exc: # pragma: no cover - vizibil doar la regresie + errors.append(exc) + + threads = [threading.Thread(target=_run) for _ in range(2)] + for t in threads: + t.start() + for t in threads: + t.join(timeout=5) + + assert not errors + assert backend.max_active <= 1 # lock-ul de modul serializeaza cele doua treceri + + from app.embedding_cache import load_cached_vectors, text_hash + cached = load_cached_vectors(conn, emb.FASTEMBED_MODEL, [text_hash(d) for d in denumiri]) + assert len(cached) == 2 # niciun rand nou nu a fost purjat fals de trecerea concurenta + + +# --------------------------------------------------------------------------- # +# Echivalenta ranking cu toleranta float32 (T6/A5) # +# --------------------------------------------------------------------------- # + +def test_block_false_nu_asteapta_dupa_warmup_in_curs(conn): + """Calea de request (block=False) nu trebuie sa blocheze cat warmup-ul (block=True) + tine lock-ul -- trebuie sa iasa imediat, nu sa astepte pana termina warmup-ul.""" + warmup_poate_continua = threading.Event() + warmup_a_intrat_in_embed = threading.Event() + + class SlowBackend: + def embed(self, texts): + warmup_a_intrat_in_embed.set() + warmup_poate_continua.wait(timeout=5) + return [_det_vector(t) for t in texts] + + emb = _inject_engine(SlowBackend()) + _seed_silver(conn, [("SCHIMB ULEI MOTOR", "OE-1", 0)]) + + from app.db import get_connection + from app.mapping import ensure_embeddings_corpus + + warmup_thread = threading.Thread( + target=lambda: ensure_embeddings_corpus(get_connection(), block=True) + ) + warmup_thread.start() + assert warmup_a_intrat_in_embed.wait(timeout=5), "warmup trebuia sa ajunga in embed()" + + t0 = time.monotonic() + ensure_embeddings_corpus(conn, block=False) # nu trebuie sa astepte lock-ul + durata_request = time.monotonic() - t0 + assert durata_request < 1.0, "block=False nu are voie sa astepte warmup-ul in curs" + assert not emb.has_corpus() # warmup-ul nu a terminat inca, request-ul a iesit fara sa faca nimic + + warmup_poate_continua.set() + warmup_thread.join(timeout=5) + assert emb.has_corpus() # warmup-ul a terminat normal, neblocat de request + + +def test_indexare_esuata_nu_purjeaza_cache_ul_vechi(conn): + """Cand embed-ul unui text NOU esueaza, indexarea nu se termina cu succes -> + purge_stale nu trebuie sa ruleze, altfel randuri inca valide (disparute doar + din setul CERUT, nu esecul lor) ar fi sterse fals din cache.""" + denumiri_initiale = ["SCHIMB ULEI MOTOR", "INLOCUIT PLACUTE FRANA"] + _seed_silver(conn, [(d, "OE-1", 0) for d in denumiri_initiale]) + _inject_engine(CountingMockBackend()) + + from app.mapping import ensure_embeddings_corpus + ensure_embeddings_corpus(conn) + + from app.embedding_cache import load_cached_vectors, text_hash + hashes_initiale = [text_hash(d) for d in denumiri_initiale] + assert len(load_cached_vectors(conn, "sentence-transformers/paraphrase-multilingual-MiniLM-L12-v2", hashes_initiale)) == 2 + + # Corpusul cerut se schimba: "INLOCUIT PLACUTE FRANA" dispare, apare un text nou + # a carui vectorizare va esua -- indexarea intreaga trebuie sa rateze. + conn.execute( + "DELETE FROM mapping_suggestions WHERE denumire_normalizata=?", + ("INLOCUIT PLACUTE FRANA",), + ) + conn.commit() + _seed_silver(conn, [("TEXT NOU CARE ESUEAZA", "OE-2", 0)]) + + class BrokenOnNewText: + def embed(self, texts): + raise RuntimeError("model indisponibil pentru text nou") + + emb = _inject_engine(BrokenOnNewText()) + ensure_embeddings_corpus(conn) # esueaza intern, prins de degradarea gratioasa + + # Randul disparut din corpusul cerut RAMANE in cache -- purge nu a rulat. + out = load_cached_vectors(conn, emb.FASTEMBED_MODEL, hashes_initiale) + assert len(out) == 2 + + +def test_ranking_echivalent_index_direct_vs_din_cache_float32(conn): + from app.embedding_cache import sync_corpus_vectors + from app.embeddings import EmbeddingEngine + + corpus = [ + {"denumire": "SCHIMB ULEI MOTOR", "cod": "OE-3"}, + {"denumire": "REPARATIE CUTIE VITEZE", "cod": "OE-1"}, + {"denumire": "VERIFICARE DIRECTIE VOLAN", "cod": "OE-4"}, + {"denumire": "INLOCUIT PLACUTE FRANA", "cod": "OE-2"}, + ] + texts = [item["denumire"] for item in corpus] + backend = CountingMockBackend() + + engine_direct = EmbeddingEngine(backend=backend) + engine_direct.index_corpus(corpus) # embed complet, vectori float64 in RAM + + vectors_cache = sync_corpus_vectors(conn, "model-test-ranking", texts, backend.embed) + engine_cache = EmbeddingEngine(backend=backend) + engine_cache.index_corpus(corpus, vectors=vectors_cache) # round-trip float32 din cache + + for query in ("SCHIMB ULEI", "VERIFICARE DIRECTIE"): + r_direct = engine_direct.suggest_nearest(query, top_k=len(corpus)) + r_cache = engine_cache.suggest_nearest(query, top_k=len(corpus)) + assert [r["cod"] for r in r_direct] == [r["cod"] for r in r_cache]