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
3 changes: 3 additions & 0 deletions src/everos/config/default.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
9 changes: 9 additions & 0 deletions src/everos/config/settings.py
Original file line number Diff line number Diff line change
Expand Up @@ -518,6 +518,15 @@ 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. 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):
Expand Down
96 changes: 95 additions & 1 deletion src/everos/core/persistence/lancedb/base.py
Original file line number Diff line number Diff line change
Expand Up @@ -34,12 +34,28 @@ 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

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_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
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
Expand Down Expand Up @@ -177,6 +193,84 @@ 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
) -> list[str]:
"""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; 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 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
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. 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
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 = {
col: idx
for idx in await table.list_indices()
for col in (idx.columns or [])
}
touched: list[str] = []
for column in columns:
existing = indices.get(column)
if existing is not None:
stats = await table.index_stats(existing.name)
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:
continue
await table.create_index(
column,
replace=existing is not None,
config=IvfFlat(
distance_type="cosine",
num_partitions=max(1, rows // VECTOR_INDEX_ROWS_PER_PARTITION),
),
)
touched.append(column)
return touched


def touch(record: BaseLanceTable) -> BaseLanceTable:
"""Set ``record.updated_at = now`` and return the record (chainable)."""
Expand Down
33 changes: 30 additions & 3 deletions src/everos/core/persistence/lancedb/repository.py
Original file line number Diff line number Diff line change
Expand Up @@ -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."""
Expand Down Expand Up @@ -136,8 +146,9 @@
``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."""
Expand Down Expand Up @@ -639,14 +650,30 @@ 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()
)

# ── Read ───────────────────────────────────────────────────────────────

async def ensure_vector_indexes(self) -> list[str]:
"""Build the ANN index on vector columns that crossed the row
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(
table, min_rows=_vector_index_min_rows()
)

async def count(self) -> int:
"""Total row count."""
async with self._deadline(_READ_TIMEOUT_SECONDS, "count"):
Expand Down
5 changes: 5 additions & 0 deletions src/everos/infra/persistence/backends/lancedb.py
Original file line number Diff line number Diff line change
Expand Up @@ -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 (
Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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)

Expand Down
2 changes: 2 additions & 0 deletions src/everos/infra/persistence/index/protocols.py
Original file line number Diff line number Diff line change
Expand Up @@ -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]):
Expand Down
3 changes: 3 additions & 0 deletions src/everos/infra/persistence/index/router.py
Original file line number Diff line number Diff line change
Expand Up @@ -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]:
Expand Down
11 changes: 11 additions & 0 deletions src/everos/infra/persistence/lancedb/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -289,6 +289,17 @@ 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()
Expand Down
2 changes: 2 additions & 0 deletions src/everos/infra/persistence/lancedb/repos/agent_skill.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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()
Expand Down
3 changes: 3 additions & 0 deletions src/everos/infra/persistence/milvus/repository.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)

Expand Down
5 changes: 5 additions & 0 deletions src/everos/memory/cascade/worker.py
Original file line number Diff line number Diff line change
Expand Up @@ -1027,6 +1027,11 @@ async def _run_optimize_once(self, kind: str) -> None:
)
if state is not None:
state.last_prune_at = now
# 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.
Expand Down
Loading
Loading