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
1 change: 1 addition & 0 deletions docs/lessons.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
40 changes: 40 additions & 0 deletions openspec/changes/complete-operation-ledger-recovery/evidence.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Original file line number Diff line number Diff line change
Expand Up @@ -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
102 changes: 60 additions & 42 deletions src/operations.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -663,16 +665,23 @@ impl OperationService {
async fn list_nonterminal_ids(&self) -> Result<Vec<String>> {
let connection = self.ledger_connection().await?;
let db = connection.db.clone();
let rows: Vec<DbOperationId> = 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<DbOperationId> = 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<Vec<OperationEvent>> {
Expand Down Expand Up @@ -992,42 +1001,50 @@ impl OperationService {
let mut pending = VecDeque::from(initial);
let mut queued = pending.iter().cloned().collect::<HashSet<_>>();
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;
}
}
}
Expand Down Expand Up @@ -2598,7 +2615,8 @@ mod tests {
let rows: Vec<Value> = storage
.db()
.unwrap()
.query(LIST_NONTERMINAL_OPERATION_IDS_QUERY)
.query(LIST_OPERATION_IDS_BY_STATE_QUERY)
.bind(("state", "blocked"))
.await
.unwrap()
.check()
Expand Down
148 changes: 148 additions & 0 deletions tests/operation_query_deadline.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<dyn EmbeddingService> = 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::<String>();
let prerequisite_hash = Sha256::digest(serde_json::to_vec(&prerequisite_payload).unwrap())
.iter()
.map(|byte| format!("{byte:02x}"))
.collect::<String>();
let dependent_key = Sha256::digest(b"a-dependent")
.iter()
.map(|byte| format!("{byte:02x}"))
.collect::<String>();
let prerequisite_key = Sha256::digest(b"z-prerequisite")
.iter()
.map(|byte| format!("{byte:02x}"))
.collect::<String>();

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<dyn MemoryStorage>,
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");
}
Loading