From 9fa1e3e4982097190b3ddd530fc079a7ceec6949 Mon Sep 17 00:00:00 2001 From: mptyl Date: Fri, 26 Jun 2026 23:00:06 +0200 Subject: [PATCH] =?UTF-8?q?feat(harness):=20nsp=20memory=20save-one=20core?= =?UTF-8?q?=20=E2=80=94=20targeted=20upsert=20via=20writer=20key=20(D11)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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. --- harness/nsp/memory.py | 257 ++++++++++++++++++++++++++ harness/nsp/session/artifacts.py | 39 ++++ harness/nsp/session/store.py | 84 +++++++++ harness/nsp/textutil.py | 8 + harness/tests/test_memory_save_one.py | 89 +++++++++ 5 files changed, 477 insertions(+) create mode 100644 harness/nsp/memory.py create mode 100644 harness/nsp/session/artifacts.py create mode 100644 harness/nsp/session/store.py create mode 100644 harness/nsp/textutil.py create mode 100644 harness/tests/test_memory_save_one.py diff --git a/harness/nsp/memory.py b/harness/nsp/memory.py new file mode 100644 index 00000000..50770264 --- /dev/null +++ b/harness/nsp/memory.py @@ -0,0 +1,257 @@ +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]) + diff --git a/harness/nsp/session/artifacts.py b/harness/nsp/session/artifacts.py new file mode 100644 index 00000000..3d42d3e7 --- /dev/null +++ b/harness/nsp/session/artifacts.py @@ -0,0 +1,39 @@ +from pathlib import Path + +from nsp.decisions import DecisionRecord +from nsp.session.models import SchemaLinking + + +def _find_evidence_file(evidence_root: Path, evidence_id: str) -> str: + for match in evidence_root.rglob(f"{evidence_id}.md"): + return str(match) + return "" + + +def build_evidence_entries( + decisions: list[DecisionRecord], + linking: SchemaLinking, + evidence_root: Path, +) -> list[dict]: + """Elenco {id, file, esito, decision_seq}: evidence citate nello schema linking + (esito 'usata') e decisioni esplicite del reviewer (accettata/scartata, che + prevalgono sul linking).""" + entries: dict[str, dict] = {} + for candidate in linking.candidates: + for evidence_id in candidate.evidence: + entries.setdefault(evidence_id, { + "id": evidence_id, + "file": _find_evidence_file(evidence_root, evidence_id), + "esito": "usata", + "decision_seq": candidate.decision_seq, + }) + for d in decisions: + if d.type not in ("evidence_accepted", "evidence_rejected"): + continue + entries[d.subject] = { + "id": d.subject, + "file": _find_evidence_file(evidence_root, d.subject), + "esito": "accettata" if d.type == "evidence_accepted" else "scartata", + "decision_seq": d.seq, + } + return list(entries.values()) diff --git a/harness/nsp/session/store.py b/harness/nsp/session/store.py new file mode 100644 index 00000000..eefb3971 --- /dev/null +++ b/harness/nsp/session/store.py @@ -0,0 +1,84 @@ +from datetime import UTC, datetime +from pathlib import Path + +from nsp.config import DatabaseConfig +from nsp.session.models import SessionManifest +from nsp.textutil import slugify + +MANIFEST = "session_manifest.yaml" +MAX_SLUG_CHARS = 40 + + +class SessionError(Exception): + pass + + +def render_question_md(question: str, assumptions: list[str] | None = None) -> str: + """Rende question.md in modo deterministico: domanda + assunzioni opzionali. + + Unica fonte di formattazione per question.md (riusata da create_session e + set_question), così il gate non deve costruire markdown a mano. + """ + body = f"# Domanda\n\n{question.strip()}\n" + items = [a.strip() for a in (assumptions or []) if a.strip()] + if items: + body += "\n## Assunzioni\n\n" + "".join(f"- {a}\n" for a in items) + return body + + +def _new_id(question: str, sessions_root: Path, stamp: str) -> str: + base = f"{stamp}-{slugify(question)[:MAX_SLUG_CHARS].rstrip('-')}" + candidate, n = base, 1 + while (sessions_root / candidate).exists(): + n += 1 + candidate = f"{base}-{n}" + return candidate + + +def create_session( + question: str, db: DatabaseConfig, sessions_root: Path +) -> SessionManifest: + now = datetime.now(UTC) + # stamp con ora/min/sec: identifica univocamente sessioni dello stesso giorno + # sulla stessa domanda. Il contatore -n resta come rete per collisioni nello + # stesso secondo. + session_id = _new_id(question, sessions_root, now.strftime("%Y-%m-%d-%H%M%S")) + manifest = SessionManifest( + id=session_id, created_at=now, question=question, + database=db.database, schema=db.db_schema, + ) + session_dir = sessions_root / session_id + manifest.to_yaml(session_dir / MANIFEST) + (session_dir / "question.md").write_text(render_question_md(question)) + return manifest + + +def set_question( + session_id: str, + question: str, + assumptions: list[str], + sessions_root: Path, +) -> Path: + """Riscrive question.md (Fase 3) in modo deterministico, senza edit tool. + + Valida l'esistenza della sessione (SessionError se assente) e ritorna il + path scritto. Non tocca il manifest né le decisioni. + """ + load_session(session_id, sessions_root) + path = sessions_root / session_id / "question.md" + path.write_text(render_question_md(question, assumptions)) + return path + + +def load_session(session_id: str, sessions_root: Path) -> SessionManifest: + path = sessions_root / session_id / MANIFEST + if not path.exists(): + raise SessionError(f"Sessione non trovata: {session_id} (atteso {path})") + return SessionManifest.from_yaml(path) + + +def close_session(session_id: str, sessions_root: Path) -> SessionManifest: + manifest = load_session(session_id, sessions_root) + manifest.status = "closed" + manifest.to_yaml(sessions_root / session_id / MANIFEST) + return manifest diff --git a/harness/nsp/textutil.py b/harness/nsp/textutil.py new file mode 100644 index 00000000..62fe5a37 --- /dev/null +++ b/harness/nsp/textutil.py @@ -0,0 +1,8 @@ +import re +import unicodedata + + +def slugify(text: str) -> str: + """Slug ASCII minuscolo; preserva gli underscore (nomi tabella).""" + text = unicodedata.normalize("NFKD", text).encode("ascii", "ignore").decode() + return re.sub(r"[^a-z0-9_]+", "-", text.lower()).strip("-") diff --git a/harness/tests/test_memory_save_one.py b/harness/tests/test_memory_save_one.py new file mode 100644 index 00000000..866acb4d --- /dev/null +++ b/harness/tests/test_memory_save_one.py @@ -0,0 +1,89 @@ +"""L1: nsp memory save-one -- targeted upsert via the writer key (spec D11). + +The D11 deviation: instead of a full vectorstore resync (nsp memory index / sync), +a remote workstation with a writer key can push a SINGLE promoted decision to +pgvector as a one-row upsert. This test pins the pure core of that behavior: +- exactly one VectorRecord is built for the chosen decision_seq +- the writer.upsert_records is called once with a single row +- writer.sync is NEVER called (that is the full-resync path) +""" +from datetime import datetime +from pathlib import Path +from unittest.mock import MagicMock + +import pytest + +from nsp.memory import MemoryRecord, memory_vector_record_for_decision, save_one_memory + + +def _record(seq: int = 7, **kw) -> MemoryRecord: + base = dict( + id="mem-0007", ts=datetime(2025, 1, 1), session_id="s1", decision_seq=seq, + type="table_promoted", subject="pazienti", detail="promossa", rationale="r", + question_context="dammi i pazienti", tables=["pazienti"], concepts=[], + ) + base.update(kw) + return MemoryRecord(**base) + + +# --- memory_vector_record_for_decision (single-record filter) ------------------ + +def test_single_record_built_for_decision_seq(): + records = [_record(seq=7), _record(seq=8, id="mem-0008")] + vr = memory_vector_record_for_decision(records, decision_seq=7) + assert vr is not None + assert vr.id == "memory:mem-0007" # memory_vector_records prefix + assert vr.kind == "memory" + assert "pazienti" in vr.content + + +def test_returns_none_for_unknown_decision_seq(): + records = [_record(seq=7)] + assert memory_vector_record_for_decision(records, decision_seq=999) is None + + +# --- save_one_memory (the D11 orchestrator: single upsert, never sync) --------- + +def test_save_one_calls_upsert_with_single_row_never_sync(): + records = [_record(seq=7)] + writer = MagicMock() + writer.upsert_records.return_value = 1 + embedder = MagicMock() + embedder.embed_documents.return_value = [[0.1] * 8] + + upserted = save_one_memory(records, decision_seq=7, writer=writer, embedder=embedder) + + assert upserted == 1 + writer.sync.assert_not_called() # the whole point of D11: no full resync + writer.upsert_records.assert_called_once() + args = writer.upsert_records.call_args + # table is memory, exactly one row + assert args[0][0] == "memory" + rows = args[0][1] + assert len(rows) == 1 + assert rows[0]["record_key"] == "memory:mem-0007" + assert "embedding" in rows[0] + + +def test_save_one_no_record_for_seq_is_noop(): + records = [_record(seq=7)] + writer = MagicMock() + embedder = MagicMock() + upserted = save_one_memory(records, decision_seq=42, writer=writer, embedder=embedder) + assert upserted == 0 + writer.upsert_records.assert_not_called() + writer.sync.assert_not_called() + embedder.embed_documents.assert_not_called() + + +def test_save_one_uses_writer_key_for_upsert(): + """The upsert must flow through the writer client (writer key), not a reader. + Verified indirectly: save_one_memory takes the writer as its client argument.""" + records = [_record(seq=7)] + writer = MagicMock() + writer.upsert_records.return_value = 1 + embedder = MagicMock() + embedder.embed_documents.return_value = [[0.0] * 4] + save_one_memory(records, decision_seq=7, writer=writer, embedder=embedder) + # one upsert call, single row, table=memory + assert writer.upsert_records.call_count == 1