From 3ded93498601b8a929657df11c0289f4e51afa90 Mon Sep 17 00:00:00 2001 From: Travis James Date: Mon, 21 Sep 2026 00:14:17 -0500 Subject: [PATCH] fix: bound ledger backlog reconciliation --- docs/lessons.md | 1 + .../evidence.md | 40 +++++ .../spec.md | 23 +++ src/operations.rs | 102 +++++++----- tests/operation_query_deadline.rs | 148 ++++++++++++++++++ 5 files changed, 272 insertions(+), 42 deletions(-) diff --git a/docs/lessons.md b/docs/lessons.md index 20c2088..b42e74e 100644 --- a/docs/lessons.md +++ b/docs/lessons.md @@ -40,6 +40,7 @@ Format: `YYYY-MM-DD — Rule — *(context: what went wrong)*` - 2026-08-26 — When adding a typed field to a struct that derives Serde, derive the same Serde direction on the field type before the first compile. *(Context: `LocalEmbeddingBackend` was added to deserializable `Config` without `Deserialize`.)* - 2026-08-26 — Do not run the aggregate `prometheus-rust-auditor audit` command in this repository: its pipeline generates a GitHub Actions workflow, which violates the local-only validation policy. Run the read-only enforcement, format, dependency, inventory, and partition checks individually. - 2026-09-21 — Never launch `cargo fmt` and `cargo test` in parallel in this repository; they share one target directory, so even a read-only format check violates the single-writer build rule when paired with compilation. +- 2026-09-21 — A bounded reconciliation query can still fail systematically when it bypasses the available state index, and rescanning it after every commit makes backlog recovery quadratic. Query indexed states under separate deadlines and rescan once per dependency wave. - 2026-08-26 — A copied SwiftPM executable is not self-contained when a C target ships resources: install `mlx-swift_Cmlx.bundle` beside the executable and smoke-test the copied path, because the build-tree binary can hide a missing `default.metallib` deployment. - 2026-08-26 — Persisted executor generations span server processes, but a newly pre-warmed child starts with process-local numbering. Adopt the healthy child into the durable generation sequence; do not kill it merely because its initial number is lower, or queue recovery turns into an unnecessary cold model launch. - 2026-08-26 — A liveness heartbeat must run independently of the work it supervises. An async Swift `Task` heartbeat can be starved while MLX/Metal synchronously occupies a cooperative executor; use a dedicated Dispatch timer and synchronously drain it before writing the terminal protocol message. diff --git a/openspec/changes/complete-operation-ledger-recovery/evidence.md b/openspec/changes/complete-operation-ledger-recovery/evidence.md index 994afb0..8a29ee1 100644 --- a/openspec/changes/complete-operation-ledger-recovery/evidence.md +++ b/openspec/changes/complete-operation-ledger-recovery/evidence.md @@ -64,3 +64,43 @@ The repaired source is not yet the installed binary. Deployment certification requires the merged commit to be built, signed, installed at both owned paths, and observed draining every accepted receipt before `prometheus doctor --json` can pass. + +## Installed-runtime follow-up + +The first installed build remained ready but its startup reconciliation timed +out twice at the 10-second discovery deadline. Direct HTTP measurements against +the same database showed the cause: the `state NOT IN [...]` discovery query +took 8.74 seconds, while separate state-index equality queries took 1.20–2.22 +seconds each. The coordinator also repeated the full discovery query after +every commit, making backlog recovery quadratic. + +The follow-up implementation queries `accepted`, `validated`, `blocked`, and +`processing` separately under the configured per-query deadline, then sorts +and deduplicates the combined identities. The drain loop now rescans once after +a wave that committed work, preserving dependency recovery without one full +ledger scan per committed operation. + +- `cargo fmt --all --check` and `git diff --check` — exit 0. +- `openspec validate complete-operation-ledger-recovery --strict` — exit 0. +- `openspec validate bound-operation-query-deadlines --strict` — exit 0. +- `openspec validate bound-operation-reconciliation-projection --strict` — + exit 0. +- `RUSTC_WRAPPER= cargo test --locked --test operation_query_deadline --no-default-features --features server-only -- --nocapture` + — exit 0; 2 passed, 0 failed. The added real-server case proves an + earlier-sorted dependent operation commits in a second drain wave after its + later-sorted prerequisite commits. +- The first run of that integration target produced 1 pass and 1 failure + because the new fixture used literal record keys instead of the production + SHA-256 key derivation. Correcting the fixture made the unchanged production + path pass. +- `RUSTC_WRAPPER= cargo test --locked --test executor_recovery -- --nocapture` + — exit 0; 4 passed, 0 failed. +- `prometheus-rust-auditor enforce`, `format`, and `inventory` — exit 0 with + no findings. `partition` exited 0 with the same six informational + AI-loop-pending rows. + +The task's isolated review budget was already consumed by the two recorded +rounds, ending in PASS. The installed-runtime regression therefore proceeds to +the same local gates and a fresh deployment proof rather than a third critic +round. Task 3.1 remains open until the new build drains accepted receipts to +zero and `prometheus doctor --json` exits successfully. diff --git a/openspec/changes/complete-operation-ledger-recovery/specs/operation-ledger-connection-recovery/spec.md b/openspec/changes/complete-operation-ledger-recovery/specs/operation-ledger-connection-recovery/spec.md index e03c04e..2a65be4 100644 --- a/openspec/changes/complete-operation-ledger-recovery/specs/operation-ledger-connection-recovery/spec.md +++ b/openspec/changes/complete-operation-ledger-recovery/specs/operation-ledger-connection-recovery/spec.md @@ -76,3 +76,26 @@ database errors MUST NOT enter this retry path. - **WHEN** startup discovery returns the typed recovered stale-ledger condition - **THEN** discovery is attempted once on the replacement generation - **AND** any other startup error is reported without this retry + +### Requirement: Backlog discovery and dependency rescans remain bounded + +The operation coordinator SHALL discover nonterminal work through the indexed +state values `accepted`, `validated`, `blocked`, and `processing`. Each state +query SHALL have its own database deadline, and the combined identities SHALL +be deduplicated and sorted before processing. After processing a drain wave, +the coordinator SHALL rescan at most once when that wave committed work, so +newly unblocked dependencies are recovered without one full-ledger query per +commit. + +#### Scenario: A large terminal history shares the ledger + +- **WHEN** startup reconciliation runs beside many committed or rejected rows +- **THEN** discovery queries only the four indexed nonterminal state values +- **AND** every discovered operation identity is processed in deterministic order + +#### Scenario: A dependency commits during a drain wave + +- **WHEN** an operation was blocked earlier in the wave and its dependency later commits +- **THEN** the coordinator performs one nonterminal rescan after the wave +- **AND** the newly unblocked operation is processed in the next wave +- **AND** the coordinator does not rescan the full ledger after every individual commit diff --git a/src/operations.rs b/src/operations.rs index 349d1a1..4ad9d4c 100644 --- a/src/operations.rs +++ b/src/operations.rs @@ -185,7 +185,9 @@ struct DbOperationId { operation_id: String, } -const LIST_NONTERMINAL_OPERATION_IDS_QUERY: &str = "SELECT operation_id FROM memory_operation WHERE state NOT IN ['committed', 'rejected'] ORDER BY operation_id ASC"; +const LIST_OPERATION_IDS_BY_STATE_QUERY: &str = + "SELECT operation_id FROM memory_operation WHERE state = $state ORDER BY operation_id ASC"; +const NONTERMINAL_OPERATION_STATES: [&str; 4] = ["accepted", "validated", "blocked", "processing"]; #[derive(Debug, Deserialize, SurrealValue)] struct DbOperationReceipt { @@ -663,16 +665,23 @@ impl OperationService { async fn list_nonterminal_ids(&self) -> Result> { let connection = self.ledger_connection().await?; let db = connection.db.clone(); - let rows: Vec = self - .await_database( - "reconciliation discovery", - connection.generation, - db.query(LIST_NONTERMINAL_OPERATION_IDS_QUERY), - ) - .await? - .check()? - .take(0)?; - Ok(rows.into_iter().map(|row| row.operation_id).collect()) + let mut operation_ids = Vec::new(); + for state in NONTERMINAL_OPERATION_STATES { + let rows: Vec = self + .await_database( + "reconciliation discovery", + connection.generation, + db.query(LIST_OPERATION_IDS_BY_STATE_QUERY) + .bind(("state", state)), + ) + .await? + .check()? + .take(0)?; + operation_ids.extend(rows.into_iter().map(|row| row.operation_id)); + } + operation_ids.sort_unstable(); + operation_ids.dedup(); + Ok(operation_ids) } async fn events_after(&self, operation_id: &str, sequence: u64) -> Result> { @@ -992,42 +1001,50 @@ impl OperationService { let mut pending = VecDeque::from(initial); let mut queued = pending.iter().cloned().collect::>(); let mut retried = HashSet::new(); + let mut rescan_after_wave = false; - while let Some(operation_id) = pending.pop_front() { - queued.remove(&operation_id); - if let Err(error) = self.process_pending(&operation_id).await { - if is_recovered_ledger_deadline(&error) - && retried.insert(operation_id.clone()) - && queued.insert(operation_id.clone()) - { - tracing::warn!(%operation_id, "retrying stale-ledger interruption once"); - pending.push_back(operation_id); + loop { + while let Some(operation_id) = pending.pop_front() { + queued.remove(&operation_id); + if let Err(error) = self.process_pending(&operation_id).await { + if is_recovered_ledger_deadline(&error) + && retried.insert(operation_id.clone()) + && queued.insert(operation_id.clone()) + { + tracing::warn!(%operation_id, "retrying stale-ledger interruption once"); + pending.push_back(operation_id); + continue; + } + tracing::error!(%operation_id, %error, "operation processing paused"); + self.record_processing_error(&operation_id, &error).await; continue; } - tracing::error!(%operation_id, %error, "operation processing paused"); - self.record_processing_error(&operation_id, &error).await; - continue; + let committed_now = match self.get(&operation_id).await { + Ok(Some(receipt)) => receipt.state == OperationState::Committed, + Ok(None) => false, + Err(error) => { + tracing::error!(%operation_id, %error, "operation post-process read failed"); + false + } + }; + rescan_after_wave |= committed_now; } - let committed_now = match self.get(&operation_id).await { - Ok(Some(receipt)) => receipt.state == OperationState::Committed, - Ok(None) => false, - Err(error) => { - tracing::error!(%operation_id, %error, "operation post-process read failed"); - false - } - }; - if committed_now { - match self.list_nonterminal_ids().await { - Ok(ids) => { - for id in ids { - if queued.insert(id.clone()) { - pending.push_back(id); - } + + if !rescan_after_wave { + break; + } + rescan_after_wave = false; + match self.list_nonterminal_ids().await { + Ok(ids) => { + for id in ids { + if queued.insert(id.clone()) { + pending.push_back(id); } } - Err(error) => { - tracing::error!(%error, "dependent operation reconciliation failed") - } + } + Err(error) => { + tracing::error!(%error, "dependent operation reconciliation failed"); + break; } } } @@ -2598,7 +2615,8 @@ mod tests { let rows: Vec = storage .db() .unwrap() - .query(LIST_NONTERMINAL_OPERATION_IDS_QUERY) + .query(LIST_OPERATION_IDS_BY_STATE_QUERY) + .bind(("state", "blocked")) .await .unwrap() .check() diff --git a/tests/operation_query_deadline.rs b/tests/operation_query_deadline.rs index fff379a..b9249b6 100644 --- a/tests/operation_query_deadline.rs +++ b/tests/operation_query_deadline.rs @@ -245,3 +245,151 @@ async fn concurrent_receipt_timeouts_leave_the_same_coordinator_able_to_commit_l .expect("the original coordinator commits later work"); assert_eq!(receipt["operation_id"], "deadline-probe"); } + +#[tokio::test] +async fn startup_reconciliation_processes_a_dependency_in_a_later_drain_wave() { + let (_server, endpoint) = start_server().await; + let embedder: Arc = Arc::new(NoOpEmbedder); + let storage = Arc::new( + SurrealStorage::new( + &SurrealConfig { + mode: SurrealMode::Server, + endpoint: Some(endpoint), + embedded_path: None, + username: None, + password: None, + namespace: format!("reconcile_{}", uuid::Uuid::new_v4().simple()), + database: "operations".to_owned(), + retry: RetryConfig { + max_connect_retries: 0, + query_timeout_ms: 10_000, + ..RetryConfig::default() + }, + }, + Arc::clone(&embedder), + ) + .await + .expect("isolated server-mode SurrealStorage"), + ); + let database = storage.db().unwrap().clone(); + let dependent_payload = json!({ + "name": "dependent", + "description": "must wait for the later-sorted prerequisite", + "agent_id": null, + "user_id": "test" + }); + let prerequisite_payload = json!({ + "name": "prerequisite", + "description": "commits in the first drain wave", + "agent_id": null, + "user_id": "test" + }); + let dependent_hash = Sha256::digest(serde_json::to_vec(&dependent_payload).unwrap()) + .iter() + .map(|byte| format!("{byte:02x}")) + .collect::(); + let prerequisite_hash = Sha256::digest(serde_json::to_vec(&prerequisite_payload).unwrap()) + .iter() + .map(|byte| format!("{byte:02x}")) + .collect::(); + let dependent_key = Sha256::digest(b"a-dependent") + .iter() + .map(|byte| format!("{byte:02x}")) + .collect::(); + let prerequisite_key = Sha256::digest(b"z-prerequisite") + .iter() + .map(|byte| format!("{byte:02x}")) + .collect::(); + + for (key, operation_id, dependencies, payload_hash, payload) in [ + ( + dependent_key, + "a-dependent", + vec!["z-prerequisite"], + dependent_hash, + dependent_payload, + ), + ( + prerequisite_key, + "z-prerequisite", + Vec::new(), + prerequisite_hash, + prerequisite_payload, + ), + ] { + database + .query( + "CREATE type::record('memory_operation', $key) CONTENT { + operation_id: $id, + schema_version: 2, + kind: 'create_task_stream', + dependencies: $dependencies, + payload_hash: $payload_hash, + payload: $payload, + state: 'accepted', + blocked_by: [], + result: NONE, + error: NONE, + executor_generation: 0, + executor_progress_seq: 0, + executor_exit_count: 0, + executor_last_exit: NONE, + executor_error: NONE, + progress_seq: 1, + created_at: time::now(), + updated_at: time::now() + }", + ) + .bind(("key", key)) + .bind(("id", operation_id)) + .bind(("dependencies", dependencies)) + .bind(("payload_hash", payload_hash)) + .bind(("payload", payload)) + .await + .unwrap() + .check() + .unwrap(); + } + + let router = api::build_router_with_query_timeout( + Arc::clone(&storage) as Arc, + embedder, + Duration::from_secs(10), + ); + + let receipts = tokio::time::timeout(Duration::from_secs(5), async { + loop { + let mut committed = Vec::new(); + for operation_id in ["a-dependent", "z-prerequisite"] { + let response = router + .clone() + .oneshot( + Request::builder() + .uri(format!("/api/v2/operations/{operation_id}")) + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + assert_eq!(response.status(), StatusCode::OK); + let body: Value = serde_json::from_slice( + &to_bytes(response.into_body(), usize::MAX).await.unwrap(), + ) + .unwrap(); + committed.push(body); + } + if committed + .iter() + .all(|receipt| receipt["state"] == "committed") + { + break committed; + } + tokio::time::sleep(Duration::from_millis(10)).await; + } + }) + .await + .expect("startup reconciliation commits both drain waves"); + + assert_eq!(receipts[0]["operation_id"], "a-dependent"); + assert_eq!(receipts[1]["operation_id"], "z-prerequisite"); +}