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
7 changes: 7 additions & 0 deletions docs/configuration.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
7 changes: 7 additions & 0 deletions src/brainlayer/cli/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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(
Expand Down
38 changes: 36 additions & 2 deletions src/brainlayer/scrub_at_rest.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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"}
| {
Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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:
Expand All @@ -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)
Comment thread
coderabbitai[bot] marked this conversation as resolved.
with _maintenance_lock(path):
path = assert_not_live_db(path, allow_live=allow_live_db)
if allow_live_db:
Expand Down
141 changes: 141 additions & 0 deletions tests/test_scrub_at_rest_command.py
Original file line number Diff line number Diff line change
Expand Up @@ -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):

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

`test_backup_wait_before_lock_and_quiesce` has a cyclomatic complexity of 16 with "high" risk


A function with high cyclomatic complexity can be hard to understand and
maintain. Cyclomatic complexity is a software metric that measures the number of
independent paths through a function. A higher cyclomatic complexity indicates
that the function has more decision points and is more complex.

Copy link
Copy Markdown
Owner Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Addressed at 9ab89d2 by making three scenario selections data-driven, removing four decision points while retaining every case and assertion. No analyzer suppression or threshold change. All 330 focused tests and the full filtered pre-push gate pass at this head.

— brainlayerCodex-2c941903 (worker) · codex/gpt-6.1-sol

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

`test_backup_wait_before_lock_and_quiesce` has a cyclomatic complexity of 18 with "high" risk


A function with high cyclomatic complexity can be hard to understand and
maintain. Cyclomatic complexity is a software metric that measures the number of
independent paths through a function. A higher cyclomatic complexity indicates
that the function has more decision points and is more complex.

Copy link
Copy Markdown
Owner Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Addressed at 9ab89d2 by making three scenario selections data-driven, removing four decision points while retaining every case and assertion. No analyzer suppression or threshold change. All 330 focused tests and the full filtered pre-push gate pass at this head.

— brainlayerCodex-2c941903 (worker) · codex/gpt-6.1-sol

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)
Loading