From b844b8e3312fce1f5f55f770068398a27801083b Mon Sep 17 00:00:00 2001 From: Travis James Date: Sun, 20 Sep 2026 23:56:31 -0500 Subject: [PATCH 1/2] fix: isolate operation ledger recovery --- Cargo.lock | 1 + Cargo.toml | 1 + crates/surreal-memory/src/storage/surreal.rs | 15 + docs/lessons.md | 1 + .../.openspec.yaml | 2 + .../design.md | 95 ++++ .../proposal.md | 41 ++ .../spec.md | 78 +++ .../tasks.md | 11 + src/operations.rs | 511 +++++++++++++++++- tests/operation_query_deadline.rs | 58 +- 11 files changed, 765 insertions(+), 49 deletions(-) create mode 100644 openspec/changes/complete-operation-ledger-recovery/.openspec.yaml create mode 100644 openspec/changes/complete-operation-ledger-recovery/design.md create mode 100644 openspec/changes/complete-operation-ledger-recovery/proposal.md create mode 100644 openspec/changes/complete-operation-ledger-recovery/specs/operation-ledger-connection-recovery/spec.md create mode 100644 openspec/changes/complete-operation-ledger-recovery/tasks.md diff --git a/Cargo.lock b/Cargo.lock index 569d541..5bfe251 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -6147,6 +6147,7 @@ name = "surreal-memory-server" version = "1.8.0" dependencies = [ "anyhow", + "arc-swap", "async-stream", "async-trait", "axum", diff --git a/Cargo.toml b/Cargo.toml index 819e04e..5c9551d 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -65,6 +65,7 @@ serde_json = { workspace = true } surrealdb = { workspace = true } surrealdb_types = { workspace = true } anyhow = { workspace = true } +arc-swap = { workspace = true } thiserror = { workspace = true } sha2 = { workspace = true } tracing = { workspace = true } diff --git a/crates/surreal-memory/src/storage/surreal.rs b/crates/surreal-memory/src/storage/surreal.rs index 84dad01..194b5be 100644 --- a/crates/surreal-memory/src/storage/surreal.rs +++ b/crates/surreal-memory/src/storage/surreal.rs @@ -924,6 +924,21 @@ DEFINE INDEX IF NOT EXISTS memory_embedding_hnsw self.live_db() } + /// Return a connection for the durable operation ledger. + /// + /// Server mode opens an independent transport so cancelling a ledger query + /// cannot strand unrelated storage work on the same WebSocket. Embedded + /// mode clones the in-process SDK handle because RocksDB cannot be opened a + /// second time at the same path. + pub async fn operation_ledger_connection(&self) -> Result> { + match self.connection_info.config.mode { + SurrealMode::Embedded => self.live_db(), + SurrealMode::Server => { + Self::connect_with_attempts(&self.connection_info.config, 0).await + } + } + } + /// Persist a fully planned and embedded logical memory under a stable key. /// /// The durable-operation coordinator derives `record_key` from its caller diff --git a/docs/lessons.md b/docs/lessons.md index c4ef556..20c2088 100644 --- a/docs/lessons.md +++ b/docs/lessons.md @@ -39,6 +39,7 @@ Format: `YYYY-MM-DD — Rule — *(context: what went wrong)*` - 2026-08-26 — Never stack a PR onto another PR's branch when the base will merge first. PR #9 targeted PR #8's branch; #8 merged to main, then #9 merged into an orphaned branch and GitHub reported "MERGED" while main had none of it. Target `main` and rebase, or the fix silently never ships. - 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-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/.openspec.yaml b/openspec/changes/complete-operation-ledger-recovery/.openspec.yaml new file mode 100644 index 0000000..cbd245e --- /dev/null +++ b/openspec/changes/complete-operation-ledger-recovery/.openspec.yaml @@ -0,0 +1,2 @@ +schema: spec-driven +created: 2026-09-20 diff --git a/openspec/changes/complete-operation-ledger-recovery/design.md b/openspec/changes/complete-operation-ledger-recovery/design.md new file mode 100644 index 0000000..9f99d4e --- /dev/null +++ b/openspec/changes/complete-operation-ledger-recovery/design.md @@ -0,0 +1,95 @@ +## Context + +See `proposal.md` for the failure. `OperationService` currently stores a +generation-tagged SDK handle, but it seeds generation zero by cloning the +general `SurrealStorage` handle. Its timeout error is also erased into +`anyhow::Error`, so the coordinator retries every processing failure rather +than the single recovery case the spec permits. + +The Surreal SDK handle is clone-safe, but server-mode clones multiplex one +physical WebSocket. An outer Tokio timeout can drop a query future while the +SDK is still unwinding that transport. The server-mode ledger therefore needs +a separately opened SDK connection. Embedded mode is different: opening the +same RocksDB path twice is invalid, and its clone is an in-process handle rather +than a shared remote transport. + +## Goals / Non-Goals + +**Goals:** + +- Make lazy first-use initialization and later replacement use the same + serialized independent-connection path. +- Preserve the existing public HTTP deadline error text. +- Keep retry classification inspectable after errors acquire context. +- Prove the retry and generation decisions deterministically, with a real + server integration proving isolation from general storage. + +**Non-Goals:** + +- Change the Surreal SDK's internal query timeout or general storage reconnect + policy. +- Retry executor work, payload validation, or arbitrary database failures. +- Add a third connection pool or configurable recovery knobs. + +## Decisions + +### Lazy independent initialization + +`OperationService` starts with no ledger connection. Its async accessor locks +the existing replacement mutex, checks again after acquiring the lock, asks +`SurrealStorage` for an operation-ledger handle, and publishes generation zero. +That storage method opens a fresh authenticated transport in server mode and +clones the live in-process handle in embedded mode. This keeps the synchronous +router builder unchanged while ensuring a server query cannot fall back to the +shared WebSocket. + +Making the whole router builder async was rejected because it would broaden +every construction site and still require concurrency control for the first +request. Pre-opening a connection by blocking inside the synchronous builder +was rejected because it can deadlock a Tokio runtime. + +### One typed deadline error with recovery state + +The deadline error retains the stage, elapsed milliseconds, and whether a +newer ledger generation was available after replacement. Its display remains +the existing API error. The coordinator inspects this type through the error +chain and retries only when recovery succeeded. + +String matching was rejected because repository lessons prohibit retry +classification from rendered messages. Retrying every `anyhow::Error` was +rejected because it repeats executor work. + +### Generation check under one async mutex + +Initialization and replacement share one Tokio mutex. A replacement checks the +currently published generation after locking; if another caller already +advanced it, the replacement succeeds without opening or publishing another +connection. This prevents an older recovery from overwriting a newer handle. + +### Evidence split + +Pure decision helpers and generation transitions receive focused unit tests, +including connector failure and non-ledger errors. The server integration +holds a large ledger receipt query past the application deadline while a +general storage health query completes, then proves the same production +coordinator accepts and commits later work. + +## Risks / Trade-offs + +- **[Risk] The first server-mode operation request now pays one connection + setup.** → The setup happens once and is serialized; embedded mode retains + its existing in-process cost. +- **[Risk] A failed initial connection leaves the slot empty.** → A later call + may attempt initialization again; no shared fallback is allowed. +- **[Risk] Real-server timing evidence can be affected by host pressure.** → + The integration uses an isolated fixture and deterministic large result, + while retry classification and generation behavior use non-timing tests. + +## Migration Plan + +1. Merge and build the repaired binary from a clean committed tree. +2. Install and sign both owned binary copies. +3. Restart the database, memory server, and learning worker in dependency order. +4. Observe accepted receipts reach zero and run `prometheus doctor --json`. +5. Roll back to the preceding signed binary if readiness or backlog progress + regresses; stored receipts require no schema migration. diff --git a/openspec/changes/complete-operation-ledger-recovery/proposal.md b/openspec/changes/complete-operation-ledger-recovery/proposal.md new file mode 100644 index 0000000..353b5ce --- /dev/null +++ b/openspec/changes/complete-operation-ledger-recovery/proposal.md @@ -0,0 +1,41 @@ +## Why + +The first operation-ledger query still uses the shared storage connection, so +cancelling that query can poison unrelated memory persistence before rotation +begins. The current recovery retry also catches executor failures, which can +repeat non-database work and violates the intended retry boundary. + +## What Changes + +- Open the server-mode operation ledger on an independent connection before + its first query, including startup reconciliation and concurrent API + requests; retain one clone-safe in-process handle for embedded mode, where a + second RocksDB open on the same path is invalid. +- Serialize connection initialization and replacement so an older recovery + cannot overwrite a newer healthy generation. +- Classify a database deadline as retryable only after the stale ledger + generation has been replaced. +- Retry startup discovery and an interrupted coordinator operation once only + for that typed stale-ledger condition. +- Add deterministic regressions for initialization, overlapping replacement, + replacement failure, retry classification, and startup recovery. + +## Capabilities + +### New Capabilities + +- `operation-ledger-connection-recovery`: Independent connection ownership, + generation-safe replacement, and narrowly classified retry behavior for the + durable operation ledger. + +### Modified Capabilities + +None. + +## Impact + +This changes `src/operations.rs`, the `SurrealStorage` independent-connection +API, focused operation tests, and the installed server runtime. The +uncomfortable constraint is that the deployed backlog cannot be certified +until the repaired binary is installed and processes every accepted receipt; +passing isolated tests alone does not close the release task. 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 new file mode 100644 index 0000000..e03c04e --- /dev/null +++ b/openspec/changes/complete-operation-ledger-recovery/specs/operation-ledger-connection-recovery/spec.md @@ -0,0 +1,78 @@ +## Purpose + +Keep durable operation-ledger cancellation and recovery isolated from ordinary +memory storage while preserving bounded, generation-safe retry behavior. + +## ADDED Requirements + +### Requirement: The server-backed operation ledger owns an independent transport + +Before executing its first ledger query, the operation service SHALL open a +server-mode ledger connection independently of the general storage transport. +Concurrent first use MUST converge on one published ledger generation. In +embedded mode it SHALL instead clone the in-process SDK handle because a second +RocksDB open on the same path is invalid; the clone MUST NOT publish or mutate +general storage connection state. + +#### Scenario: First ledger query overlaps ordinary storage work + +- **WHEN** the first operation-ledger query is still running +- **THEN** ordinary server-mode storage work remains able to complete on its own transport +- **AND** cancellation of the ledger query does not change general storage connection state + +#### Scenario: Concurrent callers initialize the ledger + +- **WHEN** multiple operation requests arrive before a ledger connection has been published +- **THEN** initialization is serialized +- **AND** every caller observes the same published ledger generation + +#### Scenario: Embedded ledger initializes without reopening RocksDB + +- **WHEN** the operation service uses an embedded database +- **THEN** ledger initialization clones the existing in-process SDK handle +- **AND** it does not attempt a second open of the embedded database path + +### Requirement: Ledger replacement is generation safe + +The operation service SHALL serialize ledger replacement and SHALL replace a +connection only while the caller's generation is still current. A replacement +started for an older generation MUST NOT overwrite a newer published +connection. + +#### Scenario: Timed-out queries overlap + +- **WHEN** multiple queries from the same ledger generation time out concurrently +- **THEN** at most one replacement connection becomes the next generation +- **AND** every later replacement attempt observes that newer generation + +#### Scenario: Replacement cannot connect + +- **WHEN** a timed-out ledger query cannot establish a replacement connection within its bound +- **THEN** the request returns the database deadline error +- **AND** the failure is not classified as a recovered stale-ledger interruption +- **AND** non-database operation work is not retried + +### Requirement: Retry is limited to recovered stale-ledger interruption + +The operation coordinator SHALL retry work at most once only when a database +deadline interrupted a stale ledger generation and a replacement generation +was successfully published. Executor, validation, payload, and ordinary +database errors MUST NOT enter this retry path. + +#### Scenario: Coordinator work is interrupted by a stale ledger generation + +- **WHEN** a coordinator operation returns the typed recovered stale-ledger condition +- **THEN** the same operation is attempted once on the replacement generation +- **AND** a second failure is recorded without another retry + +#### Scenario: Executor work fails + +- **WHEN** a coordinator operation fails in the executor or another non-ledger stage +- **THEN** the failure is recorded +- **AND** the coordinator does not retry it through ledger recovery + +#### Scenario: Startup discovery is interrupted by a stale ledger generation + +- **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 diff --git a/openspec/changes/complete-operation-ledger-recovery/tasks.md b/openspec/changes/complete-operation-ledger-recovery/tasks.md new file mode 100644 index 0000000..6481e8f --- /dev/null +++ b/openspec/changes/complete-operation-ledger-recovery/tasks.md @@ -0,0 +1,11 @@ +## 1. Independent ledger recovery + +- [x] 1.1 Open the initial server-mode ledger transport independently while preserving embedded initialization, serialize generation-safe replacement, classify recovered deadline errors, and prove initialization, overlap, replacement failure, stale-only retry, startup retry, and general-storage isolation with focused tests. + +## 2. Review and publication + +- [ ] 2.1 Pass formatting, compilation, focused tests, strict OpenSpec validation, the repository Rust format, enforcement, and inventory audit gates; record dependency and partition baseline findings; pass a fresh isolated critic review; commit, push, and merge a reviewable PR. + +## 3. Installed runtime certification + +- [ ] 3.1 Build and sign the merged binary, reinstall and restart managed services, drain accepted receipts to zero, and prove service health plus `prometheus doctor --json` success. diff --git a/src/operations.rs b/src/operations.rs index 04e3c10..349d1a1 100644 --- a/src/operations.rs +++ b/src/operations.rs @@ -8,13 +8,14 @@ use std::{ collections::{HashSet, VecDeque}, convert::Infallible, - future::IntoFuture, + future::{Future, IntoFuture}, pin::Pin, sync::Arc, time::Duration, }; use anyhow::{Context, Result}; +use arc_swap::ArcSwapOption; use axum::{ Json, Router, extract::{Path, Query, State}, @@ -32,6 +33,7 @@ use surreal_memory::{ embeddings::{EmbeddingPlanPart, ExecutorEvent, ExecutorEventKind, ExecutorSnapshot}, }; use surrealdb::types::{Datetime, RecordId}; +use surrealdb::{Surreal, engine::any::Any}; use surrealdb_types::SurrealValue; use tokio::sync::{broadcast, mpsc}; use tokio_stream::StreamExt; @@ -324,8 +326,52 @@ pub struct OperationService { storage: Arc, embedding_service: Arc, query_timeout: Duration, + ledger_connector: LedgerConnector, + ledger_connection: Arc>, + connection_replacement: Arc>, wake_tx: mpsc::Sender, events_tx: broadcast::Sender, + #[cfg(test)] + process_override: Option, +} + +struct LedgerConnection { + generation: u64, + db: Surreal, +} + +type LedgerConnectFuture = Pin>> + Send>>; +type LedgerConnector = Arc LedgerConnectFuture + Send + Sync>; +#[cfg(test)] +type ProcessOverrideFuture = Pin> + Send>>; +#[cfg(test)] +type ProcessOverride = Arc ProcessOverrideFuture + Send + Sync>; + +#[derive(Debug, thiserror::Error)] +#[error("operation database {stage} timed out after {timeout_ms}ms")] +struct LedgerDeadlineExceeded { + stage: &'static str, + timeout_ms: u128, + recovered: bool, +} + +fn is_recovered_ledger_deadline(error: &anyhow::Error) -> bool { + error.chain().any(|cause| { + cause + .downcast_ref::() + .is_some_and(|deadline| deadline.recovered) + }) +} + +async fn retry_recovered_ledger_once(mut operation: F) -> (Result, bool) +where + F: FnMut() -> Fut, + Fut: std::future::Future>, +{ + match operation().await { + Err(error) if is_recovered_ledger_deadline(&error) => (operation().await, true), + result => (result, false), + } } #[derive(Debug)] @@ -360,12 +406,28 @@ impl OperationService { ) -> Self { let (wake_tx, wake_rx) = mpsc::channel(wake_capacity); let (events_tx, _) = broadcast::channel(event_capacity); + let connector_storage = Arc::clone(&storage); + let ledger_connector: LedgerConnector = Arc::new(move || { + let storage = Arc::clone(&connector_storage); + Box::pin(async move { + let surreal = storage + .as_any() + .downcast_ref::() + .context("durable operations require SurrealStorage")?; + surreal.operation_ledger_connection().await + }) + }); let service = Self { storage, embedding_service, query_timeout, + ledger_connector, + ledger_connection: Arc::new(ArcSwapOption::empty()), + connection_replacement: Arc::new(tokio::sync::Mutex::new(())), wake_tx, events_tx, + #[cfg(test)] + process_override: None, }; if let Some(executor_events) = service.embedding_service.subscribe_executor_events() { let journal = service.clone(); @@ -383,20 +445,80 @@ impl OperationService { .context("durable operations require SurrealStorage") } - async fn await_database(&self, stage: &'static str, future: F) -> Result + async fn ledger_connection(&self) -> Result> { + if let Some(connection) = self.ledger_connection.load_full() { + return Ok(connection); + } + let _replacement = self.connection_replacement.lock().await; + if let Some(connection) = self.ledger_connection.load_full() { + return Ok(connection); + } + let db = (self.ledger_connector)().await?; + let connection = Arc::new(LedgerConnection { generation: 0, db }); + self.ledger_connection.store(Some(Arc::clone(&connection))); + Ok(connection) + } + + async fn replace_ledger_connection(&self, stale_generation: u64) -> Result<()> { + let _replacement = self.connection_replacement.lock().await; + if self + .ledger_connection + .load_full() + .is_some_and(|connection| connection.generation != stale_generation) + { + return Ok(()); + } + let db = (self.ledger_connector)().await?; + self.ledger_connection + .store(Some(Arc::new(LedgerConnection { + generation: stale_generation.saturating_add(1), + db, + }))); + Ok(()) + } + + async fn await_database( + &self, + stage: &'static str, + generation: u64, + future: F, + ) -> Result where F: IntoFuture>, anyhow::Error: From, { - tokio::time::timeout(self.query_timeout, future.into_future()) - .await - .with_context(|| { - format!( - "operation database {stage} timed out after {}ms", - self.query_timeout.as_millis() + match tokio::time::timeout(self.query_timeout, future.into_future()).await { + Ok(result) => result.map_err(anyhow::Error::from), + Err(_) => { + let replacement_timeout = self.query_timeout.max(Duration::from_secs(1)); + let replacement = tokio::time::timeout( + replacement_timeout, + self.replace_ledger_connection(generation), ) - })? - .map_err(anyhow::Error::from) + .await; + let recovered = match replacement { + Ok(Ok(())) => true, + Ok(Err(error)) => { + tracing::error!(%error, stage, "operation database connection replacement failed"); + false + } + Err(_) => { + tracing::error!( + stage, + timeout_ms = replacement_timeout.as_millis(), + "operation database connection replacement timed out" + ); + false + } + }; + Err(LedgerDeadlineExceeded { + stage, + timeout_ms: self.query_timeout.as_millis(), + recovered, + } + .into()) + } + } } pub async fn submit( @@ -453,14 +575,15 @@ impl OperationService { }; let key = record_key(&request.operation_id); let event_key = format!("{key}-0000000000000001"); - let db = self - .surreal() - .map_err(SubmitError::Storage)? - .db() + let connection = self + .ledger_connection() + .await .map_err(SubmitError::Storage)?; + let db = connection.db.clone(); let response = self .await_database( "submit", + connection.generation, db.query( "BEGIN TRANSACTION;\n\ CREATE type::record('memory_operation', $key) CONTENT $operation;\n\ @@ -487,7 +610,7 @@ impl OperationService { } return Err(SubmitError::Conflict(Box::new(existing))); } - return Err(SubmitError::Storage(error.into())); + return Err(SubmitError::Storage(error)); } let receipt = self @@ -511,10 +634,12 @@ impl OperationService { } pub async fn get(&self, operation_id: &str) -> Result> { - let db = self.surreal()?.db()?; + let connection = self.ledger_connection().await?; + let db = connection.db.clone(); let mut rows: Vec = self .await_database( "receipt lookup", + connection.generation, db.query(GET_OPERATION_RECEIPT_QUERY) .bind(("id", operation_id.to_owned())), ) @@ -525,19 +650,23 @@ impl OperationService { } async fn get_db(&self, operation_id: &str) -> Result> { - let db = self.surreal()?.db()?; + let connection = self.ledger_connection().await?; + let db = connection.db.clone(); self.await_database( "operation lookup", + connection.generation, db.select(("memory_operation", record_key(operation_id))), ) .await } async fn list_nonterminal_ids(&self) -> Result> { - let db = self.surreal()?.db()?; + 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? @@ -547,10 +676,12 @@ impl OperationService { } async fn events_after(&self, operation_id: &str, sequence: u64) -> Result> { - let db = self.surreal()?.db()?; + let connection = self.ledger_connection().await?; + let db = connection.db.clone(); let rows: Vec = self .await_database( "event history lookup", + connection.generation, db.query( "SELECT * FROM memory_operation_event WHERE operation_id = $id AND sequence > $sequence ORDER BY sequence ASC", ) @@ -590,9 +721,11 @@ impl OperationService { occurred_at: now, }; let event_key = format!("{}-{sequence:016}", record_key(operation_id)); - let db = self.surreal()?.db()?; + let connection = self.ledger_connection().await?; + let db = connection.db.clone(); self.await_database( "state transition", + connection.generation, db.query( "BEGIN TRANSACTION;\n\ UPDATE memory_operation SET state = $state, blocked_by = $blocked_by, result = $result, error = $error, progress_seq = $sequence, updated_at = $now WHERE operation_id = $id;\n\ @@ -640,10 +773,12 @@ impl OperationService { } async fn operation_parts(&self, operation_id: &str) -> Result> { - let db = self.surreal()?.db()?; + let connection = self.ledger_connection().await?; + let db = connection.db.clone(); let parts: Vec = self .await_database( "part lookup", + connection.generation, db.query( "SELECT * FROM memory_operation_part WHERE operation_id = $id ORDER BY part_index ASC", ) @@ -671,7 +806,8 @@ impl OperationService { return Ok(()); } - let db = self.surreal()?.db()?; + let connection = self.ledger_connection().await?; + let db = connection.db.clone(); let rows = plan .iter() .map(|part| { @@ -693,6 +829,7 @@ impl OperationService { .collect::>(); self.await_database( "plan persistence", + connection.generation, db.query( "BEGIN TRANSACTION;\n\ INSERT INTO memory_operation_part $parts;\n\ @@ -711,9 +848,11 @@ impl OperationService { part_index: u64, embedding: Vec, ) -> Result<()> { - let db = self.surreal()?.db()?; + let connection = self.ledger_connection().await?; + let db = connection.db.clone(); self.await_database( "part persistence", + connection.generation, db.query( "UPDATE memory_operation_part SET state = 'indexed', embedding = $embedding, updated_at = $now WHERE operation_id = $id AND part_index = $index", ) @@ -799,9 +938,11 @@ impl OperationService { event.generation, event.progress_seq ); - let db = self.surreal()?.db()?; + let connection = self.ledger_connection().await?; + let db = connection.db.clone(); self.await_database( "executor event persistence", + connection.generation, db.query("CREATE type::record('memory_executor_event', $key) CONTENT $event") .bind(("key", key)) .bind(("event", row)), @@ -818,9 +959,11 @@ impl OperationService { else { return Ok(()); }; - let db = self.surreal()?.db()?; + let connection = self.ledger_connection().await?; + let db = connection.db.clone(); self.await_database( "executor snapshot persistence", + connection.generation, db.query( "UPDATE memory_operation SET executor_generation = $generation, executor_progress_seq = $progress, executor_exit_count = $exit_count, executor_last_exit = $last_exit, executor_error = $executor_error, updated_at = $now WHERE operation_id = $id", ) @@ -837,13 +980,30 @@ impl OperationService { Ok(()) } + async fn process_pending(&self, operation_id: &str) -> Result<()> { + #[cfg(test)] + if let Some(process) = &self.process_override { + return process(operation_id.to_owned()).await; + } + self.process(operation_id).await + } + async fn drain_pending(&self, initial: Vec) { let mut pending = VecDeque::from(initial); let mut queued = pending.iter().cloned().collect::>(); + let mut retried = HashSet::new(); while let Some(operation_id) = pending.pop_front() { queued.remove(&operation_id); - if let Err(error) = self.process(&operation_id).await { + 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; @@ -881,8 +1041,16 @@ impl OperationService { } async fn run(self, mut wake_rx: mpsc::Receiver) { - match self.reconcile_nonterminal().await { + let (reconciliation, retried) = + retry_recovered_ledger_once(|| self.reconcile_nonterminal()).await; + match reconciliation { + Ok(_) if retried => { + tracing::warn!("operation startup reconciliation recovered after one retry") + } Ok(_) => {} + Err(error) if retried => { + tracing::error!(%error, "operation startup reconciliation retry failed"); + } Err(error) => tracing::error!(%error, "operation startup reconciliation failed"), } @@ -1593,6 +1761,295 @@ mod tests { } } + fn service_without_coordinator( + storage: Arc, + embedder: Arc, + ledger_connector: LedgerConnector, + initial_connection: Option, + query_timeout: Duration, + ) -> OperationService { + let (wake_tx, _) = mpsc::channel(4); + let (events_tx, _) = broadcast::channel(4); + OperationService { + storage: storage as Arc, + embedding_service: embedder, + query_timeout, + ledger_connector, + ledger_connection: Arc::new(ArcSwapOption::new(initial_connection.map(Arc::new))), + connection_replacement: Arc::new(tokio::sync::Mutex::new(())), + wake_tx, + events_tx, + process_override: None, + } + } + + fn deadline_error(recovered: bool) -> anyhow::Error { + LedgerDeadlineExceeded { + stage: "fixture", + timeout_ms: 1, + recovered, + } + .into() + } + + #[tokio::test] + async fn concurrent_first_use_publishes_one_ledger_generation() { + let embedder: Arc = Arc::new(NoOpEmbedder); + let storage = Arc::new( + SurrealStorage::new_mem(Arc::clone(&embedder)) + .await + .expect("in-memory SurrealStorage"), + ); + let database = storage.db().unwrap(); + let connect_calls = Arc::new(AtomicUsize::new(0)); + let ledger_connector: LedgerConnector = { + let connect_calls = Arc::clone(&connect_calls); + Arc::new(move || { + let database = database.clone(); + let connect_calls = Arc::clone(&connect_calls); + Box::pin(async move { + connect_calls.fetch_add(1, Ordering::SeqCst); + tokio::task::yield_now().await; + Ok(database) + }) + }) + }; + let service = service_without_coordinator( + Arc::clone(&storage), + embedder, + ledger_connector, + None, + Duration::from_secs(1), + ); + + let connections = + futures_util::future::join_all((0..4).map(|_| service.ledger_connection())) + .await + .into_iter() + .collect::>>() + .unwrap(); + + assert_eq!(connect_calls.load(Ordering::SeqCst), 1); + assert!( + connections + .iter() + .all(|connection| connection.generation == 0) + ); + assert!( + connections + .iter() + .all(|connection| Arc::ptr_eq(connection, &connections[0])) + ); + } + + #[tokio::test] + async fn overlapping_replacements_publish_only_the_next_generation() { + let embedder: Arc = Arc::new(NoOpEmbedder); + let storage = Arc::new( + SurrealStorage::new_mem(Arc::clone(&embedder)) + .await + .expect("in-memory SurrealStorage"), + ); + let database = storage.db().unwrap(); + let connect_calls = Arc::new(AtomicUsize::new(0)); + let ledger_connector: LedgerConnector = { + let connect_calls = Arc::clone(&connect_calls); + let replacement = database.clone(); + Arc::new(move || { + let replacement = replacement.clone(); + let connect_calls = Arc::clone(&connect_calls); + Box::pin(async move { + connect_calls.fetch_add(1, Ordering::SeqCst); + tokio::time::sleep(Duration::from_millis(10)).await; + Ok(replacement) + }) + }) + }; + let service = service_without_coordinator( + Arc::clone(&storage), + embedder, + ledger_connector, + Some(LedgerConnection { + generation: 0, + db: database, + }), + Duration::from_secs(1), + ); + + let replacements = + futures_util::future::join_all((0..4).map(|_| service.replace_ledger_connection(0))) + .await; + + assert!(replacements.into_iter().all(|result| result.is_ok())); + assert_eq!(connect_calls.load(Ordering::SeqCst), 1); + assert_eq!(service.ledger_connection().await.unwrap().generation, 1); + } + + #[tokio::test] + async fn replacement_failure_is_not_a_recovered_ledger_deadline() { + let embedder: Arc = Arc::new(NoOpEmbedder); + let storage = Arc::new( + SurrealStorage::new_mem(Arc::clone(&embedder)) + .await + .expect("in-memory SurrealStorage"), + ); + let database = storage.db().unwrap(); + let connect_calls = Arc::new(AtomicUsize::new(0)); + let ledger_connector: LedgerConnector = { + let connect_calls = Arc::clone(&connect_calls); + Arc::new(move || { + let connect_calls = Arc::clone(&connect_calls); + Box::pin(async move { + connect_calls.fetch_add(1, Ordering::SeqCst); + anyhow::bail!("fixture replacement failed") + }) + }) + }; + let service = service_without_coordinator( + Arc::clone(&storage), + embedder, + ledger_connector, + Some(LedgerConnection { + generation: 0, + db: database, + }), + Duration::from_millis(1), + ); + + let error = service + .await_database( + "fixture", + 0, + std::future::pending::>(), + ) + .await + .unwrap_err(); + + assert_eq!( + error.to_string(), + "operation database fixture timed out after 1ms" + ); + assert!(!is_recovered_ledger_deadline(&error)); + assert_eq!(connect_calls.load(Ordering::SeqCst), 1); + assert_eq!(service.ledger_connection().await.unwrap().generation, 0); + } + + #[tokio::test] + async fn startup_retry_runs_once_only_for_a_recovered_ledger_deadline() { + let recovered_calls = Arc::new(AtomicUsize::new(0)); + let recovered_counter = Arc::clone(&recovered_calls); + let (recovered_result, retried) = retry_recovered_ledger_once(move || { + let attempt = recovered_counter.fetch_add(1, Ordering::SeqCst); + async move { + if attempt == 0 { + Err(deadline_error(true)) + } else { + Ok("recovered") + } + } + }) + .await; + assert_eq!(recovered_result.unwrap(), "recovered"); + assert!(retried); + assert_eq!(recovered_calls.load(Ordering::SeqCst), 2); + + for first_error in [deadline_error(false), anyhow::anyhow!("executor failed")] { + let calls = Arc::new(AtomicUsize::new(0)); + let counter = Arc::clone(&calls); + let mut first_error = Some(first_error); + let (result, retried) = retry_recovered_ledger_once(move || { + counter.fetch_add(1, Ordering::SeqCst); + let error = first_error.take().expect("only one attempt"); + async move { Err::<(), _>(error) } + }) + .await; + assert!(result.is_err()); + assert!(!retried); + assert_eq!(calls.load(Ordering::SeqCst), 1); + } + + let repeated_calls = Arc::new(AtomicUsize::new(0)); + let repeated_counter = Arc::clone(&repeated_calls); + let (result, retried) = retry_recovered_ledger_once(move || { + repeated_counter.fetch_add(1, Ordering::SeqCst); + async { Err::<(), _>(deadline_error(true)) } + }) + .await; + assert!(result.is_err()); + assert!(retried); + assert_eq!(repeated_calls.load(Ordering::SeqCst), 2); + } + + #[tokio::test] + async fn coordinator_drain_retries_only_a_recovered_ledger_deadline_once() { + let embedder: Arc = Arc::new(NoOpEmbedder); + let storage = Arc::new( + SurrealStorage::new_mem(Arc::clone(&embedder)) + .await + .expect("in-memory SurrealStorage"), + ); + let database = storage.db().unwrap(); + + let make_service = |process_override: ProcessOverride| { + let replacement = database.clone(); + let ledger_connector: LedgerConnector = Arc::new(move || { + let replacement = replacement.clone(); + Box::pin(async move { Ok(replacement) }) + }); + let mut service = service_without_coordinator( + Arc::clone(&storage), + Arc::clone(&embedder), + ledger_connector, + Some(LedgerConnection { + generation: 0, + db: database.clone(), + }), + Duration::from_secs(1), + ); + service.process_override = Some(process_override); + service + }; + + let recovered_then_executor_calls = Arc::new(AtomicUsize::new(0)); + let recovered_counter = Arc::clone(&recovered_then_executor_calls); + let recovered_then_executor: ProcessOverride = Arc::new(move |_| { + let attempt = recovered_counter.fetch_add(1, Ordering::SeqCst); + Box::pin(async move { + if attempt == 0 { + Err(deadline_error(true)) + } else { + anyhow::bail!("fixture executor failure") + } + }) + }); + make_service(recovered_then_executor) + .drain_pending(vec!["recovered-then-executor".to_owned()]) + .await; + assert_eq!(recovered_then_executor_calls.load(Ordering::SeqCst), 2); + + let repeated_deadline_calls = Arc::new(AtomicUsize::new(0)); + let repeated_counter = Arc::clone(&repeated_deadline_calls); + let repeated_deadline: ProcessOverride = Arc::new(move |_| { + repeated_counter.fetch_add(1, Ordering::SeqCst); + Box::pin(async { Err(deadline_error(true)) }) + }); + make_service(repeated_deadline) + .drain_pending(vec!["repeated-deadline".to_owned()]) + .await; + assert_eq!(repeated_deadline_calls.load(Ordering::SeqCst), 2); + + let executor_calls = Arc::new(AtomicUsize::new(0)); + let executor_counter = Arc::clone(&executor_calls); + let executor_failure: ProcessOverride = Arc::new(move |_| { + executor_counter.fetch_add(1, Ordering::SeqCst); + Box::pin(async { anyhow::bail!("fixture executor failure") }) + }); + make_service(executor_failure) + .drain_pending(vec!["executor-only".to_owned()]) + .await; + assert_eq!(executor_calls.load(Ordering::SeqCst), 1); + } + async fn test_router() -> (Router, OperationService) { let embedder: Arc = Arc::new(NoOpEmbedder); let storage: Arc = Arc::new( diff --git a/tests/operation_query_deadline.rs b/tests/operation_query_deadline.rs index e0b347b..fff379a 100644 --- a/tests/operation_query_deadline.rs +++ b/tests/operation_query_deadline.rs @@ -77,7 +77,7 @@ impl EmbeddingService for NoOpEmbedder { } #[tokio::test] -async fn receipt_timeout_leaves_the_same_coordinator_able_to_commit_later_work() { +async fn concurrent_receipt_timeouts_leave_the_same_coordinator_able_to_commit_later_work() { let (_server, endpoint) = start_server().await; let embedder: Arc = Arc::new(NoOpEmbedder); let storage = Arc::new( @@ -92,7 +92,7 @@ async fn receipt_timeout_leaves_the_same_coordinator_able_to_commit_later_work() database: "operations".to_owned(), retry: RetryConfig { max_connect_retries: 0, - query_timeout_ms: 1_000, + query_timeout_ms: 10_000, ..RetryConfig::default() }, }, @@ -146,31 +146,45 @@ async fn receipt_timeout_leaves_the_same_coordinator_able_to_commit_later_work() let router = api::build_router_with_query_timeout( Arc::clone(&storage) as Arc, embedder, - Duration::from_millis(50), + Duration::from_millis(10), ); - let timeout_response = router - .clone() - .oneshot( - Request::builder() - .uri("/api/v2/operations/large-receipt") - .body(Body::empty()) + let timeout_batch = futures_util::future::join_all((0..4).map(|_| { + let router = router.clone(); + async move { + router + .oneshot( + Request::builder() + .uri("/api/v2/operations/large-receipt") + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap() + } + })); + let general_storage_health = async { + tokio::time::sleep(Duration::from_millis(1)).await; + tokio::time::timeout(Duration::from_millis(500), storage.health_check()) + .await + .expect("general storage remains responsive during ledger cancellation") + .expect("general storage health query succeeds") + }; + let (timeout_responses, storage_healthy) = tokio::join!(timeout_batch, general_storage_health); + assert!(storage_healthy); + for timeout_response in timeout_responses { + assert_eq!(timeout_response.status(), StatusCode::INTERNAL_SERVER_ERROR); + let timeout_body: Value = serde_json::from_slice( + &to_bytes(timeout_response.into_body(), usize::MAX) + .await .unwrap(), ) - .await .unwrap(); - assert_eq!(timeout_response.status(), StatusCode::INTERNAL_SERVER_ERROR); - let timeout_body: Value = serde_json::from_slice( - &to_bytes(timeout_response.into_body(), usize::MAX) - .await - .unwrap(), - ) - .unwrap(); - assert_eq!( - timeout_body["error"], - "operation database receipt lookup timed out after 50ms" - ); - tokio::time::sleep(Duration::from_millis(1_100)).await; + assert_eq!( + timeout_body["error"], + "operation database receipt lookup timed out after 10ms" + ); + } let payload = json!({ "name": "deadline-probe", From 406200e643f03c99297c95a5dcb2565a91235fcc Mon Sep 17 00:00:00 2001 From: Travis James Date: Sun, 20 Sep 2026 23:56:57 -0500 Subject: [PATCH 2/2] docs: record operation recovery evidence --- .../evidence.md | 54 +++++++++++++++++ .../evidence.md | 38 ++++++++++++ .../evidence.md | 59 +++++++++++++++++++ 3 files changed, 151 insertions(+) create mode 100644 openspec/changes/bound-operation-query-deadlines/evidence.md create mode 100644 openspec/changes/bound-operation-reconciliation-projection/evidence.md create mode 100644 openspec/changes/complete-operation-ledger-recovery/evidence.md diff --git a/openspec/changes/bound-operation-query-deadlines/evidence.md b/openspec/changes/bound-operation-query-deadlines/evidence.md new file mode 100644 index 0000000..25fc0ce --- /dev/null +++ b/openspec/changes/bound-operation-query-deadlines/evidence.md @@ -0,0 +1,54 @@ +# Verification evidence + +Date: 2026-09-21 + +## Source and review + +- PR #20: https://github.com/Prometheus-AGS/surreal-memory-server/pull/20 + - source: `93372f820438dc485277530e8d5b23c21711090d` + - merge: `b000bced73f4f626625b4379d26cda657999e080` +- Critic round 1 returned BLOCK because the recovery assertion used a second + coordinator. The regression was repaired to use the production coordinator. +- Critic round 2 returned BLOCK because the fixture could race with coordinator + processing. The terminal receipt is now seeded before the coordinator starts. +- Follow-up change `complete-operation-ledger-recovery` addressed the recorded + blocker in a new review cycle. Its round 1 critic required a direct + coordinator retry regression and removal of a redundant legacy-spec edit; + both were repaired. Round 2 returned PASS with no blocking findings. + +## Local verification + +- `RUSTC_WRAPPER= cargo check --locked --package surreal-memory-server --no-default-features --features server-only` + — exit 0; compilation completed in 1 minute 16 seconds. +- `RUSTC_WRAPPER= cargo test --locked --test operation_query_deadline --no-default-features --features server-only` + — exit 0; 1 test passed. A real isolated SurrealDB server held a large + receipt query past the 50 ms application deadline, and the same production + coordinator then committed a subsequent operation. +- `RUSTC_WRAPPER= cargo test --locked --test executor_recovery` — exit 0; 4 + tests passed. +- `RUSTFLAGS='-Dwarnings' cargo build --release --locked --no-default-features --features embedded,metal,local-embeddings` + — exit 0; release build completed in 37.71 seconds. +- `cargo fmt --all --check` — exit 0 on the merged tree. +- `openspec validate bound-operation-query-deadlines --strict` — exit 0; the + change is valid on the merged tree. +- `RUSTC_WRAPPER= cargo test --locked --lib operations::tests:: -- --nocapture` + — exit 0; 17 tests passed, including independent initialization, + generation-safe overlap, replacement failure, startup retry, and the actual + coordinator drain retry branch. +- `prometheus-rust-auditor enforce`, `format`, and `inventory` — exit 0; no + findings. The dependency phase remains a repository baseline limitation: + `cargo-deny` is not installed, while `cargo audit` reports the same three + advisories in both the committed and working lockfiles. + +## Deployment state + +The signed installed binary has SHA-256 +`63d4b297a9e4b1ebd2cb61ebac0752a00532eaf611b5b0a967ed9dc875942046` +at both owned install paths, and both copies pass `codesign --verify`. + +The deployed server recorded +`operation database reconciliation discovery timed out after 10000ms` instead +of leaving the startup future pending indefinitely. The service stayed ready, +and the durable queue continued to advance after the learning worker loaded. +Backlog recovery is still running, so task 1.3 remains open until the accepted +count reaches zero and the final doctor check exits successfully. diff --git a/openspec/changes/bound-operation-reconciliation-projection/evidence.md b/openspec/changes/bound-operation-reconciliation-projection/evidence.md new file mode 100644 index 0000000..c613e9b --- /dev/null +++ b/openspec/changes/bound-operation-reconciliation-projection/evidence.md @@ -0,0 +1,38 @@ +# Verification evidence + +Date: 2026-09-21 + +## Source and review + +- PR #18: https://github.com/Prometheus-AGS/surreal-memory-server/pull/18 + - source: `16db67a94230f93a73be6a608cb4b357ec0231c4` + - merge: `199651d8ec632a1918f64d3e9f015b5a18f84111` +- PR #19: https://github.com/Prometheus-AGS/surreal-memory-server/pull/19 + - source: `9e9c180d4f9bc8c88b96770ea413e9a1c2bcc233` + - merge: `c88144158d8e7bfc6f690d88586e5e68183521d7` +- The PR #19 isolated artifact critic returned PASS. PR #18 also passed the + repository Rust audit gate with no blocking finding. + +## Local verification + +- `RUSTC_WRAPPER= cargo test --locked operations::tests::startup_reconciliation -- --nocapture` + — exit 0; 2 tests passed. +- `RUSTC_WRAPPER= cargo test --locked operations::tests::startup_reconciliation_query_projects_only_operation_identity -- --nocapture` + — exit 0; 1 test passed. +- `cargo fmt --all --check` — exit 0 on the merged tree. +- `openspec validate bound-operation-reconciliation-projection --strict` — + exit 0; the change is valid on the merged tree. + +## Deployment state + +The projection changes are installed through the signed +`surreal-memory-server` binary whose SHA-256 is +`63d4b297a9e4b1ebd2cb61ebac0752a00532eaf611b5b0a967ed9dc875942046` +at both `~/.local/bin` and `/usr/local/bin`. Both copies pass +`codesign --verify`. + +Backlog recovery is still running. The installed service reduced the local +durable queue from 271 accepted and 2,286 completed operations to 261 accepted +and 2,296 completed operations while both database health and memory readiness +returned HTTP 200. Task 1.3 remains open until the accepted count reaches zero +and the final doctor check exits successfully. diff --git a/openspec/changes/complete-operation-ledger-recovery/evidence.md b/openspec/changes/complete-operation-ledger-recovery/evidence.md new file mode 100644 index 0000000..d444f26 --- /dev/null +++ b/openspec/changes/complete-operation-ledger-recovery/evidence.md @@ -0,0 +1,59 @@ +# Verification evidence + +Date: 2026-09-21 + +## Behavior + +- Server-mode operation-ledger initialization opens a separately authenticated + transport before the first query; embedded mode clones its in-process handle + and does not reopen RocksDB. +- Initialization and replacement share one async mutex and check the published + generation after acquiring it, so overlapping stale callers publish only one + next generation. +- A database deadline is retryable only when replacement succeeded or another + caller already published a newer generation. Replacement failure preserves + the existing API error but is not classified as recovered. +- Startup discovery and the coordinator drain each retry at most once only for + that typed recovered condition. Executor and ordinary errors are not retried. + +## Local verification + +- `cargo fmt --all --check` — exit 0. +- `git diff --check` — exit 0. +- `RUSTC_WRAPPER= cargo check --locked --package surreal-memory-server --no-default-features --features server-only` + — exit 0. +- `RUSTC_WRAPPER= cargo test --locked --lib operations::tests:: -- --nocapture` + — exit 0; 17 passed, 0 failed. +- `RUSTC_WRAPPER= cargo test --locked --test operation_query_deadline --no-default-features --features server-only -- --nocapture` + — exit 0; 1 passed, 0 failed. The real isolated server kept general storage + responsive while four ledger receipt queries timed out, then the same + production coordinator committed later work. +- `RUSTC_WRAPPER= cargo test --locked --test executor_recovery -- --nocapture` + — exit 0; 4 passed, 0 failed. +- `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. +- `prometheus-rust-auditor format`, `enforce`, and `inventory` — exit 0; no + findings. `partition` exited 0 with six informational AI-loop-pending rows. +- `prometheus-rust-auditor deps` could not pass because `cargo-deny` is absent + and the committed lockfile already contains the same three `cargo audit` + advisories as the working lockfile: RUSTSEC-2026-0235, RUSTSEC-2023-0071, + and RUSTSEC-2026-0285. This change adds only the already workspace-pinned + `arc-swap` package dependency and introduces none of those advisory paths. + +## Isolated review + +- Round 1: BLOCK. It required exercising the actual coordinator drain retry + branch and removing the duplicate mutation to the prior deadline spec. +- Repairs: added a deterministic `drain_pending` regression for recovered, + repeated, and executor errors; restored the prior change's spec to committed + bytes. +- Round 2: PASS. No blocking findings. + +## Deployment state + +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.