diff --git a/docs/configuration.md b/docs/configuration.md index 009b14be..ef4b299c 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -188,6 +188,13 @@ if the current total differs. The guarded mode holds the maintenance lock, requires an active enrichment pause sentinel and a verified backup receipt for that DB no older than 24 hours, and reuses VACUUM's quiet-window (04:00–06:00 local), idle-queue and writer gates. +`--wait-for-backup-seconds N` (default 0) lets guarded live apply poll at intervals of up to 30 seconds +for a qualifying receipt before taking the lock or quiescing services. Dry runs and +offline applies ignore it. Poll intervals shorten near the deadline. Waiting ends +at N seconds or when only 20 minutes remain in the quiet window, reserving that +time for the guarded run. A refusal after sleeping reports `verified-backup-timeout`; +if no wait is possible before the first sleep, it reports `verified-backup-required`, +as does the default. Receipt freshness continues to use `attempted_at`. It quiesces the fleet/throughput/tier-0 watchdogs, health-check healer, BrainBar UI/daemon, hotlane, watcher, drain, index, tier-3 ingest, decay and enrichment using maintenance's service helpers. It checks that jobs remain unloaded, BrainBar processes are gone and no writable database diff --git a/src/brainlayer/cli/__init__.py b/src/brainlayer/cli/__init__.py index 5cb94394..369346d8 100644 --- a/src/brainlayer/cli/__init__.py +++ b/src/brainlayer/cli/__init__.py @@ -2883,6 +2883,12 @@ def scrub_at_rest_command( expect_rows: int | None = typer.Option( None, "--expect-rows", min=0, help="Required guarded apply total from the preceding dry run." ), + wait_for_backup_seconds: int = typer.Option( + 0, + "--wait-for-backup-seconds", + min=0, + help="Wait for a verified backup before guarded apply; dry runs ignore this option.", + ), ) -> None: """Redact provider credentials on a copy or a guarded runtime DB; print counts only.""" from ..scrub_at_rest import scrub_at_rest @@ -2895,6 +2901,7 @@ def scrub_at_rest_command( batch_size=batch_size, allow_live_db=allow_live_db, expect_rows=expect_rows, + wait_for_backup_seconds=wait_for_backup_seconds, ) except Exception as exc: typer.echo( diff --git a/src/brainlayer/scrub_at_rest.py b/src/brainlayer/scrub_at_rest.py index d18b8035..982feb01 100644 --- a/src/brainlayer/scrub_at_rest.py +++ b/src/brainlayer/scrub_at_rest.py @@ -13,7 +13,7 @@ from .chunk_origin_wipe import assert_not_live_db from .chunk_write import canonical_content_hash from .dedupe import BUSY_RETRY_ATTEMPTS, _busy_retry_delay, compute_dedupe_fields -from .maintenance import _maintenance_lock +from .maintenance import _maintenance_lock, _remaining_quiet_window_seconds from .pipeline.secret_scrub import scrub_secrets from .runtime_store import ReadonlyStore, WriterRuntimeStore from .vector_store import value_free_sqlite_logging @@ -52,6 +52,7 @@ BRAINBAR_EXIT_TIMEOUT_SECONDS = 30.0 +BACKUP_WAIT_SAFE_MARGIN_SECONDS = 20 * 60.0 _QUIESCE_DETAILS = frozenset( {"quiesce-services", "brainbar-process-probe", "lsof-writers", "process:BrainBar", "process:BrainBarDaemon"} | { @@ -285,6 +286,34 @@ def _live_requirements(config): raise ScrubAtRestError("verified backup within 24 hours required", reason="verified-backup-required") +def _wait_for_verified_backup(path, timeout_seconds): + """Wait without holding a maintenance lock or stopping any writer services.""" + config = maintenance.MaintenanceConfig(db_path=path, backup_reuse_max_age_hours=24) + deadline = time.monotonic() + timeout_seconds + slept = False + while True: + try: + _live_requirements(config) + return + except ScrubAtRestError as exc: + if exc.reason != "verified-backup-required": + raise + remaining = min( + deadline - time.monotonic(), + _remaining_quiet_window_seconds(config) - BACKUP_WAIT_SAFE_MARGIN_SECONDS, + ) + if remaining <= 0: + if not slept: + raise ScrubAtRestError("verified backup within 24 hours required", reason="verified-backup-required") + raise ScrubAtRestError("verified backup wait timed out", reason="verified-backup-timeout") + time.sleep(min(30.0, remaining / 2)) + slept = True + # Inspect the final receipt only within the deadline and window reserve. + window_left = _remaining_quiet_window_seconds(config) - BACKUP_WAIT_SAFE_MARGIN_SECONDS + if time.monotonic() > deadline or window_left <= 0: + raise ScrubAtRestError("verified backup wait timed out", reason="verified-backup-timeout") + + def _check_no_brainbar_processes(): # An unregistered UI can open the daemon bundle even after its job is booted out. try: @@ -400,12 +429,13 @@ def scrub_at_rest( providers: str = "google_oauth", allow_live_db: bool = False, expect_rows: int | None = None, + wait_for_backup_seconds: int = 0, ) -> dict: """Read-only surveys need no opt-in; guarded applies require a current row total.""" failure = None try: with value_free_sqlite_logging(): - if providers not in PROVIDER_MODES or not 1 <= batch_size <= 1000: + if providers not in PROVIDER_MODES or not 1 <= batch_size <= 1000 or wait_for_backup_seconds < 0: raise ValueError("invalid scrub options") selected = PROVIDER_MODES[providers] try: @@ -422,6 +452,10 @@ def scrub_at_rest( if dry_run: with ReadonlyStore(path) as store: return _run(store, True, batch_size, selected) + if allow_live_db and wait_for_backup_seconds > 0: + if expect_rows is None or expect_rows < 0: + raise ScrubAtRestError("expected row count required", reason="expected-row-count-required") + _wait_for_verified_backup(path, wait_for_backup_seconds) with _maintenance_lock(path): path = assert_not_live_db(path, allow_live=allow_live_db) if allow_live_db: diff --git a/tests/test_scrub_at_rest_command.py b/tests/test_scrub_at_rest_command.py index b6a297cc..cf73490e 100644 --- a/tests/test_scrub_at_rest_command.py +++ b/tests/test_scrub_at_rest_command.py @@ -1064,3 +1064,144 @@ def fail(*args, **kwargs): module.scrub_at_rest(db.db_path, allow_live_db=True, expect_rows=live_guard.total) assert error.value.detail == "state:com.etanhey.brainlayer-fleet-watchdog" assert "private-probe-value" not in str(error.value.__cause__) + + +@pytest.mark.parametrize("mode", ["google_oauth", "context7", "exa_labeled"]) +@pytest.mark.parametrize( + "outcome", ["appears", "final-poll", "oversleep", "timeout", "window", "zero", "dry-run", "missing-count"] +) +def test_backup_wait_before_lock_and_quiesce(db, live_guard, monkeypatch, mode, outcome): + from contextlib import contextmanager + + from brainlayer import maintenance + from brainlayer import scrub_at_rest as module + + live_guard.backup.unlink() + monkeypatch.setenv("BRAINLAYER_MCP_SOCKET", str(db.db_path.parent / "absent.sock")) + monkeypatch.setenv("BRAINLAYER_FORBID_BRAINBAR_SOCKET", "1") + elapsed, sleeps, locks = [0.0], [], [] + factory = maintenance.MaintenanceConfig + + def config(**kwargs): + result = factory(**kwargs) + start = {"window": live_guard.now.replace(hour=5, minute=39, second=59)}.get(outcome, live_guard.now) + result.now_fn = lambda: start + dt.timedelta(seconds=elapsed[0]) + return result + + monkeypatch.setattr(maintenance, "MaintenanceConfig", config) + + def sleep(seconds): + assert not locks + assert not any(isinstance(e, tuple) for e in live_guard.events) + elapsed[0] += seconds + {"final-poll": 0.01, "window": 1.0}.get(outcome, 0) + elapsed[0] = {("oversleep", 2): 91}.get((outcome, len(sleeps)), elapsed[0]) + sleeps.append(seconds) + arrival = ( + elapsed[0] >= 85 + if outcome == "final-poll" + else len(sleeps) == {"appears": 2, "oversleep": 3, "window": 1}.get(outcome) + ) + if arrival: + live_guard.backup.write_text(json.dumps(live_guard.receipt) + "\n") + + monkeypatch.setattr(module, "time", SimpleNamespace(monotonic=lambda: elapsed[0], sleep=sleep)) + lock = module._maintenance_lock + + @contextmanager + def tracked_lock(path): + locks.append(path) + with lock(path): + yield + + monkeypatch.setattr(module, "_maintenance_lock", tracked_lock) + total = sum(t["rows"] for t in module._run(db, True, 100, module.PROVIDER_MODES[mode])["tables"].values()) + flags = ["--dry-run"] if outcome == "dry-run" else ["--allow-live-db", "--expect-rows", str(total)] + if outcome == "missing-count": + flags = ["--allow-live-db"] + result = CliRunner().invoke( + app, + [ + "scrub-at-rest", + "--db", + str(db.db_path), + "--providers", + mode, + "--wait-for-backup-seconds", + "0" if outcome == "zero" else "90", + *flags, + ], + ) + if outcome in {"appears", "final-poll", "dry-run"}: + assert result.exit_code == 0, result.output + assert len(sleeps) == {"appears": 2, "final-poll": 5, "dry-run": 0}[outcome] + assert len(locks) == (0 if outcome == "dry-run" else 1) + else: + assert result.exit_code == 1, result.output + expected = "verified-backup-required" if outcome == "zero" else "verified-backup-timeout" + if outcome == "missing-count": + expected = "expected-row-count-required" + assert json.loads(result.stdout)["reason"] == expected + if outcome != "zero": + assert not locks + assert not any(isinstance(e, tuple) for e in live_guard.events) + assert elapsed[0] == {"timeout": 90, "oversleep": 91, "window": 1.5, "zero": 0, "missing-count": 0}[outcome] + + +@pytest.mark.parametrize("guarded", [False, True]) +def test_backup_wait_does_not_sleep_for_offline_apply_or_existing_receipt(db, live_guard, monkeypatch, guarded): + from brainlayer import chunk_origin_wipe + from brainlayer import scrub_at_rest as module + + if not guarded: + live_guard.backup.unlink() + monkeypatch.setattr(chunk_origin_wipe, "_live_db_candidates", lambda: []) + monkeypatch.setattr( + module, + "time", + SimpleNamespace(monotonic=lambda: 0, sleep=lambda _: pytest.fail("unexpected backup wait")), + ) + result = module.scrub_at_rest( + db.db_path, + allow_live_db=guarded, + expect_rows=live_guard.total, + wait_for_backup_seconds=90, + ) + assert result["tables"]["chunks"]["rows"] == 1 + + +@pytest.mark.parametrize("state", ["before-window", "after-window", "within-reserve", "missing-pause"]) +def test_backup_wait_refuses_without_sleep_for_unavailable_window_or_other_gate(db, live_guard, monkeypatch, state): + from brainlayer import maintenance + from brainlayer import scrub_at_rest as module + + live_guard.backup.unlink() + if state == "missing-pause": + live_guard.pause.unlink() + monkeypatch.setenv("BRAINLAYER_MCP_SOCKET", str(db.db_path.parent / "absent.sock")) + monkeypatch.setenv("BRAINLAYER_FORBID_BRAINBAR_SOCKET", "1") + start = { + "before-window": live_guard.now.replace(hour=3, minute=50), + "after-window": live_guard.now.replace(hour=7, minute=50), + "within-reserve": live_guard.now.replace(hour=5, minute=50), + }.get(state, live_guard.now) + factory = maintenance.MaintenanceConfig + elapsed, sleeps = [0.0], [] + + def config(**kwargs): + result = factory(**kwargs) + result.now_fn = lambda: start + dt.timedelta(seconds=elapsed[0]) + return result + + def sleep(seconds): + sleeps.append(seconds) + elapsed[0] += 1801 + + monkeypatch.setattr(maintenance, "MaintenanceConfig", config) + monkeypatch.setattr(module, "time", SimpleNamespace(monotonic=lambda: elapsed[0], sleep=sleep)) + monkeypatch.setattr(module, "_maintenance_lock", lambda _: pytest.fail("lock taken before refusal")) + reason = "enrichment-pause-required" if state == "missing-pause" else "verified-backup-required" + with pytest.raises(module.ScrubAtRestError) as error: + module.scrub_at_rest(db.db_path, allow_live_db=True, expect_rows=live_guard.total, wait_for_backup_seconds=1800) + assert error.value.reason == reason + assert sleeps == [] + assert not any(isinstance(e, tuple) for e in live_guard.events)