fix(maria): lacat intre reindexari, scriere atomica a indexului, XML invalid raportat

Descoperit la prima sincronizare reala din Drive: `maria-sync.timer` a pornit
peste rularea manuala si doua procese faceau embeddings in paralel pe acelasi
Ollama, ambele urmand sa scrie acelasi rag_index.json. Embedding-ul a incetinit
de la ~7s la ~20s din concurenta, iar ultimul care termina ar fi suprascris
munca celuilalt.

- config.exclusive(): lacat `flock` intre procese, luat la intrarea in sync.py si
  indexer.py. Nu asteapta — a doua rulare iese curat cu "o reindexare e deja in
  curs", fiindca ar reface exact acelasi lucru. Verificat pe procese reale.
- indexer scrie indexul atomic (tmp + os.replace): consumer-ul reciteste fisierul
  la 30s si putea prinde un JSON pe jumatate scris.
- build() intoarce `warnings` pentru XML-urile care nu se pot parsa, iar rularea
  din linia de comanda le scrie in stderr. Pana acum, un XML invalid se indexa
  tacut ca text simplu, cu o singura linie pierduta in log.

Context de performanta, masurat pe LXC 171 fara alta incarcare: un embedding
`nomic-embed-text` ia ~7,3s, deci o reindexare completa a celor 173 de chunk-uri
dureaza ~21 de minute — mai mult decat intervalul timer-ului. Nu e o problema
practica (amprenta reindexeaza doar la schimbare, iar lacatul opreste
suprapunerea), dar explica de ce prima rulare pare blocata.

docs/rclone-google-drive-headless.md: procedura de conectare a unui container
headless la Drive prin `rclone authorize`, cu transcriptul rularii reale de pe
Windows, capcanele (sync e distructiv pe destinatie, connection string in loc de
cale pe nume, unde stau secretele) si de ce nu contul de serviciu. Indexata in
CLAUDE.md.

6 teste noi (lacat, eliberare la exceptie, scriere atomica, avertismente).

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01Q4uzvgm7AyJch5WH8QHRhY
This commit is contained in:
Claude Agent
2026-08-31 18:42:01 +00:00
parent f3f5eedb79
commit 3b7069e954
6 changed files with 243 additions and 5 deletions

View File

@@ -7,6 +7,8 @@ model de citit pentru cine intretine ambele punti pe acest container.
from __future__ import annotations
import contextlib
import fcntl
import os
import pathlib
@@ -16,6 +18,7 @@ _DEFAULT_DIR = pathlib.Path.home() / ".maria-bridge"
STATE_DIR: pathlib.Path = pathlib.Path(os.environ.get("MARIA_BRIDGE_DIR") or _DEFAULT_DIR)
DOCS_DIR: pathlib.Path = STATE_DIR / "documents"
INDEX_FILE: pathlib.Path = STATE_DIR / "rag_index.json"
LOCK_FILE: pathlib.Path = STATE_DIR / ".rag.lock"
LOG_DIR: pathlib.Path = STATE_DIR / "logs"
ENV_FILE: pathlib.Path = STATE_DIR / "env"
AUTH_DIR: pathlib.Path = STATE_DIR / "whatsapp-auth"
@@ -70,12 +73,13 @@ def parse_env(text: str) -> dict[str, str]:
def reload(base_dir: str | os.PathLike | None = None) -> dict[str, str]:
"""Recalculeaza caile si reciteste env-ul. Returneaza dictionarul incarcat."""
global STATE_DIR, DOCS_DIR, INDEX_FILE, LOG_DIR, ENV_FILE, AUTH_DIR, _env
global STATE_DIR, DOCS_DIR, INDEX_FILE, LOCK_FILE, LOG_DIR, ENV_FILE, AUTH_DIR, _env
if base_dir is None:
base_dir = os.environ.get("MARIA_BRIDGE_DIR") or _DEFAULT_DIR
STATE_DIR = pathlib.Path(base_dir)
DOCS_DIR = STATE_DIR / "documents"
INDEX_FILE = STATE_DIR / "rag_index.json"
LOCK_FILE = STATE_DIR / ".rag.lock"
LOG_DIR = STATE_DIR / "logs"
ENV_FILE = STATE_DIR / "env"
AUTH_DIR = STATE_DIR / "whatsapp-auth"
@@ -102,3 +106,30 @@ def get_int(key: str, default: int) -> int:
reload()
class Busy(RuntimeError):
"""Alta reindexare e deja in curs."""
@contextlib.contextmanager
def exclusive():
"""Lacat intre PROCESE pentru reindexare (timer vs dashboard vs rulare manuala).
Reindexarea dureaza minute (embeddings pe CPU). Fara lacat, `maria-sync.timer`
poate porni peste o sincronizare manuala si ambele scriu acelasi rag_index.json.
Nu asteptam: a doua rulare se anuleaza curat, fiindca oricum ar reface acelasi
lucru imediat dupa.
"""
STATE_DIR.mkdir(parents=True, exist_ok=True)
fh = open(LOCK_FILE, "w", encoding="utf-8")
try:
try:
fcntl.flock(fh, fcntl.LOCK_EX | fcntl.LOCK_NB)
except OSError:
raise Busy("o reindexare e deja in curs")
fh.write(str(os.getpid()))
fh.flush()
yield
finally:
fh.close()

