diff --git a/src/everos/core/persistence/lancedb/repository.py b/src/everos/core/persistence/lancedb/repository.py index de79702c0..0fd59fe81 100644 --- a/src/everos/core/persistence/lancedb/repository.py +++ b/src/everos/core/persistence/lancedb/repository.py @@ -208,6 +208,23 @@ def _remove_empty_index_dirs( return removed +_TRANSIENT_EXECUTION_MARKERS = ("Spill has sent an error",) +"""Substrings of lance error messages that name a query-execution failure a +retry clears. ``LanceError(IO): Execution error: Spill has sent an error`` is +DataFusion's sort / merge spill to the OS temp dir failing mid-query; lancedb +raises it as a bare ``RuntimeError``. Seen only on the Windows soak box under +nine concurrent clients (187 / 180 / 34 times over three runs), never on an +idle box, and the same row projected fine on the next attempt — yet the worker +filed every one as unrecoverable, so ~200 md files per run needed a manual +``cascade fix``. Match the exact phrase: a generic IO error (disk full, file +gone) must stay permanent.""" + + +def _is_transient_execution_error(exc: BaseException) -> bool: + text = str(exc) + return any(marker in text for marker in _TRANSIENT_EXECUTION_MARKERS) + + class LanceRepoBase[T: BaseLanceTable]: """Generic CRUD repository for one LanceDB table. @@ -296,6 +313,19 @@ async def _deadline(self, budget: float, op: str) -> AsyncIterator[None]: raise VectorStoreBusyError( f"{op} on table {self.table_name!r} exceeded its {budget:g}s deadline" ) from exc + except RuntimeError as exc: + if not _is_transient_execution_error(exc): + raise + logger.warning( + "lancedb_transient_execution_error", + table=self.table_name, + op=op, + error=str(exc)[:200], + ) + raise VectorStoreBusyError( + f"{op} on table {self.table_name!r} hit a transient lance " + f"execution error: {exc}" + ) from exc @asynccontextmanager async def _locked(self, budget: float, op: str) -> AsyncIterator[None]: diff --git a/tests/unit/test_core/test_persistence/test_lancedb/test_transient_execution_errors.py b/tests/unit/test_core/test_persistence/test_lancedb/test_transient_execution_errors.py new file mode 100644 index 000000000..800241968 --- /dev/null +++ b/tests/unit/test_core/test_persistence/test_lancedb/test_transient_execution_errors.py @@ -0,0 +1,47 @@ +"""A lance query-execution failure that a retry clears must reach the cascade +worker as :class:`VectorStoreBusyError` (retried with backoff), not as a bare +``RuntimeError`` (filed as unrecoverable, needing a manual ``cascade fix``). +""" + +from __future__ import annotations + +from typing import ClassVar + +import pytest + +from everos.core.errors import VectorStoreBusyError +from everos.core.persistence.lancedb import BaseLanceTable, LanceRepoBase +from everos.core.persistence.lancedb import repository as repo_mod + +_SPILL = ( + "lance error: LanceError(IO): Execution error: Spill has sent an error, " + "C:\\Users\\x\\everos-src\\.venv\\Lib\\site-packages\\lance\\..." +) + + +class _Note(BaseLanceTable): + TABLE_NAME: ClassVar[str] = "_note" + id: str + + +class _Repo(LanceRepoBase[_Note]): + schema = _Note + + +def test_only_the_spill_phrase_counts_as_transient() -> None: + assert repo_mod._is_transient_execution_error(RuntimeError(_SPILL)) + assert not repo_mod._is_transient_execution_error( + RuntimeError("lance error: LanceError(IO): No space left on device") + ) + + +async def test_spill_failure_inside_a_read_becomes_a_busy_error() -> None: + with pytest.raises(VectorStoreBusyError, match="transient lance execution"): + async with _Repo()._deadline(1.0, "find_where"): + raise RuntimeError(_SPILL) + + +async def test_other_runtime_errors_still_propagate_unchanged() -> None: + with pytest.raises(RuntimeError, match="No space left"): + async with _Repo()._deadline(1.0, "find_where"): + raise RuntimeError("lance error: LanceError(IO): No space left on device")