289 lines
10 KiB
Python
289 lines
10 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_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)],
|
|
)
|