From 2a0545ec7fab976788767d2b91f6f655c0ce771e Mon Sep 17 00:00:00 2001 From: zhanghui Date: Thu, 24 Sep 2026 15:16:26 +0800 Subject: [PATCH 1/5] feat(lancedb): build an IVF_FLAT index on vector columns MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit EverOS built FTS indexes on its LanceDB tables and nothing on the vector columns, so every `nearest_to` was a brute-force scan of the whole column: linear in rows and bytes. Measured on the Windows soak box (1024-dim float32): 10k rows = 42 MB = 155 ms per scan, 27k rows = 112 MB = 590 ms; a hybrid search issues two or three scans, which is where its 0.5–1.6 s idle latency went and why it climbed through the 10-hour soak. `BaseLanceTable.ensure_vector_indexes` creates an IVF_FLAT (cosine, to match `dense_search`) index on each vector column once it holds `[lancedb] vector_index_min_rows` (default 2000) non-null vectors; all-null columns (a Tier 1 store) and small stores are left alone. It runs at startup with the FTS pass, on the cascade's heavy maintenance beat (so a table that crosses the threshold while the server runs gets indexed), and with `replace=True` on the rebuild cadence. `rebuild_indexes` used to drop every index not on a BM25 column — it now keeps the vector indexes it would otherwise have removed every 12 h. LanceDB's `optimize()` merges new rows into an existing index, so the unindexed tail stays small between beats. Tests: columns detected from the Arrow schema; nothing below the threshold or on all-null vectors; IvfFlat built at the threshold; idempotent unless replaced; indexed top-5 equals the brute-force top-5. Co-Authored-By: Claude Fable 5.1 --- src/everos/config/default.toml | 3 + src/everos/config/settings.py | 7 + src/everos/core/persistence/lancedb/base.py | 46 ++++++- .../core/persistence/lancedb/repository.py | 26 +++- .../infra/persistence/lancedb/__init__.py | 3 + src/everos/memory/cascade/worker.py | 6 + .../test_lancedb/test_vector_index.py | 125 ++++++++++++++++++ 7 files changed, 214 insertions(+), 2 deletions(-) create mode 100644 tests/unit/test_core/test_persistence/test_lancedb/test_vector_index.py diff --git a/src/everos/config/default.toml b/src/everos/config/default.toml index 3dceee2f4..092d18d0e 100644 --- a/src/everos/config/default.toml +++ b/src/everos/config/default.toml @@ -49,6 +49,9 @@ cache_size_kb = 2048 # >0 -> eventual (interval seconds between checks) # Uncomment to override: # read_consistency_seconds = 5.0 +# Rows (with a non-null vector) before a table's vector columns get an +# IVF_FLAT index. Below it every vector query scans the whole column. +# vector_index_min_rows = 2000 [index] # Rebuildable vector/BM25 index. Markdown remains the source of truth. diff --git a/src/everos/config/settings.py b/src/everos/config/settings.py index fc1be21ad..cde390ac5 100644 --- a/src/everos/config/settings.py +++ b/src/everos/config/settings.py @@ -518,6 +518,13 @@ class LanceDBSettings(BaseModel): read_consistency_seconds: float | None = None index_cache_size_bytes: int = 16 * 1024 * 1024 + vector_index_min_rows: int = Field(default=2000, ge=1) + """Rows (with a non-null vector) a table needs before its vector columns + get an ANN index. Below this a brute-force scan is cheaper than the + index; above it the scan grows linearly with the table — 27k rows of + 1024-dim vectors was 112 MB and ~0.6 s per query on a laptop SSD, and + a hybrid search runs two or three of them. Checked at startup, on the + cascade's heavy maintenance beat, and rebuilt on the rebuild cadence.""" class CascadeSettings(BaseModel): diff --git a/src/everos/core/persistence/lancedb/base.py b/src/everos/core/persistence/lancedb/base.py index 09f685d61..87389b471 100644 --- a/src/everos/core/persistence/lancedb/base.py +++ b/src/everos/core/persistence/lancedb/base.py @@ -34,7 +34,7 @@ class Episode(BaseLanceTable): import pyarrow as pa from lancedb import AsyncTable -from lancedb.index import FTS +from lancedb.index import FTS, IvfFlat from lancedb.pydantic import LanceModel from pydantic import Field @@ -177,6 +177,50 @@ async def ensure_fts_indexes( ), ) + @classmethod + def vector_columns(cls) -> list[str]: + """Names of the schema's vector columns (Arrow fixed-size lists).""" + return [ + field.name + for field in cls.to_arrow_schema() + if pa.types.is_fixed_size_list(field.type) + ] + + @classmethod + async def ensure_vector_indexes( + cls, table: AsyncTable, *, min_rows: int, replace: bool = False + ) -> list[str]: + """Create an IVF_FLAT (cosine) index on each vector column that has + at least ``min_rows`` non-null vectors; return the columns indexed. + + Without an index LanceDB answers ``nearest_to`` with a brute-force + scan of the whole column — linear in rows and in bytes (27k rows of + 1024-dim float32 is 112 MB and ~0.6 s per query on a laptop SSD), + and a hybrid search issues two or three of them. IVF_FLAT keeps + exact distances inside the probed partitions, so at these sizes the + recall cost is small; ``cosine`` matches the query side + (``distance_type("cosine")`` in ``dense_search``). Columns that are + still all-null (a Tier 1 store has no embeddings) are skipped: there + is nothing to train on. Idempotent unless ``replace``; the rebuild + cadence passes ``replace=True`` to retrain in place. + """ + columns = cls.vector_columns() + if not columns: + return [] + indices = await table.list_indices() + indexed = {col for idx in indices for col in (idx.columns or [])} + built: list[str] = [] + for column in columns: + if column in indexed and not replace: + continue + if await table.count_rows(f"{column} IS NOT NULL") < min_rows: + continue + await table.create_index( + column, replace=replace, config=IvfFlat(distance_type="cosine") + ) + built.append(column) + return built + def touch(record: BaseLanceTable) -> BaseLanceTable: """Set ``record.updated_at = now`` and return the record (chainable).""" diff --git a/src/everos/core/persistence/lancedb/repository.py b/src/everos/core/persistence/lancedb/repository.py index de79702c0..331b539e3 100644 --- a/src/everos/core/persistence/lancedb/repository.py +++ b/src/everos/core/persistence/lancedb/repository.py @@ -52,6 +52,16 @@ for a wedged table.""" _REBUILD_TIMEOUT_SECONDS = 300.0 + + +def _vector_index_min_rows() -> int: + """``[lancedb] vector_index_min_rows`` — imported lazily so this module + stays free of the settings import at load time (same as MemoryRoot).""" + from everos.config.settings import load_settings + + return load_settings().lancedb.vector_index_min_rows + + """Index rebuild (drop + recreate every index) — the one genuinely slow critical section, measured at ~0.3s per 50k rows per indexed column, so 5 minutes covers a multi-million-row table with wide headroom.""" @@ -639,14 +649,28 @@ async def rebuild_indexes(self) -> None: # ``return_exceptions``, so the whole search request 500s. Only # indexes on columns that are no longer indexed at all get dropped; # nothing queries those, so their drop opens no window. - wanted = set(self.schema.BM25_FIELDS or ()) + wanted = set(self.schema.BM25_FIELDS or ()) | set( + self.schema.vector_columns() + ) for idx in await table.list_indices(): if not wanted.intersection(idx.columns or ()): await table.drop_index(idx.name) await self.schema.ensure_fts_indexes(table, replace=True) + await self.schema.ensure_vector_indexes( + table, min_rows=_vector_index_min_rows(), replace=True + ) # ── Read ─────────────────────────────────────────────────────────────── + async def ensure_vector_indexes(self) -> list[str]: + """Build the ANN index on vector columns that crossed the row + threshold since startup — the cascade's heavy beat calls this.""" + async with self._locked(_REBUILD_TIMEOUT_SECONDS, "ensure_vector_indexes"): + table = await self._table() + return await self.schema.ensure_vector_indexes( + table, min_rows=_vector_index_min_rows() + ) + async def count(self) -> int: """Total row count.""" async with self._deadline(_READ_TIMEOUT_SECONDS, "count"): diff --git a/src/everos/infra/persistence/lancedb/__init__.py b/src/everos/infra/persistence/lancedb/__init__.py index 78ea873ed..0a48ba232 100644 --- a/src/everos/infra/persistence/lancedb/__init__.py +++ b/src/everos/infra/persistence/lancedb/__init__.py @@ -27,6 +27,7 @@ import contextlib import datetime as dt +from everos.config.settings import load_settings from everos.core.observability.logging import get_logger from everos.core.persistence import BaseLanceTable, MemoryRoot, memory_root_lock @@ -292,9 +293,11 @@ async def ensure_business_indexes() -> None: """ await migrate_table_schemas() await migrate_fts_indexes() + min_rows = load_settings().lancedb.vector_index_min_rows for schema in _BUSINESS_SCHEMAS: table = await get_table(schema.TABLE_NAME, schema) await schema.ensure_fts_indexes(table) + await schema.ensure_vector_indexes(table, min_rows=min_rows) async def verify_business_schemas() -> None: diff --git a/src/everos/memory/cascade/worker.py b/src/everos/memory/cascade/worker.py index 7ab050a05..d122fd83e 100644 --- a/src/everos/memory/cascade/worker.py +++ b/src/everos/memory/cascade/worker.py @@ -1027,6 +1027,12 @@ async def _run_optimize_once(self, kind: str) -> None: ) if state is not None: state.last_prune_at = now + # A table that crossed the ANN threshold since startup gets its + # vector index on the same heavy beat; a no-op otherwise. + ensure_vectors = getattr(repo, "ensure_vector_indexes", None) + if ensure_vectors is not None: + async with asyncio.timeout(_MAINTENANCE_TASK_TIMEOUT_SECONDS): + await ensure_vectors() else: # Light beat: lock-free compaction. A commit conflict here # is benign — handled below. diff --git a/tests/unit/test_core/test_persistence/test_lancedb/test_vector_index.py b/tests/unit/test_core/test_persistence/test_lancedb/test_vector_index.py new file mode 100644 index 000000000..5de64e1c8 --- /dev/null +++ b/tests/unit/test_core/test_persistence/test_lancedb/test_vector_index.py @@ -0,0 +1,125 @@ +"""``BaseLanceTable.ensure_vector_indexes`` — the ANN index on vector columns. + +Without it every ``nearest_to`` is a brute-force scan of the column, linear +in rows; the soak measured 0.6 s per scan at 27k rows of 1024-dim vectors. +The index is built only once a column holds ``min_rows`` non-null vectors, +so a Tier 1 store (no embeddings) and a small store are left alone. +""" + +from __future__ import annotations + +import random +from collections.abc import AsyncIterator +from pathlib import Path +from typing import ClassVar + +import lancedb +import pytest +from lancedb import AsyncTable +from lancedb.pydantic import Vector + +from everos.core.persistence.lancedb import BaseLanceTable + +_DIM = 8 + + +class _VecSpec(BaseLanceTable): + TABLE_NAME: ClassVar[str] = "vec_probe" + BM25_FIELDS: ClassVar[list[str]] = ["body"] + + id: str + body: str + vector: Vector(_DIM) | None = None # type: ignore[valid-type] + + +@pytest.fixture +async def vec_table(tmp_path: Path) -> AsyncIterator[AsyncTable]: + conn = await lancedb.connect_async(str(tmp_path / "lancedb")) + yield await conn.create_table(_VecSpec.TABLE_NAME, schema=_VecSpec) + + +def _rows(n: int, *, with_vectors: bool = True) -> list[_VecSpec]: + rng = random.Random(7) + return [ + _VecSpec( + id=f"r{i}", + body=f"row {i}", + vector=[rng.random() for _ in range(_DIM)] if with_vectors else None, + ) + for i in range(n) + ] + + +async def _vector_indices(table: AsyncTable) -> list[tuple[str, str]]: + return [ + (idx.columns[0], idx.index_type) + for idx in await table.list_indices() + if idx.columns and idx.columns[0] == "vector" + ] + + +def test_vector_columns_are_the_fixed_size_list_fields() -> None: + assert _VecSpec.vector_columns() == ["vector"] + + +async def test_below_the_row_threshold_no_index_is_built( + vec_table: AsyncTable, +) -> None: + await vec_table.add(_rows(40)) + assert await _VecSpec.ensure_vector_indexes(vec_table, min_rows=50) == [] + assert await _vector_indices(vec_table) == [] + + +async def test_at_the_threshold_an_ivf_flat_index_is_built( + vec_table: AsyncTable, +) -> None: + await vec_table.add(_rows(60)) + built = await _VecSpec.ensure_vector_indexes(vec_table, min_rows=50) + assert built == ["vector"] + assert await _vector_indices(vec_table) == [("vector", "IvfFlat")] + + +async def test_ensure_is_idempotent_unless_replace(vec_table: AsyncTable) -> None: + await vec_table.add(_rows(60)) + await _VecSpec.ensure_vector_indexes(vec_table, min_rows=50) + assert await _VecSpec.ensure_vector_indexes(vec_table, min_rows=50) == [] + rebuilt = await _VecSpec.ensure_vector_indexes(vec_table, min_rows=50, replace=True) + assert rebuilt == ["vector"] + assert len(await _vector_indices(vec_table)) == 1 + + +async def test_all_null_vectors_are_skipped_even_above_the_threshold( + vec_table: AsyncTable, +) -> None: + """A Tier 1 store has rows but no embeddings; there is nothing to train.""" + await vec_table.add(_rows(60, with_vectors=False)) + assert await _VecSpec.ensure_vector_indexes(vec_table, min_rows=50) == [] + assert await _vector_indices(vec_table) == [] + + +async def test_indexed_search_returns_the_same_neighbours_as_the_scan( + vec_table: AsyncTable, +) -> None: + """IVF_FLAT keeps exact distances inside the probed partitions: at this + size the top-5 must match the brute-force answer.""" + rows = _rows(300) + await vec_table.add(rows) + q = rows[17].vector + + async def _top5(bypass: bool) -> list[str]: + query = ( + vec_table.query() + .nearest_to(q) + .column("vector") + .distance_type("cosine") + .limit(5) + ) + if bypass: + query = query.bypass_vector_index() + return [r["id"] for r in await query.to_list()] + + flat = await _top5(bypass=True) + await _VecSpec.ensure_vector_indexes(vec_table, min_rows=50) + indexed = await _top5(bypass=False) + assert indexed[0] == "r17" + assert set(indexed) == set(flat) From c966667a51dc5ad713f286c5b528692d479a97c6 Mon Sep 17 00:00:00 2001 From: zhanghui Date: Thu, 24 Sep 2026 16:06:28 +0800 Subject: [PATCH 2/5] fix(lancedb): collapse vector index deltas on the heavy beat Review of the first version (adversarial subagent, probes on lancedb 0.34) found two blockers: * Every optimize() on a table with new rows appends a delta to the IVF index and never merges it (num_indices +1 per light beat, confirmed locally). A query probes every delta, so latency climbed with the beats since the last rebuild: 27k x 1024 rows measured 5.9 ms at 0 deltas, 25.7 ms at 100, 303 ms at 400 -- worse than the 24 ms scan the index replaces. Only create_index(replace=True) collapses them (optimize(retrain=True) does not); ensure_vector_indexes now retrains a column whose index has num_indices > 1, and the heavy beat (300 s) calls it, so at most ~30 deltas accumulate under sustained writes. * The heavy-beat hook was dead code: the worker holds the routed repository, which did not forward ensure_vector_indexes, and the getattr guard skipped silently. The method is now on the IndexRepository protocol, forwarded by the router and the LanceDB backend, a no-op on Milvus, and the worker calls it unguarded. Also: pin the IVF partition size (4096 rows) and nprobes (32) as one pair of constants used by every nearest_to, so the search is exact up to ~130k rows instead of depending on lance's changing defaults; drop the replace= parameter (the conditional covers the rebuild cadence and removes the retrain-at-startup the review also flagged); fix the read-timeout docstring that still said no ANN index is built. Tests: delta collapse (real LanceDB; precondition asserts one delta per beat), nprobes accepted on an unindexed column, every repository class exposes the maintenance methods and the router forwards them, the heavy beat calls ensure_vector_indexes. Each fails under its own mutation. Co-Authored-By: Claude Fable 5.1 --- src/everos/core/persistence/lancedb/base.py | 76 ++++++++++++++----- .../core/persistence/lancedb/repository.py | 11 ++- .../infra/persistence/backends/lancedb.py | 5 ++ .../infra/persistence/index/protocols.py | 2 + src/everos/infra/persistence/index/router.py | 3 + .../persistence/lancedb/repos/agent_skill.py | 2 + .../infra/persistence/milvus/repository.py | 3 + src/everos/memory/cascade/worker.py | 11 ++- .../test_lancedb/test_vector_index.py | 47 +++++++++++- .../test_index_router_maintenance.py | 60 +++++++++++++++ .../test_memory/test_cascade/test_worker.py | 25 ++++++ 11 files changed, 213 insertions(+), 32 deletions(-) create mode 100644 tests/unit/test_infra/test_index_router_maintenance.py diff --git a/src/everos/core/persistence/lancedb/base.py b/src/everos/core/persistence/lancedb/base.py index 87389b471..3380182ee 100644 --- a/src/everos/core/persistence/lancedb/base.py +++ b/src/everos/core/persistence/lancedb/base.py @@ -40,6 +40,16 @@ class Episode(BaseLanceTable): from everos.component.utils.datetime import get_utc_now +VECTOR_INDEX_ROWS_PER_PARTITION = 4096 +"""IVF partition size at build time. Pinned so :data:`VECTOR_QUERY_NPROBES` +means something: lance's default partition count has changed across releases.""" + +VECTOR_QUERY_NPROBES = 32 +"""Partitions probed per vector query. 32 partitions of 4096 rows cover the whole +column, i.e. exact search, up to ~130k rows; past that the unprobed partitions +are skipped and recall degrades gradually. Every ``nearest_to`` in the tree sets +it, so the index and the query agree on what "exact" costs.""" + class BaseLanceTable(LanceModel): """Pydantic / LanceDB base with ``created_at`` / ``updated_at`` and @@ -188,38 +198,68 @@ def vector_columns(cls) -> list[str]: @classmethod async def ensure_vector_indexes( - cls, table: AsyncTable, *, min_rows: int, replace: bool = False + cls, table: AsyncTable, *, min_rows: int ) -> list[str]: - """Create an IVF_FLAT (cosine) index on each vector column that has - at least ``min_rows`` non-null vectors; return the columns indexed. + """Keep one IVF_FLAT (cosine) index per vector column that holds at + least ``min_rows`` non-null vectors; return the columns touched. Without an index LanceDB answers ``nearest_to`` with a brute-force scan of the whole column — linear in rows and in bytes (27k rows of 1024-dim float32 is 112 MB and ~0.6 s per query on a laptop SSD), and a hybrid search issues two or three of them. IVF_FLAT keeps - exact distances inside the probed partitions, so at these sizes the - recall cost is small; ``cosine`` matches the query side - (``distance_type("cosine")`` in ``dense_search``). Columns that are - still all-null (a Tier 1 store has no embeddings) are skipped: there - is nothing to train on. Idempotent unless ``replace``; the rebuild - cadence passes ``replace=True`` to retrain in place. + exact distances inside the probed partitions; with + :data:`VECTOR_INDEX_ROWS_PER_PARTITION` rows per partition and + :data:`VECTOR_QUERY_NPROBES` probes the search stays exact up to + ~130k rows. ``cosine`` matches the query side. Columns that are + still all-null (a Tier 1 store has no embeddings) are skipped: + there is nothing to train on. + + Two cases do work; everything else is a no-op: + + * no index yet and the column crossed ``min_rows`` -> build one; + * the index has grown delta indices -> retrain it in place. Every + ``optimize()`` on a table with new rows appends one *delta* index + instead of merging (``num_indices`` +1 per light beat, never + collapsing on its own), and a query probes every delta, so latency + climbs with the beats since the last rebuild: 27k x 1024 rows + measured 5.9 ms at 0 deltas, 25.7 ms at 100, 303 ms at 400 — + worse than the 24 ms scan the index replaces. The cascade's heavy + beat (300 s) calls this, so at most ~30 deltas accumulate under + sustained writes. + + ponytail: retraining (``create_index(replace=True)``) rewrites the + whole index once per heavy beat under load; merging the deltas + instead needs pylance's ``optimize_indices``, which is not a + dependency. Revisit when a table passes ~500k rows. """ columns = cls.vector_columns() if not columns: return [] - indices = await table.list_indices() - indexed = {col for idx in indices for col in (idx.columns or [])} - built: list[str] = [] + indices = { + col: idx + for idx in await table.list_indices() + for col in (idx.columns or []) + } + touched: list[str] = [] for column in columns: - if column in indexed and not replace: - continue - if await table.count_rows(f"{column} IS NOT NULL") < min_rows: + existing = indices.get(column) + if existing is not None: + stats = await table.index_stats(existing.name) + if stats is None or stats.num_indices <= 1: + continue + rows = await table.count_rows(f"{column} IS NOT NULL") + if rows < min_rows: continue await table.create_index( - column, replace=replace, config=IvfFlat(distance_type="cosine") + column, + replace=existing is not None, + config=IvfFlat( + distance_type="cosine", + num_partitions=max(1, rows // VECTOR_INDEX_ROWS_PER_PARTITION), + ), ) - built.append(column) - return built + touched.append(column) + return touched def touch(record: BaseLanceTable) -> BaseLanceTable: diff --git a/src/everos/core/persistence/lancedb/repository.py b/src/everos/core/persistence/lancedb/repository.py index 331b539e3..9652ad8be 100644 --- a/src/everos/core/persistence/lancedb/repository.py +++ b/src/everos/core/persistence/lancedb/repository.py @@ -146,8 +146,9 @@ def _vector_index_min_rows() -> int: ``processing`` forever with nothing logged (a hang raises nothing, so the drain-failure counter stays at zero and ``/health`` keeps reporting healthy). Same last-resort shape as :data:`_COMPACT_TIMEOUT_SECONDS`, and generous by -design: everos builds no vector ANN index, so reads are flat scans — measured -~62ms over 117k rows, i.e. 60s is ~1000x headroom and never fires normally. On +design: a vector read is an IVF probe, or a flat scan below the index +threshold — measured ~62ms over 117k unindexed rows, i.e. 60s is ~1000x +headroom and never fires normally. On expiry the caller gets a retryable :class:`VectorStoreBusyError`, so a drain row is retried and a search request fails with a structured error rather than hanging the request.""" @@ -657,14 +658,16 @@ async def rebuild_indexes(self) -> None: await table.drop_index(idx.name) await self.schema.ensure_fts_indexes(table, replace=True) await self.schema.ensure_vector_indexes( - table, min_rows=_vector_index_min_rows(), replace=True + table, min_rows=_vector_index_min_rows() ) # ── Read ─────────────────────────────────────────────────────────────── async def ensure_vector_indexes(self) -> list[str]: """Build the ANN index on vector columns that crossed the row - threshold since startup — the cascade's heavy beat calls this.""" + threshold since startup, and retrain one whose delta indices piled + up — the cascade's heavy beat calls this; the cases are spelled out + on :meth:`BaseLanceTable.ensure_vector_indexes`.""" async with self._locked(_REBUILD_TIMEOUT_SECONDS, "ensure_vector_indexes"): table = await self._table() return await self.schema.ensure_vector_indexes( diff --git a/src/everos/infra/persistence/backends/lancedb.py b/src/everos/infra/persistence/backends/lancedb.py index 784d013b4..0513ec2ba 100644 --- a/src/everos/infra/persistence/backends/lancedb.py +++ b/src/everos/infra/persistence/backends/lancedb.py @@ -18,6 +18,7 @@ from everos.component.utils.datetime import ensure_utc, to_iso_format from everos.core.persistence import LanceRepoBase +from everos.core.persistence.lancedb.base import VECTOR_QUERY_NPROBES from everos.infra.persistence import lancedb as _lancedb from ..predicate import ( @@ -220,6 +221,7 @@ async def dense_search( .nearest_to(list(vector)) .column(vector_field) .distance_type("cosine") + .nprobes(VECTOR_QUERY_NPROBES) ) expression = _render_optional(where) if expression: @@ -254,6 +256,9 @@ async def prune(self, older_than: dt.timedelta) -> None: async def rebuild_indexes(self) -> None: await self._repo.rebuild_indexes() + async def ensure_vector_indexes(self) -> None: + await self._repo.ensure_vector_indexes() + async def find_by_owner(self, owner_id: str, *, limit: int = 100) -> list[T]: return await self.find_where(eq("owner_id", owner_id), limit=limit) diff --git a/src/everos/infra/persistence/index/protocols.py b/src/everos/infra/persistence/index/protocols.py index 8a9967bcf..db6d4a912 100644 --- a/src/everos/infra/persistence/index/protocols.py +++ b/src/everos/infra/persistence/index/protocols.py @@ -124,6 +124,8 @@ async def prune(self, older_than: dt.timedelta) -> None: ... async def rebuild_indexes(self) -> None: ... + async def ensure_vector_indexes(self) -> None: ... + @runtime_checkable class EpisodeIndexRepository(IndexRepository[T], Protocol[T]): diff --git a/src/everos/infra/persistence/index/router.py b/src/everos/infra/persistence/index/router.py index 83480c04d..93180ba9c 100644 --- a/src/everos/infra/persistence/index/router.py +++ b/src/everos/infra/persistence/index/router.py @@ -195,6 +195,9 @@ async def prune(self, older_than: dt.timedelta) -> None: async def rebuild_indexes(self) -> None: await self._repo().rebuild_indexes() + async def ensure_vector_indexes(self) -> None: + await self._repo().ensure_vector_indexes() + class RoutedEpisodeRepository(RoutedIndexRepository[Any]): def _repo(self) -> EpisodeIndexRepository[Any]: diff --git a/src/everos/infra/persistence/lancedb/repos/agent_skill.py b/src/everos/infra/persistence/lancedb/repos/agent_skill.py index 215ff6c26..37f377860 100644 --- a/src/everos/infra/persistence/lancedb/repos/agent_skill.py +++ b/src/everos/infra/persistence/lancedb/repos/agent_skill.py @@ -7,6 +7,7 @@ from lancedb import AsyncTable from everos.core.persistence.lancedb import LanceRepoBase +from everos.core.persistence.lancedb.base import VECTOR_QUERY_NPROBES from ..lancedb_manager import get_table from ..tables.agent_skill import AgentSkill @@ -59,6 +60,7 @@ async def find_topk_relevant_in_cluster( table.query() .nearest_to(list(query_vector)) .distance_type("cosine") + .nprobes(VECTOR_QUERY_NPROBES) .where(_in_cluster(owner_id, cluster_id)) .limit(top_k) .to_list() diff --git a/src/everos/infra/persistence/milvus/repository.py b/src/everos/infra/persistence/milvus/repository.py index 54774c10a..1204934cd 100644 --- a/src/everos/infra/persistence/milvus/repository.py +++ b/src/everos/infra/persistence/milvus/repository.py @@ -458,6 +458,9 @@ async def prune(self, older_than: dt.timedelta) -> None: async def rebuild_indexes(self) -> None: """Milvus AUTOINDEX maintenance is service-managed.""" + async def ensure_vector_indexes(self) -> None: + """Milvus builds its vector index at collection creation.""" + async def count(self) -> int: return await self._count_where(None) diff --git a/src/everos/memory/cascade/worker.py b/src/everos/memory/cascade/worker.py index d122fd83e..c06411e76 100644 --- a/src/everos/memory/cascade/worker.py +++ b/src/everos/memory/cascade/worker.py @@ -1027,12 +1027,11 @@ async def _run_optimize_once(self, kind: str) -> None: ) if state is not None: state.last_prune_at = now - # A table that crossed the ANN threshold since startup gets its - # vector index on the same heavy beat; a no-op otherwise. - ensure_vectors = getattr(repo, "ensure_vector_indexes", None) - if ensure_vectors is not None: - async with asyncio.timeout(_MAINTENANCE_TASK_TIMEOUT_SECONDS): - await ensure_vectors() + # Same heavy beat: a table that crossed the ANN threshold since + # startup gets its vector index, and one whose light beats piled + # up delta indices is retrained; a no-op otherwise. + async with asyncio.timeout(_MAINTENANCE_TASK_TIMEOUT_SECONDS): + await repo.ensure_vector_indexes() else: # Light beat: lock-free compaction. A commit conflict here # is benign — handled below. diff --git a/tests/unit/test_core/test_persistence/test_lancedb/test_vector_index.py b/tests/unit/test_core/test_persistence/test_lancedb/test_vector_index.py index 5de64e1c8..6949af9a3 100644 --- a/tests/unit/test_core/test_persistence/test_lancedb/test_vector_index.py +++ b/tests/unit/test_core/test_persistence/test_lancedb/test_vector_index.py @@ -19,6 +19,7 @@ from lancedb.pydantic import Vector from everos.core.persistence.lancedb import BaseLanceTable +from everos.core.persistence.lancedb.base import VECTOR_QUERY_NPROBES _DIM = 8 @@ -79,15 +80,53 @@ async def test_at_the_threshold_an_ivf_flat_index_is_built( assert await _vector_indices(vec_table) == [("vector", "IvfFlat")] -async def test_ensure_is_idempotent_unless_replace(vec_table: AsyncTable) -> None: +async def _num_indices(table: AsyncTable) -> int: + (name,) = [ + idx.name + for idx in await table.list_indices() + if idx.columns and idx.columns[0] == "vector" + ] + stats = await table.index_stats(name) + assert stats is not None + return stats.num_indices + + +async def test_delta_indexes_left_by_optimize_are_collapsed( + vec_table: AsyncTable, +) -> None: + """Every ``optimize()`` on a table with new rows appends a delta index and + a query probes them all; the heavy beat folds them back into one index + and otherwise leaves a healthy index alone.""" await vec_table.add(_rows(60)) - await _VecSpec.ensure_vector_indexes(vec_table, min_rows=50) + assert await _VecSpec.ensure_vector_indexes(vec_table, min_rows=50) == ["vector"] assert await _VecSpec.ensure_vector_indexes(vec_table, min_rows=50) == [] - rebuilt = await _VecSpec.ensure_vector_indexes(vec_table, min_rows=50, replace=True) - assert rebuilt == ["vector"] + for _ in range(2): + await vec_table.add(_rows(5)) + await vec_table.optimize() + assert await _num_indices(vec_table) == 3, "precondition: one delta per beat" + assert await _VecSpec.ensure_vector_indexes(vec_table, min_rows=50) == ["vector"] + assert await _num_indices(vec_table) == 1 assert len(await _vector_indices(vec_table)) == 1 +async def test_nprobes_is_accepted_on_an_unindexed_column( + vec_table: AsyncTable, +) -> None: + """``dense_search`` always sets ``nprobes``; a table below the index + threshold (Tier 1, a fresh store) must still answer.""" + await vec_table.add(_rows(10)) + rows = await ( + vec_table.query() + .nearest_to(_rows(1)[0].vector) + .column("vector") + .distance_type("cosine") + .nprobes(VECTOR_QUERY_NPROBES) + .limit(3) + .to_list() + ) + assert len(rows) == 3 + + async def test_all_null_vectors_are_skipped_even_above_the_threshold( vec_table: AsyncTable, ) -> None: diff --git a/tests/unit/test_infra/test_index_router_maintenance.py b/tests/unit/test_infra/test_index_router_maintenance.py new file mode 100644 index 000000000..f86c5c46f --- /dev/null +++ b/tests/unit/test_infra/test_index_router_maintenance.py @@ -0,0 +1,60 @@ +"""The maintenance surface the cascade worker calls must exist on every +repository it can be handed. + +The worker holds ``RoutedEpisodeRepository`` and friends, not the LanceDB +repository. A method added to the LanceDB layer but not forwarded here is +unreachable in production: the first ``ensure_vector_indexes`` hook was +guarded by ``getattr(repo, ..., None)`` and silently never ran. +""" + +from __future__ import annotations + +import pytest + +from everos.infra.persistence.backends.lancedb import LanceIndexRepository +from everos.infra.persistence.index.protocols import IndexRepository +from everos.infra.persistence.index.router import RoutedIndexRepository + +_MAINTENANCE = ("optimize", "prune", "rebuild_indexes", "ensure_vector_indexes") + + +class _Stub: + schema = object() + table_name = "stub" + + def __init__(self) -> None: + self.calls: list[str] = [] + + async def ensure_vector_indexes(self) -> None: + self.calls.append("ensure_vector_indexes") + + +@pytest.mark.parametrize("name", _MAINTENANCE) +def test_every_repository_class_exposes_the_maintenance_method(name: str) -> None: + for cls in (LanceIndexRepository, RoutedIndexRepository): + assert callable(getattr(cls, name, None)), f"{cls.__name__}.{name} missing" + assert name in IndexRepository.__protocol_attrs__ # type: ignore[attr-defined] + + +def test_milvus_repository_exposes_the_maintenance_methods() -> None: + pytest.importorskip("pymilvus") + from everos.infra.persistence.milvus import repository as milvus + + classes = [ + c + for c in vars(milvus).values() + if isinstance(c, type) + and c.__module__ == milvus.__name__ + and callable(getattr(c, "rebuild_indexes", None)) + ] + assert classes, "no Milvus repository class found" + for cls in classes: + for name in _MAINTENANCE: + assert callable(getattr(cls, name, None)), f"{cls.__name__}.{name}" + + +async def test_router_forwards_ensure_vector_indexes_to_the_lance_repo() -> None: + stub = _Stub() + routed = RoutedIndexRepository(stub, milvus_repo_name="unused") # type: ignore[arg-type] + await routed.ensure_vector_indexes() + assert stub.calls == ["ensure_vector_indexes"] diff --git a/tests/unit/test_memory/test_cascade/test_worker.py b/tests/unit/test_memory/test_cascade/test_worker.py index e7953c34a..fd95e0d77 100644 --- a/tests/unit/test_memory/test_cascade/test_worker.py +++ b/tests/unit/test_memory/test_cascade/test_worker.py @@ -306,6 +306,7 @@ def __init__( self.prune_calls: list[float] = [] self.prune_args: list[dt.timedelta] = [] self.rebuild_calls: list[float] = [] + self.ensure_calls: list[float] = [] self.optimize_delay = optimize_delay self.rebuild_delay = rebuild_delay self.rebuild_raises = rebuild_raises @@ -333,6 +334,9 @@ async def rebuild_indexes(self) -> None: raise RuntimeError("rebuild boom") self.rebuild_calls.append(time.monotonic()) + async def ensure_vector_indexes(self) -> None: + self.ensure_calls.append(time.monotonic()) + class _OkHandlerWithRepo(_OkHandler): """OK handler exposing a fake ``index_repo`` for scheduler tests.""" @@ -593,6 +597,27 @@ async def test_optimize_prunes_on_first_call_then_throttles( assert len(fake.optimize_calls) == 1, "second beat is the light path" +async def test_heavy_beat_ensures_vector_indexes(patched_repo: _FakeRepo) -> None: + """Vector index build / delta retrain rides the heavy beat only, so a + table crossing the row threshold mid-run gets its index on the next + prune cadence, not at the 12 h rebuild; light beats leave indexes alone.""" + fake = _FakeLanceRepo() + w = CascadeWorker( + {"episode": _OkHandlerWithRepo(fake)}, + retry_backoff_seconds=0, + optimize_min_interval_seconds=0.01, + optimize_prune_interval_seconds=10.0, + optimize_prune_retention_seconds=45.0, + ) + w._schedule_optimize("episode") + await w._flush_optimizers() + assert len(fake.ensure_calls) == 1, "heavy beat must ensure vector indexes" + await asyncio.sleep(0.02) + w._schedule_optimize("episode") + await w._flush_optimizers() + assert len(fake.ensure_calls) == 1, "light beat must not touch indexes" + + async def test_failed_prune_backs_off_a_cadence_and_keeps_health_signal( patched_repo: _FakeRepo, ) -> None: From 0bfdc22999608c584753c5b6ba643e2cdda281cd Mon Sep 17 00:00:00 2001 From: zhanghui Date: Thu, 24 Sep 2026 16:20:45 +0800 Subject: [PATCH 3/5] fix(lancedb): retrain a vector index only past 16 delta indices Second review pass: with the cap at one delta, a trickle writer (one small upsert per 5-minute window) paid a full index rewrite every heavy beat -- 107 MB per column at 27k x 1024 rows, 391 MB at 100k -- while 16 deltas cost 1-3 ms extra per query. Under sustained writes the light beat adds ~30 deltas per heavy beat, so the cap changes nothing there; it only spares light users the write amplification. Also cover the gap the first review hit one layer down: the LanceDB backend forwarding is asserted with a stub (a 'pass' body satisfied the protocol check alone), and the docstring names the cross-process optimize() commit conflict that can preempt a retrain (same exposure as prune; retried on the next heavy beat). Tests: the cap test is monkeypatched to 2 and asserts both sides (2 deltas left alone, 3 collapsed); it fails when the cap is ignored. The forwarding test fails with a 'pass' body. Co-Authored-By: Claude Fable 5.1 --- src/everos/core/persistence/lancedb/base.py | 16 ++++++++++--- .../test_lancedb/test_vector_index.py | 24 ++++++++++++------- .../test_index_router_maintenance.py | 9 +++++++ 3 files changed, 37 insertions(+), 12 deletions(-) diff --git a/src/everos/core/persistence/lancedb/base.py b/src/everos/core/persistence/lancedb/base.py index 3380182ee..7a32749b5 100644 --- a/src/everos/core/persistence/lancedb/base.py +++ b/src/everos/core/persistence/lancedb/base.py @@ -44,6 +44,12 @@ class Episode(BaseLanceTable): """IVF partition size at build time. Pinned so :data:`VECTOR_QUERY_NPROBES` means something: lance's default partition count has changed across releases.""" +VECTOR_INDEX_MAX_DELTAS = 16 +"""Delta indices a vector column may accumulate before the heavy beat retrains +it. Each light beat with new rows adds one; probing 16 of them cost ~1-3 ms extra +at 27k-100k rows, while a retrain rewrites the whole index (107 MB at 27k x 1024, +391 MB at 100k), so a trickle writer must not pay that every 300 s.""" + VECTOR_QUERY_NPROBES = 32 """Partitions probed per vector query. 32 partitions of 4096 rows cover the whole column, i.e. exact search, up to ~130k rows; past that the unprobed partitions @@ -217,7 +223,8 @@ async def ensure_vector_indexes( Two cases do work; everything else is a no-op: * no index yet and the column crossed ``min_rows`` -> build one; - * the index has grown delta indices -> retrain it in place. Every + * the index has grown more than :data:`VECTOR_INDEX_MAX_DELTAS` delta + indices -> retrain it in place. Every ``optimize()`` on a table with new rows appends one *delta* index instead of merging (``num_indices`` +1 per light beat, never collapsing on its own), and a query probes every delta, so latency @@ -225,7 +232,10 @@ async def ensure_vector_indexes( measured 5.9 ms at 0 deltas, 25.7 ms at 100, 303 ms at 400 — worse than the 24 ms scan the index replaces. The cascade's heavy beat (300 s) calls this, so at most ~30 deltas accumulate under - sustained writes. + sustained writes. The retrain is one atomic index swap (searches + never see the column unindexed); a concurrent ``optimize()`` from + another process can preempt it with a benign commit conflict, in + which case the next heavy beat retries — the same exposure prune has. ponytail: retraining (``create_index(replace=True)``) rewrites the whole index once per heavy beat under load; merging the deltas @@ -245,7 +255,7 @@ async def ensure_vector_indexes( existing = indices.get(column) if existing is not None: stats = await table.index_stats(existing.name) - if stats is None or stats.num_indices <= 1: + if stats is None or stats.num_indices <= VECTOR_INDEX_MAX_DELTAS: continue rows = await table.count_rows(f"{column} IS NOT NULL") if rows < min_rows: diff --git a/tests/unit/test_core/test_persistence/test_lancedb/test_vector_index.py b/tests/unit/test_core/test_persistence/test_lancedb/test_vector_index.py index 6949af9a3..cc4e9e859 100644 --- a/tests/unit/test_core/test_persistence/test_lancedb/test_vector_index.py +++ b/tests/unit/test_core/test_persistence/test_lancedb/test_vector_index.py @@ -18,7 +18,7 @@ from lancedb import AsyncTable from lancedb.pydantic import Vector -from everos.core.persistence.lancedb import BaseLanceTable +from everos.core.persistence.lancedb import BaseLanceTable, base from everos.core.persistence.lancedb.base import VECTOR_QUERY_NPROBES _DIM = 8 @@ -91,19 +91,25 @@ async def _num_indices(table: AsyncTable) -> int: return stats.num_indices -async def test_delta_indexes_left_by_optimize_are_collapsed( - vec_table: AsyncTable, +async def test_delta_indexes_left_by_optimize_are_collapsed_past_the_cap( + vec_table: AsyncTable, monkeypatch: pytest.MonkeyPatch ) -> None: """Every ``optimize()`` on a table with new rows appends a delta index and - a query probes them all; the heavy beat folds them back into one index - and otherwise leaves a healthy index alone.""" + a query probes them all. Up to the cap they are tolerated (a retrain + rewrites the whole index); past it the heavy beat folds them into one.""" + monkeypatch.setattr(base, "VECTOR_INDEX_MAX_DELTAS", 2) await vec_table.add(_rows(60)) assert await _VecSpec.ensure_vector_indexes(vec_table, min_rows=50) == ["vector"] assert await _VecSpec.ensure_vector_indexes(vec_table, min_rows=50) == [] - for _ in range(2): - await vec_table.add(_rows(5)) - await vec_table.optimize() - assert await _num_indices(vec_table) == 3, "precondition: one delta per beat" + await vec_table.add(_rows(5)) + await vec_table.optimize() + assert await _num_indices(vec_table) == 2, "precondition: one delta per beat" + assert await _VecSpec.ensure_vector_indexes(vec_table, min_rows=50) == [], ( + "within the cap the index is left alone" + ) + await vec_table.add(_rows(5)) + await vec_table.optimize() + assert await _num_indices(vec_table) == 3 assert await _VecSpec.ensure_vector_indexes(vec_table, min_rows=50) == ["vector"] assert await _num_indices(vec_table) == 1 assert len(await _vector_indices(vec_table)) == 1 diff --git a/tests/unit/test_infra/test_index_router_maintenance.py b/tests/unit/test_infra/test_index_router_maintenance.py index f86c5c46f..cb2fdef19 100644 --- a/tests/unit/test_infra/test_index_router_maintenance.py +++ b/tests/unit/test_infra/test_index_router_maintenance.py @@ -58,3 +58,12 @@ async def test_router_forwards_ensure_vector_indexes_to_the_lance_repo() -> None routed = RoutedIndexRepository(stub, milvus_repo_name="unused") # type: ignore[arg-type] await routed.ensure_vector_indexes() assert stub.calls == ["ensure_vector_indexes"] + + +async def test_lance_backend_forwards_ensure_vector_indexes_to_the_repo() -> None: + """Attribute presence is not enough: a ``pass`` body would satisfy the + protocol and silently never build an index.""" + stub = _Stub() + backend = LanceIndexRepository(stub, schema=object) # type: ignore[arg-type] + await backend.ensure_vector_indexes() + assert stub.calls == ["ensure_vector_indexes"] From 669b23db906fb1cdb80b700413999296a44787d3 Mon Sep 17 00:00:00 2001 From: zhanghui Date: Thu, 24 Sep 2026 16:37:49 +0800 Subject: [PATCH 4/5] fix(lancedb): leave vector index builds to the cascade worker ensure_business_indexes runs in every process that opens the root: the server lifespan and the CLI's _runtime. With the vector step in it, a read-only 'everos cascade status' trained an IVF index on the running server's table the moment it crossed the row threshold and failed with 'Retryable commit conflict' (Windows soak, final run: 10/9 storm errors in 15 minutes; the first two runs of this PR had none because the CLI storms ran before any table reached 2000 rows). Vector indexes belong to the cascade worker alone: its first rebuild sweep at server start builds a missing one (rebuild_indexes -> ensure_vector_ indexes) and the heavy beat keeps it healthy. FTS stays in the startup pass because a search on a column without its inverted index raises. Test: a spy on BaseLanceTable.ensure_vector_indexes must see no call from ensure_business_indexes; re-adding the call fails it. Co-Authored-By: Claude Fable 5.1 --- src/everos/config/settings.py | 6 ++- .../infra/persistence/lancedb/__init__.py | 14 +++++-- ...startup_leaves_vector_indexes_to_worker.py | 39 +++++++++++++++++++ 3 files changed, 54 insertions(+), 5 deletions(-) create mode 100644 tests/unit/test_infra/test_lancedb/test_startup_leaves_vector_indexes_to_worker.py diff --git a/src/everos/config/settings.py b/src/everos/config/settings.py index cde390ac5..b9df27266 100644 --- a/src/everos/config/settings.py +++ b/src/everos/config/settings.py @@ -523,8 +523,10 @@ class LanceDBSettings(BaseModel): get an ANN index. Below this a brute-force scan is cheaper than the index; above it the scan grows linearly with the table — 27k rows of 1024-dim vectors was 112 MB and ~0.6 s per query on a laptop SSD, and - a hybrid search runs two or three of them. Checked at startup, on the - cascade's heavy maintenance beat, and rebuilt on the rebuild cadence.""" + a hybrid search runs two or three of them. Applied by the cascade worker + only: its first rebuild sweep after server start builds a missing index + and the heavy maintenance beat keeps it healthy. The CLI never builds + one (it would race the running server's commits).""" class CascadeSettings(BaseModel): diff --git a/src/everos/infra/persistence/lancedb/__init__.py b/src/everos/infra/persistence/lancedb/__init__.py index 0a48ba232..e16d512bf 100644 --- a/src/everos/infra/persistence/lancedb/__init__.py +++ b/src/everos/infra/persistence/lancedb/__init__.py @@ -27,7 +27,6 @@ import contextlib import datetime as dt -from everos.config.settings import load_settings from everos.core.observability.logging import get_logger from everos.core.persistence import BaseLanceTable, MemoryRoot, memory_root_lock @@ -290,14 +289,23 @@ async def ensure_business_indexes() -> None: Adding a new business table = adding it to ``_BUSINESS_SCHEMAS``; everything else (table name, columns to index) reads off the schema's ClassVars. + + Vector (ANN) indexes are deliberately **not** built here. This runs in + every process that opens the root — the server lifespan and the CLI's + ``_runtime`` — and a CLI command training an index on a live server's + table races its commits (soak: ``cascade status`` storms failed with + ``Retryable commit conflict`` the moment a table crossed the row + threshold). Vector indexes belong to the cascade worker alone: its + first rebuild sweep at server start builds a missing one and the heavy + beat maintains it (:meth:`BaseLanceTable.ensure_vector_indexes`). FTS + stays here because a search on a column without its inverted index + raises instead of degrading. """ await migrate_table_schemas() await migrate_fts_indexes() - min_rows = load_settings().lancedb.vector_index_min_rows for schema in _BUSINESS_SCHEMAS: table = await get_table(schema.TABLE_NAME, schema) await schema.ensure_fts_indexes(table) - await schema.ensure_vector_indexes(table, min_rows=min_rows) async def verify_business_schemas() -> None: diff --git a/tests/unit/test_infra/test_lancedb/test_startup_leaves_vector_indexes_to_worker.py b/tests/unit/test_infra/test_lancedb/test_startup_leaves_vector_indexes_to_worker.py new file mode 100644 index 000000000..088e921f3 --- /dev/null +++ b/tests/unit/test_infra/test_lancedb/test_startup_leaves_vector_indexes_to_worker.py @@ -0,0 +1,39 @@ +"""``ensure_business_indexes`` runs in every process that opens the root — the +server lifespan and the CLI's ``_runtime`` — so it must never train an ANN +index: a CLI command doing that races the running server's commits (soak: +``cascade status`` storms failed with ``Retryable commit conflict`` the moment +a table crossed the row threshold). Vector indexes are the cascade worker's: +first rebuild sweep at server start, then the heavy beat. +""" + +from __future__ import annotations + +from pathlib import Path + +import pytest + +from everos.core.persistence import BaseLanceTable +from everos.infra.persistence.lancedb import ensure_business_indexes, lancedb_manager + + +@pytest.fixture(autouse=True) +async def _isolated_root(tmp_path: Path, monkeypatch: pytest.MonkeyPatch): + monkeypatch.setenv("EVEROS_ROOT", str(tmp_path)) + lancedb_manager._conn = None + lancedb_manager._tables.clear() + yield + await lancedb_manager.dispose_connection() + + +async def test_startup_pass_never_trains_a_vector_index( + monkeypatch: pytest.MonkeyPatch, +) -> None: + calls: list[str] = [] + + async def spy(cls, table, *, min_rows): # type: ignore[no-untyped-def] + calls.append(cls.__name__) + return [] + + monkeypatch.setattr(BaseLanceTable, "ensure_vector_indexes", classmethod(spy)) + await ensure_business_indexes() + assert calls == [], f"startup pass trained vector indexes on {calls}" From c0e6eeb6805972f922680a40ec267e1efb296bac Mon Sep 17 00:00:00 2001 From: zhanghui Date: Thu, 24 Sep 2026 17:13:38 +0800 Subject: [PATCH 5/5] test(cascade): a failing vector index build stays inside the heavy beat Pins what the code already does: an exception from ensure_vector_indexes on the heavy beat is caught with the other maintenance failures, counted toward optimize_failures (so the health signal can show a streak), leaves the prune that ran before it credited, and the next beat still runs. Co-Authored-By: Claude Fable 5.1 --- .../test_memory/test_cascade/test_worker.py | 31 +++++++++++++++++++ 1 file changed, 31 insertions(+) diff --git a/tests/unit/test_memory/test_cascade/test_worker.py b/tests/unit/test_memory/test_cascade/test_worker.py index fd95e0d77..09801e628 100644 --- a/tests/unit/test_memory/test_cascade/test_worker.py +++ b/tests/unit/test_memory/test_cascade/test_worker.py @@ -301,6 +301,7 @@ def __init__( optimize_delay: float = 0.0, rebuild_delay: float = 0.0, rebuild_raises: bool = False, + ensure_raises: bool = False, ) -> None: self.optimize_calls: list[float] = [] self.prune_calls: list[float] = [] @@ -310,6 +311,7 @@ def __init__( self.optimize_delay = optimize_delay self.rebuild_delay = rebuild_delay self.rebuild_raises = rebuild_raises + self.ensure_raises = ensure_raises @property def beats(self) -> list[float]: @@ -336,6 +338,8 @@ async def rebuild_indexes(self) -> None: async def ensure_vector_indexes(self) -> None: self.ensure_calls.append(time.monotonic()) + if self.ensure_raises: + raise RuntimeError("KMeans cannot train: simulated create_index failure") class _OkHandlerWithRepo(_OkHandler): @@ -618,6 +622,33 @@ async def test_heavy_beat_ensures_vector_indexes(patched_repo: _FakeRepo) -> Non assert len(fake.ensure_calls) == 1, "light beat must not touch indexes" +async def test_vector_index_failure_on_the_heavy_beat_is_contained( + patched_repo: _FakeRepo, +) -> None: + """A failing index build / retrain must not escape the maintenance beat + or stop maintenance: it is counted like any other optimize failure + (visible in the health signal), the prune that ran before it stays + credited, and the next beats keep coming.""" + fake = _FakeLanceRepo(ensure_raises=True) + w = CascadeWorker( + {"episode": _OkHandlerWithRepo(fake)}, + retry_backoff_seconds=0, + optimize_min_interval_seconds=0.01, + optimize_prune_interval_seconds=10.0, + optimize_prune_retention_seconds=45.0, + ) + w._schedule_optimize("episode") + await w._flush_optimizers() # raises nothing + state = w._optimizer_states["episode"] + assert len(fake.ensure_calls) == 1 + assert state.optimize_failures == 1, "counted, so /health can see a streak" + assert state.last_prune_at > 0, "the prune before it stays credited" + await asyncio.sleep(0.02) + w._schedule_optimize("episode") + await w._flush_optimizers() + assert len(fake.optimize_calls) == 1, "the next (light) beat still runs" + + async def test_failed_prune_backs_off_a_cadence_and_keeps_health_signal( patched_repo: _FakeRepo, ) -> None: