Files
ThothII/harness/nsp/memory.py
T
marcopan 9fa1e3e498 feat(harness): nsp memory save-one core — targeted upsert via writer key (D11)
The D11 deviation is a single-row pgvector upsert, not a full vectorstore resync.
Ports memory.py + session/{store,artifacts} + textutil (deps of memory), renamed
psdwp3->nsp. Decision import paths rewired from nsp.session.decisions to nsp.decisions
(our A4 port lives at the top level). session/models.py left UNCHANGED to preserve
the Phase-A ThothII additions (D12/D15 author/summary, D14a grounded_values,
D14b concept_formulas).

New in memory.py:
- memory_vector_record_for_decision(records, decision_seq): the single VectorRecord
  for a chosen decision (reuses memory_vector_records, filtered to one).
- save_one_memory(records, decision_seq, writer, embedder): embeds one record and
  calls writer.upsert_records('memory', [row]) -- NEVER writer.sync (that's the
  full-resync, server-side-only path). Returns the upsert count.

L1: test_memory_save_one (5 tests) pins the contract -- single row, one upsert
call, sync never called, None/0 for unknown seq.

Deferred: the full nsp memory save-one CLI command (config/session loading + the
workstation write-guard) lands when memory_cmd.py is ported alongside the other
CLI commands. The pure D11 core is what L1 can honestly cover here.
2026-06-26 23:00:06 +02:00

258 lines
8.9 KiB
Python

import os
import re
import tempfile
from datetime import UTC, datetime
from pathlib import Path
from pydantic import BaseModel
from nsp.decisions import DecisionRecord, DecisionType, list_decisions
from nsp.session.models import SessionManifest
from nsp.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
`nsp 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]:
decisions = list_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,
},
)
)
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).
`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 nsp.vectorstore.rest_writer import pack_metadata
from nsp.vectorstore.store import content_hash
record = memory_vector_record_for_decision(records, decision_seq)
if record is None:
return 0
embedding = embedder.embed_documents([record.content])[0]
row = {
"record_key": record.id,
"kind": record.kind,
"content_hash": content_hash(record.content),
"metadata": pack_metadata(record),
"embedding": embedding,
}
return writer.upsert_records("memory", [row])