Files
ThothII/harness/tht/memory.py
T
marcopanandClaude Opus 4.8 c4d130828f fix(harness): remediation difetti review — gate↔CLI, D15, D7/D6, D14, robustezza
Implementazione del piano di remediation progressiva sui difetti emersi
dall'analisi dell'harness. Tutto verificato: 214 test Python (incl. L0 su
Postgres reale), 14 test JS del gate, ruff pulito.

Blocco 1 (CRITICA, integrazione gate↔CLI):
- phase advance: gate usa --auto + exit 6; reviewer_confirm kind:phase fa
  advance esplicito che applica i prerequisiti (prima non avanzava per le
  fasi a conferma umana).
- cte plan riceve i --name dal gate (param names); set-question con id
  posizionale; skill `tht search find`; nuovo comando `tht memory save-one`
  con dedup hash client-side in save_one_memory.

Blocco 2 (D15, stato post-rollback):
- campo `phase` su DecisionRecord + effective_decisions phase-aware per i
  subject "a nome" (cte_approved ecc.); _compute_promotions e finalize sulla
  vista effective; finalize confronta col piano CTE effettivo, non glob;
  `decision add --retracts` + comando `decision retract`.

Blocco 3 (D7 read-only + D6 manifest):
- assert_read_only su tutti e quattro i codepath (direct + REST);
- manifest author/summary/updated_at/updated_by/schema_version popolati +
  helper touch_manifest sulle mutazioni.

Blocco 4-5 (D14a/D14b):
- decision_min_phase data-driven via `emits:` in workflow.yaml;
- formula evidence: status auto, search_formulas, gruppo CLI `tht formula`,
  `search find --kind formula`, load_evidence_dir salta i .sql.md.

Blocco 6 (robustezza):
- taskdoc slice promoted_tables + bound enforced; report escaping/bound +
  rsplit note; filtro kind reader REST/direct; conteggio upserted robusto;
  guard REST run_query non-list; LSH disallineato -> LshIndexError.

Blocco 7 (pulizia):
- dead code gate e KIND_TO_TABLE morto rimossi; doc Postgres-only
  (README + connection.py).

Blocco 0 (parziale): test di compatibilità firma gate↔CLI
(tests/integration). Rinviati: fake-Pi runtime completo, artifact-gate da
disco (#23), parità eligibility REST/direct (#28), unificazione
reserved-labels (#30), memory_rejected da deselezione (#33).

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-27 17:16:51 +02:00

271 lines
9.5 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
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, SENZA
cap: il chiamante applica MAX_PROMOTION_CANDIDATES e segnala il troncamento."""
cand = _compute_promotions(
session_dir, manifest, seqs=None, existing=load_registry(registry_path)
)
return [c for c in cand if c.type in REUSABLE_TYPES]
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, *, writer, 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.
`writer` is a VectorRestClient (writer key); `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.vectorstore.rest_writer import pack_metadata
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 = writer.existing_hashes("memory", ["memory"])
if existing.get(record.id) == new_hash:
return 0 # unchanged: skip embedding + upsert
embedding = embedder.embed_documents([record.content])[0]
row = {
"record_key": record.id,
"kind": record.kind,
"content_hash": new_hash,
"metadata": pack_metadata(record),
"embedding": embedding,
}
return writer.upsert_records("memory", [row])