fix(vector): separate write transport fields
This commit is contained in:
@@ -163,7 +163,8 @@ git commit -m "refactor(dwh): adapt direct and REST transports"
|
|||||||
- Modify: `harness/tht/vectorstore/reader.py`
|
- Modify: `harness/tht/vectorstore/reader.py`
|
||||||
|
|
||||||
**Interfaces:**
|
**Interfaces:**
|
||||||
- Produces: `VectorStore`, `VectorCapabilities`, `VectorHealth`, `VectorRecord`, `VectorHit`, `ThothHttpVectorStore`.
|
- Produces: `VectorStore`, `VectorCapabilities`, `VectorHealth`, `VectorRecord`,
|
||||||
|
`VectorWriteRecord`, `VectorHit`, `ThothHttpVectorStore`.
|
||||||
- Preserves: current `VectorRestClient`, `DirectSearcher`, and `RestSearcher` behavior behind wrappers.
|
- Preserves: current `VectorRestClient`, `DirectSearcher`, and `RestSearcher` behavior behind wrappers.
|
||||||
|
|
||||||
- [ ] **Step 1: Write read/write capability and dual-credential tests**
|
- [ ] **Step 1: Write read/write capability and dual-credential tests**
|
||||||
@@ -193,9 +194,13 @@ class VectorStore(Protocol):
|
|||||||
def search(self, collections: list[str], embedding: list[float], *, limit: int,
|
def search(self, collections: list[str], embedding: list[float], *, limit: int,
|
||||||
kinds: list[str] | None = None) -> list[VectorHit]: ...
|
kinds: list[str] | None = None) -> list[VectorHit]: ...
|
||||||
def existing_hashes(self, collection: str, kinds: list[str]) -> dict[str, str]: ...
|
def existing_hashes(self, collection: str, kinds: list[str]) -> dict[str, str]: ...
|
||||||
def upsert(self, collection: str, records: list[VectorRecord]) -> int: ...
|
def upsert(self, collection: str, records: list[VectorWriteRecord]) -> int: ...
|
||||||
```
|
```
|
||||||
|
|
||||||
|
`VectorWriteRecord` is the transport-neutral write envelope: it contains the canonical
|
||||||
|
`VectorRecord`, a precomputed embedding, and a content hash. Adapters must preserve
|
||||||
|
`VectorRecord.metadata` unchanged, including semantic keys named `embedding` or `content_hash`.
|
||||||
|
|
||||||
- [ ] **Step 4: Run vector regression tests**
|
- [ ] **Step 4: Run vector regression tests**
|
||||||
|
|
||||||
Run: `cd harness && .venv/bin/pytest tests/test_vector_port_contract.py tests/test_vector_dual_key.py tests/test_search_similar_kinds.py tests/test_memory_save_one.py tests/test_solved_question.py -q`
|
Run: `cd harness && .venv/bin/pytest tests/test_vector_port_contract.py tests/test_vector_dual_key.py tests/test_search_similar_kinds.py tests/test_memory_save_one.py tests/test_solved_question.py -q`
|
||||||
|
|||||||
@@ -3,12 +3,15 @@ from unittest.mock import MagicMock
|
|||||||
import pytest
|
import pytest
|
||||||
|
|
||||||
from tht.adapters.vector.thoth_http import ThothHttpVectorStore
|
from tht.adapters.vector.thoth_http import ThothHttpVectorStore
|
||||||
|
from tht.evidence.model import EvidenceDoc
|
||||||
from tht.ports.vector import (
|
from tht.ports.vector import (
|
||||||
VectorHit,
|
VectorHit,
|
||||||
VectorRecord,
|
VectorRecord,
|
||||||
VectorStore,
|
VectorStore,
|
||||||
|
VectorWriteRecord,
|
||||||
VectorWriteUnavailable,
|
VectorWriteUnavailable,
|
||||||
)
|
)
|
||||||
|
from tht.vectorstore.records import evidence_records
|
||||||
|
|
||||||
|
|
||||||
def test_http_store_reports_reader_without_writer():
|
def test_http_store_reports_reader_without_writer():
|
||||||
@@ -67,13 +70,16 @@ def test_http_store_keeps_reader_and_writer_operations_separate():
|
|||||||
writer.existing_hashes.assert_called_once_with("memory", ["memory"])
|
writer.existing_hashes.assert_called_once_with("memory", ["memory"])
|
||||||
|
|
||||||
records = [
|
records = [
|
||||||
VectorRecord(
|
VectorWriteRecord(
|
||||||
|
record=VectorRecord(
|
||||||
id="m1",
|
id="m1",
|
||||||
kind="memory",
|
kind="memory",
|
||||||
ref="session:s1",
|
ref="session:s1",
|
||||||
title="Choice",
|
title="Choice",
|
||||||
content="Use the curated table",
|
content="Use the curated table",
|
||||||
metadata={"embedding": [0.1, 0.2], "content_hash": "abc"},
|
),
|
||||||
|
embedding=[0.1, 0.2],
|
||||||
|
content_hash="abc",
|
||||||
)
|
)
|
||||||
]
|
]
|
||||||
assert store.upsert("memory", records) == 1
|
assert store.upsert("memory", records) == 1
|
||||||
@@ -81,11 +87,64 @@ def test_http_store_keeps_reader_and_writer_operations_separate():
|
|||||||
reader.upsert_records.assert_not_called()
|
reader.upsert_records.assert_not_called()
|
||||||
|
|
||||||
|
|
||||||
|
def test_http_upsert_serializes_a_canonical_builder_record():
|
||||||
|
record = evidence_records(
|
||||||
|
[EvidenceDoc(id="joins", title="Join guidance", body="Use the curated join")],
|
||||||
|
max_chunk_chars=1000,
|
||||||
|
)[0]
|
||||||
|
writer = MagicMock()
|
||||||
|
writer.upsert_records.return_value = 1
|
||||||
|
store = ThothHttpVectorStore(reader=MagicMock(), writer=writer)
|
||||||
|
|
||||||
|
assert store.upsert(
|
||||||
|
"evidence",
|
||||||
|
[VectorWriteRecord(record=record, embedding=[0.2, 0.3], content_hash="digest")],
|
||||||
|
) == 1
|
||||||
|
row = writer.upsert_records.call_args.args[1][0]
|
||||||
|
assert row["record_key"] == "evidence:joins:0"
|
||||||
|
assert row["metadata"]["status"] == "reviewed"
|
||||||
|
assert row["embedding"] == [0.2, 0.3]
|
||||||
|
assert row["content_hash"] == "digest"
|
||||||
|
|
||||||
|
|
||||||
|
def test_http_upsert_preserves_metadata_named_like_transport_fields():
|
||||||
|
record = VectorRecord(
|
||||||
|
id="collision",
|
||||||
|
kind="memory",
|
||||||
|
ref="session:s1",
|
||||||
|
title="Collision",
|
||||||
|
content="Semantic metadata must survive",
|
||||||
|
metadata={"embedding": "semantic embedding", "content_hash": "semantic hash"},
|
||||||
|
)
|
||||||
|
writer = MagicMock()
|
||||||
|
store = ThothHttpVectorStore(reader=MagicMock(), writer=writer)
|
||||||
|
|
||||||
|
store.upsert(
|
||||||
|
"memory",
|
||||||
|
[VectorWriteRecord(record=record, embedding=[0.4], content_hash="transport hash")],
|
||||||
|
)
|
||||||
|
row = writer.upsert_records.call_args.args[1][0]
|
||||||
|
assert row["embedding"] == [0.4]
|
||||||
|
assert row["content_hash"] == "transport hash"
|
||||||
|
assert row["metadata"]["embedding"] == "semantic embedding"
|
||||||
|
assert row["metadata"]["content_hash"] == "semantic hash"
|
||||||
|
|
||||||
|
|
||||||
def test_http_store_is_runtime_vector_store():
|
def test_http_store_is_runtime_vector_store():
|
||||||
store = ThothHttpVectorStore(reader=MagicMock(), writer=None)
|
store = ThothHttpVectorStore(reader=MagicMock(), writer=None)
|
||||||
assert isinstance(store, VectorStore)
|
assert isinstance(store, VectorStore)
|
||||||
|
|
||||||
|
|
||||||
|
def test_vector_contract_is_exported_from_public_packages():
|
||||||
|
from tht.adapters.vector import ThothHttpVectorStore as PublicHttpStore
|
||||||
|
from tht.ports import VectorStore as PublicVectorStore
|
||||||
|
from tht.ports import VectorWriteRecord as PublicVectorWriteRecord
|
||||||
|
|
||||||
|
assert PublicHttpStore is ThothHttpVectorStore
|
||||||
|
assert PublicVectorStore is VectorStore
|
||||||
|
assert PublicVectorWriteRecord is VectorWriteRecord
|
||||||
|
|
||||||
|
|
||||||
def test_http_health_uses_reader_list_tables_and_reports_failure():
|
def test_http_health_uses_reader_list_tables_and_reports_failure():
|
||||||
reader = MagicMock()
|
reader = MagicMock()
|
||||||
store = ThothHttpVectorStore(reader=reader, writer=None)
|
store = ThothHttpVectorStore(reader=reader, writer=None)
|
||||||
|
|||||||
@@ -2,7 +2,12 @@
|
|||||||
|
|
||||||
from sqlalchemy import Engine
|
from sqlalchemy import Engine
|
||||||
|
|
||||||
from tht.ports.vector import VectorCapabilities, VectorHealth, VectorRecord, VectorWriteUnavailable
|
from tht.ports.vector import (
|
||||||
|
VectorCapabilities,
|
||||||
|
VectorHealth,
|
||||||
|
VectorWriteRecord,
|
||||||
|
VectorWriteUnavailable,
|
||||||
|
)
|
||||||
from tht.vectorstore.store import VectorHit, VectorStore as TableVectorStore
|
from tht.vectorstore.store import VectorHit, VectorStore as TableVectorStore
|
||||||
|
|
||||||
|
|
||||||
@@ -43,5 +48,5 @@ class LegacyDirectVectorStore:
|
|||||||
def existing_hashes(self, collection: str, kinds: list[str]) -> dict[str, str]:
|
def existing_hashes(self, collection: str, kinds: list[str]) -> dict[str, str]:
|
||||||
raise VectorWriteUnavailable("Legacy direct reader has no writer interface")
|
raise VectorWriteUnavailable("Legacy direct reader has no writer interface")
|
||||||
|
|
||||||
def upsert(self, collection: str, records: list[VectorRecord]) -> int:
|
def upsert(self, collection: str, records: list[VectorWriteRecord]) -> int:
|
||||||
raise VectorWriteUnavailable("Legacy direct reader has no writer interface")
|
raise VectorWriteUnavailable("Legacy direct reader has no writer interface")
|
||||||
|
|||||||
@@ -4,12 +4,11 @@ from tht.ports.vector import (
|
|||||||
VectorCapabilities,
|
VectorCapabilities,
|
||||||
VectorHealth,
|
VectorHealth,
|
||||||
VectorHit,
|
VectorHit,
|
||||||
VectorRecord,
|
VectorWriteRecord,
|
||||||
VectorStoreError,
|
|
||||||
VectorWriteUnavailable,
|
VectorWriteUnavailable,
|
||||||
)
|
)
|
||||||
from tht.vectorstore.rest_client import VectorRestClient
|
from tht.vectorstore.rest_client import VectorRestClient
|
||||||
from tht.vectorstore.store import content_hash, hit_from_metadata
|
from tht.vectorstore.store import hit_from_metadata
|
||||||
|
|
||||||
|
|
||||||
def _merge(hits: list[VectorHit], limit: int) -> list[VectorHit]:
|
def _merge(hits: list[VectorHit], limit: int) -> list[VectorHit]:
|
||||||
@@ -63,31 +62,26 @@ class ThothHttpVectorStore:
|
|||||||
def existing_hashes(self, collection: str, kinds: list[str]) -> dict[str, str]:
|
def existing_hashes(self, collection: str, kinds: list[str]) -> dict[str, str]:
|
||||||
return self._require_writer().existing_hashes(collection, kinds)
|
return self._require_writer().existing_hashes(collection, kinds)
|
||||||
|
|
||||||
def upsert(self, collection: str, records: list[VectorRecord]) -> int:
|
def upsert(self, collection: str, records: list[VectorWriteRecord]) -> int:
|
||||||
writer = self._require_writer()
|
writer = self._require_writer()
|
||||||
rows = [self._row(record) for record in records]
|
rows = [self._row(record) for record in records]
|
||||||
return writer.upsert_records(collection, rows)
|
return writer.upsert_records(collection, rows)
|
||||||
|
|
||||||
@staticmethod
|
@staticmethod
|
||||||
def _row(record: VectorRecord) -> dict:
|
def _row(write_record: VectorWriteRecord) -> dict:
|
||||||
extra = dict(record.metadata)
|
record = write_record.record
|
||||||
try:
|
|
||||||
embedding = extra.pop("embedding")
|
|
||||||
except KeyError as exc:
|
|
||||||
raise VectorStoreError(f"Vector record {record.id!r} has no embedding") from exc
|
|
||||||
digest = extra.pop("content_hash", content_hash(record.content))
|
|
||||||
metadata = {
|
metadata = {
|
||||||
"kind": record.kind,
|
"kind": record.kind,
|
||||||
"ref": record.ref,
|
"ref": record.ref,
|
||||||
"record_key": record.id,
|
"record_key": record.id,
|
||||||
"title": record.title,
|
"title": record.title,
|
||||||
"content": record.content,
|
"content": record.content,
|
||||||
**extra,
|
**record.metadata,
|
||||||
}
|
}
|
||||||
return {
|
return {
|
||||||
"record_key": record.id,
|
"record_key": record.id,
|
||||||
"kind": record.kind,
|
"kind": record.kind,
|
||||||
"content_hash": digest,
|
"content_hash": write_record.content_hash,
|
||||||
"metadata": metadata,
|
"metadata": metadata,
|
||||||
"embedding": embedding,
|
"embedding": write_record.embedding,
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -7,6 +7,16 @@ from tht.ports.dwh import (
|
|||||||
DistinctValues,
|
DistinctValues,
|
||||||
UnsupportedCapability,
|
UnsupportedCapability,
|
||||||
)
|
)
|
||||||
|
from tht.ports.vector import (
|
||||||
|
VectorCapabilities,
|
||||||
|
VectorHealth,
|
||||||
|
VectorHit,
|
||||||
|
VectorRecord,
|
||||||
|
VectorStore,
|
||||||
|
VectorStoreError,
|
||||||
|
VectorWriteRecord,
|
||||||
|
VectorWriteUnavailable,
|
||||||
|
)
|
||||||
|
|
||||||
__all__ = [
|
__all__ = [
|
||||||
"DwhAdapter",
|
"DwhAdapter",
|
||||||
@@ -14,4 +24,12 @@ __all__ = [
|
|||||||
"DwhHealth",
|
"DwhHealth",
|
||||||
"DistinctValues",
|
"DistinctValues",
|
||||||
"UnsupportedCapability",
|
"UnsupportedCapability",
|
||||||
|
"VectorCapabilities",
|
||||||
|
"VectorHealth",
|
||||||
|
"VectorHit",
|
||||||
|
"VectorRecord",
|
||||||
|
"VectorStore",
|
||||||
|
"VectorStoreError",
|
||||||
|
"VectorWriteRecord",
|
||||||
|
"VectorWriteUnavailable",
|
||||||
]
|
]
|
||||||
|
|||||||
@@ -20,6 +20,15 @@ class VectorHealth:
|
|||||||
detail: str | None = None
|
detail: str | None = None
|
||||||
|
|
||||||
|
|
||||||
|
@dataclass(frozen=True)
|
||||||
|
class VectorWriteRecord:
|
||||||
|
"""A canonical record plus transport-neutral, precomputed vector data."""
|
||||||
|
|
||||||
|
record: VectorRecord
|
||||||
|
embedding: list[float]
|
||||||
|
content_hash: str
|
||||||
|
|
||||||
|
|
||||||
class VectorStoreError(Exception):
|
class VectorStoreError(Exception):
|
||||||
"""Base error exposed by vector adapters."""
|
"""Base error exposed by vector adapters."""
|
||||||
|
|
||||||
@@ -46,7 +55,7 @@ class VectorStore(Protocol):
|
|||||||
|
|
||||||
def existing_hashes(self, collection: str, kinds: list[str]) -> dict[str, str]: ...
|
def existing_hashes(self, collection: str, kinds: list[str]) -> dict[str, str]: ...
|
||||||
|
|
||||||
def upsert(self, collection: str, records: list[VectorRecord]) -> int: ...
|
def upsert(self, collection: str, records: list[VectorWriteRecord]) -> int: ...
|
||||||
|
|
||||||
|
|
||||||
__all__ = [
|
__all__ = [
|
||||||
@@ -56,5 +65,6 @@ __all__ = [
|
|||||||
"VectorRecord",
|
"VectorRecord",
|
||||||
"VectorStore",
|
"VectorStore",
|
||||||
"VectorStoreError",
|
"VectorStoreError",
|
||||||
|
"VectorWriteRecord",
|
||||||
"VectorWriteUnavailable",
|
"VectorWriteUnavailable",
|
||||||
]
|
]
|
||||||
|
|||||||
Reference in New Issue
Block a user