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
45 changes: 28 additions & 17 deletions docs/cascade_runbook.md
Original file line number Diff line number Diff line change
Expand Up @@ -109,9 +109,10 @@ naturally.

## One-shot replay: `everos cascade sync [PATH]`

Use this when the watcher missed an event (WSL mount, network share,
external editor with no inotify) or when you want a deterministic
flush before, say, a smoke test:
Use this when no server is running and you want the index caught up
with the markdown — after batch edits, before a smoke test, or on a
mount where the watcher misses events (WSL mount, network share,
external editor with no inotify) while the daemon is down:

```bash
everos cascade sync # drain everything pending
Expand All @@ -120,10 +121,16 @@ everos cascade sync users/u1/episodes/X.md # re-enqueue + drain

The CLI builds the same `CascadeOrchestrator` as the daemon but only
calls `sync_once` / `drain_once` — no watcher / scanner background task.
Its drain still runs the same compaction + version-cleanup (`prune`) as
the daemon, but `prune` uses `delete_unverified=False`, so it never
deletes a file another process may be mid-commit on. Safe to run in
parallel with a live `everos server`.
It holds the OME lock for the whole run and **refuses to start (exit code
3) while a server holds it**: two processes writing the same LanceDB
tables cannot see each other's snapshot and both insert the row (4–5 %
duplicate rows after a 10-hour soak with two concurrent `sync` processes
next to a server). The same rule applies to `cascade fix --apply` and
`cascade rebuild`; `cascade status` and `cascade fix` (listing) are
read-only and work alongside a server. A running server projects every
markdown change itself, so nothing is lost by waiting for it — unless it
was started with `EVEROS_DISABLE_CASCADE=1` or has been quiesced, in
which case stop it before syncing.

## Rebuild the index: `everos cascade rebuild`

Expand All @@ -135,12 +142,12 @@ everos cascade rebuild # prompts for confirmation
everos cascade rebuild --yes # non-interactive
```

> **Stop the `everos server` first.** Unlike `cascade sync`, rebuild
> **drops and recreates** the active backend's tables or collections. A running
> daemon holds
> cached table handles that would keep pointing at (and writing to) the
> dropped dataset, corrupting the rebuild. This is the one cascade
> command that is **not** safe to run alongside a live server.
> **Stop the `everos server` first.** Like every index-writing cascade
> command, rebuild refuses to run while a server holds the memory root
> (exit code 3) — and it has the strongest reason: it **drops and
> recreates** the active backend's tables or collections. A running daemon
> holds cached table handles that would keep pointing at (and writing to)
> the dropped dataset, corrupting the rebuild.

What it does, in order:

Expand Down Expand Up @@ -231,7 +238,10 @@ Workarounds:
- Rely on the scanner — at default 30 s interval, throughput is
bounded but eventually-consistent.
- Drop the scan interval to ~5 s if the memory root is small.
- Run `everos cascade sync` explicitly after batch edits.
- With no server running, run `everos cascade sync` explicitly after batch
edits. A running server picks them up itself, and `sync` refuses to run
next to it (exit code 3): two processes writing the same index insert
rows twice.

### Daemon process crash mid-batch

Expand Down Expand Up @@ -373,6 +383,7 @@ is a deployment-side change with no schema work.
in the entry inline. Tracked separately.
- **Reference-file change detection (agent_skill)**: edits to
`references/*.md` siblings won't trigger a re-index — only changes
to `SKILL.md` itself fire the watcher. Workaround: run
`everos cascade sync agents/<a>/skills/skill_<n>/SKILL.md` after
editing references.
to `SKILL.md` itself fire the watcher. Workaround: touch or re-save
`SKILL.md` so the watcher fires; with the server stopped,
`everos cascade sync agents/<a>/skills/skill_<n>/SKILL.md` re-enqueues
it directly.
4 changes: 2 additions & 2 deletions docs/how-memory-works.md
Original file line number Diff line number Diff line change
Expand Up @@ -263,8 +263,8 @@ Two paths, two guarantees:

