Files
ThothII/harness/tht/vectorstore/rest_writer.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

85 lines
3.0 KiB
Python

"""Scrittura controllata del pgvector via REST.
Usata dalle postazioni remote solo quando e' configurata una seconda API key di scrittura.
Mantiene l'upsert incrementale del VectorStore diretto, ma non esegue delete/clear: le
operazioni distruttive restano solo-server via connessione Postgres diretta.
"""
from tht.vectorstore.records import VectorRecord
from tht.vectorstore.rest_client import VectorRestClient
from tht.vectorstore.store import SyncStats, content_hash
TABLE_TO_KINDS = {
"schema_records": {"schema_table", "schema_column"},
"evidence": {"evidence"},
"memory": {"memory"},
}
def pack_metadata(record: VectorRecord) -> dict:
"""Impacchetta nel metadata tutta la semantica letta poi da `search_similar`."""
return {
"kind": record.kind,
"ref": record.ref,
"record_key": record.id,
"title": record.title,
"content": record.content,
**record.metadata,
}
class RestVectorWriter:
"""Writer table-scoped via RPC REST allowlist.
Il metodo `sync` e' volutamente upsert-only: aggiorna/aggiunge record, conta gli stale,
ma non li elimina. Per cleanup completo usare i comandi server-side con `vector_db`.
"""
def __init__(self, client: VectorRestClient, table: str):
if table not in TABLE_TO_KINDS:
raise ValueError(f"Tabella vector non supportata per scrittura REST: {table}")
self.client = client
self.table = table
def existing_hashes(self, kinds: set[str]) -> dict[str, str]:
allowed = TABLE_TO_KINDS[self.table]
bad = kinds - allowed
if bad:
raise ValueError(
f"Kind non ammessi per vectors.{self.table}: {', '.join(sorted(bad))}"
)
return self.client.existing_hashes(self.table, sorted(kinds))
def sync(self, records: list[VectorRecord], embedder, kinds: set[str]) -> SyncStats:
stats = SyncStats()
existing = self.existing_hashes(kinds)
to_embed: list[VectorRecord] = []
for record in records:
h = content_hash(record.content)
if record.id not in existing:
to_embed.append(record)
stats.added += 1
elif existing[record.id] != h:
to_embed.append(record)
stats.updated += 1
else:
stats.unchanged += 1
stats.deleted = 0
vectors = embedder.embed_documents([r.content for r in to_embed]) if to_embed else []
rows = [
{
"record_key": record.id,
"kind": record.kind,
"content_hash": content_hash(record.content),
"metadata": pack_metadata(record),
"embedding": vector,
}
for record, vector in zip(to_embed, vectors)
]
if rows:
self.client.upsert_records(self.table, rows)
# Gli stale non vengono cancellati in REST writer: restano responsabilita' server-side.
return stats