import os import re import tempfile from datetime import UTC, datetime from pathlib import Path from pydantic import BaseModel from tht.decisions import DecisionRecord, DecisionType from tht.session.models import SessionManifest from tht.vectorstore.records import VectorRecord class MemoryRecord(BaseModel): id: str ts: datetime session_id: str decision_seq: int type: DecisionType subject: str detail: str = "" rationale: str = "" question_context: str = "" tables: list[str] = [] concepts: list[str] = [] _MEM_ID_RE = re.compile(r"\bmem-\d{4,}\b") def decided_memory_ids(decisions: list[DecisionRecord]) -> set[str]: """Id delle memorie gia' DECISE nella sessione: rifiutate (decisione `memory_rejected`, id nel subject) o applicate (citate come `mem-XXXX` nel rationale, per convenzione della checklist memorie). Servono a `tht memory search --session` per non riproporre cio' che il reviewer ha gia' deciso (fix: memorie scartate riproposte).""" out: set[str] = set() for d in decisions: if d.type == "memory_rejected" and d.subject: out.add(d.subject) out.update(_MEM_ID_RE.findall(d.rationale or "")) return out _DECLINED_SEQ_RE = re.compile(r"\bseq:(\d+)\b") def declined_promotion_seqs(decisions: list[DecisionRecord]) -> set[int]: """decision_seq dei candidati che il reviewer ha rifiutato al gate di promozione (F8, `memory_promotion_declined` con detail "seq:"): una riapertura del gate non deve riproporli. I promossi sono gia' dedupati dal registro.""" out: set[int] = set() for d in decisions: if d.type != "memory_promotion_declined": continue m = _DECLINED_SEQ_RE.search(d.detail or "") if m: out.add(int(m.group(1))) return out def load_registry(registry_path: Path) -> list[MemoryRecord]: if not registry_path.exists(): return [] return [ MemoryRecord.model_validate_json(line) for line in registry_path.read_text().splitlines() if line.strip() ] def save_registry(records: list[MemoryRecord], path: Path) -> None: """Riscrive il registro in modo atomico (tmp nella stessa dir + os.replace).""" path.parent.mkdir(parents=True, exist_ok=True) fd, tmp = tempfile.mkstemp(dir=path.parent, prefix=".registry-", suffix=".tmp") try: with os.fdopen(fd, "w") as f: for r in records: f.write(r.model_dump_json() + "\n") os.replace(tmp, path) except BaseException: if os.path.exists(tmp): os.unlink(tmp) raise class MemoryNotFound(Exception): """Sollevata quando un id memoria non esiste nel registro.""" EDITABLE_FIELDS = frozenset( {"subject", "type", "detail", "rationale", "question_context", "tables", "concepts"} ) def update_record(path: Path, mem_id: str, fields: dict) -> MemoryRecord: records = load_registry(path) for i, r in enumerate(records): if r.id == mem_id: data = r.model_dump() data.update({k: v for k, v in fields.items() if k in EDITABLE_FIELDS}) records[i] = MemoryRecord.model_validate(data) save_registry(records, path) return records[i] raise MemoryNotFound(mem_id) def delete_record(path: Path, mem_id: str) -> MemoryRecord: records = load_registry(path) for i, r in enumerate(records): if r.id == mem_id: removed = records.pop(i) save_registry(records, path) return removed raise MemoryNotFound(mem_id) def _default_tables(decision: DecisionRecord) -> list[str]: if decision.type in ("table_promoted", "table_excluded"): return [decision.subject] if decision.type in ("column_corrected", "join_modified") and "." in decision.subject: return [decision.subject.split(".")[0]] return [] def _default_concepts(decision: DecisionRecord) -> list[str]: if decision.type == "concept_clarified": return [decision.subject] return [] def question_context(decisions: list[DecisionRecord], manifest: SessionManifest) -> str: rewritten = [d for d in decisions if d.type == "question_rewritten"] return rewritten[-1].detail if rewritten else manifest.question def _next_id_num(existing: list[MemoryRecord]) -> int: nums = [int(r.id.split("-", 1)[1]) for r in existing if r.id.startswith("mem-")] return (max(nums) + 1) if nums else 1 # Solo questi tipi di decisione sono concetti riusabili in altre generazioni # (scelta reviewer): il resto e' query-specifico e non va proposto in promozione. REUSABLE_TYPES = frozenset({"concept_clarified", "table_promoted", "table_excluded"}) MAX_PROMOTION_CANDIDATES = 5 def _compute_promotions( session_dir: Path, manifest: SessionManifest, *, seqs: list[int] | None, existing: list[MemoryRecord], ) -> list[MemoryRecord]: # Vista effective (D15, §4.8): non promuovere decisioni stale dopo un rollback. from tht.phase import effective_decisions decisions = effective_decisions(session_dir) selected = decisions if seqs is None else [d for d in decisions if d.seq in seqs] already = {(r.session_id, r.decision_seq) for r in existing} n = _next_id_num(existing) context = question_context(decisions, manifest) out: list[MemoryRecord] = [] for d in selected: if (manifest.id, d.seq) in already: continue out.append( MemoryRecord( id=f"mem-{n:04d}", ts=datetime.now(UTC), session_id=manifest.id, decision_seq=d.seq, type=d.type, subject=d.subject, detail=d.detail, rationale=d.rationale, question_context=context, tables=_default_tables(d), concepts=_default_concepts(d), ) ) n += 1 return out def promote( session_dir: Path, manifest: SessionManifest, *, seqs: list[int] | None, registry_path: Path, ) -> list[MemoryRecord]: """Copia le decisioni indicate (tutte se seqs=None) nel registro globale. Salta quelle gia' promosse (chiave: session_id + decision_seq).""" existing = load_registry(registry_path) promoted = _compute_promotions(session_dir, manifest, seqs=seqs, existing=existing) if promoted: save_registry(existing + promoted, registry_path) return promoted def reusable_promotions( session_dir: Path, manifest: SessionManifest, registry_path: Path ) -> list[MemoryRecord]: """Candidati riusabili (tipi in REUSABLE_TYPES) non ancora promossi ne' rifiutati al gate, SENZA cap: il chiamante applica MAX_PROMOTION_CANDIDATES e segnala il troncamento.""" from tht.phase import effective_decisions cand = _compute_promotions( session_dir, manifest, seqs=None, existing=load_registry(registry_path) ) declined = declined_promotion_seqs(effective_decisions(session_dir)) return [ c for c in cand if c.type in REUSABLE_TYPES and c.decision_seq not in declined ] def preview_promotions( session_dir: Path, manifest: SessionManifest, registry_path: Path ) -> list[MemoryRecord]: """Candidati promuovibili: riusabili e troncati a MAX_PROMOTION_CANDIDATES.""" return reusable_promotions(session_dir, manifest, registry_path)[ :MAX_PROMOTION_CANDIDATES ] def memory_vector_records(records: list[MemoryRecord]) -> list[VectorRecord]: out: list[VectorRecord] = [] for r in records: lines = [ f"Decisione {r.type}: {r.subject}", r.detail, r.rationale, f"Domanda di contesto: {r.question_context}", ] if r.tables: lines.append("Tabelle: " + ", ".join(r.tables)) if r.concepts: lines.append("Concetti: " + ", ".join(r.concepts)) out.append( VectorRecord( id=f"memory:{r.id}", kind="memory", ref=r.id, title=f"{r.type}: {r.subject}", content="\n".join(filter(None, lines)), metadata={ "type": r.type, "session_id": r.session_id, "tables": r.tables, "concepts": r.concepts, "subject": r.subject, "detail": r.detail, "rationale": r.rationale, }, ) ) return out def memory_vector_record_for_decision( records: list[MemoryRecord], decision_seq: int ) -> VectorRecord | None: """The single VectorRecord for `decision_seq`, or None if no memory matches. D11 save-one builds only this one record (not the full memory_vector_records list) so the remote upsert is a single row. """ match = [r for r in records if r.decision_seq == decision_seq] if not match: return None return memory_vector_records(match)[0] def save_one_memory( records: list[MemoryRecord], decision_seq: int, *, store, embedder ) -> int: """Targeted one-row upsert of a promoted decision to pgvector via the writer key (spec D11). This is NOT a full vectorstore resync: it embeds and pushes a single record, so a workstation with a writer key can publish one memory without rebuilding the index. Returns the upsert count (0 if no record matched or the record is already up to date). Hash dedup client-side (spec §5.4): the SHA-256 of the content is compared with the writer's existing_vector_hashes; the embedding (Ollama round-trip) and the upsert are skipped when the content is unchanged. Idempotent by construction. `store` is the configured writable VectorStore; `embedder` an embeddings client. The destructive cleanup (sync's delete-stale step) is intentionally absent: it remains a server-side-only operation via the direct vectordb connection. """ from tht.ports.vector import VectorWriteRecord from tht.vectorstore.store import content_hash record = memory_vector_record_for_decision(records, decision_seq) if record is None: return 0 new_hash = content_hash(record.content) existing = store.existing_hashes("memory", ["memory"]) if existing.get(record.id) == new_hash: return 0 # unchanged: skip embedding + upsert embedding = embedder.embed_documents([record.content])[0] return store.upsert( "memory", [VectorWriteRecord(record=record, embedding=embedding, content_hash=new_hash)], )