So a `/search` immediately after the `/flush` that produced a record may
miss it. The markdown is durable regardless; index lag never loses data. If
you need read-your-write, retry with backoff, or force the queue with
`everos cascade sync`.
you need read-your-write, retry with backoff (the running server is the
only index writer; `everos cascade sync` is for when no server is running).

Integrity is anchored by a few invariants (details in
[storage_layout.md](storage_layout.md)): the frontmatter `id` /
Expand Down
79 changes: 58 additions & 21 deletions src/everos/entrypoints/cli/commands/cascade.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,7 @@

- ``cascade sync [PATH]`` — flush the work queue. With ``PATH`` the
command first force-enqueues that single file (used after a manual
md edit when waiting for the watcher is impractical), then drains.
md edit with no server running), then drains.
- ``cascade status`` — print the queue + LSN summary that the daemon
sees right now.
- ``cascade fix`` — list every ``failed`` row. With ``--apply``, also
Expand All @@ -30,6 +30,7 @@
from __future__ import annotations

import asyncio
import contextlib
import enum
import os
from collections.abc import AsyncIterator
Expand All @@ -47,6 +48,7 @@
from everos.core.persistence import MemoryRoot
from everos.entrypoints.cli._log_setup import configure_cli_logging
from everos.entrypoints.cli.commands._backfill_cmd import run_backfill
from everos.infra.ome.exceptions import EngineLockHeldError
from everos.infra.persistence.index import (
connect,
drop_business_tables,
Expand All @@ -61,6 +63,7 @@
)
from everos.memory.cascade import (
CascadeOrchestrator,
hold_ome_lock,
match_kind,
ome_lock_is_free,
)
Expand Down Expand Up @@ -136,15 +139,21 @@ def _apply_verbose_logging(verbose: bool | None) -> None:


@asynccontextmanager
async def _runtime(*, verify: bool = True, ensure: bool = True) -> AsyncIterator[None]:
async def _runtime(
*, verify: bool = True, ensure: bool = True, exclusive: bool = False
) -> AsyncIterator[None]:
"""Stand up sqlite + lancedb the same way the API lifespan would.

The CLI uses the same lazy, process-wide singletons the API lifespan
does. They are **per-process**: a running daemon has its own
connection and table-handle cache, so read/write traffic interleaves
safely, but a change to the table *set* made here (drop / recreate)
is invisible to the daemon's cached handles — which is why
``rebuild`` refuses to run while a server holds the OME lock.
connection and table-handle cache, and reads its own LanceDB
snapshot. Reads interleave safely; writes do not — a second process
upserting the same table cannot see what the daemon just committed
(nor the other way round) and both insert the row, so every command
that writes the index (``sync``, ``fix --apply``, ``rebuild``) passes
``exclusive=True`` and holds the OME lock for its whole run. Refused
with exit code 3 while a server (or another exclusive CLI phase) holds
it; a server starting meanwhile fails at its own lock instead.

``verify=False`` skips :func:`verify_business_schemas` — required by
``cascade rebuild``, whose whole purpose is to recover from a table
Expand All @@ -159,19 +168,39 @@ async def _runtime(*, verify: bool = True, ensure: bool = True) -> AsyncIterator
Rebuild recreates the tables and their indexes itself after dropping,
so skipping the pre-drop pass loses nothing.
"""
engine = get_engine()
async with engine.begin() as conn:
await conn.run_sync(SQLModel.metadata.create_all)
await connect()
if verify:
await verify_business_schemas()
if ensure:
await ensure_business_indexes()
lock = hold_ome_lock() if exclusive else contextlib.nullcontext()
try:
yield
lock.__enter__()
except EngineLockHeldError:
typer.echo(
"error: another process holds this memory root's OME lock — a "
"running `everos server`\n"
" (or another exclusive CLI phase). Two processes writing the "
"same index insert rows\n"
" twice, so this command needs the root to itself. A server "
"projects markdown changes\n"
" on its own unless it was started with EVEROS_DISABLE_CASCADE=1 "
"or quiesced; stop it\n"
" first, then re-run.",
err=True,
)
raise typer.Exit(code=3) from None
try:
engine = get_engine()
async with engine.begin() as conn:
await conn.run_sync(SQLModel.metadata.create_all)
await connect()
if verify:
await verify_business_schemas()
if ensure:
await ensure_business_indexes()
try:
yield
finally:
await shutdown()
await dispose_engine()
finally:
await shutdown()
await dispose_engine()
lock.__exit__(None, None, None)


