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 <noreply@anthropic.com>
This commit is contained in:
Claude Agent
2026-07-06 21:57:49 +00:00
parent a46e364594
commit 705ad030fe
11 changed files with 1437 additions and 35 deletions

210
app/embedding_cache.py Normal file
View File

@@ -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]

View File

@@ -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]:

View File

@@ -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

View File

@@ -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

View File

@@ -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),