Files

331 lines
12 KiB
Python

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:<n>"): 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_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 i concetti chiariti sono riusabili tra domande. Le scelte sulle tabelle
# dipendono dallo schema-linking della singola domanda e non sono memory.
REUSABLE_TYPES = frozenset({"concept_clarified"})
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 d.type not in REUSABLE_TYPES:
continue
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=[], 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 promote_snapshot(snapshot, *, seqs: list[int] | None, registry_path: Path) -> list[MemoryRecord]:
existing = load_registry(registry_path)
promoted = _compute_promotions(snapshot, snapshot.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))
reusable = [
c for c in cand if c.type in REUSABLE_TYPES and c.decision_seq not in declined
]
return _dedupe_reusable_promotions(reusable)
def reusable_promotions_snapshot(snapshot, registry_path: Path) -> list[MemoryRecord]:
from tht.phase import effective_decisions
cand = _compute_promotions(snapshot, snapshot.manifest, seqs=None, existing=load_registry(registry_path))
declined = declined_promotion_seqs(effective_decisions(snapshot))
reusable = [c for c in cand if c.type in REUSABLE_TYPES and c.decision_seq not in declined]
return _dedupe_reusable_promotions(reusable)
def _dedupe_reusable_promotions(records: list[MemoryRecord]) -> list[MemoryRecord]:
"""Keep the first proposal for identical reviewer-visible memory content."""
seen: set[tuple[str, str, str, str, str]] = set()
out: list[MemoryRecord] = []
for record in records:
key = tuple(
" ".join(value.split()).casefold()
for value in (
record.type,
record.subject,
record.detail,
record.rationale,
record.question_context,
)
)
if key in seen:
continue
seen.add(key)
out.append(record)
return out
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 preview_promotions_snapshot(snapshot, registry_path: Path) -> list[MemoryRecord]:
return reusable_promotions_snapshot(snapshot, registry_path)[:MAX_PROMOTION_CANDIDATES]
def memory_vector_records(records: list[MemoryRecord]) -> list[VectorRecord]:
out: list[VectorRecord] = []
for r in records:
if r.type not in REUSABLE_TYPES:
continue
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.
"""
vectors = memory_vector_records(
[r for r in records if r.decision_seq == decision_seq]
)
if not vectors:
return None
return vectors[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)],
)