def _build_orchestrator() -> CascadeOrchestrator:
Expand Down Expand Up @@ -215,12 +244,20 @@ def sync(
typer.Option("--verbose", "-v", help=_VERBOSE_OPTION_HELP),
] = None,
) -> None:
"""Drain the cascade queue (and optionally re-enqueue a path first)."""
"""Drain the cascade queue (and optionally re-enqueue a path first).

Holds the OME lock for the run (see :func:`_runtime`): a second process
writing the same LanceDB tables inserts rows twice — the server reads
its own table snapshot and cannot see what the CLI process just
committed, so both decide the row is new (4-5 % duplicate rows after a
10-hour soak with two concurrent ``cascade sync`` processes). Refused
with exit code 3 while a server holds the root.
"""
_apply_root_env(root)
_apply_verbose_logging(verbose)

async def _run() -> None:
async with _runtime():
async with _runtime(exclusive=True):
orchestrator = _build_orchestrator()
if path is not None:
rel = _resolve_relative(path)
Expand Down Expand Up @@ -315,7 +352,7 @@ def fix(
_apply_verbose_logging(verbose)

async def _run() -> None:
async with _runtime():
async with _runtime(exclusive=apply):
rows = await md_change_state_repo.list_failed()
if not rows:
typer.echo("no failed rows")
Expand Down Expand Up @@ -495,7 +532,7 @@ async def _run() -> None:
# to fix; the startup guard would abort before we could rebuild.
# ensure=False: the pre-drop migration pass would raise on exactly
# the damage we are here to repair (see _runtime).
async with _runtime(verify=False, ensure=False):
async with _runtime(verify=False, ensure=False, exclusive=True):
# Reset the queue FIRST so every crash window converges on
# "queue pending → next scan re-indexes". Doing it after the
# drop leaves a window where a crash yields empty tables with
Expand Down
2 changes: 2 additions & 0 deletions src/everos/memory/cascade/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@
from ._backfill import BackfillPhase as BackfillPhase
from ._backfill import BackfillPresenter as BackfillPresenter
from ._backfill import NullBackfillPresenter as NullBackfillPresenter
from ._backfill import hold_ome_lock as hold_ome_lock
from ._backfill import ome_lock_is_free as ome_lock_is_free
from .orchestrator import CascadeConfig as CascadeConfig
from .orchestrator import CascadeHealth as CascadeHealth
Expand All @@ -38,6 +39,7 @@
"CascadeOrchestrator",
"KindSpec",
"NullBackfillPresenter",
"hold_ome_lock",
"match_kind",
"ome_lock_is_free",
]
29 changes: 28 additions & 1 deletion src/everos/memory/cascade/_backfill.py
Original file line number Diff line number Diff line change
Expand Up @@ -25,9 +25,10 @@
from __future__ import annotations

import asyncio
import contextlib
import dataclasses
import datetime as dt
from collections.abc import Callable
from collections.abc import Callable, Iterator
from pathlib import Path
from typing import Any, Protocol
from uuid import uuid4
Expand Down Expand Up @@ -1095,6 +1096,32 @@ def _probe_ome_lock_available() -> bool:
handle.close()


@contextlib.contextmanager
def hold_ome_lock() -> Iterator[None]:
"""Hold the OME jobstore lock for the duration of a CLI write phase.

Same file and flags as :meth:`OfflineEngine._acquire_lock`, so a server
that starts meanwhile fails at startup with :class:`EngineLockHeldError`
instead of becoming a second index writer. Raises
:class:`EngineLockHeldError` when another process already holds it.
"""
root = MemoryRoot.resolve()
lock_path = Path(str(root.ome_db) + ".lock")
lock_path.parent.mkdir(parents=True, exist_ok=True)
handle = open(lock_path, "a+") # noqa: SIM115
try:
try:
portalocker.lock(handle, portalocker.LOCK_EX | portalocker.LOCK_NB)
except portalocker.LockException as exc:
raise EngineLockHeldError(f"another process holds {lock_path}") from exc
try:
yield
finally:
portalocker.unlock(handle)
finally:
handle.close()


