mirror of
https://github.com/BerriAI/litellm.git
synced 2026-09-12 23:01:41 +00:00
fix(vector_stores): turn MongoDB's silent misconfiguration failures into errors
Driving the sad path against a live Atlas cluster showed four cases returning an empty result set instead of failing: a missing index, a missing database, a missing collection, and the async path for all three. $vectorSearch reports none of these as errors, so a misconfigured store looked exactly like a query that matched nothing, which is the worst shape for this to fail in. An empty result set is now checked against the index catalogue, which does report all three correctly, and a store that cannot work says so. The check costs one extra round trip and only on the empty path, so a search that returned hits is unaffected. Atlas also reports a wrong vector path and a dimension mismatch under the same error code. Both previously surfaced as "index not found", which sent the reader looking in the wrong place; they are now told apart and each names the setting that is actually wrong.
This commit is contained in:
parent
800cd17d17
commit
85bda43d63
3 changed files with 210 additions and 10 deletions
|
|
@ -105,6 +105,24 @@ def _index_hint(index_name: str, database: str, collection: str) -> str:
|
|||
)
|
||||
|
||||
|
||||
def missing_index_error(index_name: str, database: str, collection: str) -> ValueError:
|
||||
"""$vectorSearch against a missing index, database or collection returns zero documents
|
||||
instead of failing, so an empty result set is checked against the index catalogue and
|
||||
turned into this rather than being reported as 'no matches'."""
|
||||
return ValueError(
|
||||
f"{_index_hint(index_name, database, collection)} A vector search against a database, "
|
||||
"collection or index that does not exist returns no results rather than an error, so this "
|
||||
"was reported as an empty result set by MongoDB."
|
||||
)
|
||||
|
||||
|
||||
def index_not_ready_error(index_name: str, database: str, collection: str, status: str) -> ValueError:
|
||||
return ValueError(
|
||||
f"The Atlas Vector Search index '{index_name}' on '{database}.{collection}' is not queryable "
|
||||
f"yet; its status is {status}. Searches against it return no results until the build finishes."
|
||||
)
|
||||
|
||||
|
||||
def translate_mongo_error(error: Exception, index_name: str, database: str, collection: str) -> Exception:
|
||||
"""Turn a driver failure into a message that names the misconfiguration, never a silent empty result.
|
||||
|
||||
|
|
@ -136,14 +154,19 @@ def translate_mongo_error(error: Exception, index_name: str, database: str, coll
|
|||
f"lacks read access to '{database}.{collection}'. Driver detail: {error.details}"
|
||||
)
|
||||
detail: Final = str(error).lower()
|
||||
if "index" in detail and ("not found" in detail or "does not exist" in detail or "unknown" in detail):
|
||||
return ValueError(f"{_index_hint(index_name, database, collection)} Driver detail: {error}")
|
||||
if "dimension" in detail or "numdimensions" in detail or "queryvector" in detail:
|
||||
if "dimension" in detail:
|
||||
return ValueError(
|
||||
"The query embedding does not match the vector dimensions the Atlas index was built for. "
|
||||
"litellm_embedding_model must be the same model that produced the stored vectors. "
|
||||
f"Driver detail: {error}"
|
||||
)
|
||||
if "is not indexed as vector" in detail:
|
||||
return ValueError(
|
||||
"mongodb_embedding_field names a field the Atlas Vector Search index does not cover. "
|
||||
f"It must match the 'path' the index '{index_name}' was created on. Driver detail: {error}"
|
||||
)
|
||||
if "index" in detail and ("not found" in detail or "does not exist" in detail or "unknown" in detail):
|
||||
return ValueError(f"{_index_hint(index_name, database, collection)} Driver detail: {error}")
|
||||
return ValueError(
|
||||
f"MongoDB rejected the vector search against '{database}.{collection}' using index "
|
||||
f"'{index_name}'. Driver detail: {error}"
|
||||
|
|
|
|||
|
|
@ -26,6 +26,8 @@ from litellm.llms.mongodb.common_utils import (
|
|||
MongoClientKey,
|
||||
get_async_client,
|
||||
get_sync_client,
|
||||
index_not_ready_error,
|
||||
missing_index_error,
|
||||
translate_mongo_error,
|
||||
)
|
||||
from litellm.types.utils import EmbeddingResponse
|
||||
|
|
@ -249,6 +251,19 @@ class MongoDBVectorStoreConfig(BaseDirectVectorStoreConfig):
|
|||
data=[cls._to_result(document, text_field) for document in documents],
|
||||
)
|
||||
|
||||
@staticmethod
|
||||
def _raise_for_unusable_index(
|
||||
catalogue: Sequence[Mapping[str, object]], index_name: str, database: str, collection: str
|
||||
) -> None:
|
||||
"""An empty result set is ambiguous: Atlas returns zero documents both for a query that
|
||||
genuinely matched nothing and for a missing database, collection or index. Only the second
|
||||
is a misconfiguration, so the index catalogue decides which one happened."""
|
||||
if not catalogue:
|
||||
raise missing_index_error(index_name, database, collection)
|
||||
entry: Final = catalogue[0]
|
||||
if not entry.get("queryable"):
|
||||
raise index_not_ready_error(index_name, database, collection, str(entry.get("status") or "unknown"))
|
||||
|
||||
@staticmethod
|
||||
def _embedding_vector(embedding_response: EmbeddingResponse) -> Sequence[float]:
|
||||
data: Final = embedding_response.data
|
||||
|
|
@ -284,12 +299,21 @@ class MongoDBVectorStoreConfig(BaseDirectVectorStoreConfig):
|
|||
)
|
||||
|
||||
client: Final = self.sync_client_factory(key)
|
||||
target: Final = client[database][collection] # pyright: ignore[reportIndexIssue] # factory is typed as returning object so injected doubles are accepted
|
||||
try:
|
||||
documents: Final = list(client[database][collection].aggregate(pipeline)) # pyright: ignore[reportIndexIssue] # factory is typed as returning object so injected doubles are accepted
|
||||
documents: Final = list(target.aggregate(pipeline))
|
||||
except Exception as e:
|
||||
raise translate_mongo_error(
|
||||
e, index_name=vector_store_id, database=database, collection=collection
|
||||
) from e
|
||||
if not documents:
|
||||
try:
|
||||
catalogue: Final = list(target.list_search_indexes(vector_store_id))
|
||||
except Exception as e:
|
||||
raise translate_mongo_error(
|
||||
e, index_name=vector_store_id, database=database, collection=collection
|
||||
) from e
|
||||
self._raise_for_unusable_index(catalogue, vector_store_id, database, collection)
|
||||
return self._to_response(documents, query_text, params.text_field)
|
||||
|
||||
async def aexecute_search_vector_store_request(
|
||||
|
|
@ -317,13 +341,23 @@ class MongoDBVectorStoreConfig(BaseDirectVectorStoreConfig):
|
|||
)
|
||||
|
||||
client: Final = self.async_client_factory(key)
|
||||
target: Final = client[database][collection] # pyright: ignore[reportIndexIssue] # factory is typed as returning object so injected doubles are accepted
|
||||
try:
|
||||
cursor: Final = await client[database][collection].aggregate(pipeline) # pyright: ignore[reportIndexIssue] # factory is typed as returning object so injected doubles are accepted
|
||||
cursor: Final = await target.aggregate(pipeline)
|
||||
documents: Final = [document async for document in cursor]
|
||||
except Exception as e:
|
||||
raise translate_mongo_error(
|
||||
e, index_name=vector_store_id, database=database, collection=collection
|
||||
) from e
|
||||
if not documents:
|
||||
try:
|
||||
index_cursor: Final = await target.list_search_indexes(vector_store_id)
|
||||
catalogue: Final = [entry async for entry in index_cursor]
|
||||
except Exception as e:
|
||||
raise translate_mongo_error(
|
||||
e, index_name=vector_store_id, database=database, collection=collection
|
||||
) from e
|
||||
self._raise_for_unusable_index(catalogue, vector_store_id, database, collection)
|
||||
return self._to_response(documents, query_text, params.text_field)
|
||||
|
||||
def transform_create_vector_store_request(
|
||||
|
|
|
|||
|
|
@ -30,11 +30,16 @@ BASE_PARAMS = {
|
|||
}
|
||||
|
||||
|
||||
READY_INDEX = [{"name": INDEX, "status": "READY", "queryable": True}]
|
||||
|
||||
|
||||
class FakeCollection:
|
||||
def __init__(self, documents, error=None):
|
||||
def __init__(self, documents, error=None, search_indexes=None):
|
||||
self.documents = documents
|
||||
self.error = error
|
||||
self.search_indexes = READY_INDEX if search_indexes is None else search_indexes
|
||||
self.pipeline = None
|
||||
self.listed_indexes = []
|
||||
|
||||
def aggregate(self, pipeline):
|
||||
self.pipeline = pipeline
|
||||
|
|
@ -42,6 +47,10 @@ class FakeCollection:
|
|||
raise self.error
|
||||
return iter(self.documents)
|
||||
|
||||
def list_search_indexes(self, name):
|
||||
self.listed_indexes.append(name)
|
||||
return iter(self.search_indexes)
|
||||
|
||||
|
||||
class FakeAsyncCollection(FakeCollection):
|
||||
async def aggregate(self, pipeline):
|
||||
|
|
@ -55,6 +64,15 @@ class FakeAsyncCollection(FakeCollection):
|
|||
|
||||
return cursor()
|
||||
|
||||
async def list_search_indexes(self, name):
|
||||
self.listed_indexes.append(name)
|
||||
|
||||
async def cursor():
|
||||
for entry in self.search_indexes:
|
||||
yield entry
|
||||
|
||||
return cursor()
|
||||
|
||||
|
||||
class FakeDatabase:
|
||||
def __init__(self, collection):
|
||||
|
|
@ -92,8 +110,8 @@ class FakeAsyncEmbeddingFn(FakeEmbeddingFn):
|
|||
return SimpleNamespace(data=[{"embedding": self.embedding}] if self.embedding is not None else [])
|
||||
|
||||
|
||||
def _config(documents=(), embedding=(0.1, 0.2, 0.3), error=None):
|
||||
collection = FakeCollection(list(documents), error)
|
||||
def _config(documents=(), embedding=(0.1, 0.2, 0.3), error=None, search_indexes=None):
|
||||
collection = FakeCollection(list(documents), error, search_indexes)
|
||||
client = FakeClient(collection)
|
||||
config = MongoDBVectorStoreConfig(
|
||||
embedding_fn=FakeEmbeddingFn(list(embedding) if embedding is not None else None),
|
||||
|
|
@ -102,8 +120,8 @@ def _config(documents=(), embedding=(0.1, 0.2, 0.3), error=None):
|
|||
return config, client, collection
|
||||
|
||||
|
||||
def _async_config(documents=(), embedding=(0.1, 0.2, 0.3), error=None):
|
||||
collection = FakeAsyncCollection(list(documents), error)
|
||||
def _async_config(documents=(), embedding=(0.1, 0.2, 0.3), error=None, search_indexes=None):
|
||||
collection = FakeAsyncCollection(list(documents), error, search_indexes)
|
||||
client = FakeClient(collection)
|
||||
config = MongoDBVectorStoreConfig(
|
||||
aembedding_fn=FakeAsyncEmbeddingFn(list(embedding) if embedding is not None else None),
|
||||
|
|
@ -642,3 +660,128 @@ class TestMissingDriver:
|
|||
|
||||
with patch.dict(sys.modules, {"pymongo.errors": None}):
|
||||
assert translate_mongo_error(original, INDEX, "db", "col") is original
|
||||
|
||||
|
||||
class TestEmptyResultsAreDisambiguated:
|
||||
"""$vectorSearch returns zero documents for a missing database, collection or index just as it
|
||||
does for a query that matched nothing, so an empty result set is checked against the index
|
||||
catalogue before it is reported as 'no matches'."""
|
||||
|
||||
def test_a_missing_index_becomes_an_error_rather_than_an_empty_page(self):
|
||||
config, _, collection = _config(documents=[], search_indexes=[])
|
||||
|
||||
with pytest.raises(ValueError, match="No queryable Atlas Vector Search index"):
|
||||
_search(config)
|
||||
|
||||
assert collection.listed_indexes == [INDEX]
|
||||
|
||||
def test_the_missing_index_error_explains_why_mongodb_reported_no_results(self):
|
||||
config, _, _ = _config(documents=[], search_indexes=[])
|
||||
|
||||
with pytest.raises(ValueError, match="returns no results rather than an error"):
|
||||
_search(config)
|
||||
|
||||
def test_an_index_still_building_becomes_an_error_naming_its_status(self):
|
||||
config, _, _ = _config(
|
||||
documents=[], search_indexes=[{"name": INDEX, "status": "PENDING", "queryable": False}]
|
||||
)
|
||||
|
||||
with pytest.raises(ValueError, match="not queryable yet; its status is PENDING"):
|
||||
_search(config)
|
||||
|
||||
def test_a_genuine_no_match_against_a_ready_index_returns_an_empty_page(self):
|
||||
config, _, collection = _config(documents=[])
|
||||
|
||||
response = _search(config)
|
||||
|
||||
assert response["data"] == []
|
||||
assert response["object"] == "vector_store.search_results.page"
|
||||
assert collection.listed_indexes == [INDEX]
|
||||
|
||||
def test_the_catalogue_is_not_consulted_when_the_search_returned_hits(self):
|
||||
config, _, collection = _config(documents=[{"_id": 1, "text": "hit", "score": 0.9}])
|
||||
|
||||
_search(config)
|
||||
|
||||
assert collection.listed_indexes == []
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_async_missing_index_becomes_an_error_rather_than_an_empty_page(self):
|
||||
config, _, collection = _async_config(documents=[], search_indexes=[])
|
||||
|
||||
with pytest.raises(ValueError, match="No queryable Atlas Vector Search index"):
|
||||
await _asearch(config)
|
||||
|
||||
assert collection.listed_indexes == [INDEX]
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_async_index_still_building_becomes_an_error_naming_its_status(self):
|
||||
config, _, _ = _async_config(
|
||||
documents=[], search_indexes=[{"name": INDEX, "status": "PENDING", "queryable": False}]
|
||||
)
|
||||
|
||||
with pytest.raises(ValueError, match="not queryable yet; its status is PENDING"):
|
||||
await _asearch(config)
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_async_genuine_no_match_returns_an_empty_page(self):
|
||||
config, _, _ = _async_config(documents=[])
|
||||
|
||||
response = await _asearch(config)
|
||||
|
||||
assert response["data"] == []
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_async_catalogue_is_not_consulted_when_the_search_returned_hits(self):
|
||||
config, _, collection = _async_config(documents=[{"_id": 1, "text": "hit", "score": 0.9}])
|
||||
|
||||
await _asearch(config)
|
||||
|
||||
assert collection.listed_indexes == []
|
||||
|
||||
def test_a_failure_while_checking_the_catalogue_is_translated_too(self):
|
||||
from pymongo.errors import OperationFailure
|
||||
|
||||
class ExplodingCollection(FakeCollection):
|
||||
def list_search_indexes(self, name):
|
||||
raise OperationFailure("not authorized", code=13)
|
||||
|
||||
collection = ExplodingCollection([], None, [])
|
||||
config = MongoDBVectorStoreConfig(
|
||||
embedding_fn=FakeEmbeddingFn([0.1]),
|
||||
sync_client_factory=lambda key: FakeClient(collection),
|
||||
)
|
||||
|
||||
with pytest.raises(ValueError, match="lacks read access"):
|
||||
_search(config)
|
||||
|
||||
|
||||
class TestAtlasPlanExecutorErrors:
|
||||
"""Atlas reports a wrong vector path and a dimension mismatch through the same error code, so
|
||||
each one has to be told apart by its message or both come back as a generic index failure."""
|
||||
|
||||
def _translate(self, message):
|
||||
from pymongo.errors import OperationFailure
|
||||
|
||||
return translate_mongo_error(
|
||||
OperationFailure(message, code=8),
|
||||
index_name=INDEX,
|
||||
database="sample_mflix",
|
||||
collection="embedded_movies",
|
||||
)
|
||||
|
||||
def test_a_wrong_vector_path_points_at_the_embedding_field_setting(self):
|
||||
translated = self._translate(
|
||||
"PlanExecutor error during aggregation :: caused by :: nope is not indexed as vector"
|
||||
)
|
||||
|
||||
assert "mongodb_embedding_field names a field" in str(translated)
|
||||
|
||||
def test_a_dimension_mismatch_is_not_reported_as_a_wrong_path(self):
|
||||
translated = self._translate(
|
||||
"PlanExecutor error during aggregation :: caused by :: vector field is indexed with "
|
||||
"1536 dimensions but queried with 3072"
|
||||
)
|
||||
|
||||
assert "does not match the vector dimensions" in str(translated)
|
||||
assert "mongodb_embedding_field" not in str(translated)
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue