Files
ThothII/harness/nsp/vectorstore/rest_writer.py
T
marcopan 796a39d893 feat(harness): port vectorstore dual-key + reader RPC (D11, §5.4)
Ports vectorstore/{rest_client,rest_writer,store,reader,embeddings,records},
evidence/model (leaf dep of records), and cli/_guards (require_vector_write_allowed
workstation write-guard). Renamed psdwp3->nsp, verbatim.

VectorRestClient gains an api_key property so reader/writer clients carry their
distinct keys visibly (spec D11: vector_reader / vector_writer on the same endpoint).

scripts/create_vector_reader_rpc.sql is NEW: the reader RPCs (search_similar,
list_tables) lived server-side in Supabase and were never versioned. Authored now
mirroring the writer allowlist pattern (table allowlist, security definer, revoke
from anon/authenticated, grant to vector_reader only). Writer RPC ported verbatim.

L1: test_vector_dual_key (7 tests) pins the dual-key construction + the workstation
write-guard (exit 4 without writer key).
2026-06-26 22:55:40 +02:00

91 lines
3.1 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 nsp.vectorstore.records import VectorRecord
from nsp.vectorstore.rest_client import VectorRestClient
from nsp.vectorstore.store import SyncStats, content_hash
KIND_TO_TABLE = {
"schema_table": "schema_records",
"schema_column": "schema_records",
"evidence": "evidence",
"memory": "memory",
}
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