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
Original file line number Diff line number Diff line change
@@ -0,0 +1,30 @@
## Why

Startup reconciliation selects every field from every nonterminal operation,
including the full memory payload. The deployed server stalled before resuming
305 accepted operations while the same records remained individually readable.
The coordinator needs only each operation identity to schedule recovery.

## What Changes

- Project only `operation_id` when listing nonterminal work.
- Retain deterministic ordering and the existing state filter.
- Add a database-backed regression proving large payload bytes do not enter the
reconciliation list result.

## Capabilities

### New Capabilities

- `operation-reconciliation-projection`: Bounds startup work discovery to the
identity data required by the coordinator.

### Modified Capabilities

None.

## Impact

This changes one internal query and its row type in `src/operations.rs`. The
operation schema, state machine, processing order, HTTP contract, and stored
payloads remain unchanged.
Original file line number Diff line number Diff line change
@@ -0,0 +1,21 @@
## Purpose

Bound durable operation recovery discovery independently of stored payload size.

## ADDED Requirements

### Requirement: Reconciliation discovery projects only operation identity

The durable coordinator SHALL list nonterminal operations by projecting only
the stable operation identity required to schedule processing. It MUST preserve
deterministic operation-ID ordering and MUST NOT read stored payload, result, or
executor fields into the reconciliation list result.

#### Scenario: Many nonterminal operations contain large payloads
- **WHEN** startup reconciliation discovers accepted, blocked, or processing operations
- **THEN** its list result contains only each `operation_id` in deterministic order
- **AND** payload size does not increase the returned reconciliation data

#### Scenario: Terminal operations share the ledger
- **WHEN** committed or rejected operations exist beside nonterminal operations
- **THEN** the reconciliation list excludes terminal identities without loading their payloads
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
## 1. Bounded reconciliation discovery

- [x] 1.1 Project only nonterminal operation identities in deterministic order.
- [x] 1.2 Add a database-backed large-payload projection regression.
- [ ] 1.3 Pass formatting, focused operations tests, compilation, strict OpenSpec validation, and deployed backlog recovery.
62 changes: 60 additions & 2 deletions src/operations.rs
Original file line number Diff line number Diff line change
Expand Up @@ -176,6 +176,13 @@ struct DbOperation {
updated_at: Datetime,
}

#[derive(Debug, Deserialize, SurrealValue)]
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";

#[derive(Debug, Clone, Serialize, Deserialize, SurrealValue)]
struct DbOperationEvent {
#[serde(default)]
Expand Down Expand Up @@ -435,8 +442,8 @@ impl OperationService {

async fn list_nonterminal_ids(&self) -> Result<Vec<String>> {
let db = self.surreal()?.db()?;
let rows: Vec<DbOperation> = db
.query("SELECT * FROM memory_operation WHERE state NOT IN ['committed', 'rejected'] ORDER BY operation_id ASC")
let rows: Vec<DbOperationId> = db
.query(LIST_NONTERMINAL_OPERATION_IDS_QUERY)
.await?
.check()?
.take(0)?;
Expand Down Expand Up @@ -1981,6 +1988,57 @@ mod tests {
assert_eq!(service.reconcile_nonterminal().await.unwrap(), 12);
}

#[tokio::test]
async fn startup_reconciliation_query_projects_only_operation_identity() {
let embedder: Arc<dyn EmbeddingService> = Arc::new(NoOpEmbedder);
let storage = Arc::new(
SurrealStorage::new_mem(Arc::clone(&embedder))
.await
.expect("in-memory SurrealStorage"),
);
let service = OperationService::start_with_capacities(
Arc::clone(&storage) as Arc<dyn MemoryStorage>,
embedder,
4,
16,
);
let sentinel = "payload-must-not-enter-reconciliation-list";
let payload = json!({
"name": "projected-reconciliation",
"description": sentinel.repeat(16_384),
"agent_id": null,
"user_id": "test"
});
service
.submit(OperationRequest {
operation_id: "projected-reconciliation".to_owned(),
schema_version: OPERATION_SCHEMA_VERSION,
kind: "create_task_stream".to_owned(),
dependencies: vec!["missing-prerequisite".to_owned()],
payload_hash: payload_hash(&payload).unwrap(),
payload,
})
.await
.unwrap();

let rows: Vec<Value> = storage
.db()
.unwrap()
.query(LIST_NONTERMINAL_OPERATION_IDS_QUERY)
.await
.unwrap()
.check()
.unwrap()
.take(0)
.unwrap();
let row = rows
.iter()
.find(|row| row["operation_id"] == "projected-reconciliation")
.expect("nonterminal operation identity is projected");
assert_eq!(row.as_object().unwrap().len(), 1);
assert!(!serde_json::to_string(&rows).unwrap().contains(sentinel));
}

#[tokio::test]
async fn lagged_event_subscriber_backfills_every_durable_sequence() {
let embedder: Arc<dyn EmbeddingService> = Arc::new(NoOpEmbedder);
Expand Down
Loading