View File

@@ -13,6 +13,7 @@ Doua strategii de taiere in chunk-uri:
from __future__ import annotations
import json
import os
import re
import sys
import xml.etree.ElementTree as ET
@@ -132,19 +133,36 @@ def embed(text: str) -> list[float]:
def build() -> dict:
entries = []
warnings: list[str] = []
docs = store.documents_for_index()
for doc in docs:
text = store.read_document(doc["name"])
if doc["name"].endswith(".xml") and chunk_xml(text) is None:
warnings.append(f"{doc['name']}: XML invalid, indexat ca text simplu")
for i, chunk in enumerate(chunk_document(doc["name"], text)):
vec = embed(chunk)
entries.append({"source": doc["name"], "chunk": i, "text": chunk, "embedding": vec})
config.STATE_DIR.mkdir(parents=True, exist_ok=True)
config.INDEX_FILE.write_text(json.dumps(entries, ensure_ascii=False), encoding="utf-8")
return {"documents": len(docs), "chunks": len(entries)}
# Scriere atomica: consumer-ul reciteste fisierul la 30s si ar putea prinde
# un JSON pe jumatate scris daca am scrie direct peste el.
tmp = config.INDEX_FILE.with_suffix(".json.tmp")
tmp.write_text(json.dumps(entries, ensure_ascii=False), encoding="utf-8")
os.replace(tmp, config.INDEX_FILE)
out = {"documents": len(docs), "chunks": len(entries)}
if warnings:
out["warnings"] = warnings
return out
if __name__ == "__main__":
result = build()
try:
with config.exclusive():
result = build()
except config.Busy as exc:
print(f"[indexer] {exc}, ies fara sa fac nimic", file=sys.stderr)
raise SystemExit(0)
for w in result.get("warnings", []):
print(f"[indexer] ATENTIE {w}", file=sys.stderr)
print(
f"[indexer] {result['documents']} documente, {result['chunks']} chunk-uri -> {config.INDEX_FILE}",
file=sys.stderr,

View File

@@ -89,5 +89,9 @@ def sync_and_reindex(force: bool = False) -> dict:
if __name__ == "__main__":
out = sync_and_reindex(force="--force" in sys.argv[1:])
try:
with config.exclusive():
out = sync_and_reindex(force="--force" in sys.argv[1:])
except config.Busy as exc:
out = {"skipped": str(exc)}
print(json.dumps(out, ensure_ascii=False), file=sys.stderr)

View File

@@ -112,3 +112,70 @@ def test_md_foloseste_taierea_pe_paragrafe():
text = "primul paragraf\n\n" + "al doilea paragraf " * 20
assert len(indexer.chunk_document("x.md", text)) >= 1
assert indexer.chunk_document("x.md", text) == indexer.chunk_text(text)
# --------------------------------------------- concurenta intre reindexari
def test_lacatul_refuza_a_doua_reindexare():
"""Timer-ul nu trebuie sa porneasca peste o sincronizare manuala."""
import config
with config.exclusive():
with pytest.raises(config.Busy):
with config.exclusive():
pass
def test_lacatul_se_elibereaza_dupa_iesire():
import config
with config.exclusive():
pass
with config.exclusive(): # trebuie sa mearga din nou
pass
def test_lacatul_se_elibereaza_si_la_exceptie():
import config
with pytest.raises(ValueError):
with config.exclusive():
raise ValueError("ceva")
with config.exclusive():
pass
def test_indexul_se_scrie_atomic(write, monkeypatch):
"""Consumer-ul reciteste indexul la 30s; nu are voie sa prinda JSON pe jumatate."""
import config
import indexer
write("x.md", "un paragraf oarecare")
monkeypatch.setattr(indexer, "embed", lambda text: [0.1, 0.2])
vazute = []
real_replace = indexer.os.replace
def spion(src, dst):
vazute.append((str(src), str(dst)))
return real_replace(src, dst)
monkeypatch.setattr(indexer.os, "replace", spion)
indexer.build()
assert vazute and vazute[0][1] == str(config.INDEX_FILE)
assert not config.INDEX_FILE.with_suffix(".json.tmp").exists()
def test_xml_invalid_e_raportat_ca_avertisment(write, monkeypatch):
import indexer
write("stricat.xml", "<root><a>x</a></root><in-plus/>")
monkeypatch.setattr(indexer, "embed", lambda text: [0.0])
out = indexer.build()
assert any("stricat.xml" in w for w in out.get("warnings", []))
def test_xml_valid_nu_produce_avertismente(write, monkeypatch):
import indexer
write("bun.xml", "<root><a>x</a></root>")
monkeypatch.setattr(indexer, "embed", lambda text: [0.0])
assert "warnings" not in indexer.build()