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>
This commit is contained in:
Rajat Ahuja 2026-04-20 16:56:51 -04:00 committed by GitHub
parent 1c3e3f8816
commit 7fae16b351
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194
10 changed files with 211 additions and 128 deletions

View File

@ -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)

View File

@ -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.",

View File

@ -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)

View File

@ -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",
]

View File

@ -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}"

View File

@ -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}"

View File

@ -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

View File

@ -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

View File

@ -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)

View File

@ -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()