diff --git a/loopx/capabilities/decision_context/README.md b/loopx/capabilities/decision_context/README.md index 44311fc2ec..fe51f17da3 100644 --- a/loopx/capabilities/decision_context/README.md +++ b/loopx/capabilities/decision_context/README.md @@ -312,10 +312,9 @@ substitute capture cursors for that file or manually manufacture reviewed cursor If several settlements occur between ticks, their intermediate transitions may be unobservable; ambiguous batches are retained, not inferred to be reviewed. An older spool without review observations is baselined without retiring rows. -For either hold, reconcile against actual review evidence explicitly; if starting -a new spool after a current-source rebase, retain the old spool as a private -checkpoint. This conservative protocol does not promise automatic queue drainage -after skipped review transitions. +For either hold, reconcile against actual review evidence explicitly, or use the +guarded recovery below. This conservative protocol does not promise automatic +queue drainage after skipped review transitions. This is a change-reference spool, **not a lossless source archive**. First-scan history, pagination, late edits, deletion visibility and deadlines remain provider @@ -329,6 +328,68 @@ Existing reviewable batches remain private and can still be prepared. To roll back to an older release, also remove the three new automation fields; retain the spool as a private checkpoint rather than deleting unreviewed work. +#### Recover an unreplayable source without discarding history + +Recovery belongs to this capability's existing private SQLite spool, not Core +Goal lifecycle. The local host/CLI is the operator surface; no provider, model, +remote write permission, scheduler or dashboard setting is added. A trusted +host may call `diagnose_capture_source` / `recover_capture_source` from +`loopx.capabilities.decision_context.capture_recovery` with its existing +`source_provider_overrides`. + +1. Run `capture-diagnose` with the same `--goal-id`, `--agent-id`, `--profile`, + `--spool`, `--cursor-state` arguments and an explicit `--source-id`. + Metadata-only diagnostics never contact providers. Add `--probe` to perform + one bounded transient replay/exact-read check, without semantic review or + pending settlement. Results distinguish `replay_not_checked`, `replayable`, + `revision_unavailable`, `binding_changed`, `cursor_diverged`, + `probe_unavailable`, `state_changed`, `empty` and `acquisition_held`. +2. Preview `capture-recovery --action hold` with that same scope. It shows + affected counts and an opaque `preview_token`. Apply only with explicit + operator authorization, `--execute --expected-token `. + All pending references for that source move atomically to **held, unresolved + history**, byte-for-byte, and acquisition for that source pauses. They are + not marked reviewed. Other sources can use the released active capacity. +3. When ready to read current material, preview and apply `--action restart` + with a fresh token. This rebinds acquisition to the current profile and the + **unchanged settlement-owned reviewed cursor**, clears the source's interval + wait, and removes its acquisition hold. The next ordinary capture produces + a fresh batch for `prepare-captured` and normal `settle-review`. Older held + references remain unresolved, even after the new batch is reviewed. +4. `--action rollback --recovery-id ` also requires preview + and explicit apply. It restores that operation's prior source state only if + the source scope and profile/reviewed files still match its receipt and the + active capacity permits restoration. After new capture/review, it fails + closed rather than overwriting progress. The applied/rolled-back receipts + remain in the private spool; copy/export the spool securely for inspection. + +All apply tokens bind action, source, profile/binding, spool identity and full +queue frontier, reviewed-file digest/identity and audit state. A concurrent +capture or review, duplicate apply, rebind or file ABA requires a fresh preview. +SQLite serializes capture/recovery; apply uses the settlement cursor lock. +Arbitrary manual edits are unsupported. Recovery never writes reviewed cursors, +so it does not invalidate or supersede a legitimate separate review settlement. +It is an acquisition restart, **not a claim to have created a new review epoch**. + +The existing `max_pending_batches=N` still bounds active batches. Held history +has a separate cap of N; new hold/restart operations stop after 2N audit records +(at most one rollback per applicable receipt). No automatic eviction, compaction +or repeated capacity increase is performed. At that bound, preserve/export the +private spool and make an explicit retention decision; increasing the limit is +not evidence consumption. A hold is explicit, not an automatic fairness policy. +It can isolate a noisy source, but exhaustion can recur if other sources are +not reviewed. Backpressured sources now retry on the next tick when capacity is +available instead of waiting an additional scan interval. + +`capture-status` separates active `pending_batch_count`, unresolved +`held_batch_count`, per-source `acquisition_held` and +`semantic_review_completion=not_inferred_from_capture`. `last_checked_at` is +the last attempt, not necessarily a successful scan; host service liveness and +successful-scan timestamps remain separate. No status-only call proves historical +replay or complete decision coverage. Disable capture using the existing profile +switch; stop the scheduler before downgrading, since older runtimes do not honor +recovery holds. Retain the spool/receipts rather than treating downgrade as rollback. + ## Relationship To Other Capabilities | Capability | Primary question | Relationship | diff --git a/loopx/capabilities/decision_context/README.zh-CN.md b/loopx/capabilities/decision_context/README.zh-CN.md index 572f21828c..c480df80e6 100644 --- a/loopx/capabilities/decision_context/README.zh-CN.md +++ b/loopx/capabilities/decision_context/README.zh-CN.md @@ -273,8 +273,8 @@ loopx decision-context prepare-captured --goal-id --agent-id --agent-id ` 显式应用。该来源所有待审阅引用 + 原样转入 **held、未解决历史**,暂停其采集,释放活跃队列容量给其他来源。 + 这不是审阅完成,不修改 reviewed cursor,也不丢弃旧证据。 +3. 准备读取当前材料时,重新预览并应用 `--action restart`。它以当前 profile + binding 和 **未改动的审阅游标** 重启采集,清除此来源的扫描间隔等待。 + 下一次正常 capture 产生新批次,再走 `prepare-captured` / `settle-review`。 + 新批次被审阅,也不会把旧 held 引用变成已审阅。 +4. `--action rollback --recovery-id <已应用回执 ID>` 同样需要预览和显式应用。 + 只有来源状态、profile/reviewed 文件仍匹配回执且恢复后不超过容量上限时, + 才恢复操作前状态;新采集或审阅后拒绝覆盖进展。回执保留在私有 spool 中。 + +预览令牌绑定 action、source、profile/binding、spool 身份、完整队列前沿、 +审阅文件内容与文件身份,以及审计记录。并发采集/审阅、重复应用、重绑或文件 ABA +必须重新预览。SQLite 串行化采集与恢复写入,恢复使用与 settlement 相同的游标锁。 +不支持手工改库绕过门禁;恢复不替代合法的独立审阅结算,也不宣称创建新审阅 epoch。 + +`max_pending_batches=N` 继续限制活跃批次;另最多保留 N 条未解决历史, +2N 条审计记录后停止新增 hold/restart(每条适用回执最多再 rollback 一次)。 +不会自动删除、压缩或无限扩容。达到上限需保留/导出私有 spool 后明确处理保留策略。 +这是显式来源隔离,不是默认公平调度;若其他来源长期不被审阅,仍可能再次背压。 +行为变化:曾背压的来源在容量释放后的下一 tick 可重试,不再多等一个扫描间隔。 + +`capture-status` 分开报告 active pending、held 历史、每来源 acquisition hold, +并明确 `semantic_review_completion=not_inferred_from_capture`。`last_checked_at` +是尝试时间,不保证成功;服务存活和最近成功扫描时间仍由 host 独立报告。 +仅看 status 不能证明历史可重放或决策覆盖完整。停用仍使用原 profile 开关; +降级旧版本前必须停止调度器,因为旧运行时不认识 recovery hold。 +保留 spool 与回执,不能把软件降级当成状态回滚。 + ## 与其他能力的关系 | 能力 | 核心问题 | 与 Decision Context 的关系 | diff --git a/loopx/capabilities/decision_context/capture.py b/loopx/capabilities/decision_context/capture.py index 6de880f727..6b3c64d747 100644 --- a/loopx/capabilities/decision_context/capture.py +++ b/loopx/capabilities/decision_context/capture.py @@ -28,6 +28,14 @@ from .sources import DecisionSourceProvider, DecisionSourceSpec +class CaptureReplayError(ValueError): + """Typed recovery diagnosis; never classify provider exception prose.""" + + def __init__(self, reason: str, message: str): + self.reason = reason + super().__init__(message) + + def _open_spool(path: Path, *, goal_id: str, agent_id: str) -> sqlite3.Connection: path.parent.mkdir(parents=True, exist_ok=True) descriptor = os.open( @@ -74,6 +82,10 @@ def _binding_digest(profile: DecisionContextProfile, source: DecisionSourceSpec) def _status(db: sqlite3.Connection, source_ids: tuple[str, ...]) -> dict[str, Any]: + has_recovery = ( + db.execute("SELECT 1 FROM sqlite_master WHERE name='capture_holds'").fetchone() + is not None + ) rows = [] for source_id in source_ids: source = db.execute( @@ -89,12 +101,30 @@ def _status(db: sqlite3.Connection, source_ids: tuple[str, ...]) -> dict[str, An "status": source["status"] if source else "never_checked", "pending_batch_count": pending[0], "next_batch_id": pending[1], + "held_batch_count": db.execute( + "SELECT count(*) FROM held_batches WHERE source_id=?", (source_id,) + ).fetchone()[0] + if has_recovery + else 0, + "acquisition_held": bool( + db.execute( + "SELECT 1 FROM capture_holds WHERE source_id=?", (source_id,) + ).fetchone() + ) + if has_recovery + else False, } ) return { "schema_version": "decision_context_capture_status_v0", "sources": rows, "pending_batch_count": db.execute("SELECT count(*) FROM batches").fetchone()[0], + "held_batch_count": db.execute("SELECT count(*) FROM held_batches").fetchone()[ + 0 + ] + if has_recovery + else 0, + "semantic_review_completion": "not_inferred_from_capture", "raw_content_captured": False, "decision_cursors_mutated": False, "external_writes_performed": False, @@ -180,6 +210,15 @@ def capture_profile_sources( try: reviewed = load_private_decision_cursors(cursor_path, profile=profile) for source in sources: + if ( + db.execute( + "SELECT 1 FROM sqlite_master WHERE name='capture_holds'" + ).fetchone() + and db.execute( + "SELECT 1 FROM capture_holds WHERE source_id=?", (source.source_id,) + ).fetchone() + ): + continue binding = _binding_digest(profile, source) row = db.execute( "SELECT * FROM sources WHERE source_id=?", (source.source_id,) @@ -217,6 +256,7 @@ def capture_profile_sources( if ( row and row["checked_at"] + and row["status"] != "backpressure" and (now - datetime.fromisoformat(row["checked_at"])).total_seconds() < profile.capture_interval_seconds ): @@ -333,10 +373,12 @@ def assemble_captured_decision_evidence( None, ) if source is None or binding != _binding_digest(profile, source): - raise ValueError("capture source binding changed") + raise CaptureReplayError("binding_changed", "capture source binding changed") reviewed = load_private_decision_cursors(cursor_path, profile=profile) if reviewed.get(source.source_id) != batch["cursor_before"]: - raise ValueError("capture batch must follow the reviewed cursor") + raise CaptureReplayError( + "cursor_diverged", "capture batch must follow the reviewed cursor" + ) expected = json.loads(batch["receipt"]) def checked_rebase(collection: DecisionEvidenceCollection): @@ -353,8 +395,9 @@ def changes(value): or receipt["status"] != expected["status"] or receipt["exact_read_count"] != expected["changed_count"] ): - raise ValueError( - "captured revision unavailable; explicit source rebase required" + raise CaptureReplayError( + "revision_unavailable", + "captured revision unavailable; explicit source rebase required", ) return rebase(collection) diff --git a/loopx/capabilities/decision_context/capture_recovery.py b/loopx/capabilities/decision_context/capture_recovery.py new file mode 100644 index 0000000000..dc3cb71590 --- /dev/null +++ b/loopx/capabilities/decision_context/capture_recovery.py @@ -0,0 +1,473 @@ +"""Explicit, reference-preserving capture recovery, not semantic settlement. + +The private spool is the existing owner. Holds free active queue capacity but +retain unresolved references; restart uses the *unchanged* reviewed cursor. +Neither action acknowledges historical evidence or creates review authority. +""" + +from __future__ import annotations + +import hashlib +import json +import sqlite3 +from contextlib import closing, nullcontext +from collections.abc import Mapping +from datetime import datetime, timezone +from enum import Enum +from pathlib import Path +from typing import Any +from uuid import uuid4 + +from ...file_lock import exclusive_file_lock +from .assembler import DecisionEvidenceRecords +from .capture import ( + CaptureReplayError, + _binding_digest, + _open_spool, + assemble_captured_decision_evidence, +) +from .private_state import load_private_decision_cursors, private_file_digest +from .profile import DecisionContextProfile, resolve_decision_context_activation +from .sources import DecisionSourceProvider + + +class RecoveryAction(str, Enum): + HOLD = "hold" + RESTART = "restart" + ROLLBACK = "rollback" + + +def _hash(value: Any) -> str: + return hashlib.sha256(json.dumps(value, sort_keys=True).encode()).hexdigest() + + +def _file_fence(path: Path | None) -> dict[str, Any] | None: + if path is None: + return None + path = path.expanduser().resolve() + if not path.exists(): + return {"path": _hash(str(path)), "exists": False} + stat = path.stat() + return { + "path": _hash(str(path)), + "identity": [stat.st_dev, stat.st_ino], + "mtime_ns": stat.st_mtime_ns, + "ctime_ns": stat.st_ctime_ns, + "digest": private_file_digest(path), + } + + +def _rows( + db: sqlite3.Connection, table: str, source_id: str | None = None +) -> list[dict[str, Any]]: + # All table names are module-owned literals, never caller input. + if not db.execute("SELECT 1 FROM sqlite_master WHERE name=?", (table,)).fetchone(): + return [] + query = f"SELECT * FROM {table}" + parameters: tuple[str, ...] = () + if source_id is not None: + query += " WHERE source_id=?" + parameters = (source_id,) + return [dict(row) for row in db.execute(query + " ORDER BY rowid", parameters)] + + +def _scope(db: sqlite3.Connection, source_id: str) -> dict[str, Any]: + return { + name: _rows(db, name, source_id) + for name in ( + "sources", + "batches", + "review_observations", + "capture_holds", + "held_batches", + ) + } + + +def _fence( + db: sqlite3.Connection, + spool_path: Path, + profile_path: Path, + cursor_path: Path | None, +) -> str: + stat = spool_path.stat() + return _hash( + { + "spool_identity": [str(spool_path.resolve()), stat.st_dev, stat.st_ino], + "profile": _file_fence(profile_path), + "reviewed": _file_fence(cursor_path), + "tables": { + name: _rows(db, name) + for name in ( + "identity", + "sources", + "batches", + "review_observations", + "capture_holds", + "held_batches", + "capture_recoveries", + "sqlite_sequence", + ) + }, + } + ) + + +def _schema(db: sqlite3.Connection) -> None: + # Do not executescript here: it would commit the caller's fence transaction. + db.execute( + "CREATE TABLE IF NOT EXISTS capture_holds " + "(source_id TEXT PRIMARY KEY, recovery_id TEXT NOT NULL)" + ) + db.execute( + "CREATE TABLE IF NOT EXISTS held_batches " + "(id INTEGER PRIMARY KEY, source_id TEXT NOT NULL, cursor_before TEXT, " + "cursor_after TEXT NOT NULL, before_time TEXT NOT NULL, receipt TEXT NOT NULL, " + "recovery_id TEXT NOT NULL)" + ) + db.execute( + "CREATE TABLE IF NOT EXISTS capture_recoveries " + "(id TEXT PRIMARY KEY, source_id TEXT NOT NULL, action TEXT NOT NULL, " + "recorded_at TEXT NOT NULL, before_state TEXT NOT NULL, " + "after_digest TEXT NOT NULL, file_fence TEXT NOT NULL, " + "rollback_of TEXT, rolled_back INTEGER NOT NULL DEFAULT 0)" + ) + + +def _resolve( + goal_id: str, + agent_id: str, + profile_path: Path, + overrides: Mapping[str, DecisionSourceProvider], +) -> tuple[DecisionContextProfile, dict[str, Any] | None]: + before = _file_fence(profile_path) + activation, profile = resolve_decision_context_activation( + goal_id=goal_id, + agent_id=agent_id, + profile_path=profile_path, + available_source_provider_ids=overrides, + ) + if ( + profile is None + or not profile.enabled + or not profile.automatic_capture + or not activation["configured_for_agent"] + ): + raise ValueError("capture recovery requires an enabled capture profile") + if _file_fence(profile_path) != before: + raise ValueError("capture profile changed during recovery activation") + return profile, before + + +def _read_db(spool_path: Path, goal_id: str, agent_id: str) -> sqlite3.Connection: + db = sqlite3.connect(spool_path.resolve().as_uri() + "?mode=ro", uri=True) + db.row_factory = sqlite3.Row + try: + db.execute("BEGIN") + row = db.execute("SELECT goal, agent FROM identity").fetchone() + if row is None or tuple(row) != (goal_id, agent_id): + raise ValueError("capture spool goal/agent mismatch") + return db + except BaseException: + db.close() + raise + + +def diagnose_capture_source( + *, + goal_id: str, + agent_id: str, + profile_path: Path, + spool_path: Path, + cursor_path: Path | None, + source_id: str, + probe: bool = False, + source_provider_overrides: Mapping[str, DecisionSourceProvider] | None = None, +) -> dict[str, Any]: + """Metadata by default; explicit probe exact-reads transiently, never settles.""" + overrides = source_provider_overrides or {} + profile, profile_fence = _resolve(goal_id, agent_id, profile_path, overrides) + cursor_fence = _file_fence(cursor_path) + source = next((s for s in profile.sources if s.source_id == source_id), None) + if source is None or source_id not in profile.capture_source_ids: + raise ValueError("source is not enrolled for capture") + reviewed = load_private_decision_cursors(cursor_path, profile=profile) + with closing(_read_db(spool_path, goal_id, agent_id)) as db: + state = _scope(db, source_id) + fence = _fence(db, spool_path, profile_path, cursor_path) + oldest = state["batches"][0] if state["batches"] else None + reason = "empty" + if state["capture_holds"]: + reason = "acquisition_held" + elif oldest: + reason = "replay_not_checked" + if state["sources"][0]["binding_digest"] != _binding_digest(profile, source): + reason = "binding_changed" + elif reviewed.get(source_id) != oldest["cursor_before"]: + reason = "cursor_diverged" + elif probe: + try: + assemble_captured_decision_evidence( + goal_id=goal_id, + agent_id=agent_id, + profile_path=profile_path, + spool_path=spool_path, + cursor_path=cursor_path, + batch_id=oldest["id"], + decision_id="capture-recovery-probe", + rebase=lambda _: DecisionEvidenceRecords(), + source_provider_overrides=overrides, + ) + reason = "replayable" + except CaptureReplayError as exc: + reason = exc.reason + except Exception: + reason = "probe_unavailable" # No private provider exception text. + with closing(_read_db(spool_path, goal_id, agent_id)) as db: + if ( + _fence(db, spool_path, profile_path, cursor_path) != fence + or _file_fence(profile_path) != profile_fence + or _file_fence(cursor_path) != cursor_fence + ): + reason = "state_changed" + return { + "schema_version": "decision_capture_diagnosis_v0", + "source_id": source_id, + "diagnosis": reason, + "probe_requested": probe, + "next_batch_id": oldest["id"] if oldest else None, + "pending_batch_count": len(state["batches"]), + "held_batch_count": len(state["held_batches"]), + "decision_cursors_mutated": False, + "raw_content_captured": False, + } + + +def recover_capture_source( + *, + goal_id: str, + agent_id: str, + profile_path: Path, + spool_path: Path, + cursor_path: Path | None, + source_id: str, + action: RecoveryAction | str, + execute: bool = False, + expected_token: str | None = None, + recovery_id: str | None = None, + source_provider_overrides: Mapping[str, DecisionSourceProvider] | None = None, +) -> dict[str, Any]: + """Preview then explicitly apply one source operation with exact CAS fencing. + + Rollback is permitted only while the affected scope and review/profile files + still match the applied receipt. Other sources can progress independently. + This trusted local-host API is not a remote authorization endpoint. + """ + action = RecoveryAction(action) + if (action == RecoveryAction.ROLLBACK) != (recovery_id is not None): + raise ValueError("recovery-id is required only for rollback") + overrides = source_provider_overrides or {} + profile_path, spool_path = profile_path.expanduser(), spool_path.expanduser() + if not spool_path.is_file(): + raise ValueError("capture recovery requires an existing spool") + if cursor_path is not None: + cursor_path = cursor_path.expanduser() + paths = [profile_path.resolve(), spool_path.resolve()] + if cursor_path is not None: + paths.append(cursor_path.resolve()) + if len(set(paths)) != len(paths): + raise ValueError( + "recovery profile, spool and reviewed cursors must be separate" + ) + profile, profile_fence = _resolve(goal_id, agent_id, profile_path, overrides) + source = next((s for s in profile.sources if s.source_id == source_id), None) + if source is None or source_id not in profile.capture_source_ids: + raise ValueError("source is not enrolled for capture") + if execute and not expected_token: + raise ValueError("recovery execute requires an exact preview token") + # Same lock as settlement. SQLite serializes capture/recovery writes. A + # preview is fully read-only: no lock files, migrations or permission writes. + lock = ( + exclusive_file_lock(cursor_path) if execute and cursor_path else nullcontext() + ) + with lock: + db = ( + _open_spool(spool_path, goal_id=goal_id, agent_id=agent_id) + if execute + else _read_db(spool_path, goal_id, agent_id) + ) + with closing(db): + spool_stat = spool_path.stat() + spool_identity = (spool_stat.st_dev, spool_stat.st_ino) + files = { + "profile": _file_fence(profile_path), + "reviewed": _file_fence(cursor_path), + } + if files["profile"] != profile_fence: + raise ValueError("capture profile changed during recovery") + before = _scope(db, source_id) + reviewed = load_private_decision_cursors(cursor_path, profile=profile) + token = _hash( + { + "action": action.value, + "source_id": source_id, + "recovery_id": recovery_id, + "state": _fence(db, spool_path, profile_path, cursor_path), + } + ) + if execute and expected_token != token: + raise ValueError("capture recovery preview is stale; preview again") + old_receipt = None + pending_total = len(_rows(db, "batches")) + held_total = len(_rows(db, "held_batches")) + cap = profile.capture_max_pending_batches + if action == RecoveryAction.HOLD: + if not before["batches"] or before["capture_holds"]: + raise ValueError( + "hold requires pending batches on a non-held source" + ) + if held_total + len(before["batches"]) > cap: + raise ValueError( + "retained-history capacity reached; preserve/export spool before recovery" + ) + elif action == RecoveryAction.RESTART: + if not before["capture_holds"] or before["batches"]: + raise ValueError( + "restart requires a held source without active batches" + ) + else: + old_receipt = next( + ( + r + for r in _rows(db, "capture_recoveries") + if r["id"] == recovery_id and r["source_id"] == source_id + ), + None, + ) + if ( + not old_receipt + or old_receipt["rolled_back"] + or old_receipt["action"] == "rollback" + ): + raise ValueError("recovery receipt is not rollbackable") + if ( + old_receipt["after_digest"] != _hash(before) + or json.loads(old_receipt["file_fence"]) != files + ): + raise ValueError( + "source or reviewed/profile state changed since recovery" + ) + restored = json.loads(old_receipt["before_state"]) + if pending_total + len(restored["batch_ids"]) > cap: + raise ValueError("rollback would exceed active queue capacity") + # Bound audit metadata as well as active and unresolved references. + if ( + len(_rows(db, "capture_recoveries")) >= 2 * cap + and action != RecoveryAction.ROLLBACK + ): + raise ValueError( + "recovery audit capacity reached; preserve/export spool" + ) + result = { + "schema_version": "decision_capture_recovery_v0", + "source_id": source_id, + "action": action.value, + "executed": False, + "preview_token": token, + "pending_batch_count": len(before["batches"]), + "held_batch_count": len(before["held_batches"]), + "decision_cursors_mutated": False, + "historical_batches_reviewed": False, + "raw_content_captured": False, + "external_writes_performed": False, + } + if not execute: + return result + _schema(db) + new_id = str(uuid4()) + saved = { + key: before[key] + for key in ("sources", "review_observations", "capture_holds") + } + saved["batch_ids"] = [r["id"] for r in before["batches"]] + if action == RecoveryAction.HOLD: + db.execute( + "INSERT INTO held_batches SELECT id,source_id,cursor_before,cursor_after,before_time,receipt,? FROM batches WHERE source_id=?", + (new_id, source_id), + ) + db.execute("DELETE FROM batches WHERE source_id=?", (source_id,)) + db.execute( + "INSERT INTO capture_holds VALUES (?, ?)", (source_id, new_id) + ) + db.execute( + "UPDATE sources SET status='recovery_held' WHERE source_id=?", + (source_id,), + ) + elif action == RecoveryAction.RESTART: + db.execute("DELETE FROM capture_holds WHERE source_id=?", (source_id,)) + db.execute( + "UPDATE sources SET binding_digest=?, cursor=?, checked_at=NULL, status='recovery_restarted' WHERE source_id=?", + ( + _binding_digest(profile, source), + reviewed.get(source_id), + source_id, + ), + ) + db.execute( + "INSERT OR REPLACE INTO review_observations VALUES (?, ?)", + (source_id, reviewed.get(source_id)), + ) + else: + assert old_receipt is not None + restored = json.loads(old_receipt["before_state"]) + for table in ("sources", "review_observations", "capture_holds"): + db.execute(f"DELETE FROM {table} WHERE source_id=?", (source_id,)) + for row in restored[table]: + db.execute( + f"INSERT INTO {table} ({','.join(row)}) VALUES ({','.join('?' for _ in row)})", + tuple(row.values()), + ) + if restored["batch_ids"]: + db.execute( + "INSERT INTO batches SELECT id,source_id,cursor_before,cursor_after,before_time,receipt FROM held_batches WHERE recovery_id=?", + (recovery_id,), + ) + db.execute( + "DELETE FROM held_batches WHERE recovery_id=?", (recovery_id,) + ) + db.execute( + "UPDATE capture_recoveries SET rolled_back=1 WHERE id=?", + (recovery_id,), + ) + db.execute( + "INSERT INTO capture_recoveries VALUES (?,?,?,?,?,?,?,?,0)", + ( + new_id, + source_id, + action.value, + datetime.now(timezone.utc).isoformat(), + json.dumps(saved), + _hash(_scope(db, source_id)), + json.dumps(files), + recovery_id, + ), + ) + if { + "profile": _file_fence(profile_path), + "reviewed": _file_fence(cursor_path), + } != files: + raise ValueError( + "capture recovery profile/reviewed state changed during apply" + ) + spool_stat = spool_path.stat() + if (spool_stat.st_dev, spool_stat.st_ino) != spool_identity: + raise ValueError("capture spool replaced during recovery") + db.commit() + return { + **result, + "executed": True, + "recovery_id": new_id, + "pending_batch_count": len(_rows(db, "batches", source_id)), + "held_batch_count": len(_rows(db, "held_batches", source_id)), + "acquisition_held": bool(_rows(db, "capture_holds", source_id)), + } diff --git a/loopx/capabilities/decision_context/cli.py b/loopx/capabilities/decision_context/cli.py index 40da9cfc7e..fa4e5f6208 100644 --- a/loopx/capabilities/decision_context/cli.py +++ b/loopx/capabilities/decision_context/cli.py @@ -75,6 +75,35 @@ def register_decision_context_commands( help="Inspect the provider-neutral decision evidence and outcome contract.", ) commands = parser.add_subparsers(dest="decision_context_command", required=True) + for name in ("capture-diagnose", "capture-recovery"): + recovery = commands.add_parser( + name, help="Diagnose or explicitly recover one private capture source." + ) + add_subcommand_format(recovery) + for field in ( + "goal-id", + "agent-id", + "profile", + "spool", + "cursor-state", + "source-id", + ): + recovery.add_argument("--" + field, required=True) + if name == "capture-diagnose": + recovery.add_argument( + "--probe", + action="store_true", + help="Exact-read transiently to check replay; never settle.", + ) + else: + recovery.add_argument( + "--action", choices=("hold", "restart", "rollback"), required=True + ) + recovery.add_argument("--expected-token") + recovery.add_argument( + "--recovery-id", help="Applied receipt id to roll back." + ) + recovery.add_argument("--execute", action="store_true") for name in ("capture", "capture-status", "prepare-captured"): capture = commands.add_parser( name, help="Operate the opt-in private source-reference spool." @@ -236,7 +265,28 @@ def handle_decision_context_command( ) -> int | None: if args.command != "decision-context": return None - if args.decision_context_command in { + if args.decision_context_command in {"capture-diagnose", "capture-recovery"}: + from .capture_recovery import diagnose_capture_source, recover_capture_source + + recovery_args = dict( + goal_id=args.goal_id, + agent_id=args.agent_id, + profile_path=Path(args.profile), + spool_path=Path(args.spool), + cursor_path=Path(args.cursor_state), + source_id=args.source_id, + ) + if args.decision_context_command == "capture-diagnose": + payload = diagnose_capture_source(**recovery_args, probe=args.probe) + else: + payload = recover_capture_source( + **recovery_args, + action=args.action, + execute=args.execute, + expected_token=args.expected_token, + recovery_id=args.recovery_id, + ) + elif args.decision_context_command in { "capture", "capture-status", "prepare-captured", diff --git a/tests/capabilities/test_decision_context_capture_recovery.py b/tests/capabilities/test_decision_context_capture_recovery.py new file mode 100644 index 0000000000..2a7de72ae6 --- /dev/null +++ b/tests/capabilities/test_decision_context_capture_recovery.py @@ -0,0 +1,432 @@ +from __future__ import annotations + +import json +import shutil +import sqlite3 + +import pytest + +from loopx.capabilities.decision_context.capture import ( + capture_profile_sources, + assemble_captured_decision_evidence, +) +from loopx.capabilities.decision_context.capture_recovery import ( + recover_capture_source, + diagnose_capture_source, +) +from loopx.capabilities.decision_context.assembler import DecisionEvidenceRecords +from loopx.cli import main +from test_decision_context_capture import setup as capture_setup, settle_batch + + +@pytest.fixture +def setup(tmp_path): + return capture_setup.__wrapped__(tmp_path) + + +def recover(args, source_id, action, **kwargs): + preview = recover_capture_source( + **args, source_id=source_id, action=action, **kwargs + ) + return recover_capture_source( + **args, + source_id=source_id, + action=action, + execute=True, + expected_token=preview["preview_token"], + **kwargs, + ) + + +def source_id(payload): + return payload["sources"][0]["source_id"] + + +def rows(args, table): + with sqlite3.connect(args["spool_path"]) as db: + return db.execute(f"SELECT * FROM {table} ORDER BY rowid").fetchall() + + +def due(args): + with sqlite3.connect(args["spool_path"]) as db: + db.execute("UPDATE sources SET checked_at=NULL") + + +def test_hold_restart_exact_review_preserves_unresolved_history(setup): + args, payload, authority = setup + sid = source_id(payload) + capture_profile_sources(**args, execute=True) + original = rows(args, "batches") + authority.write_text("new-private-body") + assert ( + diagnose_capture_source(**args, source_id=sid)["diagnosis"] + == "replay_not_checked" + ) + assert ( + diagnose_capture_source(**args, source_id=sid, probe=True)["diagnosis"] + == "revision_unavailable" + ) + held = recover(args, sid, "hold") + assert held["held_batch_count"] == 1 and held["acquisition_held"] + assert [r[:-1] for r in rows(args, "held_batches")] == original + assert not args["cursor_path"].exists() + assert capture_profile_sources(**args, execute=True)["pending_batch_count"] == 0 + recover(args, sid, "restart") + captured = capture_profile_sources(**args, execute=True) + assert captured["pending_batch_count"] == 1 and captured["held_batch_count"] == 1 + assert not args["cursor_path"].exists() + with pytest.raises(ValueError, match="batch unavailable"): + assemble_captured_decision_evidence( + **args, + batch_id=1, + decision_id="old", + rebase=lambda _: DecisionEvidenceRecords(), + ) + settle_batch(args, captured["sources"][0]["next_batch_id"]) + after = capture_profile_sources(**args, execute=True) + assert after["pending_batch_count"] == 0 and after["held_batch_count"] == 1 + assert [r[:-1] for r in rows(args, "held_batches")] == original + assert b"new-private-body" not in args["spool_path"].read_bytes() + + +def test_one_held_source_no_longer_starves_unrelated_source(setup, tmp_path): + args, payload, authority = setup + sid = source_id(payload) + second_file = tmp_path / "second.txt" + second_file.write_text("independent") + second = dict( + payload["sources"][0], source_id="second", private_locator=str(second_file) + ) + payload["sources"].append(second) + payload["automation"].update(source_ids=[sid, "second"], max_pending_batches=1) + args["profile_path"].write_text(json.dumps(payload)) + captured = capture_profile_sources(**args, execute=True) + assert captured["sources"][1]["status"] == "backpressure" + authority.write_text("changed") + held = recover(args, sid, "hold") + after = capture_profile_sources(**args, execute=True) + assert after["pending_batch_count"] == 1 and after["held_batch_count"] == 1 + assert after["sources"][1]["pending_batch_count"] == 1 + assert after["sources"][0]["acquisition_held"] + assert not args["cursor_path"].exists() + with pytest.raises(ValueError, match="exceed active queue capacity"): + recover(args, sid, "rollback", recovery_id=held["recovery_id"]) + # No silent deletion or unbounded second archive when retention is full. + with pytest.raises(ValueError, match="retained-history capacity"): + recover(args, "second", "hold") + + +def test_preview_is_read_only_and_scoped_rollback_restores_exact_rows(setup): + args, payload, _ = setup + sid = source_id(payload) + capture_profile_sources(**args, execute=True) + original = rows(args, "batches") + original_sources = rows(args, "sources") + before = args["spool_path"].read_bytes() + preview = recover_capture_source(**args, source_id=sid, action="hold") + assert before == args["spool_path"].read_bytes() + assert not args["cursor_path"].with_suffix(".json.lock").exists() + held = recover_capture_source( + **args, + source_id=sid, + action="hold", + execute=True, + expected_token=preview["preview_token"], + ) + restored = recover(args, sid, "rollback", recovery_id=held["recovery_id"]) + assert not restored["acquisition_held"] + assert rows(args, "batches") == original + assert rows(args, "sources") == original_sources + assert rows(args, "held_batches") == [] + with pytest.raises(ValueError, match="not rollbackable"): + recover(args, sid, "rollback", recovery_id=held["recovery_id"]) + + +def test_restart_rollback_and_then_hold_rollback(setup): + args, payload, _ = setup + sid = source_id(payload) + capture_profile_sources(**args, execute=True) + held = recover(args, sid, "hold") + restarted = recover(args, sid, "restart") + recover(args, sid, "rollback", recovery_id=restarted["recovery_id"]) + assert capture_profile_sources(**args)["sources"][0]["acquisition_held"] + recover(args, sid, "rollback", recovery_id=held["recovery_id"]) + assert len(rows(args, "batches")) == 1 + + +@pytest.mark.parametrize( + "mutation", ["capture", "review", "aba", "profile", "rebind", "spool"] +) +def test_stale_preview_rejected_without_evidence_loss(setup, mutation, tmp_path): + args, payload, authority = setup + sid = source_id(payload) + capture_profile_sources(**args, execute=True) + preview = recover_capture_source(**args, source_id=sid, action="hold") + if mutation == "capture": + authority.write_text("next") + due(args) + capture_profile_sources(**args, execute=True) + elif mutation in {"review", "aba"}: + args["cursor_path"].write_text(json.dumps({sid: "B"})) + if mutation == "aba": + args["cursor_path"].write_text("{}") + elif mutation in {"profile", "rebind"}: + if mutation == "rebind": + payload["sources"][0]["private_locator"] += ".new" + else: + payload["automation"]["interval_seconds"] = 42 + args["profile_path"].write_text(json.dumps(payload)) + else: + replacement = tmp_path / "replacement.sqlite" + shutil.copyfile(args["spool_path"], replacement) + replacement.replace(args["spool_path"]) + before = rows(args, "batches") + with pytest.raises(ValueError, match="stale"): + recover_capture_source( + **args, + source_id=sid, + action="hold", + execute=True, + expected_token=preview["preview_token"], + ) + assert rows(args, "batches") == before + + +def test_profile_and_cursor_aba_returning_to_identical_bytes_is_rejected(setup): + args, payload, _ = setup + sid = source_id(payload) + capture_profile_sources(**args, execute=True) + args["cursor_path"].write_text("{}") + for path in (args["profile_path"], args["cursor_path"]): + preview = recover_capture_source(**args, source_id=sid, action="hold") + original = path.read_bytes() + path.write_text("changed") + path.write_bytes(original) + with pytest.raises(ValueError, match="stale"): + recover_capture_source( + **args, + source_id=sid, + action="hold", + execute=True, + expected_token=preview["preview_token"], + ) + + +def test_duplicate_apply_disabled_profile_and_absent_token_fail_closed(setup): + args, payload, _ = setup + sid = source_id(payload) + capture_profile_sources(**args, execute=True) + with pytest.raises(ValueError, match="preview token"): + recover_capture_source(**args, source_id=sid, action="hold", execute=True) + preview = recover_capture_source(**args, source_id=sid, action="hold") + recover_capture_source( + **args, + source_id=sid, + action="hold", + execute=True, + expected_token=preview["preview_token"], + ) + with pytest.raises(ValueError, match="stale"): + recover_capture_source( + **args, + source_id=sid, + action="hold", + execute=True, + expected_token=preview["preview_token"], + ) + payload["automation"]["automatic_capture"] = False + args["profile_path"].write_text(json.dumps(payload)) + with pytest.raises(ValueError, match="enabled capture"): + recover(args, sid, "restart") + + +@pytest.mark.parametrize("mutation", ["profile", "review", "capture"]) +def test_rollback_after_progress_rejected(setup, mutation): + args, payload, _ = setup + sid = source_id(payload) + capture_profile_sources(**args, execute=True) + held = recover(args, sid, "hold") + if mutation == "profile": + args["profile_path"].write_text(json.dumps(payload, indent=2)) + elif mutation == "review": + args["cursor_path"].write_text(json.dumps({sid: "new"})) + else: + recover(args, sid, "restart") + capture_profile_sources(**args, execute=True) + with pytest.raises(ValueError, match="state changed"): + recover(args, sid, "rollback", recovery_id=held["recovery_id"]) + assert len(rows(args, "held_batches")) == 1 + + +def test_binding_cursor_deletion_and_provider_failure_diagnostics(setup): + args, payload, authority = setup + sid = source_id(payload) + capture_profile_sources(**args, execute=True) + assert ( + diagnose_capture_source(**args, source_id=sid, probe=True)["diagnosis"] + == "replayable" + ) + authority.unlink() + assert ( + diagnose_capture_source(**args, source_id=sid, probe=True)["diagnosis"] + == "revision_unavailable" + ) + args["cursor_path"].write_text(json.dumps({sid: "unknown"})) + assert ( + diagnose_capture_source(**args, source_id=sid)["diagnosis"] == "cursor_diverged" + ) + payload["sources"][0]["private_locator"] += ".other" + args["profile_path"].write_text(json.dumps(payload)) + assert ( + diagnose_capture_source(**args, source_id=sid)["diagnosis"] == "binding_changed" + ) + recover(args, sid, "hold") + recover(args, sid, "restart") + assert ( + capture_profile_sources(**args, execute=True)["sources"][0]["status"] + != "binding_changed" + ) + + +def test_transaction_contention_and_mid_apply_mutation_rollback(setup, monkeypatch): + from loopx.capabilities.decision_context import capture_recovery as module + + args, payload, _ = setup + sid = source_id(payload) + capture_profile_sources(**args, execute=True) + preview = recover_capture_source(**args, source_id=sid, action="hold") + with sqlite3.connect(args["spool_path"]) as db: + db.execute("BEGIN IMMEDIATE") + with pytest.raises(sqlite3.OperationalError, match="locked"): + recover_capture_source( + **args, + source_id=sid, + action="hold", + execute=True, + expected_token=preview["preview_token"], + ) + original_schema = module._schema + + def racing_schema(db): + original_schema(db) + payload["automation"]["automatic_capture"] = False + args["profile_path"].write_text(json.dumps(payload)) + + monkeypatch.setattr(module, "_schema", racing_schema) + with pytest.raises(ValueError, match="changed during apply"): + recover_capture_source( + **args, + source_id=sid, + action="hold", + execute=True, + expected_token=preview["preview_token"], + ) + assert len(rows(args, "batches")) == 1 + + +def test_cli_recovery_preview_apply_diagnose_and_private_boundary(setup, capsys): + args, payload, _ = setup + sid = source_id(payload) + capture_profile_sources(**args, execute=True) + common = [ + "--goal-id", + args["goal_id"], + "--agent-id", + args["agent_id"], + "--profile", + str(args["profile_path"]), + "--spool", + str(args["spool_path"]), + "--cursor-state", + str(args["cursor_path"]), + "--source-id", + sid, + ] + prefix = ["--format", "json", "decision-context"] + assert main(prefix + ["capture-diagnose", *common, "--probe"]) == 0 + assert json.loads(capsys.readouterr().out)["diagnosis"] == "replayable" + assert main(prefix + ["capture-recovery", *common, "--action", "hold"]) == 0 + preview = json.loads(capsys.readouterr().out) + assert ( + main( + prefix + + [ + "capture-recovery", + *common, + "--action", + "hold", + "--execute", + "--expected-token", + preview["preview_token"], + ] + ) + == 0 + ) + output = capsys.readouterr().out + assert json.loads(output)["held_batch_count"] == 1 + assert "private-body" not in output and str(args["profile_path"]) not in output + assert not args["cursor_path"].exists() + + +def test_missing_spool_and_wrong_scope_do_not_create_recovery_state(setup): + args, payload, _ = setup + sid = source_id(payload) + with pytest.raises(ValueError, match="existing spool"): + recover_capture_source( + **args, + source_id=sid, + action="hold", + execute=True, + expected_token="not-a-preview", + ) + assert not args["spool_path"].exists() + capture_profile_sources(**args, execute=True) + with pytest.raises(ValueError, match="not enrolled"): + recover_capture_source(**args, source_id="other", action="hold") + with pytest.raises(ValueError, match="separate"): + recover_capture_source( + **{**args, "cursor_path": args["spool_path"]}, source_id=sid, action="hold" + ) + assert len(rows(args, "batches")) == 1 + + +def test_review_cursor_lock_excludes_concurrent_settlement(setup): + import subprocess + import sys + from loopx.file_lock import exclusive_file_lock + + args, payload, _ = setup + sid = source_id(payload) + capture_profile_sources(**args, execute=True) + preview = recover_capture_source(**args, source_id=sid, action="hold") + command = [ + sys.executable, + "-c", + "from loopx.cli import main; raise SystemExit(main())", + "--format", + "json", + "decision-context", + "capture-recovery", + "--action", + "hold", + "--execute", + "--expected-token", + preview["preview_token"], + "--source-id", + sid, + "--goal-id", + args["goal_id"], + "--agent-id", + args["agent_id"], + "--profile", + str(args["profile_path"]), + "--spool", + str(args["spool_path"]), + "--cursor-state", + str(args["cursor_path"]), + ] + with exclusive_file_lock(args["cursor_path"]): + result = subprocess.run(command, capture_output=True, text=True, timeout=20) + assert result.returncode != 0 + assert "lock" in result.stdout + result.stderr + assert len(rows(args, "batches")) == 1