"""Recovery drills using installed packages and real bwrap, without live services. Run with the selected runtime's python -I -B. Queue and rein authoring are fixtures. This does not prove production API expiry or paid model execution. """ from __future__ import annotations import argparse from dataclasses import replace from datetime import UTC, datetime, timedelta import hashlib import importlib import json import os from pathlib import Path import selectors import signal import shutil import subprocess import sys import tempfile from unittest.mock import patch from glas_harness.contract import ExecutionSummary, OperationalReadiness, Rein, ToolResult from glas_harness.gateway import run_execution from glas_harness.profiles import ProfileCatalog from glas_harness.transport import transport_from_sandbox from rein_aharness.claim_loop import process_one, run_claim_loop from rein_aharness.close_outbox import CloseOutbox from rein_aharness.execution_cancel import active_cancel from rein_aharness.metrics import external_metrics_dir from rein_aharness.ops_run_client import OpsRun, OpsRunConfig, OpsRunError from rein_aharness.repository_grant import RepositoryGrant from rein_aharness.repository_transaction import RepositoryBusyError, RepositoryTransaction from sandboxer.core.manager import SandboxManager from sandboxer.extensions.runtime import runtime_digest from sandboxer.extensions.registry import extensions_dir from sandboxer.lifecycle.store import SandboxStore from sandboxer.profiles.loader import profiles_dir def git(repo, *args): return subprocess.check_output(["git", "-C", str(repo), *args], text=True).strip() def repository(root): root.mkdir() git(root, "init", "-q") (root / "README.md").write_text("Disposable recovery fixture.\n") git(root, "add", "README.md") git(root, "-c", "user.name=Fixture", "-c", "user.email=fixture@example.invalid", "commit", "-qm", "baseline") return git(root, "rev-parse", "HEAD") class FixtureQueue: """Explicit local queue double; never creates an Activity Core client.""" def __init__(self, source, mode, profile_ref): self.mode = mode self.config = OpsRunConfig(worker_id="recovery-fixture", repo_roots=(str(source.parent),)) self.run = OpsRun( id="recovery-" + mode, activity_definition_id="disposable-recovery", idempotency_key=mode, target_repo=str(source), title="recovery fixture", description="", state="claimed", claim_owner=self.config.worker_id, attempt=1, lease_until=(datetime.now(UTC) + timedelta(minutes=5)).isoformat(), harness_profile_ref=profile_ref, repository_grant=RepositoryGrant("1", ("result.txt",), 1, 1, False), ) self.claims = self.heartbeats = self.closes = 0 self.reject_heartbeat = False self.closed = None def claim(self, **kwargs): self.claims += 1 return [self.run] if self.claims == 1 else [] def list_open(self, **kwargs): return [] def heartbeat(self, *args, **kwargs): self.heartbeats += 1 if self.reject_heartbeat: raise OpsRunError("fixture lease expired", status_code=409, action="heartbeat", code="expired_lease") return self.run def _close(self, action, result): self.closes += 1 intent = {"action": action, "result": result} if self.closed is not None: assert self.closed == intent, "replay changed the terminal intent" self.closed = intent if self.mode == "lost-close" and self.closes == 1: raise OpsRunError("fixture response lost after apply", action=action) return replace(self.run, state="succeeded" if action == "complete" else "failed", close_disposition="reconciled" if self.closes > 1 else "applied") def complete(self, run_id, *, result): return self._close("complete", result) def fail(self, run_id, *, error, reopen, result): assert not reopen return self._close("fail", result) class FixtureRein(Rein): def __init__(self, source, baseline, queue): self.source, self.baseline, self.queue = source, baseline, queue self.calls = 0 self.workspace = None self.head = None def start_session(self, profile, inputs, sandbox): transport = transport_from_sandbox(sandbox) self.workspace = Path(transport.workspace) return transport def dispatch_tool(self, transport, tool_call): self.calls += 1 if self.queue.mode == "worker-sigkill": # Kill only this disposable child after the worker has acquired its # repository lock and Glas has persisted a ready sandbox. os.kill(os.getpid(), signal.SIGKILL) if self.queue.mode == "execution-failure": result = transport.run(["python3", "-c", "raise SystemExit(7)"], timeout=10) assert result.returncode == 7 raise RuntimeError("synthetic child execution failure") if self.queue.mode == "lease-loss": self.queue.reject_heartbeat = True assert active_cancel().wait(5), "periodic heartbeat did not cancel execution" elif self.queue.mode == "sigterm-before-import": # run_claim_loop installed the production signal handler in this # disposable proof process. No signal targets the standing worker. assert callable(signal.getsignal(signal.SIGTERM)) os.kill(os.getpid(), signal.SIGTERM) assert active_cancel().cancelled # Deliberately return a late successful commit after cancellation. The # worker must reject it and leave the source baseline untouched. result = transport.run(["python3", "-c", """from pathlib import Path import subprocess Path('result.txt').write_text('disposable recovery result\\n') def git(*args): return subprocess.check_output(['git', *args], text=True).strip() git('add', 'result.txt') git('-c', 'user.name=Fixture', '-c', 'user.email=fixture@example.invalid', 'commit', '-qm', 'result') print(git('rev-parse', 'HEAD')) """], timeout=30) assert result.returncode == 0, "sandbox authoring failed" self.head = result.stdout.strip() assert git(self.source, "rev-parse", "HEAD") == self.baseline return ToolResult(ok=True, output="disposable fixture", tokens_spent=0) def end_session(self, session): return ExecutionSummary(committed=True, commit_sha=self.head, outcome="succeeded", tokens_spent=0, cost_usd=0.0) def cleanup_session(self, session): assert self.workspace.exists() assert git(self.source, "rev-parse", "HEAD") == self.baseline def case(root, mode, profile_ref): root.mkdir() source = root / "source" baseline = repository(source) os.environ["REIN_AHARNESS_STATE_DIR"] = str(root / "state") os.environ["XDG_DATA_HOME"] = str(root / "data") queue = FixtureQueue(source, mode, profile_ref) catalog = ProfileCatalog() profile, _ = catalog.resolve(profile_ref) assert profile.sandbox_profile == "profile.bwrap-local", "requires empty-egress bwrap fixture" catalog.profiles()[(profile.id, profile.version)] = profile.model_copy(update={ "operational_readiness": OperationalReadiness(status="ready", reason="disposable fixture", owner="tests", evidence_ref="test:recovery"), }) store = SandboxStore(root / "sandboxes.json") manager = SandboxManager(store=store) rein = FixtureRein(source, baseline, queue) outbox = CloseOutbox(state_dir=root / "state") def gateway(request, **kwargs): return run_execution(request, catalog=catalog, rein=rein, manager=manager, **kwargs) try: with patch("glas_harness.gateway.run_execution", gateway), patch( "rein_aharness.claim_loop._heartbeat_interval", return_value=0.05 ): if mode == "sigterm-before-import": results = [] def record_cycle(*args, **kwargs): result = process_one(*args, outbox=outbox, **kwargs) results.append(result) return result handlers = {sig: signal.getsignal(sig) for sig in (signal.SIGINT, signal.SIGTERM)} try: with patch("rein_aharness.claim_loop.ActivityCoreOpsClient", return_value=queue), patch( "rein_aharness.claim_loop.process_one", record_cycle ): assert run_claim_loop(once=True, report_to_hub=False) == 1 finally: for sig, handler in handlers.items(): signal.signal(sig, handler) assert len(results) == 1 result = results[0] else: result = process_one(queue, outbox=outbox, report_to_hub=False) assert rein.calls == 1 and rein.workspace is not None, result.reason assert not rein.workspace.exists(), "sandbox workspace remains" statuses = store.list_all() assert len(statuses) == 1 and statuses[0].state.value == "destroyed" assert not git(source, "status", "--porcelain"), "source left dirty" with RepositoryTransaction(source) as transaction: assert transaction.locked, "repository lock leaked" metrics = [json.loads(line) for line in (external_metrics_dir(source, "rein-aharness") / "executions.jsonl").read_text().splitlines()] assert len(metrics) == 1 if mode == "lost-close": assert result.reason == "close evidence remains pending", result.reason assert outbox.status()["pending"] == 1 and metrics[0]["success"] assert git(source, "rev-parse", "HEAD") == rein.head # Construct a fresh store, as after worker restart. Replay must not # dispatch anything and must retain the exact accepted commit. restarted = CloseOutbox(state_dir=root / "state") with patch("rein_aharness.claim_loop.execute_profiled_run", side_effect=AssertionError("replayed workload")): replay = process_one(queue, outbox=restarted, report_to_hub=False) assert replay.empty and queue.closes == 2 and rein.calls == 1 assert restarted.status() == {"pending": 0, "delivered": 1, "quarantined": 0} assert git(source, "rev-list", "--count", "HEAD") == "2" else: assert result.ok is False and not metrics[0]["success"] assert git(source, "rev-parse", "HEAD") == baseline assert not (source / "result.txt").exists() if mode == "lease-loss": assert queue.closes == 0 and "lease lost" in result.reason else: assert queue.closed["action"] == "fail" assert outbox.status()["delivered"] == 1 return {"case": mode, "ok": True, "dispatches": rein.calls, "close_attempts": queue.closes, "outbox": outbox.status(), "sandbox_destroyed": True, "workspace_removed": True, "lock_reacquired": True, "source_clean": True, "source_commit_count": int(git(source, "rev-list", "--count", "HEAD")), "metrics_records": len(metrics), "metrics_success": metrics[0]["success"]} finally: # Cleanup only sandboxes created in this private fixture store. for status in store.list_all(): if status.state.value != "destroyed": manager.destroy(status.sandbox_id) def killed_lock_owner(root): source = root / "lock-source" repository(source) child = subprocess.Popen([sys.executable, "-I", "-B", "-c", """import sys,time from pathlib import Path from rein_aharness.repository_transaction import RepositoryTransaction with RepositoryTransaction(Path(sys.argv[1])): print('locked', flush=True) time.sleep(30) """, str(source)], stdout=subprocess.PIPE, text=True) try: with selectors.DefaultSelector() as selector: selector.register(child.stdout, selectors.EVENT_READ) assert selector.select(10), "lock child did not start" assert child.stdout.readline().strip() == "locked" try: with RepositoryTransaction(source): raise AssertionError("second process acquired held lock") except RepositoryBusyError: pass child.kill() child.wait(timeout=5) with RepositoryTransaction(source) as transaction: assert transaction.locked return {"ok": True, "held_lock_refused": True, "sigkill_lock_reacquired": True} finally: if child.poll() is None: child.kill() child.wait(timeout=5) child.stdout.close() def killed_worker(root, runtime, checksum, profile_ref): child_root = root / "worker-sigkill" try: child = subprocess.run([ sys.executable, "-I", "-B", str(Path(__file__).resolve()), "--runtime", str(runtime), "--sha256", checksum, "--profile-ref", profile_ref, "--crash-root", str(child_root), ], capture_output=True, text=True, timeout=60) assert child.returncode == -signal.SIGKILL, "worker did not reach the intended crash point" os.environ["XDG_DATA_HOME"] = str(child_root / "data") os.environ["REIN_AHARNESS_STATE_DIR"] = str(child_root / "state") store = SandboxStore(child_root / "sandboxes.json") manager = SandboxManager(store=store) statuses = store.list_all() assert len(statuses) == 1, "crashed worker did not persist sandbox ownership" status = statuses[0] workspace = Path(status.inputs["workspace_dir"]) # Explicit owner recovery, not a claim of automatic startup sweeping. manager.destroy(status.sandbox_id) assert store.get(status.sandbox_id).state.value == "destroyed" assert not workspace.exists() source = child_root / "source" assert not git(source, "status", "--porcelain") assert git(source, "rev-list", "--count", "HEAD") == "1" with RepositoryTransaction(source) as transaction: assert transaction.locked assert CloseOutbox(state_dir=child_root / "state").status() == { "pending": 0, "delivered": 0, "quarantined": 0, } return {"ok": True, "worker_sigkill": True, "owner_cleanup": "explicit", "sandbox_destroyed": True, "workspace_removed": True, "lock_reacquired": True, "source_clean": True, "accepted_commits": 0} finally: # A timeout or failed assertion must still clean any persisted fixture # sandbox. subprocess.run kills and reaps a timed-out child first. os.environ["XDG_DATA_HOME"] = str(child_root / "data") if (child_root / "sandboxes.json").exists(): cleanup_store = SandboxStore(child_root / "sandboxes.json") cleanup_manager = SandboxManager(store=cleanup_store) for remaining in cleanup_store.list_all(): if remaining.state.value != "destroyed": cleanup_manager.destroy(remaining.sandbox_id) def main(): parser = argparse.ArgumentParser(description=__doc__) parser.add_argument("--runtime", type=Path, required=True) parser.add_argument("--sha256", required=True) parser.add_argument("--profile-ref", required=True) parser.add_argument("--crash-root", type=Path, help=argparse.SUPPRESS) args = parser.parse_args() runtime = args.runtime.resolve() assert sys.flags.isolated and sys.dont_write_bytecode, "use candidate python -I -B" assert runtime_digest(runtime) == args.sha256, "runtime digest mismatch" imports = {} for name in ("rein_aharness", "glas_harness", "sandboxer", "llm_connect"): path = Path(importlib.import_module(name).__file__).resolve() assert path.is_relative_to(runtime), "source fallback" imports[name] = str(path.relative_to(runtime)) catalog = ProfileCatalog() for definitions in (catalog.profile_dir, catalog.rein_dir, profiles_dir(), extensions_dir()): assert Path(definitions).resolve().is_relative_to(runtime), "external definition fallback" os.environ["SANDBOXER_NO_STATE_HUB"] = "1" if args.crash_root is not None: case(args.crash_root, "worker-sigkill", args.profile_ref) raise AssertionError("crash child survived") root = Path(tempfile.mkdtemp(prefix="installed-recovery-")) try: results = [case(root / mode, mode, args.profile_ref) for mode in ("lost-close", "execution-failure", "lease-loss", "sigterm-before-import")] lock_result = killed_lock_owner(root) crash_result = killed_worker(root, runtime, args.sha256, args.profile_ref) except BaseException: print(f"Failed recovery fixture retained at {root}", file=sys.stderr) raise else: shutil.rmtree(root) assert runtime_digest(runtime) == args.sha256, "runtime mutated" print(json.dumps({"ok": True, "runtime_sha256": args.sha256, "proof_script_sha256": hashlib.sha256(Path(__file__).read_bytes()).hexdigest(), "scope": "installed worker/Glas/bwrap; synthetic queue and deterministic rein", "profile_ref": args.profile_ref, "imports": imports, "cases": results, "packaged_definitions": True, "killed_lock_owner": lock_result, "artifact_unchanged": True, "killed_worker": crash_result, "paid_requests": 0, "production_queue_calls": 0, "limitations": ["not a live Activity Core expiry proof", "hard-crash sandbox recovery is explicit owner cleanup"]}, indent=2)) if __name__ == "__main__": main()