From 7fae16b351f4db3ef5672bf51c625fa19c6fe5b7 Mon Sep 17 00:00:00 2001 From: Rajat Ahuja Date: Mon, 20 Apr 2026 16:56:51 -0400 Subject: [PATCH] handle turbopuffer server errors (#561) * fix: catch InternalServerError from turbopuffer * fix: remove unused VectorUpsertResult * fix: downgrade vector store sync errors to warnings * fix: remove upsert_with_retry * fix: (vector) add silent path and explicit path for vector db server errors --------- Co-authored-by: Vineeth Voruganti <13438633+VVoruganti@users.noreply.github.com> --- src/crud/document.py | 53 ++++++++++---- src/crud/message.py | 30 ++++++-- src/reconciler/sync_vectors.py | 21 +++++- src/vector_store/__init__.py | 22 +----- src/vector_store/lancedb.py | 8 +- src/vector_store/turbopuffer.py | 41 +++++++++-- src/vector_store/utils.py | 57 --------------- tests/conftest.py | 7 +- tests/deriver/test_vector_reconciliation.py | 19 ++--- tests/vector_store/test_turbopuffer.py | 81 +++++++++++++++++++++ 10 files changed, 211 insertions(+), 128 deletions(-) delete mode 100644 src/vector_store/utils.py create mode 100644 tests/vector_store/test_turbopuffer.py diff --git a/src/crud/document.py b/src/crud/document.py index 7de9dfd6..ba652eea 100644 --- a/src/crud/document.py +++ b/src/crud/document.py @@ -17,13 +17,16 @@ from src.crud.peer import get_peer from src.crud.session import get_session from src.dependencies import tracked_db from src.embedding_client import embedding_client -from src.exceptions import ResourceNotFoundException, ValidationException +from src.exceptions import ( + ResourceNotFoundException, + ValidationException, + VectorStoreError, +) from src.utils.filter import apply_filter from src.vector_store import ( VectorRecord, VectorStore, get_external_vector_store, - upsert_with_retry, ) logger = getLogger(__name__) @@ -560,11 +563,9 @@ async def create_documents( ) ) - # Upsert to external vector store with retry and update sync state + # Upsert to external vector store and update sync state try: - await upsert_with_retry( - external_vector_store, namespace, vector_records - ) + await external_vector_store.upsert_many(namespace, vector_records) # Success: mark as synced await db.execute( update(models.Document) @@ -577,9 +578,21 @@ async def create_documents( ) await db.commit() + except VectorStoreError: + # Vector store unavailable - increment sync_attempts for reconciliation + logger.warning("Vector store unavailable; leaving docs unsynced") + await db.execute( + update(models.Document) + .where(models.Document.id.in_(doc_ids)) + .values( + sync_attempts=models.Document.sync_attempts + 1, + last_sync_at=func.now(), + ) + ) + await db.commit() + except Exception: - # Failed after retries - increment sync_attempts for reconciliation - logger.exception("Failed to upsert vectors after retries") + logger.exception("Unexpected error upserting vectors") await db.execute( update(models.Document) .where(models.Document.id.in_(doc_ids)) @@ -846,11 +859,9 @@ async def create_observations( ) ) - # Upsert to external vector store with retry and update sync state + # Upsert to external vector store and update sync state try: - await upsert_with_retry( - external_vector_store, namespace, vector_records - ) + await external_vector_store.upsert_many(namespace, vector_records) # Success: mark as synced await db.execute( update(models.Document) @@ -863,10 +874,24 @@ async def create_observations( ) await db.commit() + except VectorStoreError: + logger.warning( + "Vector store unavailable for namespace %s; leaving observations unsynced", + namespace, + ) + await db.execute( + update(models.Document) + .where(models.Document.id.in_(doc_ids)) + .values( + sync_attempts=models.Document.sync_attempts + 1, + last_sync_at=func.now(), + ) + ) + await db.commit() + except Exception: - # Failed after retries - increment sync_attempts for reconciliation logger.exception( - f"Failed to upsert vectors for {namespace} after retries" + "Unexpected error upserting vectors for %s", namespace ) await db.execute( update(models.Document) diff --git a/src/crud/message.py b/src/crud/message.py index 08c4c861..3334c6e8 100644 --- a/src/crud/message.py +++ b/src/crud/message.py @@ -11,9 +11,10 @@ from src import models, schemas from src.config import settings from src.dependencies import tracked_db from src.embedding_client import embedding_client +from src.exceptions import VectorStoreError from src.utils.filter import apply_filter from src.utils.formatting import ILIKE_ESCAPE_CHAR, escape_ilike_pattern -from src.vector_store import VectorRecord, get_external_vector_store, upsert_with_retry +from src.vector_store import VectorRecord, get_external_vector_store from .session import get_or_create_session @@ -349,11 +350,11 @@ async def create_messages( ) ) - # Upsert to external vector store with retry and update sync state + # Upsert to external vector store and update sync state if vector_records: try: - await upsert_with_retry( - external_vector_store, namespace, vector_records + await external_vector_store.upsert_many( + namespace, vector_records ) # Success: mark as synced if we have DB rows if embedding_ids: @@ -368,10 +369,9 @@ async def create_messages( ) await db.commit() - except Exception: - # Failed after retries - increment sync_attempts for reconciliation - logger.exception( - "Failed to upsert message vectors after retries" + except VectorStoreError: + logger.warning( + "Vector store unavailable; leaving message vectors unsynced" ) if embedding_ids: await db.execute( @@ -385,6 +385,20 @@ async def create_messages( ) await db.commit() + except Exception: + logger.exception("Unexpected error upserting message vectors") + if embedding_ids: + await db.execute( + update(models.MessageEmbedding) + .where(models.MessageEmbedding.id.in_(embedding_ids)) + .values( + sync_attempts=models.MessageEmbedding.sync_attempts + + 1, + last_sync_at=func.now(), + ) + ) + await db.commit() + except Exception: logger.exception( "Failed to generate message embeddings for %s messages in workspace %s and session %s.", diff --git a/src/reconciler/sync_vectors.py b/src/reconciler/sync_vectors.py index 32de0d19..4a17e40e 100644 --- a/src/reconciler/sync_vectors.py +++ b/src/reconciler/sync_vectors.py @@ -19,6 +19,7 @@ from src import models from src.config import settings from src.dependencies import tracked_db from src.embedding_client import embedding_client +from src.exceptions import VectorStoreError from src.vector_store import VectorRecord, VectorStore, get_external_vector_store logger = logging.getLogger(__name__) @@ -259,8 +260,16 @@ async def _sync_documents( .values(sync_state="synced", last_sync_at=func.now(), sync_attempts=0) ) synced_count += len(docs_to_sync) + except VectorStoreError: + logger.warning( + "Vector store unavailable while syncing namespace %s", namespace + ) + await _bump_document_sync_attempts(db, docs_to_sync) + failed_count += len(docs_to_sync) except Exception: - logger.exception("Failed to sync documents to namespace %s", namespace) + logger.exception( + "Unexpected error syncing documents to namespace %s", namespace + ) await _bump_document_sync_attempts(db, docs_to_sync) failed_count += len(docs_to_sync) @@ -399,9 +408,17 @@ async def _sync_message_embeddings( .values(sync_state="synced", last_sync_at=func.now(), sync_attempts=0) ) synced_count += len(embs_to_sync) + except VectorStoreError: + logger.warning( + "Vector store unavailable while syncing message embeddings to namespace %s", + namespace, + ) + await _bump_message_embedding_sync_attempts(db, embs_to_sync) + failed_count += len(embs_to_sync) except Exception: logger.exception( - "Failed to sync message embeddings to namespace %s", namespace + "Unexpected error syncing message embeddings to namespace %s", + namespace, ) await _bump_message_embedding_sync_attempts(db, embs_to_sync) failed_count += len(embs_to_sync) diff --git a/src/vector_store/__init__.py b/src/vector_store/__init__.py index c77064a7..5a22abd9 100644 --- a/src/vector_store/__init__.py +++ b/src/vector_store/__init__.py @@ -50,17 +50,6 @@ class VectorQueryResult(BaseModel): metadata: dict[str, Any] = Field(default_factory=dict) -class VectorUpsertResult(BaseModel): - """Result for a vector upsert operation.""" - - model_config: ClassVar[ConfigDict] = ConfigDict( - extra="forbid", - frozen=True, - ) - - ok: bool - - class VectorStore(ABC): """ Abstract base class for vector store implementations. @@ -123,7 +112,7 @@ class VectorStore(ABC): self, namespace: str, vectors: list[VectorRecord], - ) -> VectorUpsertResult: + ) -> None: """ Upsert multiple vectors into the store. @@ -131,8 +120,8 @@ class VectorStore(ABC): namespace: The namespace to store the vectors in vectors: List of VectorRecord objects to upsert - Returns: - Result describing primary/secondary outcomes. + Raises: + Exception: If the write fails. """ ... @@ -192,9 +181,6 @@ class VectorStore(ABC): ... -from src.vector_store.utils import upsert_with_retry # noqa: E402 - - def _create_store_by_type(store_type: str) -> VectorStore: """Create a vector store instance by type name.""" if store_type == "turbopuffer": @@ -251,9 +237,7 @@ __all__ = [ "VectorStore", "VectorRecord", "VectorQueryResult", - "VectorUpsertResult", "get_external_vector_store", "close_external_vector_store", - "upsert_with_retry", "_hash_namespace_components", ] diff --git a/src/vector_store/lancedb.py b/src/vector_store/lancedb.py index 77c3c6cd..f63b8cfd 100644 --- a/src/vector_store/lancedb.py +++ b/src/vector_store/lancedb.py @@ -17,7 +17,7 @@ from lancedb import AsyncConnection, AsyncTable from src.config import settings from src.exceptions import VectorStoreError -from . import VectorQueryResult, VectorRecord, VectorStore, VectorUpsertResult +from . import VectorQueryResult, VectorRecord, VectorStore logger = logging.getLogger(__name__) @@ -156,7 +156,7 @@ class LanceDBVectorStore(VectorStore): self, namespace: str, vectors: list[VectorRecord], - ) -> VectorUpsertResult: + ) -> None: """ Upsert multiple vectors into LanceDB. @@ -165,7 +165,7 @@ class LanceDBVectorStore(VectorStore): vectors: List of VectorRecord objects to upsert """ if not vectors: - return VectorUpsertResult(ok=True) + return try: rows = [self._row_to_dict(v) for v in vectors] @@ -180,7 +180,7 @@ class LanceDBVectorStore(VectorStore): ) logger.debug(f"Upserted {len(vectors)} vectors to namespace {namespace}") - return VectorUpsertResult(ok=True) + return except Exception as e: logger.exception( f"Failed to upsert {len(vectors)} vectors to namespace {namespace}" diff --git a/src/vector_store/turbopuffer.py b/src/vector_store/turbopuffer.py index 39b93f5f..c84c8774 100644 --- a/src/vector_store/turbopuffer.py +++ b/src/vector_store/turbopuffer.py @@ -8,13 +8,14 @@ import logging from collections.abc import Sequence from typing import Any, Literal, cast -from turbopuffer import AsyncTurbopuffer, NotFoundError +from turbopuffer import AsyncTurbopuffer, InternalServerError, NotFoundError from turbopuffer.lib.namespace import AsyncNamespace from turbopuffer.types import Filter from src.config import settings +from src.exceptions import VectorStoreError -from . import VectorQueryResult, VectorRecord, VectorStore, VectorUpsertResult +from . import VectorQueryResult, VectorRecord, VectorStore logger = logging.getLogger(__name__) @@ -62,7 +63,7 @@ class TurbopufferVectorStore(VectorStore): self, namespace: str, vectors: list[VectorRecord], - ) -> VectorUpsertResult: + ) -> None: """ Upsert multiple vectors into Turbopuffer. @@ -71,7 +72,7 @@ class TurbopufferVectorStore(VectorStore): vectors: List of VectorRecord objects to upsert """ if not vectors: - return VectorUpsertResult(ok=True) + return ns = self._get_namespace(namespace) @@ -89,7 +90,18 @@ class TurbopufferVectorStore(VectorStore): upsert_rows=rows, distance_metric=DISTANCE_METRIC, ) - return VectorUpsertResult(ok=True) + return + except InternalServerError as exc: + # Turbopuffer unavailable. SDK implicitly retries 5xx responses, + # so raise a vector store error and let callers leave writes unsynced. + logger.warning( + "Turbopuffer unavailable for upsert to namespace %s (%s after retries)", + namespace, + exc.status_code, + ) + raise VectorStoreError( + f"Turbopuffer unavailable for upsert to namespace {namespace}" + ) from exc except Exception: logger.exception( f"Failed to upsert {len(vectors)} vectors to namespace {namespace}" @@ -183,6 +195,16 @@ class TurbopufferVectorStore(VectorStore): ) return [] + except InternalServerError as exc: + # Turbopuffer unavailable. SDK implicitly retries 5xx responses, + # so we should return []. + logger.warning( + "Turbopuffer unavailable for query on namespace %s (%s after retries), returning empty results", + namespace, + exc.status_code, + ) + return [] + except Exception: logger.exception(f"Failed to query namespace {namespace}") raise @@ -247,6 +269,15 @@ class TurbopufferVectorStore(VectorStore): except NotFoundError: # Namespace doesn't exist - nothing to delete logger.debug(f"Namespace {namespace} does not exist, nothing to delete") + except InternalServerError as exc: + logger.warning( + "Turbopuffer unavailable for delete from namespace %s (%s after retries)", + namespace, + exc.status_code, + ) + raise VectorStoreError( + f"Turbopuffer unavailable while deleting vectors in namespace {namespace}" + ) from exc except Exception: logger.exception( f"Failed to delete {len(ids)} vectors from namespace {namespace}" diff --git a/src/vector_store/utils.py b/src/vector_store/utils.py deleted file mode 100644 index ae613ada..00000000 --- a/src/vector_store/utils.py +++ /dev/null @@ -1,57 +0,0 @@ -""" -Vector store utility functions. -""" - -from __future__ import annotations - -import logging -from typing import TYPE_CHECKING - -from tenacity import ( - AsyncRetrying, - retry_if_exception_type, - stop_after_attempt, - wait_exponential, -) - -if TYPE_CHECKING: - from src.vector_store import VectorRecord, VectorStore, VectorUpsertResult - -logger = logging.getLogger(__name__) - - -async def upsert_with_retry( - vector_store: VectorStore, - namespace: str, - vector_records: list[VectorRecord], - max_attempts: int = 3, -) -> VectorUpsertResult | None: - """ - Upsert vectors with exponential backoff retry. - - Args: - vector_store: The vector store to upsert into - namespace: The namespace for the vectors - vector_records: List of VectorRecord objects to upsert - max_attempts: Maximum number of retry attempts (default 3) - - Returns: - VectorUpsertResult on success, or None if vector_records is empty - - Raises: - Exception: If all retries fail - """ - if not vector_records: - return None - - result: VectorUpsertResult | None = None - async for attempt in AsyncRetrying( - stop=stop_after_attempt(max_attempts), - wait=wait_exponential(multiplier=0.5, min=0.5, max=2.0), - retry=retry_if_exception_type(Exception), - reraise=True, - ): - with attempt: - result = await vector_store.upsert_many(namespace, vector_records) - - return result diff --git a/tests/conftest.py b/tests/conftest.py index cd739b7c..3c9b8e63 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -509,21 +509,18 @@ def mock_vector_store(request: pytest.FixtureRequest): from src.vector_store import ( VectorQueryResult, VectorRecord, - VectorUpsertResult, _hash_namespace_components, # pyright: ignore[reportPrivateUsage] ) # Create a mock vector store that stores vectors in memory vector_storage: dict[str, dict[str, tuple[list[float], dict[str, Any]]]] = {} - async def mock_upsert_many( - namespace: str, vectors: list[VectorRecord] - ) -> VectorUpsertResult: + async def mock_upsert_many(namespace: str, vectors: list[VectorRecord]) -> None: if namespace not in vector_storage: vector_storage[namespace] = {} for vector in vectors: vector_storage[namespace][vector.id] = (vector.embedding, vector.metadata) - return VectorUpsertResult(ok=True) + return async def mock_query( namespace: str, embedding: list[float], **kwargs: Any diff --git a/tests/deriver/test_vector_reconciliation.py b/tests/deriver/test_vector_reconciliation.py index 3b8e0e92..cc637597 100644 --- a/tests/deriver/test_vector_reconciliation.py +++ b/tests/deriver/test_vector_reconciliation.py @@ -27,7 +27,6 @@ from src.reconciler.sync_vectors import ( from src.vector_store import ( VectorRecord, VectorStore, - VectorUpsertResult, _hash_namespace_components, # pyright: ignore[reportPrivateUsage] ) @@ -84,9 +83,7 @@ class TestStateTransitions: mock_vector_store.get_vector_namespace = MagicMock( return_value=f"honcho.doc.{_hash_namespace_components(workspace.name, peer1.name, peer1.name)}" ) - mock_vector_store.upsert_many = AsyncMock( - return_value=VectorUpsertResult(ok=True) - ) + mock_vector_store.upsert_many = AsyncMock(return_value=None) # Run sync synced, failed = await _sync_documents(db_session, docs, mock_vector_store) @@ -309,13 +306,11 @@ class TestBatchProcessing: ) -> str: return f"honcho.doc.{_hash_namespace_components(workspace, observer, observed)}" - async def mock_upsert( - namespace: str, vectors: list[VectorRecord] - ) -> VectorUpsertResult: + async def mock_upsert(namespace: str, vectors: list[VectorRecord]) -> None: if namespace not in namespace_calls: namespace_calls[namespace] = [] namespace_calls[namespace].extend(vectors) - return VectorUpsertResult(ok=True) + return mock_vector_store.get_vector_namespace = mock_get_namespace mock_vector_store.upsert_many = mock_upsert @@ -440,9 +435,7 @@ class TestReEmbedding: mock_vector_store.get_vector_namespace = MagicMock( return_value=f"honcho.doc.{_hash_namespace_components(workspace.name, peer1.name, peer1.name)}" ) - mock_vector_store.upsert_many = AsyncMock( - return_value=VectorUpsertResult(ok=True) - ) + mock_vector_store.upsert_many = AsyncMock(return_value=None) # Run sync synced, failed = await _sync_documents(db_session, docs, mock_vector_store) @@ -512,9 +505,7 @@ class TestReEmbedding: mock_vector_store.get_vector_namespace = MagicMock( return_value=f"honcho.doc.{_hash_namespace_components(workspace.name, peer1.name, peer1.name)}" ) - mock_vector_store.upsert_many = AsyncMock( - return_value=VectorUpsertResult(ok=True) - ) + mock_vector_store.upsert_many = AsyncMock(return_value=None) # Run sync await _sync_documents(db_session, docs, mock_vector_store) diff --git a/tests/vector_store/test_turbopuffer.py b/tests/vector_store/test_turbopuffer.py new file mode 100644 index 00000000..73bf3c20 --- /dev/null +++ b/tests/vector_store/test_turbopuffer.py @@ -0,0 +1,81 @@ +"""Tests for TurbopufferVectorStore error handling on 5xx responses.""" + +from __future__ import annotations + +from unittest.mock import AsyncMock, MagicMock + +import httpx +import pytest +from turbopuffer import InternalServerError + +from src.config import settings +from src.exceptions import VectorStoreError +from src.vector_store import VectorRecord +from src.vector_store.turbopuffer import TurbopufferVectorStore + + +def _internal_server_error(status_code: int = 503) -> InternalServerError: + request = httpx.Request( + "POST", "https://api.turbopuffer.com/v2/namespaces/ns/write" + ) + response = httpx.Response(status_code, request=request) + return InternalServerError("turbopuffer unavailable", response=response, body=None) + + +@pytest.fixture +def store(monkeypatch: pytest.MonkeyPatch) -> TurbopufferVectorStore: + monkeypatch.setattr(settings.VECTOR_STORE, "TURBOPUFFER_API_KEY", "test-key") + monkeypatch.setattr(settings.VECTOR_STORE, "TURBOPUFFER_REGION", "gcp-us-east4") + return TurbopufferVectorStore() + + +@pytest.fixture +def record() -> VectorRecord: + return VectorRecord( + id="doc_1", embedding=[0.1, 0.2, 0.3, 0.4], metadata={"foo": "bar"} + ) + + +@pytest.mark.asyncio +async def test_upsert_many_raises_vector_store_error_on_5xx( + store: TurbopufferVectorStore, + record: VectorRecord, +) -> None: + namespace_mock = MagicMock() + namespace_mock.write = AsyncMock(side_effect=_internal_server_error(503)) + store._get_namespace = MagicMock(return_value=namespace_mock) # pyright: ignore[reportPrivateUsage] + + with pytest.raises(VectorStoreError) as excinfo: + await store.upsert_many("honcho.doc.test", [record]) + + assert "honcho.doc.test" in str(excinfo.value) + assert isinstance(excinfo.value.__cause__, InternalServerError) + namespace_mock.write.assert_awaited_once() + + +@pytest.mark.asyncio +async def test_upsert_many_short_circuits_on_empty( + store: TurbopufferVectorStore, +) -> None: + namespace_mock = MagicMock() + namespace_mock.write = AsyncMock() + store._get_namespace = MagicMock(return_value=namespace_mock) # pyright: ignore[reportPrivateUsage] + + await store.upsert_many("honcho.doc.test", []) + + namespace_mock.write.assert_not_awaited() + + +@pytest.mark.asyncio +async def test_upsert_many_succeeds_without_raising( + store: TurbopufferVectorStore, + record: VectorRecord, +) -> None: + namespace_mock = MagicMock() + namespace_mock.write = AsyncMock() + store._get_namespace = MagicMock(return_value=namespace_mock) # pyright: ignore[reportPrivateUsage] + + result = await store.upsert_many("honcho.doc.test", [record]) + + assert result is None + namespace_mock.write.assert_awaited_once()