Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
30 changes: 30 additions & 0 deletions src/everos/core/persistence/lancedb/repository.py
Original file line number Diff line number Diff line change
Expand Up @@ -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.

Expand Down Expand Up @@ -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]:
Expand Down
Original file line number Diff line number Diff line change
@@ -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")
Loading