def _build_cluster_engine() -> OfflineEngine:
"""Construct (but do not start) the throw-away OME engine Phase 2 drives.

Expand Down
71 changes: 71 additions & 0 deletions tests/integration/test_cascade_cli_integration.py
Original file line number Diff line number Diff line change
Expand Up @@ -15,11 +15,13 @@
from __future__ import annotations

import asyncio
import contextlib
import datetime as _dt
import re
from collections.abc import Iterator
from pathlib import Path

import portalocker
import pytest
from typer.testing import CliRunner

Expand Down Expand Up @@ -313,6 +315,75 @@ def test_rebuild_refuses_to_run_while_a_server_holds_the_lock(
assert "rebuild complete" not in combined


@contextlib.contextmanager
def _server_holds_the_lock(root: Path) -> Iterator[None]:
"""Hold the OME jobstore lock the way a running ``everos server`` does."""
lock_path = root / ".index" / "sqlite" / "ome.db.lock"
lock_path.parent.mkdir(parents=True, exist_ok=True)
with open(lock_path, "a+") as handle:
portalocker.lock(handle, portalocker.LOCK_EX | portalocker.LOCK_NB)
try:
yield
finally:
portalocker.unlock(handle)


def test_sync_refuses_to_run_while_a_server_holds_the_lock(
cli_runtime: Path, monkeypatch: pytest.MonkeyPatch
) -> None:
"""``sync`` is an index writer; next to a running server it exits 3
before opening anything (a second writer inserts rows twice).
"""

def _boom() -> None:
raise AssertionError("the runtime must not open the DB when refusing")

monkeypatch.setattr(cascade_mod, "get_engine", _boom)
with _server_holds_the_lock(cli_runtime):
result = CliRunner().invoke(cascade_mod.app, ["sync"])

assert result.exit_code == 3, result.output
assert "holds this memory root" in result.stderr.lower()
assert "stop it" in result.stderr.lower()
assert "sync complete" not in result.output


def test_fix_apply_refuses_but_fix_and_status_still_run_next_to_a_server(
cli_runtime: Path,
) -> None:
"""Only the writing commands need the root to themselves."""
with _server_holds_the_lock(cli_runtime):
applied = CliRunner().invoke(cascade_mod.app, ["fix", "--apply"])
listed = CliRunner().invoke(cascade_mod.app, ["fix"])
status = CliRunner().invoke(cascade_mod.app, ["status"])

assert applied.exit_code == 3, applied.output
assert "holds this memory root" in applied.stderr.lower()
assert listed.exit_code == 0, listed.output + listed.stderr
assert status.exit_code == 0, status.output + status.stderr


def test_sync_holds_the_lock_while_it_drains(
cli_runtime: Path, monkeypatch: pytest.MonkeyPatch
) -> None:
"""The lock is held for the whole drain, not probed and released — a
server starting meanwhile must fail at its own lock, not join in.
"""
seen: list[bool] = []

class _Orchestrator:
async def sync_once(self) -> int:
seen.append(cascade_mod.ome_lock_is_free())
return 0

monkeypatch.setattr(cascade_mod, "_build_orchestrator", lambda: _Orchestrator())
result = CliRunner().invoke(cascade_mod.app, ["sync"])

assert result.exit_code == 0, result.output + result.stderr
assert seen == [False] # the lock was ours during the drain
assert cascade_mod.ome_lock_is_free() # and released afterwards


# Reduce false negatives on date drift.
def test_resolve_relative_via_command_arg(cli_runtime: Path) -> None:
"""An absolute path under the root works through ``cascade sync <path>``."""
Expand Down
Loading