diff --git a/docs/cascade_runbook.md b/docs/cascade_runbook.md index 338d2a7cf..9d0ee527d 100644 --- a/docs/cascade_runbook.md +++ b/docs/cascade_runbook.md @@ -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 @@ -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` @@ -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: @@ -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 @@ -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//skills/skill_/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//skills/skill_/SKILL.md` re-enqueues + it directly. diff --git a/docs/how-memory-works.md b/docs/how-memory-works.md index be6fda712..7bb95acf5 100644 --- a/docs/how-memory-works.md +++ b/docs/how-memory-works.md @@ -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` / diff --git a/src/everos/entrypoints/cli/commands/cascade.py b/src/everos/entrypoints/cli/commands/cascade.py index d8095f56e..105935e7e 100644 --- a/src/everos/entrypoints/cli/commands/cascade.py +++ b/src/everos/entrypoints/cli/commands/cascade.py @@ -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 @@ -30,6 +30,7 @@ from __future__ import annotations import asyncio +import contextlib import enum import os from collections.abc import AsyncIterator @@ -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, @@ -61,6 +63,7 @@ ) from everos.memory.cascade import ( CascadeOrchestrator, + hold_ome_lock, match_kind, ome_lock_is_free, ) @@ -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 @@ -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: @@ -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) @@ -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") @@ -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 diff --git a/src/everos/memory/cascade/__init__.py b/src/everos/memory/cascade/__init__.py index 7ff6b1bbf..4c659c6dc 100644 --- a/src/everos/memory/cascade/__init__.py +++ b/src/everos/memory/cascade/__init__.py @@ -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 @@ -38,6 +39,7 @@ "CascadeOrchestrator", "KindSpec", "NullBackfillPresenter", + "hold_ome_lock", "match_kind", "ome_lock_is_free", ] diff --git a/src/everos/memory/cascade/_backfill.py b/src/everos/memory/cascade/_backfill.py index 11f2ac850..ac0516332 100644 --- a/src/everos/memory/cascade/_backfill.py +++ b/src/everos/memory/cascade/_backfill.py @@ -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 @@ -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. diff --git a/tests/integration/test_cascade_cli_integration.py b/tests/integration/test_cascade_cli_integration.py index c89ad7f84..86fc729b6 100644 --- a/tests/integration/test_cascade_cli_integration.py +++ b/tests/integration/test_cascade_cli_integration.py @@ -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 @@ -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 ``."""