fix(adapter): close final foundation review
This commit is contained in:
@@ -1,6 +1,6 @@
|
||||
"""Thoth/PostgREST implementation of the DWH port."""
|
||||
|
||||
from tht.config import DatabaseConfig, RestConfig
|
||||
from tht.config import DatabaseIdentityConfig, RestConfig
|
||||
from tht.db.introspect import introspect_rest
|
||||
from tht.db import sampling
|
||||
from tht.execute import ExecResult, ExecutionError, PlanSummary
|
||||
@@ -13,7 +13,7 @@ from tht.rest.execute import explain_rest, run_controlled_rest
|
||||
class ThothRestDwhAdapter:
|
||||
capabilities = DwhCapabilities()
|
||||
|
||||
def __init__(self, database: DatabaseConfig, rest: RestConfig):
|
||||
def __init__(self, database: DatabaseIdentityConfig, rest: RestConfig):
|
||||
self._database = database
|
||||
self._client = RestClient(rest)
|
||||
|
||||
|
||||
@@ -41,13 +41,12 @@ def build_vector_store(cfg: Config, *, require_write: bool = False) -> VectorSto
|
||||
dim=dim,
|
||||
)
|
||||
case "thoth_vector_http":
|
||||
if resource.reader is None:
|
||||
raise ConfigError("Vector reader non configurato")
|
||||
if require_write and resource.writer is None:
|
||||
raise ConfigError("Vector writer non configurato")
|
||||
return ThothHttpVectorStore(
|
||||
VectorRestClient(resource.reader),
|
||||
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 other: # pragma: no cover - Pydantic's discriminator rejects this first.
|
||||
raise ConfigError(f"Adapter vector non supportato: {other}")
|
||||
|
||||
@@ -7,6 +7,7 @@ from tht.ports.vector import (
|
||||
VectorHealth,
|
||||
VectorWriteRecord,
|
||||
VectorWriteUnavailable,
|
||||
require_positive_limit,
|
||||
)
|
||||
from tht.vectorstore.store import VectorHit, VectorStore as TableVectorStore
|
||||
|
||||
@@ -26,8 +27,20 @@ class LegacyDirectVectorStore:
|
||||
with self._engine.connect() as connection:
|
||||
connection.exec_driver_sql("SELECT 1")
|
||||
except Exception as exc:
|
||||
return VectorHealth(ok=False, detail=str(exc))
|
||||
return VectorHealth(ok=True)
|
||||
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,
|
||||
@@ -37,6 +50,7 @@ class LegacyDirectVectorStore:
|
||||
limit: int,
|
||||
kinds: list[str] | None = None,
|
||||
) -> list[VectorHit]:
|
||||
require_positive_limit(limit)
|
||||
hits: list[VectorHit] = []
|
||||
for collection in collections:
|
||||
table = TableVectorStore(
|
||||
|
||||
@@ -4,8 +4,10 @@ from tht.ports.vector import (
|
||||
VectorCapabilities,
|
||||
VectorHealth,
|
||||
VectorHit,
|
||||
VectorReadUnavailable,
|
||||
VectorWriteRecord,
|
||||
VectorWriteUnavailable,
|
||||
require_positive_limit,
|
||||
)
|
||||
from tht.vectorstore.rest_client import VectorRestClient
|
||||
from tht.vectorstore.store import hit_from_metadata
|
||||
@@ -18,21 +20,63 @@ def _merge(hits: list[VectorHit], limit: int) -> list[VectorHit]:
|
||||
class ThothHttpVectorStore:
|
||||
"""Vector port backed by the existing allowlisted REST RPCs."""
|
||||
|
||||
def __init__(self, reader: VectorRestClient, writer: VectorRestClient | None):
|
||||
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=True, existing_hashes=writable, upsert=writable)
|
||||
return VectorCapabilities(
|
||||
search=self._reader is not None, existing_hashes=writable, upsert=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:
|
||||
self._reader.list_tables()
|
||||
return True, None, client.list_tables()
|
||||
except Exception as exc:
|
||||
return VectorHealth(ok=False, detail=str(exc))
|
||||
return VectorHealth(ok=True)
|
||||
return False, str(exc), []
|
||||
|
||||
def search(
|
||||
self,
|
||||
@@ -42,6 +86,9 @@ class ThothHttpVectorStore:
|
||||
limit: int,
|
||||
kinds: list[str] | None = None,
|
||||
) -> list[VectorHit]:
|
||||
require_positive_limit(limit)
|
||||
if self._reader is None:
|
||||
raise VectorReadUnavailable("Vector reader credential is not configured")
|
||||
hits: list[VectorHit] = []
|
||||
for collection in collections:
|
||||
rows = self._reader.search_similar(collection, embedding, limit, kinds=kinds)
|
||||
|
||||
Reference in New Issue
Block a user