refactor: remove pgvector runtime

This commit is contained in:
2026-08-08 21:38:03 +02:00
parent 5c12d9bb79
commit 8826f8ac6b
34 changed files with 200 additions and 3362 deletions
+2 -61
View File
@@ -3,12 +3,10 @@
from tht.adapters.dwh import PostgresDwhAdapter, ThothRestDwhAdapter
from tht.adapters.evidence import FilesystemEvidenceSource, HttpManifestEvidenceSource
from tht.adapters.evidence.s3 import S3EvidenceSource
from tht.adapters.vector import PgVectorStore, QdrantVectorStore, ThothHttpVectorStore
from tht.adapters.vector import QdrantVectorStore
from tht.config import Config, ConfigError
from tht.db.connection import make_engine
from tht.ports.dwh import DwhAdapter
from tht.ports.vector import VectorStore
from tht.vectorstore.rest_client import VectorRestClient
def build_dwh(cfg: Config) -> DwhAdapter:
@@ -33,29 +31,6 @@ def build_vector_store(cfg: Config, *, require_write: bool = False) -> VectorSto
raise ConfigError("Risorsa vectors non configurata")
match resource.type:
case "pgvector_direct":
reader = resource.reader or resource.connection
# Legacy server workspaces use one RW `vector_db` connection. Keep
# that deployment contract without turning a workstation's legacy
# compatibility connection into an implicit writer.
writer = resource.writer or (
resource.connection if cfg.profile == "server" else None
)
if require_write and writer is None:
raise ConfigError("Vector writer non configurato per pgvector_direct")
return PgVectorStore(
reader,
writer,
expected_dimension=cfg.embeddings.dim if cfg.embeddings is not None else None,
)
case "thoth_vector_http":
if require_write and resource.writer is None:
raise ConfigError("Vector writer non configurato")
return ThothHttpVectorStore(
VectorRestClient(resource.reader) if resource.reader is not None else None,
VectorRestClient(resource.writer) if resource.writer is not None else None,
expected_dimension=cfg.embeddings.dim if cfg.embeddings is not None else None,
)
case "qdrant":
return QdrantVectorStore(
base_url=resource.base_url,
@@ -68,40 +43,6 @@ def build_vector_store(cfg: Config, *, require_write: bool = False) -> VectorSto
raise ConfigError(f"Adapter vector non supportato: {other}")
def build_vector_loader(cfg: Config, collection: str):
"""Compatibility construction for legacy collection sync commands."""
resource = cfg.vectors
if resource is None:
raise ConfigError("Risorsa vectors non configurata")
if cfg.embeddings is None:
raise ConfigError("Embeddings non configurati")
if (
resource.type == "thoth_vector_http"
and resource.writer is not None
and (cfg.profile == "workstation" or resource.direct is None)
):
from tht.vectorstore.rest_writer import RestVectorWriter
return RestVectorWriter(VectorRestClient(resource.writer), table=collection)
connection = (
resource.writer or resource.connection
if resource.type == "pgvector_direct"
else resource.direct
)
if connection is None:
raise ConfigError("Vector writer non configurato")
from tht.vectorstore.store import VectorStore as TableVectorStore
return TableVectorStore(
make_engine(connection),
schema=connection.db_schema,
table=collection,
dim=cfg.embeddings.dim,
)
def build_evidence_sources(cfg: Config):
"""Build configured Evidence sources, including the legacy curated filesystem tree."""
evidence = cfg.evidence
@@ -152,4 +93,4 @@ def build_evidence_sources(cfg: Config):
return sources
__all__ = ["build_dwh", "build_evidence_sources", "build_vector_loader", "build_vector_store"]
__all__ = ["build_dwh", "build_evidence_sources", "build_vector_store"]
+1 -4
View File
@@ -1,8 +1,5 @@
"""Vector-store adapter implementations."""
from tht.adapters.vector.legacy_direct import LegacyDirectVectorStore
from tht.adapters.vector.pgvector import PgVectorStore
from tht.adapters.vector.qdrant import QdrantVectorStore
from tht.adapters.vector.thoth_http import ThothHttpVectorStore
__all__ = ["LegacyDirectVectorStore", "PgVectorStore", "QdrantVectorStore", "ThothHttpVectorStore"]
__all__ = ["QdrantVectorStore"]
+41
View File
@@ -0,0 +1,41 @@
"""Shared collection and kind validation for vector stores."""
from __future__ import annotations
from tht.ports.vector import VectorStoreError
COLLECTION_KINDS = {
"schema_records": {"schema_table", "schema_column"},
"evidence": {"evidence"},
"memory": {"memory", "solved_question"},
}
ALLOWED_COLLECTIONS = frozenset(COLLECTION_KINDS)
ALLOWED_KINDS = frozenset().union(*COLLECTION_KINDS.values())
def validate_collection(collection: str) -> str:
if collection not in ALLOWED_COLLECTIONS:
raise VectorStoreError(f"Collection not allowed: {collection}")
return collection
def validate_collection_kinds(collection: str, kinds: list[str]) -> None:
invalid = set(kinds) - COLLECTION_KINDS[collection]
if invalid:
raise VectorStoreError(f"Kind not allowed for {collection}: {', '.join(sorted(invalid))}")
def validate_known_kinds(kinds: list[str]) -> None:
invalid = set(kinds) - ALLOWED_KINDS
if invalid:
raise VectorStoreError(f"Kind not allowed: {', '.join(sorted(invalid))}")
__all__ = [
"ALLOWED_COLLECTIONS",
"ALLOWED_KINDS",
"COLLECTION_KINDS",
"validate_collection",
"validate_collection_kinds",
"validate_known_kinds",
]
@@ -1,70 +0,0 @@
"""Compatibility adapter for the existing direct PostgreSQL vector reader."""
from sqlalchemy import Engine
from tht.ports.vector import (
VectorCapabilities,
VectorHealth,
VectorStoreError,
VectorWriteRecord,
VectorWriteUnavailable,
require_positive_limit,
)
from tht.vectorstore.store import VectorHit, VectorStore as TableVectorStore
class LegacyDirectVectorStore:
"""Read-only port wrapper around the legacy table-scoped pgvector store."""
capabilities = VectorCapabilities(search=True, existing_hashes=False, upsert=False)
def __init__(self, engine: Engine, schema: str = "vectors", dim: int = 768):
self._engine = engine
self._schema = schema
self._dim = dim
def health(self) -> VectorHealth:
try:
with self._engine.connect() as connection:
connection.exec_driver_sql("SELECT 1")
except Exception as exc:
return VectorHealth(
ok=False,
detail=str(exc),
read_configured=True,
read_reachable=False,
read_detail=str(exc),
expected_dimension=self._dim,
)
return VectorHealth(
ok=True,
read_configured=True,
read_reachable=True,
expected_dimension=self._dim,
)
def search(
self,
collections: list[str],
embedding: list[float],
*,
limit: int,
kinds: list[str] | None = None,
metadata_filter: dict[str, object] | None = None,
) -> list[VectorHit]:
require_positive_limit(limit)
if metadata_filter is not None:
raise VectorStoreError("Legacy vector store cannot enforce metadata filtering")
hits: list[VectorHit] = []
for collection in collections:
table = TableVectorStore(
self._engine, schema=self._schema, table=collection, dim=self._dim
)
hits.extend(table.search(embedding, top_n=limit, kinds=kinds))
return sorted(hits, key=lambda hit: hit.similarity, reverse=True)[:limit]
def existing_hashes(self, collection: str, kinds: list[str]) -> dict[str, str]:
raise VectorWriteUnavailable("Legacy direct reader has no writer interface")
def upsert(self, collection: str, records: list[VectorWriteRecord]) -> int:
raise VectorWriteUnavailable("Legacy direct reader has no writer interface")
-524
View File
@@ -1,524 +0,0 @@
"""Direct PostgreSQL/pgvector implementation of the vector port."""
import json
import re
from psycopg2 import Error as PsycopgError
from psycopg2 import sql
from sqlalchemy import Engine
from sqlalchemy.exc import SQLAlchemyError
from tht.config import DatabaseConfig
from tht.db.connection import make_engine
from tht.ports.vector import (
VectorCapabilities,
VectorHealth,
VectorReadUnavailable,
VectorStoreError,
VectorWriteRecord,
VectorWriteUnavailable,
require_positive_limit,
)
from tht.vectorstore.store import VectorHit, hit_from_metadata
COLLECTION_KINDS = {
"schema_records": {"schema_table", "schema_column"},
"evidence": {"evidence"},
"memory": {"memory", "solved_question"},
}
ALLOWED_COLLECTIONS = frozenset(COLLECTION_KINDS)
ALLOWED_KINDS = frozenset().union(*COLLECTION_KINDS.values())
_VECTOR_DIMENSION = re.compile(r"^(?:[a-z_][a-z0-9_]*\.)?vector\((\d+)\)$")
def _collection(schema: str, name: str) -> sql.Identifier:
if name not in ALLOWED_COLLECTIONS:
raise VectorStoreError(f"Collection not allowed: {name}")
return sql.Identifier(schema, name)
def _vector_literal(values: list[float]) -> str:
return "[" + ",".join(str(float(value)) for value in values) + "]"
def _vector_type(schema: str) -> sql.Identifier:
return sql.Identifier(schema, "vector")
def _cosine_operator(schema: str) -> sql.Composed:
return sql.SQL("OPERATOR({}.<=>)").format(sql.Identifier(schema))
def _vector_sql_names(cursor, table_schema: str, collection: str) -> tuple[str, str]:
"""Discover pgvector type and operator namespaces from the embedding column."""
cursor.execute(
"""SELECT type_ns.nspname, operator_ns.nspname
FROM pg_catalog.pg_attribute attribute
JOIN pg_catalog.pg_class table_class
ON table_class.oid = attribute.attrelid
JOIN pg_catalog.pg_namespace table_ns
ON table_ns.oid = table_class.relnamespace
JOIN pg_catalog.pg_type vector_type
ON vector_type.oid = attribute.atttypid
JOIN pg_catalog.pg_namespace type_ns
ON type_ns.oid = vector_type.typnamespace
JOIN pg_catalog.pg_operator cosine
ON cosine.oprname = %s
AND cosine.oprleft = vector_type.oid
AND cosine.oprright = vector_type.oid
JOIN pg_catalog.pg_namespace operator_ns
ON operator_ns.oid = cosine.oprnamespace
WHERE table_ns.nspname = %s
AND table_class.relname = %s
AND attribute.attname = %s
AND NOT attribute.attisdropped
ORDER BY cosine.oid
LIMIT 1""",
("<=>", table_schema, collection, "embedding"),
)
row = cursor.fetchone()
if row is None:
raise VectorStoreError(f"Collection {collection} has no usable pgvector embedding")
return row[0], row[1]
def _validate_collection_kinds(collection: str, kinds: list[str]) -> None:
invalid = set(kinds) - COLLECTION_KINDS[collection]
if invalid:
raise VectorStoreError(f"Kind not allowed for {collection}: {', '.join(sorted(invalid))}")
def _validate_known_kinds(kinds: list[str]) -> None:
invalid = set(kinds) - ALLOWED_KINDS
if invalid:
raise VectorStoreError(f"Kind not allowed: {', '.join(sorted(invalid))}")
class PgVectorStore:
"""Direct store with independent reader and writer database credentials."""
def __init__(
self,
read_config: DatabaseConfig | None,
write_config: DatabaseConfig | None = None,
*,
expected_dimension: int | None = None,
):
self._reader = make_engine(read_config) if read_config is not None else None
self._writer = make_engine(write_config) if write_config is not None else None
config = read_config or write_config
self._schema = config.db_schema if config is not None else "vectors"
if read_config and write_config and read_config.db_schema != write_config.db_schema:
raise VectorStoreError("Reader and writer vector schemas must match")
self._expected_dimension = expected_dimension
@property
def capabilities(self) -> VectorCapabilities:
writable = self._writer is not None
return VectorCapabilities(
search=self._reader is not None,
existing_hashes=writable,
upsert=writable,
metadata_filter=self._reader is not None,
delete_generation=writable,
list_evidence_generations=writable,
)
def _probe(
self, engine: Engine | None, *, writable: bool
) -> tuple[bool | None, str | None, set[int]]:
if engine is None:
return None, None, set()
try:
raw = engine.raw_connection()
try:
with raw.cursor() as cursor:
cursor.execute("SELECT 1")
cursor.execute(
"SELECT has_schema_privilege(current_user, %s, 'USAGE')",
(self._schema,),
)
schema_usage = bool(cursor.fetchone()[0])
if not schema_usage:
return False, "vector schema incomplete: missing schema usage", set()
cursor.execute(
"""SELECT c.relname, format_type(a.atttypid, a.atttypmod),
has_table_privilege(current_user, c.oid, 'SELECT'),
has_table_privilege(current_user, c.oid, 'INSERT'),
has_table_privilege(current_user, c.oid, 'UPDATE'),
has_column_privilege(current_user, c.oid, 'record_key', 'SELECT')
AND has_column_privilege(
current_user, c.oid, 'content_hash', 'SELECT'
)
AND has_column_privilege(current_user, c.oid, 'kind', 'SELECT'),
CASE WHEN id_attr.attname IS NOT NULL THEN
pg_get_serial_sequence(
format('%%I.%%I', n.nspname, c.relname), 'id'
)
END AS id_sequence,
CASE WHEN id_attr.attname IS NOT NULL THEN
has_sequence_privilege(
current_user,
pg_get_serial_sequence(
format('%%I.%%I', n.nspname, c.relname), 'id'
),
'USAGE'
)
END AS sequence_usage
FROM pg_class c
JOIN pg_namespace n ON n.oid = c.relnamespace
LEFT JOIN pg_attribute a ON a.attrelid = c.oid
AND a.attname = 'embedding' AND NOT a.attisdropped
LEFT JOIN pg_attribute id_attr ON id_attr.attrelid = c.oid
AND id_attr.attname = 'id' AND NOT id_attr.attisdropped
WHERE n.nspname = %s AND c.relname = ANY(%s)
AND c.relkind IN ('r', 'p')""",
(self._schema, list(ALLOWED_COLLECTIONS)),
)
rows = cursor.fetchall()
present = {row[0] for row in rows}
missing_tables = sorted(ALLOWED_COLLECTIONS - present)
missing_embeddings = sorted(row[0] for row in rows if row[1] is None)
privilege_missing = sorted(
row[0]
for row in rows
if (writable and not (row[3] and row[4] and row[5]))
or (not writable and not row[2])
)
missing_sequences = sorted(
row[0] for row in rows if writable and row[6] is None
)
sequence_privilege_missing = sorted(
row[0] for row in rows if writable and row[6] is not None and not row[7]
)
problems = []
if missing_tables:
problems.append("missing tables " + ", ".join(missing_tables))
if missing_embeddings:
problems.append(
"missing embedding columns " + ", ".join(missing_embeddings)
)
if privilege_missing:
authority = "write" if writable else "read"
problems.append(
f"missing {authority} privileges " + ", ".join(privilege_missing)
)
if missing_sequences:
problems.append("missing id sequences " + ", ".join(missing_sequences))
if sequence_privilege_missing:
problems.append(
"missing sequence privileges " + ", ".join(sequence_privilege_missing)
)
if problems:
return False, "vector schema incomplete: " + "; ".join(problems), set()
dimensions = {
int(match.group(1))
for _, type_name, *_ in rows
if (match := _VECTOR_DIMENSION.match(type_name))
}
invalid_types = sorted(
row[0]
for row in rows
if row[1] is not None and not _VECTOR_DIMENSION.match(row[1])
)
if invalid_types:
return (
False,
"vector schema incomplete: invalid embedding types "
+ ", ".join(invalid_types),
set(),
)
if self._expected_dimension is not None:
mismatches = sorted(
f"{name}={int(match.group(1))}"
for name, type_name, *_ in rows
if (match := _VECTOR_DIMENSION.match(type_name))
and int(match.group(1)) != self._expected_dimension
)
if mismatches:
return (
False,
"embedding dimension mismatch: " + ", ".join(mismatches),
dimensions,
)
return True, None, dimensions
finally:
raw.close()
except (AttributeError, TypeError, ValueError, PsycopgError, SQLAlchemyError) as exc:
return False, f"vector database probe failed: {type(exc).__name__}", set()
def health(self) -> VectorHealth:
read_ok, read_detail, read_dimensions = self._probe(self._reader, writable=False)
write_ok, write_detail, write_dimensions = self._probe(self._writer, writable=True)
dimensions = tuple(sorted(read_dimensions | write_dimensions))
compatible = (
None
if self._expected_dimension is None or not dimensions
else dimensions == (self._expected_dimension,)
)
reachable = [value for value in (read_ok, write_ok) if value is not None]
details = [value for value in (read_detail, write_detail) if value]
return VectorHealth(
ok=bool(reachable) and all(reachable) and compatible is not False,
detail="; ".join(details) or None,
read_configured=self._reader is not None,
read_reachable=read_ok,
read_detail=read_detail,
write_configured=self._writer is not None,
write_reachable=write_ok,
write_detail=write_detail,
expected_dimension=self._expected_dimension,
observed_dimensions=dimensions,
dimension_compatible=compatible,
)
def search(
self,
collections: list[str],
embedding: list[float],
*,
limit: int,
kinds: list[str] | None = None,
metadata_filter: dict[str, object] | None = None,
) -> list[VectorHit]:
require_positive_limit(limit)
if self._reader is None:
raise VectorReadUnavailable("Vector reader credential is not configured")
if self._expected_dimension is not None and len(embedding) != self._expected_dimension:
raise VectorStoreError("Query embedding dimension does not match configured dimension")
if kinds:
_validate_known_kinds(kinds)
hits: list[VectorHit] = []
raw = None
try:
raw = self._reader.raw_connection()
with raw.cursor() as cursor:
for collection in collections:
table = _collection(self._schema, collection)
type_schema, operator_schema = _vector_sql_names(
cursor, self._schema, collection
)
collection_kinds = (
sorted(set(kinds) & COLLECTION_KINDS[collection]) if kinds else None
)
if kinds and not collection_kinds:
continue
clauses = []
filter_params = []
if collection_kinds:
clauses.append(sql.SQL("kind = ANY(%s)"))
filter_params.append(collection_kinds)
if metadata_filter is not None:
if collection != "evidence" or set(metadata_filter) != {
"vector_generation", "document_ids", "workspace_id"
}:
raise VectorStoreError("Unsupported vector metadata filter")
generation = metadata_filter["vector_generation"]
document_ids = metadata_filter["document_ids"]
workspace_id = metadata_filter["workspace_id"]
if not isinstance(generation, str) or not isinstance(document_ids, list) or not isinstance(workspace_id, str):
raise VectorStoreError("Invalid vector metadata filter")
clauses.append(sql.SQL("metadata->>'vector_generation' = %s"))
clauses.append(sql.SQL("metadata->>'document_id' = ANY(%s)"))
clauses.append(sql.SQL("metadata->>'workspace_id' = %s"))
filter_params.extend((generation, document_ids, workspace_id))
where = (
sql.SQL(" WHERE ") + sql.SQL(" AND ").join(clauses)
if clauses else sql.SQL("")
)
query = sql.SQL(
"SELECT metadata, 1 - (embedding {} %s::{}) AS similarity "
"FROM {}{} ORDER BY embedding {} %s::{}, record_key LIMIT %s"
).format(
_cosine_operator(operator_schema),
_vector_type(type_schema),
table,
where,
_cosine_operator(operator_schema),
_vector_type(type_schema),
)
params = [_vector_literal(embedding)]
params.extend(filter_params)
params.extend((_vector_literal(embedding), limit))
cursor.execute(query, params)
hits.extend(hit_from_metadata(row[1], row[0]) for row in cursor.fetchall())
except VectorStoreError:
raise
except Exception as exc:
raise VectorReadUnavailable("Vector read operation unavailable") from exc
finally:
if raw is not None:
raw.close()
return sorted(hits, key=lambda hit: (-hit.similarity, hit.id))[:limit]
def _require_writer(self) -> Engine:
if self._writer is None:
raise VectorWriteUnavailable("Vector writer credential is not configured")
return self._writer
def existing_hashes(self, collection: str, kinds: list[str]) -> dict[str, str]:
engine = self._require_writer()
table = _collection(self._schema, collection)
_validate_collection_kinds(collection, kinds)
raw = None
try:
raw = engine.raw_connection()
with raw.cursor() as cursor:
cursor.execute(
sql.SQL("SELECT record_key, content_hash FROM {} WHERE kind = ANY(%s)").format(
table
),
(kinds,),
)
return dict(cursor.fetchall())
except VectorStoreError:
raise
except Exception as exc:
raise VectorWriteUnavailable("Vector write operation unavailable") from exc
finally:
if raw is not None:
raw.close()
def upsert(self, collection: str, records: list[VectorWriteRecord]) -> int:
engine = self._require_writer()
table = _collection(self._schema, collection)
for write_record in records:
_validate_collection_kinds(collection, [write_record.record.kind])
if (
self._expected_dimension is not None
and len(write_record.embedding) != self._expected_dimension
):
raise VectorStoreError("Embedding dimension does not match configured dimension")
raw = None
try:
raw = engine.raw_connection()
with raw.cursor() as cursor:
type_schema, _ = _vector_sql_names(cursor, self._schema, collection)
insert = sql.SQL(
"INSERT INTO {} (record_key, kind, content_hash, metadata, embedding) "
"VALUES (%s, %s, %s, %s::jsonb, %s::{}) "
"ON CONFLICT (record_key) DO NOTHING"
).format(table, _vector_type(type_schema))
update = sql.SQL(
"UPDATE {} SET kind = %s, content_hash = %s, metadata = %s::jsonb, "
"embedding = %s::{}, indexed_at = pg_catalog.now() WHERE record_key = %s"
).format(table, _vector_type(type_schema))
for write_record in records:
record = write_record.record
metadata = {
"kind": record.kind,
"ref": record.ref,
"record_key": record.id,
"title": record.title,
"content": record.content,
**record.metadata,
}
metadata_json = json.dumps(metadata)
vector = _vector_literal(write_record.embedding)
cursor.execute(
insert,
(record.id, record.kind, write_record.content_hash, metadata_json, vector),
)
if cursor.rowcount == 0:
cursor.execute(
update,
(
record.kind,
write_record.content_hash,
metadata_json,
vector,
record.id,
),
)
raw.commit()
except VectorStoreError:
if raw is not None:
raw.rollback()
raise
except Exception as exc:
if raw is not None:
raw.rollback()
raise VectorWriteUnavailable("Vector write operation unavailable") from exc
finally:
if raw is not None:
raw.close()
return len(records)
def delete_generation(self, collection: str, generation: str, workspace_id: str) -> int:
if collection != "evidence" or re.fullmatch(r"gen:[0-9a-f]{32}", generation) is None:
raise VectorStoreError("Only exact Evidence generations may be deleted")
if re.fullmatch(r"[a-z][a-z0-9_-]{0,63}", workspace_id) is None:
raise VectorStoreError("Invalid Evidence workspace namespace")
raw = None
try:
raw = self._require_writer().raw_connection()
with raw.cursor() as cursor:
cursor.execute(
sql.SQL(
"DELETE FROM {} WHERE kind = 'evidence' "
"AND metadata->>'vector_generation' = %s "
"AND metadata->>'workspace_id' = %s"
).format(_collection(self._schema, collection)),
(generation, workspace_id),
)
count = cursor.rowcount
raw.commit()
return count
except Exception as exc:
if raw is not None:
raw.rollback()
raise VectorWriteUnavailable("Vector generation cleanup unavailable") from exc
finally:
if raw is not None:
raw.close()
def delete_kinds(self, collection: str, kinds: list[str]) -> int:
_collection(self._schema, collection)
_validate_collection_kinds(collection, kinds)
raw = None
try:
raw = self._require_writer().raw_connection()
with raw.cursor() as cursor:
cursor.execute(
sql.SQL("DELETE FROM {} WHERE kind = ANY(%s)").format(
_collection(self._schema, collection)
),
(kinds,),
)
count = cursor.rowcount
raw.commit()
return count
except Exception as exc:
if raw is not None:
raw.rollback()
raise VectorWriteUnavailable("Vector kind cleanup unavailable") from exc
finally:
if raw is not None:
raw.close()
def list_evidence_generations(self, collection: str, workspace_id: str) -> list[str]:
if collection != "evidence":
raise VectorStoreError("Only exact Evidence generations may be listed")
if re.fullmatch(r"[a-z][a-z0-9_-]{0,63}", workspace_id) is None:
raise VectorStoreError("Invalid Evidence workspace namespace")
raw = None
try:
raw = self._require_writer().raw_connection()
with raw.cursor() as cursor:
cursor.execute(
sql.SQL(
"SELECT DISTINCT metadata->>'vector_generation' FROM {} "
"WHERE kind = 'evidence' AND metadata->>'vector_generation' "
"~ '^gen:[0-9a-f]{{32}}$' AND metadata->>'workspace_id' = %s ORDER BY 1"
).format(_collection(self._schema, collection)),
(workspace_id,),
)
return [row[0] for row in cursor.fetchall()]
except Exception as exc:
raise VectorWriteUnavailable("Vector generation inventory unavailable") from exc
finally:
if raw is not None:
raw.close()
__all__ = ["ALLOWED_COLLECTIONS", "PgVectorStore"]
+12 -12
View File
@@ -6,11 +6,11 @@ from uuid import NAMESPACE_URL, uuid5
import requests
from tht.adapters.vector.pgvector import (
from tht.adapters.vector._shared import (
COLLECTION_KINDS,
_collection,
_validate_collection_kinds,
_validate_known_kinds,
validate_collection,
validate_collection_kinds,
validate_known_kinds,
)
from tht.ports.vector import (
VectorCapabilities,
@@ -165,8 +165,8 @@ class QdrantVectorStore:
return sorted(hits, key=lambda hit: (-hit.similarity, hit.id))[:limit]
def existing_hashes(self, collection: str, kinds: list[str]) -> dict[str, str]:
_collection("vectors", collection)
_validate_collection_kinds(collection, kinds)
validate_collection(collection)
validate_collection_kinds(collection, kinds)
points = self._scroll(
[
*self._workspace_filter(),
@@ -186,11 +186,11 @@ class QdrantVectorStore:
return hashes
def upsert(self, collection: str, records: list[VectorWriteRecord]) -> int:
_collection("vectors", collection)
validate_collection(collection)
self._ensure_collection(strict=True)
points = []
for write_record in records:
_validate_collection_kinds(collection, [write_record.record.kind])
validate_collection_kinds(collection, [write_record.record.kind])
self._validate_embedding(write_record.embedding, query=False)
semantic_kind = qdrant_semantic_kind(write_record.record.kind)
points.append(
@@ -213,8 +213,8 @@ class QdrantVectorStore:
return len(records)
def delete_kinds(self, collection: str, kinds: list[str]) -> int:
_collection("vectors", collection)
_validate_collection_kinds(collection, kinds)
validate_collection(collection)
validate_collection_kinds(collection, kinds)
must = [
*self._workspace_filter(),
{"key": "record_kind", "match": {"any": sorted(kinds)}},
@@ -284,10 +284,10 @@ class QdrantVectorStore:
) -> list[str]:
selected: set[str] = set()
for collection in collections:
_collection("vectors", collection)
validate_collection(collection)
selected.update(COLLECTION_KINDS[collection])
if kinds:
_validate_known_kinds(kinds)
validate_known_kinds(kinds)
selected &= set(kinds)
return sorted(selected)
-191
View File
@@ -1,191 +0,0 @@
"""Thoth vector HTTP adapter using distinct read and write clients."""
import re
from tht.adapters.vector.pgvector import (
_collection,
_validate_collection_kinds,
_validate_known_kinds,
)
from tht.ports.vector import (
VectorCapabilities,
VectorHealth,
VectorHit,
VectorReadUnavailable,
VectorStoreError,
VectorWriteRecord,
VectorWriteUnavailable,
require_positive_limit,
)
from tht.vectorstore.rest_client import VectorRestClient, VectorRestError
from tht.vectorstore.store import hit_from_metadata
def _merge(hits: list[VectorHit], limit: int) -> list[VectorHit]:
return sorted(hits, key=lambda hit: (-hit.similarity, hit.id))[:limit]
class ThothHttpVectorStore:
"""Vector port backed by the existing allowlisted REST RPCs."""
def __init__(
self,
reader: VectorRestClient | None,
writer: VectorRestClient | None,
expected_dimension: int | None = None,
):
self._reader = reader
self._writer = writer
self._expected_dimension = expected_dimension
@property
def capabilities(self) -> VectorCapabilities:
writable = self._writer is not None
return VectorCapabilities(
search=self._reader is not None, existing_hashes=writable, upsert=writable,
metadata_filter=self._reader is not None, delete_generation=writable,
list_evidence_generations=writable,
)
def health(self) -> VectorHealth:
read_reachable, read_detail, read_tables = self._probe(self._reader)
write_reachable, write_detail, write_tables = self._probe(self._writer)
dimensions = tuple(sorted({
dimension
for row in [*read_tables, *write_tables]
if type(dimension := row.get("vector_dimensions")) is int
}))
compatible = (
None
if self._expected_dimension is None or not dimensions
else dimensions == (self._expected_dimension,)
)
reachable = [
status for status in (read_reachable, write_reachable) if status is not None
]
ok = bool(reachable) and all(reachable) and compatible is not False
details = [detail for detail in (read_detail, write_detail) if detail]
return VectorHealth(
ok=ok,
detail="; ".join(details) or None,
read_configured=self._reader is not None,
read_reachable=read_reachable,
read_detail=read_detail,
write_configured=self._writer is not None,
write_reachable=write_reachable,
write_detail=write_detail,
expected_dimension=self._expected_dimension,
observed_dimensions=dimensions,
dimension_compatible=compatible,
)
@staticmethod
def _probe(client: VectorRestClient | None) -> tuple[bool | None, str | None, list[dict]]:
if client is None:
return None, None, []
try:
return True, None, client.list_tables()
except (RuntimeError, VectorRestError) as exc:
return False, str(exc), []
def search(
self,
collections: list[str],
embedding: list[float],
*,
limit: int,
kinds: list[str] | None = None,
metadata_filter: dict[str, object] | None = None,
) -> list[VectorHit]:
require_positive_limit(limit)
if self._reader is None:
raise VectorReadUnavailable("Vector reader credential is not configured")
if self._expected_dimension is not None and len(embedding) != self._expected_dimension:
raise VectorStoreError("Query embedding dimension does not match configured dimension")
if kinds:
_validate_known_kinds(kinds)
hits: list[VectorHit] = []
for collection in collections:
_collection("vectors", collection)
try:
if metadata_filter is None:
rows = self._reader.search_similar(collection, embedding, limit, kinds=kinds)
else:
rows = self._reader.search_similar(
collection, embedding, limit, kinds=kinds,
metadata_filter=metadata_filter,
)
except VectorRestError as exc:
raise VectorStoreError(str(exc)) from exc
hits.extend(
hit_from_metadata(row.get("similarity", 0.0), row.get("metadata"))
for row in rows
)
if kinds:
allowed = set(kinds)
hits = [hit for hit in hits if hit.kind in allowed]
return _merge(hits, limit)
def _require_writer(self) -> VectorRestClient:
if self._writer is None:
raise VectorWriteUnavailable("Vector writer credential is not configured")
return self._writer
def existing_hashes(self, collection: str, kinds: list[str]) -> dict[str, str]:
_collection("vectors", collection)
_validate_collection_kinds(collection, kinds)
try:
return self._require_writer().existing_hashes(collection, kinds)
except VectorRestError as exc:
raise VectorStoreError(str(exc)) from exc
def upsert(self, collection: str, records: list[VectorWriteRecord]) -> int:
writer = self._require_writer()
_collection("vectors", collection)
for record in records:
_validate_collection_kinds(collection, [record.record.kind])
if (
self._expected_dimension is not None
and len(record.embedding) != self._expected_dimension
):
raise VectorStoreError("Embedding dimension does not match configured dimension")
rows = [self._row(record) for record in records]
try:
return writer.upsert_records(collection, rows)
except VectorRestError as exc:
raise VectorStoreError(str(exc)) from exc
def delete_generation(self, collection: str, generation: str, workspace_id: str) -> int:
if collection != "evidence" or re.fullmatch(r"gen:[0-9a-f]{32}", generation) is None:
raise VectorStoreError("Only exact Evidence generations may be deleted")
try:
return self._require_writer().delete_generation(collection, generation, workspace_id)
except VectorRestError as exc:
raise VectorStoreError(str(exc)) from exc
def list_evidence_generations(self, collection: str, workspace_id: str) -> list[str]:
if collection != "evidence":
raise VectorStoreError("Only exact Evidence generations may be listed")
try:
return self._require_writer().list_evidence_generations(collection, workspace_id)
except VectorRestError as exc:
raise VectorWriteUnavailable("Vector generation inventory unavailable") from exc
@staticmethod
def _row(write_record: VectorWriteRecord) -> dict:
record = write_record.record
metadata = {
"kind": record.kind,
"ref": record.ref,
"record_key": record.id,
"title": record.title,
"content": record.content,
**record.metadata,
}
return {
"record_key": record.id,
"kind": record.kind,
"content_hash": write_record.content_hash,
"metadata": metadata,
"embedding": write_record.embedding,
}
+10 -5
View File
@@ -2,9 +2,9 @@ from pathlib import Path
import typer
from tht.cli._guards import require_vector_write_allowed
from tht.cli.config_cmd import CONFIG_OPT
from tht.cli.schema_cmd import _load_config_or_exit
from tht.cli._guards import require_vector_write_allowed
evidence_app = typer.Typer(help="Generazione e gestione delle evidence")
@@ -56,12 +56,13 @@ def extract_cmd(config: Path = CONFIG_OPT) -> None:
@evidence_app.command("index")
def index_cmd(config: Path = CONFIG_OPT) -> None:
"""Embedda e sincronizza su pgvector tutte le evidence presenti in artifacts/."""
"""Embedda e sincronizza nel semantic store tutte le evidence presenti in artifacts/."""
from tht.adapters.factory import build_vector_store
from tht.cli.vector_cmd import (
_print_stats,
make_embedder,
open_store,
require_vector_cfg,
sync_canonical_records,
)
from tht.evidence.model import load_evidence_dir
from tht.vectorstore.records import evidence_records
@@ -71,6 +72,10 @@ def index_cmd(config: Path = CONFIG_OPT) -> None:
require_vector_cfg(cfg)
docs = load_evidence_dir(evidence_root(cfg))
records = evidence_records(docs, cfg.vector.max_chunk_chars)
store = open_store(cfg, "evidence")
stats = store.sync(records, make_embedder(cfg.embeddings), kinds={"evidence"})
stats = sync_canonical_records(
"evidence",
records,
store=build_vector_store(cfg, require_write=True),
embedder=make_embedder(cfg.embeddings),
)
_print_stats(stats)
+5 -20
View File
@@ -11,7 +11,6 @@ import typer
from sqlalchemy.exc import OperationalError, ProgrammingError
from tht.cli._guards import (
has_vector_write_rest,
require_server_profile,
require_vector_write_allowed,
)
@@ -45,15 +44,8 @@ def _resync_memory(cfg):
def clear_memory_index(cfg):
from tht.adapters.factory import build_vector_store
from tht.cli.vector_cmd import make_embedder, open_store, require_direct_vector_cfg
if cfg.vectors is not None and cfg.vectors.type == "qdrant":
return build_vector_store(cfg, require_write=True).delete_kinds("memory", ["memory"])
require_direct_vector_cfg(cfg)
legacy_store = open_store(cfg, "memory")
legacy_store.sync([], make_embedder(cfg.embeddings), kinds={"memory"})
return 0
return build_vector_store(cfg, require_write=True).delete_kinds("memory", ["memory"])
@memory_app.command("promote")
@@ -451,19 +443,13 @@ def search_cmd(
def index_solved_session(cfg, session_id: str) -> int:
"""Indicizza la coppia domanda->SQL della sessione (kind solved_question).
Solleva RuntimeError se manca la writer key e SolvedIndexError se mancano gli
artefatti: il finalize li degrada a warning, il comando CLI li converte in
errori espliciti."""
Solleva SolvedIndexError se mancano gli artefatti: il finalize lo degrada a warning,
il comando CLI lo converte in errore esplicito."""
from tht.adapters.factory import build_vector_store
from tht.cli.sql_cmd import promoted_tables_for
from tht.cli.vector_cmd import make_embedder
from tht.solved import build_solved_snapshot, save_solved_question
if not has_vector_write_rest(cfg):
raise RuntimeError(
"vector_write_rest assente: la coppia domanda->SQL si indicizza solo con la "
"writer key configurata nel workspace yaml"
)
store = build_vector_store(cfg, require_write=True)
record = build_solved_snapshot(load_snapshot_or_exit(cfg, session_id), promoted_tables_for(cfg, session_id))
return save_solved_question(
@@ -518,10 +504,9 @@ def solved_search_cmd(
from rich.table import Table
from tht.cli.vector_cmd import make_embedder, open_searcher
from tht.ports.vector import VectorReadUnavailable
from tht.ports.vector import VectorReadUnavailable, VectorStoreError
from tht.solved import SOLVED_KIND
from tht.vectorstore.embeddings import EmbeddingsError
from tht.vectorstore.rest_client import VectorRestError
cfg = _load_config_or_exit(config)
require_vector_cfg(cfg)
@@ -532,7 +517,7 @@ def solved_search_cmd(
searcher = open_searcher(cfg)
embedder = make_embedder(cfg.embeddings)
hits = searcher.search(embedder.embed_query(question), top_n=top, kinds=[SOLVED_KIND])
except (VectorRestError, VectorReadUnavailable, EmbeddingsError, OperationalError) as e:
except (VectorStoreError, VectorReadUnavailable, EmbeddingsError, OperationalError) as e:
typer.secho(
f"ATTENZIONE: exemplar non disponibili ({e}). Prosegui senza.",
fg=typer.colors.YELLOW, err=True,
+2 -3
View File
@@ -255,11 +255,10 @@ def pack_cmd(
from sqlalchemy.exc import OperationalError
from tht.cli.vector_cmd import make_embedder, open_searcher, require_vector_cfg
from tht.ports.vector import VectorReadUnavailable
from tht.ports.vector import VectorReadUnavailable, VectorStoreError
from tht.search import combined_search, schema_tables
from tht.solved import SOLVED_KIND
from tht.vectorstore.embeddings import EmbeddingsError
from tht.vectorstore.rest_client import VectorRestError
cfg = _load_config_or_exit(config)
from tht.search.evidence import validate_corpus_workspace
@@ -273,7 +272,7 @@ def pack_cmd(
evidence: list[dict] = []
solved: list[dict] = []
warnings: list[str] = []
degrade = (VectorRestError, VectorReadUnavailable, EmbeddingsError, OperationalError)
degrade = (VectorStoreError, VectorReadUnavailable, EmbeddingsError, OperationalError)
vec = None
searcher = embedder = None
+15 -40
View File
@@ -2,11 +2,7 @@ from pathlib import Path
import typer
from tht.cli._guards import (
has_vector_write_rest,
require_server_profile,
require_vector_write_allowed,
)
from tht.cli._guards import require_server_profile, require_vector_write_allowed
from tht.cli.config_cmd import CONFIG_OPT
from tht.cli.schema_cmd import _load_config_or_exit, annotations_path, physical_path
from tht.ports.vector import VectorWriteRecord
@@ -26,8 +22,8 @@ def require_vector_cfg(cfg):
missing = []
if cfg.embeddings is None:
missing.append("embeddings")
if cfg.vectors is None and cfg.vector_db is None and not has_vector_write_rest(cfg):
missing.append("vectors o vector_db o vector_write_rest")
if cfg.vectors is None:
missing.append("vectors")
if missing:
typer.secho(
f"ERRORE: sezioni mancanti nel workspace yaml: {', '.join(missing)}.",
@@ -36,27 +32,6 @@ def require_vector_cfg(cfg):
raise typer.Exit(code=1)
def require_direct_vector_cfg(cfg):
missing = [k for k in ("vector_db", "embeddings") if getattr(cfg, k) is None]
if missing:
typer.secho(
f"ERRORE: sezioni mancanti nel workspace yaml: {', '.join(missing)}.",
fg=typer.colors.RED, err=True,
)
raise typer.Exit(code=1)
def open_store(cfg, table: str):
"""Writer table-scoped per il LOADING.
Sul server preferisce la connessione diretta. In profilo workstation usa `vector_write_rest`
se configurato, con upsert remoto non distruttivo.
"""
from tht.adapters.factory import build_vector_loader
return build_vector_loader(cfg, table)
def open_searcher(cfg):
"""Searcher per la LETTURA (similarity search): via REST se `vector_rest` è configurato,
altrimenti connessione diretta (dev/test)."""
@@ -115,20 +90,20 @@ def init_cmd(
False, "--skip-ollama-check", help="Non verificare la raggiungibilita' di Ollama."
),
) -> None:
"""Crea schema e tabella pgvector (idempotente) e verifica le connessioni."""
from sqlalchemy.exc import OperationalError
"""Verifica il runtime Qdrant e la raggiungibilita' dell'embedder configurato."""
from tht.adapters.factory import build_vector_store
from tht.vectorstore.embeddings import EmbeddingsError
from tht.vectorstore.reader import ALL_TABLES
cfg = _load_config_or_exit(config)
require_server_profile(cfg, "vector init")
require_direct_vector_cfg(cfg)
try:
for table in ALL_TABLES:
open_store(cfg, table).init_schema()
except OperationalError as e:
typer.secho(f"ERRORE connessione pgvector: {e.orig}", fg=typer.colors.RED, err=True)
require_vector_cfg(cfg)
health = build_vector_store(cfg, require_write=True).health()
if not health.ok:
typer.secho(
f"ERRORE runtime vettoriale: {health.detail or 'Qdrant non raggiungibile o incompatibile'}",
fg=typer.colors.RED,
err=True,
)
raise typer.Exit(code=1)
if not skip_ollama_check:
try:
@@ -137,8 +112,8 @@ def init_cmd(
typer.secho(f"ERRORE: {e}", fg=typer.colors.RED, err=True)
raise typer.Exit(code=1)
typer.secho(
f"OK: schema {cfg.vector_db.db_schema} pronto (tabelle: {', '.join(ALL_TABLES)}) su "
f"{cfg.vector_db.host}:{cfg.vector_db.port}", fg=typer.colors.GREEN,
f"OK: runtime Qdrant pronto per la collezione {cfg.vectors.collection}",
fg=typer.colors.GREEN,
)
@@ -1,3 +0,0 @@
CREATE SCHEMA IF NOT EXISTS vectors;
REVOKE ALL ON SCHEMA vectors FROM PUBLIC;
CREATE EXTENSION IF NOT EXISTS vector WITH SCHEMA vectors;
@@ -1,32 +0,0 @@
CREATE TABLE IF NOT EXISTS vectors.schema_records (
id bigserial PRIMARY KEY,
record_key text UNIQUE NOT NULL,
kind text NOT NULL,
content_hash text NOT NULL,
metadata jsonb NOT NULL,
embedding vectors.vector(768) NOT NULL,
indexed_at timestamptz NOT NULL DEFAULT pg_catalog.now()
);
CREATE TABLE IF NOT EXISTS vectors.evidence (
id bigserial PRIMARY KEY,
record_key text UNIQUE NOT NULL,
kind text NOT NULL,
content_hash text NOT NULL,
metadata jsonb NOT NULL,
embedding vectors.vector(768) NOT NULL,
indexed_at timestamptz NOT NULL DEFAULT pg_catalog.now()
);
CREATE TABLE IF NOT EXISTS vectors.memory (
id bigserial PRIMARY KEY,
record_key text UNIQUE NOT NULL,
kind text NOT NULL,
content_hash text NOT NULL,
metadata jsonb NOT NULL,
embedding vectors.vector(768) NOT NULL,
indexed_at timestamptz NOT NULL DEFAULT pg_catalog.now()
);
REVOKE ALL ON ALL TABLES IN SCHEMA vectors FROM PUBLIC;
REVOKE ALL ON ALL SEQUENCES IN SCHEMA vectors FROM PUBLIC;
@@ -1,23 +0,0 @@
DO $roles$
BEGIN
IF NOT EXISTS (SELECT 1 FROM pg_catalog.pg_roles WHERE rolname = 'vector_reader') THEN
CREATE ROLE vector_reader NOLOGIN;
END IF;
IF NOT EXISTS (SELECT 1 FROM pg_catalog.pg_roles WHERE rolname = 'vector_writer') THEN
CREATE ROLE vector_writer NOLOGIN;
END IF;
END
$roles$;
REVOKE ALL ON SCHEMA vectors FROM vector_reader, vector_writer;
REVOKE ALL ON ALL TABLES IN SCHEMA vectors FROM vector_reader, vector_writer;
REVOKE ALL ON ALL SEQUENCES IN SCHEMA vectors FROM vector_reader, vector_writer;
GRANT USAGE ON SCHEMA vectors TO vector_reader, vector_writer;
GRANT SELECT ON ALL TABLES IN SCHEMA vectors TO vector_reader;
GRANT INSERT, UPDATE
ON vectors.schema_records, vectors.evidence, vectors.memory TO vector_writer;
GRANT SELECT (record_key, kind, content_hash)
ON vectors.schema_records, vectors.evidence, vectors.memory TO vector_writer;
GRANT USAGE ON ALL SEQUENCES IN SCHEMA vectors TO vector_writer;
@@ -1,3 +0,0 @@
-- The writer owns derived-generation reconciliation but not runtime similarity reads.
GRANT SELECT (metadata) ON vectors.evidence TO vector_writer;
GRANT DELETE ON vectors.evidence TO vector_writer;
+1 -1
View File
@@ -46,7 +46,7 @@ def _solved_hash(record: VectorRecord) -> str:
def save_solved_question(record: VectorRecord, *, store, embedder) -> int:
"""Upsert one-row della coppia domanda->SQL via writer key (stesso pattern di
save_one_memory, spec D11): hash dedup client-side, embedding solo se domanda
o SQL sono cambiati. `writer` e' un VectorRestClient (writer key). Ritorna il
o SQL sono cambiati. Ritorna il
numero di righe upsertate (0 = invariata)."""
from tht.ports.vector import VectorWriteRecord
+2 -47
View File
@@ -1,18 +1,4 @@
"""Lettura del pgvector dietro un'unica interfaccia `.search(query_vec, top_n, kinds)`, così
`search.combined_search` resta agnostico al transport. Due implementazioni:
- `RestSearcher` → produzione: similarity search via REST (`search_similar`).
- `DirectSearcher` → dev/test: connessione diretta a Postgres/pgvector.
Entrambe mappano i `kind` sulle tabelle per-dominio dello schema `vectors`.
"""
from sqlalchemy import Engine
from tht.adapters.vector.legacy_direct import LegacyDirectVectorStore
from tht.adapters.vector.thoth_http import ThothHttpVectorStore
from tht.vectorstore.rest_client import VectorRestClient
from tht.vectorstore.store import VectorHit
"""Collection mapping helpers for the workspace semantic store."""
# kind Thoth → tabella dello schema `vectors`.
KIND_TO_TABLE = {
@@ -30,35 +16,4 @@ def tables_for_kinds(kinds: list[str] | None) -> list[str]:
if not kinds:
return list(ALL_TABLES)
return sorted({KIND_TO_TABLE[k] for k in kinds if k in KIND_TO_TABLE})
class RestSearcher:
"""Similarity search via REST: una chiamata `search_similar` per tabella, poi fusione."""
def __init__(self, client: VectorRestClient):
self.client = client
self._store = ThothHttpVectorStore(reader=client, writer=None)
def search(
self, query_vec: list[float], top_n: int = 10, kinds: list[str] | None = None
) -> list[VectorHit]:
return self._store.search(
tables_for_kinds(kinds), query_vec, limit=top_n, kinds=kinds
)
class DirectSearcher:
"""Similarity search diretta su Postgres/pgvector, interrogando le tabelle per-dominio."""
def __init__(self, engine: Engine, schema: str = "vectors", dim: int = 768):
self.engine = engine
self.schema = schema
self.dim = dim
self._store = LegacyDirectVectorStore(engine, schema=schema, dim=dim)
def search(
self, query_vec: list[float], top_n: int = 10, kinds: list[str] | None = None
) -> list[VectorHit]:
return self._store.search(
tables_for_kinds(kinds), query_vec, limit=top_n, kinds=kinds
)
__all__ = ["ALL_TABLES", "KIND_TO_TABLE", "tables_for_kinds"]
-173
View File
@@ -1,173 +0,0 @@
"""Client per la similarity search del pgvector esposta via Supabase/PostgREST.
Endpoint dedicato (es. https://host/vector/v1/), distinto dal DWH. La lettura usa
`search_similar`; la scrittura remota usa RPC allowlist con una API key separata.
Errori in italiano e azionabili, stile `rest/client.py`.
"""
import re
import requests
from tht.config import RestConfig
class VectorRestError(Exception):
"""Errore di accesso al vector store via REST, con messaggio leggibile per il reviewer."""
class VectorRestClient:
def __init__(self, cfg: RestConfig):
self.cfg = cfg
self._base = cfg.base_url.rstrip("/")
@property
def api_key(self) -> str:
"""The REST API key for this client (spec D11: reader and writer carry
distinct keys against the same endpoint)."""
return self.cfg.api_key
def _post(self, fn: str, args: dict) -> requests.Response:
url = f"{self._base}/rpc/{fn}"
verify: bool | str = self.cfg.ssl_ca if self.cfg.ssl_ca else True
try:
return requests.post(
url,
json=args,
headers={"X-API-Key": self.cfg.api_key},
timeout=(self.cfg.connect_timeout, self.cfg.timeout),
verify=verify,
)
except requests.RequestException as e:
raise VectorRestError(
f"Vector REST non raggiungibile su {self.cfg.base_url} (rpc {fn}): {e}"
) from e
def _error_msg(self, fn: str, resp: requests.Response) -> str:
try:
body = resp.json()
detail = body.get("message") or body.get("details") or resp.text
except ValueError:
detail = resp.text
return f"Vector REST rpc {fn} → HTTP {resp.status_code}: {detail}"
def _call(self, fn: str, args: dict):
resp = self._post(fn, args)
if not resp.ok:
raise VectorRestError(self._error_msg(fn, resp))
if resp.status_code == 204 or not resp.text:
return None
return resp.json()
def search_similar(
self, table_name: str, query_embedding: list[float], limit_count: int,
kinds: list[str] | None = None,
metadata_filter: dict | None = None,
) -> list[dict]:
"""Ricerca per similarità coseno su `vectors.<table_name>`: ritorna le righe
`{id, similarity, metadata}` ordinate per similarity decrescente. Con `kinds`
il filtro avviene server-side nel WHERE della RPC (evita la diluizione del
top-k quando piu' kind condividono la tabella, es. memory/solved_question).
Su un server legacy senza il parametro (PostgREST 404) ritenta senza filtro:
resta il post-filter client-side di RestSearcher."""
args = {
"query_embedding": query_embedding,
"limit_count": limit_count,
"table_name": table_name,
}
if metadata_filter is not None:
# ACTIVE corpus reads must never degrade to an unfiltered legacy RPC:
# filtering after LIMIT is incomplete and could expose stale generations.
return self._call(
"search_similar",
{**args, "kinds": kinds, "metadata_filter": metadata_filter},
) or []
if kinds is not None:
try:
return self._call("search_similar", {**args, "kinds": kinds}) or []
except VectorRestError as e:
if "HTTP 404" not in str(e):
raise
# funzione a 3 argomenti (pre-migrazione kinds): fallback senza filtro
return self._call("search_similar", args) or []
def list_tables(self) -> list[dict]:
"""Tabelle vettoriali disponibili: `{table_name, vector_dimensions, …}`."""
return self._call("list_tables", {}) or []
def existing_hashes(self, table_name: str, kinds: list[str]) -> dict[str, str]:
"""Hash correnti per sync incrementale su una tabella vector allowlisted.
RPC attesa: `existing_vector_hashes(table_name, kinds)` -> righe
`{record_key, content_hash}`.
"""
rows = self._call(
"existing_vector_hashes",
{"table_name": table_name, "kinds": kinds},
) or []
return {row["record_key"]: row["content_hash"] for row in rows}
def upsert_records(self, table_name: str, rows: list[dict]) -> int:
"""Upsert controllato di record vettoriali già embeddati.
RPC attesa: `upsert_vector_records(table_name, rows)` -> `{upserted: N}` o righe.
Non espone delete/clear: il cleanup distruttivo resta solo-server.
"""
payload = self._call(
"upsert_vector_records",
{"table_name": table_name, "rows": rows},
)
if payload is None:
return len(rows)
if isinstance(payload, dict):
return int(payload.get("upserted", len(rows)))
# PostgREST puo' incapsulare uno scalar jsonb in una lista [{"upserted": N}]:
# estrai il conteggio dal primo elemento invece di restituire len(lista)=1.
if isinstance(payload, list):
if payload and isinstance(payload[0], dict) and "upserted" in payload[0]:
return int(payload[0]["upserted"])
return len(payload)
return len(rows)
def delete_generation(self, table_name: str, generation: str, workspace_id: str) -> int:
if table_name != "evidence" or re.fullmatch(r"gen:[0-9a-f]{32}", generation) is None:
raise ValueError("generation must be canonical")
if re.fullmatch(r"[a-z][a-z0-9_-]{0,63}", workspace_id) is None:
raise ValueError("workspace namespace must be canonical")
try:
payload = self._call(
"delete_vector_generation",
{"table_name": table_name, "kind": "evidence", "generation": generation,
"workspace_id": workspace_id},
)
except VectorRestError as error:
if "HTTP 404" in str(error):
raise VectorRestError(
"delete_vector_generation RPC is unavailable; deploy the cleanup migration"
) from None
raise
if isinstance(payload, dict):
return int(payload.get("deleted", 0))
return 0
def list_evidence_generations(self, table_name: str, workspace_id: str) -> list[str]:
if re.fullmatch(r"[a-z][a-z0-9_-]{0,63}", workspace_id) is None:
raise ValueError("workspace namespace must be canonical")
try:
rows = self._call(
"list_evidence_generations",
{"table_name": table_name, "kind": "evidence", "workspace_id": workspace_id},
) or []
except VectorRestError as error:
if "HTTP 404" in str(error):
raise VectorRestError(
"list_evidence_generations RPC is unavailable; deploy the cleanup migration"
) from None
raise
if not isinstance(rows, list) or any(
not isinstance(row, dict)
or re.fullmatch(r"gen:[0-9a-f]{32}", str(row.get("generation", ""))) is None
for row in rows
):
raise VectorRestError("list_evidence_generations returned malformed data")
return sorted({row["generation"] for row in rows})
-84
View File
@@ -1,84 +0,0 @@
"""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", "solved_question"},
}
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