From e53b39ffcad619f2c0776efb8ec66545c0ee4a2f Mon Sep 17 00:00:00 2001 From: serrrfirat Date: Mon, 11 May 2026 12:40:46 +0200 Subject: [PATCH 1/8] =?UTF-8?q?feat(KB-002):=20complete=20Step=201=20?= =?UTF-8?q?=E2=80=94=20extend=20TurnCheckpointRecord=20and=20define=20Loop?= =?UTF-8?q?CheckpointStore=20trait?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- crates/ironclaw_turns/src/lib.rs | 8 ++--- crates/ironclaw_turns/src/memory.rs | 6 ++++ crates/ironclaw_turns/src/store.rs | 53 +++++++++++++++++++++++++++++ 3 files changed, 63 insertions(+), 4 deletions(-) diff --git a/crates/ironclaw_turns/src/lib.rs b/crates/ironclaw_turns/src/lib.rs index c1e671772b7..b2d3ef4b9d0 100644 --- a/crates/ironclaw_turns/src/lib.rs +++ b/crates/ironclaw_turns/src/lib.rs @@ -81,8 +81,8 @@ pub use status::{ SanitizedFailure, TurnError, TurnErrorCategory, TurnRunProfile, TurnRunState, TurnStatus, }; pub use store::{ - TurnActiveLockKey, TurnActiveLockRecord, TurnCheckpointRecord, TurnIdempotencyErrorReplay, - TurnIdempotencyOperationKind, TurnIdempotencyOutcomeKind, TurnIdempotencyRecord, - TurnIdempotencyReplay, TurnLockVersion, TurnPersistenceSnapshot, TurnRecord, TurnRunRecord, - TurnStateStore, + LoopCheckpointStore, TurnActiveLockKey, TurnActiveLockRecord, TurnCheckpointRecord, + TurnIdempotencyErrorReplay, TurnIdempotencyOperationKind, TurnIdempotencyOutcomeKind, + TurnIdempotencyRecord, TurnIdempotencyReplay, TurnLockVersion, TurnPersistenceSnapshot, + TurnRecord, TurnRunRecord, TurnStateStore, }; diff --git a/crates/ironclaw_turns/src/memory.rs b/crates/ironclaw_turns/src/memory.rs index 6b68d84214c..a98564c00ff 100644 --- a/crates/ironclaw_turns/src/memory.rs +++ b/crates/ironclaw_turns/src/memory.rs @@ -1722,9 +1722,15 @@ impl Inner { self.checkpoints.push(TurnCheckpointRecord { checkpoint_id, run_id: record.run_id, + scope: Some(record.scope.clone()), sequence, status: record.status, gate_ref, + kind: crate::run_profile::LoopCheckpointKind::BeforeBlock, + // Placeholder — callers in block_run don't have the loop's actual state_ref. + // Real values will be threaded in a follow-up task. + state_ref: crate::run_profile::LoopCheckpointStateRef::new("checkpoint:block-state") + .unwrap(), created_at, }); } diff --git a/crates/ironclaw_turns/src/store.rs b/crates/ironclaw_turns/src/store.rs index d5cfb221024..32ed0bcfc81 100644 --- a/crates/ironclaw_turns/src/store.rs +++ b/crates/ironclaw_turns/src/store.rs @@ -9,6 +9,7 @@ use crate::{ TurnCheckpointId, TurnError, TurnErrorCategory, TurnId, TurnLeaseToken, TurnLifecycleEvent, TurnRunId, TurnRunProfile, TurnRunState, TurnRunnerId, TurnScope, TurnStatus, TurnTimestamp, events::EventCursor, + run_profile::{LoopCheckpointKind, LoopCheckpointStateRef}, }; #[async_trait] @@ -107,16 +108,68 @@ pub struct TurnActiveLockRecord { pub updated_at: TurnTimestamp, } +/// Serde default for `LoopCheckpointKind` — used when deserializing old +/// persisted data that predates the `kind` field. +fn default_checkpoint_kind() -> LoopCheckpointKind { + LoopCheckpointKind::BeforeBlock +} + +/// Serde default for `LoopCheckpointStateRef` — sentinel value used when +/// deserializing old persisted data that predates the `state_ref` field. +/// Real values will be threaded from loop callers in a follow-up task. +fn default_checkpoint_state_ref() -> LoopCheckpointStateRef { + // Safety: literal satisfies the "checkpoint:" prefix and length constraints. + LoopCheckpointStateRef::new("checkpoint:unknown").unwrap() +} + #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] pub struct TurnCheckpointRecord { pub checkpoint_id: TurnCheckpointId, pub run_id: TurnRunId, + /// Scope of the run that created this checkpoint. `None` for legacy records + /// persisted before scope was added to checkpoints. + #[serde(default)] + pub scope: Option, pub sequence: u64, pub status: TurnStatus, pub gate_ref: GateRef, + /// The semantic kind of checkpoint (before model, side-effect, block, final). + #[serde(default = "default_checkpoint_kind")] + pub kind: LoopCheckpointKind, + /// An opaque ref describing the loop state at the time of this checkpoint. + #[serde(default = "default_checkpoint_state_ref")] + pub state_ref: LoopCheckpointStateRef, pub created_at: TurnTimestamp, } +/// Direct single-row checkpoint read/write operations. +/// +/// Unlike the snapshot-based path (`TurnRunTransitionPort::block_run` → +/// `libsql_replace_snapshot`/`postgres_replace_snapshot`) which loads the full +/// `TurnPersistenceSnapshot`, `LoopCheckpointStore` provides targeted +/// `INSERT`/`SELECT` operations scoped by checkpoint ID and run. +/// +/// Implementations must preserve idempotent `put` semantics (re-inserting the +/// same `checkpoint_id` is not an error) and cross-run isolation (`get` returns +/// `None` when the checkpoint belongs to a different `run_id`). +#[async_trait] +pub trait LoopCheckpointStore: Send + Sync { + /// Insert a single checkpoint record scoped to the run. + /// Idempotent: re-inserting the same `checkpoint_id` is not an error. + async fn put_loop_checkpoint( + &self, + record: TurnCheckpointRecord, + ) -> Result<(), TurnError>; + + /// Direct lookup by `checkpoint_id`. Returns `None` if not found or if the + /// checkpoint belongs to a different `run_id` (cross-run rejection). + async fn get_loop_checkpoint( + &self, + checkpoint_id: TurnCheckpointId, + run_id: TurnRunId, + ) -> Result, TurnError>; +} + #[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)] #[serde(rename_all = "snake_case")] pub enum TurnIdempotencyOperationKind { From 9965120f4f6e69d3a443e2d73c33613248ba513a Mon Sep 17 00:00:00 2001 From: serrrfirat Date: Mon, 11 May 2026 12:44:48 +0200 Subject: [PATCH 2/8] =?UTF-8?q?feat(KB-002):=20complete=20Step=202=20?= =?UTF-8?q?=E2=80=94=20add=20direct=20DB=20operations=20for=20libSQL=20and?= =?UTF-8?q?=20Postgres?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- crates/ironclaw_turns/src/db.rs | 152 ++++++++++++++++++++++++++++++-- 1 file changed, 147 insertions(+), 5 deletions(-) diff --git a/crates/ironclaw_turns/src/db.rs b/crates/ironclaw_turns/src/db.rs index a32d626999f..cf5709ce569 100644 --- a/crates/ironclaw_turns/src/db.rs +++ b/crates/ironclaw_turns/src/db.rs @@ -9,9 +9,9 @@ use crate::{ PutLoopCheckpointRequest, ResolvedRunProfile, ResumeTurnRequest, ResumeTurnResponse, RunProfileResolutionError, RunProfileResolutionRequest, RunProfileResolver, SubmitTurnRequest, SubmitTurnResponse, TurnActiveLockRecord, TurnAdmissionLimitProvider, TurnAdmissionPolicy, - TurnAdmissionReservationRecord, TurnCheckpointRecord, TurnError, TurnIdempotencyRecord, - TurnLifecycleEvent, TurnPersistenceSnapshot, TurnRecord, TurnRunRecord, TurnRunState, - TurnScope, TurnStateStore, + TurnAdmissionReservationRecord, TurnCheckpointId, TurnCheckpointRecord, TurnError, + TurnIdempotencyRecord, TurnLifecycleEvent, TurnPersistenceSnapshot, TurnRecord, TurnRunId, + TurnRunRecord, TurnRunState, TurnScope, TurnStateStore, events::{EventCursor, TurnEventPage, TurnEventProjectionSource, project_turn_events}, runner::{ ApplyValidatedLoopExitRequest, BlockRunRequest, CancelRunCompletionRequest, @@ -85,6 +85,8 @@ CREATE TABLE IF NOT EXISTS turn_checkpoints ( checkpoint_id TEXT PRIMARY KEY, run_id TEXT NOT NULL, sequence INTEGER NOT NULL, + scope_key TEXT NOT NULL DEFAULT '', + kind TEXT NOT NULL DEFAULT 'before_block', payload TEXT NOT NULL ); CREATE INDEX IF NOT EXISTS idx_turn_checkpoints_run ON turn_checkpoints(run_id, sequence); @@ -165,6 +167,8 @@ CREATE TABLE IF NOT EXISTS turn_checkpoints ( checkpoint_id TEXT PRIMARY KEY, run_id TEXT NOT NULL, sequence BIGINT NOT NULL, + scope_key TEXT NOT NULL DEFAULT '', + kind TEXT NOT NULL DEFAULT 'before_block', payload JSONB NOT NULL ); CREATE INDEX IF NOT EXISTS idx_turn_checkpoints_run ON turn_checkpoints(run_id, sequence); @@ -248,6 +252,20 @@ impl LibSqlTurnStateStore { conn.execute_batch(LIBSQL_TURN_STATE_SCHEMA) .await .map_err(db_error)?; + + // Migration: add new columns to existing turn_checkpoints tables. + // For libSQL, ALTER TABLE ADD COLUMN fails if the column already exists, + // so we ignore "duplicate column name" errors. + for alter in [ + "ALTER TABLE turn_checkpoints ADD COLUMN scope_key TEXT NOT NULL DEFAULT ''", + "ALTER TABLE turn_checkpoints ADD COLUMN kind TEXT NOT NULL DEFAULT 'before_block'", + ] { + match conn.execute(alter, ()).await { + Ok(_) => {} + Err(e) if e.to_string().contains("duplicate column name") => {} + Err(e) => return Err(db_error(e)), + } + } Ok(()) } @@ -538,6 +556,56 @@ impl TurnRunTransitionPort for LibSqlTurnStateStore { } } +#[cfg(feature = "libsql")] +#[async_trait] +impl LoopCheckpointStore for LibSqlTurnStateStore { + async fn put_loop_checkpoint( + &self, + record: TurnCheckpointRecord, + ) -> Result<(), TurnError> { + let conn = self.connect().await?; + let checkpoint_scope_key = record.scope.as_ref().map(scope_key).transpose()?.unwrap_or_default(); + conn.execute( + "INSERT OR REPLACE INTO turn_checkpoints (checkpoint_id, run_id, sequence, scope_key, kind, payload) VALUES (?1, ?2, ?3, ?4, ?5, ?6)", + libsql::params![ + record.checkpoint_id.as_uuid().to_string(), + record.run_id.to_string(), + record.sequence as i64, + checkpoint_scope_key, + record.kind.as_str(), + to_json(&record)?, + ], + ) + .await + .map_err(db_error)?; + Ok(()) + } + + async fn get_loop_checkpoint( + &self, + checkpoint_id: TurnCheckpointId, + run_id: TurnRunId, + ) -> Result, TurnError> { + let conn = self.connect().await?; + let mut rows = conn + .query( + "SELECT payload FROM turn_checkpoints WHERE checkpoint_id = ?1 AND run_id = ?2", + libsql::params![ + checkpoint_id.as_uuid().to_string(), + run_id.to_string(), + ], + ) + .await + .map_err(db_error)?; + let Some(row) = rows.next().await.map_err(db_error)? else { + return Ok(None); + }; + let payload: String = row.get(0).map_err(db_error)?; + let record: TurnCheckpointRecord = serde_json::from_str(&payload).map_err(db_error)?; + Ok(Some(record)) + } +} + #[cfg(feature = "postgres")] pub struct PostgresTurnStateStore { pool: deadpool_postgres::Pool, @@ -574,6 +642,16 @@ impl PostgresTurnStateStore { .batch_execute(POSTGRES_TURN_STATE_SCHEMA) .await .map_err(db_error)?; + + // Migration: add new columns to existing turn_checkpoints tables. + // Postgres supports ADD COLUMN IF NOT EXISTS natively. + client + .batch_execute( + "ALTER TABLE turn_checkpoints ADD COLUMN IF NOT EXISTS scope_key TEXT NOT NULL DEFAULT ''; + ALTER TABLE turn_checkpoints ADD COLUMN IF NOT EXISTS kind TEXT NOT NULL DEFAULT 'before_block';", + ) + .await + .map_err(db_error)?; Ok(()) } @@ -841,6 +919,65 @@ impl TurnRunTransitionPort for PostgresTurnStateStore { } } +#[cfg(feature = "postgres")] +#[async_trait] +impl LoopCheckpointStore for PostgresTurnStateStore { + async fn put_loop_checkpoint( + &self, + record: TurnCheckpointRecord, + ) -> Result<(), TurnError> { + let client = self.client().await?; + let payload = to_json(&record)?; + let checkpoint_scope_key = record.scope.as_ref().map(scope_key).transpose()?.unwrap_or_default(); + client + .execute( + "INSERT INTO turn_checkpoints (checkpoint_id, run_id, sequence, scope_key, kind, payload) + VALUES ($1, $2, $3, $4, $5, $6::jsonb) + ON CONFLICT (checkpoint_id) DO UPDATE SET + run_id = EXCLUDED.run_id, + sequence = EXCLUDED.sequence, + scope_key = EXCLUDED.scope_key, + kind = EXCLUDED.kind, + payload = EXCLUDED.payload", + &[ + &record.checkpoint_id.as_uuid().to_string(), + &record.run_id.to_string(), + &(record.sequence as i64), + &checkpoint_scope_key, + &record.kind.as_str(), + &payload, + ], + ) + .await + .map_err(db_error)?; + Ok(()) + } + + async fn get_loop_checkpoint( + &self, + checkpoint_id: TurnCheckpointId, + run_id: TurnRunId, + ) -> Result, TurnError> { + let client = self.client().await?; + let row = client + .query_opt( + "SELECT payload::text FROM turn_checkpoints WHERE checkpoint_id = $1 AND run_id = $2", + &[ + &checkpoint_id.as_uuid().to_string(), + &run_id.to_string(), + ], + ) + .await + .map_err(db_error)?; + let Some(row) = row else { + return Ok(None); + }; + let payload: String = row.get(0); + let record: TurnCheckpointRecord = serde_json::from_str(&payload).map_err(db_error)?; + Ok(Some(record)) + } +} + #[cfg(feature = "libsql")] async fn libsql_load_payloads(conn: &libsql::Connection, sql: &str) -> Result, TurnError> where @@ -1011,11 +1148,13 @@ async fn libsql_replace_snapshot( } for record in &snapshot.checkpoints { conn.execute( - "INSERT INTO turn_checkpoints (checkpoint_id, run_id, sequence, payload) VALUES (?1, ?2, ?3, ?4)", + "INSERT INTO turn_checkpoints (checkpoint_id, run_id, sequence, scope_key, kind, payload) VALUES (?1, ?2, ?3, ?4, ?5, ?6)", libsql::params![ record.checkpoint_id.as_uuid().to_string(), record.run_id.to_string(), record.sequence as i64, + record.scope.as_ref().map(scope_key).transpose()?.unwrap_or_default(), + record.kind.as_str(), to_json(record)?, ], ) @@ -1259,12 +1398,15 @@ async fn postgres_replace_snapshot( } for record in &snapshot.checkpoints { let payload = to_json(record)?; + let checkpoint_scope_key = record.scope.as_ref().map(scope_key).transpose()?.unwrap_or_default(); txn.execute( - "INSERT INTO turn_checkpoints (checkpoint_id, run_id, sequence, payload) VALUES ($1, $2, $3, $4::jsonb)", + "INSERT INTO turn_checkpoints (checkpoint_id, run_id, sequence, scope_key, kind, payload) VALUES ($1, $2, $3, $4, $5, $6::jsonb)", &[ &record.checkpoint_id.as_uuid().to_string(), &record.run_id.to_string(), &(record.sequence as i64), + &checkpoint_scope_key, + &record.kind.as_str(), &payload, ], ) From 62d11db6a0c8f6dea565794f1daf4a82b80e8388 Mon Sep 17 00:00:00 2001 From: serrrfirat Date: Mon, 11 May 2026 12:47:10 +0200 Subject: [PATCH 3/8] =?UTF-8?q?feat(KB-002):=20complete=20Step=203=20?= =?UTF-8?q?=E2=80=94=20implement=20in-memory=20LoopCheckpointStore?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- crates/ironclaw_turns/src/memory.rs | 34 +++++++++++++++++++++++++++++ 1 file changed, 34 insertions(+) diff --git a/crates/ironclaw_turns/src/memory.rs b/crates/ironclaw_turns/src/memory.rs index a98564c00ff..802892a32c1 100644 --- a/crates/ironclaw_turns/src/memory.rs +++ b/crates/ironclaw_turns/src/memory.rs @@ -725,6 +725,40 @@ impl TurnRunTransitionPort for InMemoryTurnStateStore { } } +#[async_trait] +impl crate::LoopCheckpointStore for InMemoryTurnStateStore { + async fn put_loop_checkpoint( + &self, + record: TurnCheckpointRecord, + ) -> Result<(), crate::TurnError> { + let mut inner = self.lock_inner()?; + // Idempotent: update if same checkpoint_id exists, else push. + if let Some(existing) = inner + .checkpoints + .iter_mut() + .find(|c| c.checkpoint_id == record.checkpoint_id) + { + *existing = record; + } else { + inner.checkpoints.push(record); + } + Ok(()) + } + + async fn get_loop_checkpoint( + &self, + checkpoint_id: crate::TurnCheckpointId, + run_id: crate::TurnRunId, + ) -> Result, crate::TurnError> { + let inner = self.lock_inner()?; + Ok(inner + .checkpoints + .iter() + .find(|c| c.checkpoint_id == checkpoint_id && c.run_id == run_id) + .cloned()) + } +} + impl Inner { fn from_persistence_snapshot( snapshot: TurnPersistenceSnapshot, From 3de543232b230b4737c2778af4d15e6cdccd3e51 Mon Sep 17 00:00:00 2001 From: serrrfirat Date: Mon, 11 May 2026 12:51:52 +0200 Subject: [PATCH 4/8] =?UTF-8?q?test(KB-002):=20complete=20Step=204=20?= =?UTF-8?q?=E2=80=94=20loop=20checkpoint=20store=20contract=20tests?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../tests/loop_checkpoint_store_contract.rs | 221 ++++++++++++++++++ 1 file changed, 221 insertions(+) create mode 100644 crates/ironclaw_turns/tests/loop_checkpoint_store_contract.rs diff --git a/crates/ironclaw_turns/tests/loop_checkpoint_store_contract.rs b/crates/ironclaw_turns/tests/loop_checkpoint_store_contract.rs new file mode 100644 index 00000000000..f1656960ef3 --- /dev/null +++ b/crates/ironclaw_turns/tests/loop_checkpoint_store_contract.rs @@ -0,0 +1,221 @@ +#![cfg(any(feature = "libsql", feature = "postgres"))] + +use chrono::{TimeZone, Utc}; +use ironclaw_host_api::{AgentId, ProjectId, TenantId, ThreadId, UserId}; +use ironclaw_turns::{ + GateRef, LoopCheckpointStore, TurnCheckpointId, TurnCheckpointRecord, TurnRunId, TurnScope, + TurnStatus, + run_profile::{LoopCheckpointKind, LoopCheckpointStateRef}, +}; +use std::sync::Arc; + +#[cfg(feature = "libsql")] +use ironclaw_turns::LibSqlTurnStateStore; + +fn test_scope(thread: &str) -> TurnScope { + TurnScope::new( + TenantId::new("tenant1").unwrap(), + Some(AgentId::new("agent1").unwrap()), + Some(ProjectId::new("project1").unwrap()), + ThreadId::new(thread).unwrap(), + ) +} + +fn make_checkpoint( + run_id: TurnRunId, + sequence: u64, + kind: LoopCheckpointKind, +) -> TurnCheckpointRecord { + TurnCheckpointRecord { + checkpoint_id: TurnCheckpointId::new(), + run_id, + scope: Some(test_scope("thread-a")), + sequence, + status: TurnStatus::BlockedApproval, + gate_ref: GateRef::new("gate:test-gate").unwrap(), + kind, + state_ref: LoopCheckpointStateRef::new("checkpoint:test-state").unwrap(), + created_at: Utc.with_ymd_and_hms(2026, 5, 11, 12, 0, 0).unwrap(), + } +} + +// ── Tests against in-memory backend ────────────────────────────────────────── + +#[tokio::test] +async fn inmemory_put_get_roundtrip() { + let store = ironclaw_turns::InMemoryTurnStateStore::default(); + let run_id = TurnRunId::new(); + let record = make_checkpoint(run_id, 1, LoopCheckpointKind::BeforeModel); + let checkpoint_id = record.checkpoint_id; + + store.put_loop_checkpoint(record.clone()).await.unwrap(); + let fetched = store + .get_loop_checkpoint(checkpoint_id, run_id) + .await + .unwrap(); + + let fetched = fetched.expect("should find checkpoint"); + assert_eq!(fetched.checkpoint_id, checkpoint_id); + assert_eq!(fetched.run_id, run_id); + assert_eq!(fetched.kind, LoopCheckpointKind::BeforeModel); + assert_eq!( + fetched.state_ref, + LoopCheckpointStateRef::new("checkpoint:test-state").unwrap() + ); + assert_eq!(fetched.scope, Some(test_scope("thread-a"))); +} + +#[tokio::test] +async fn inmemory_cross_run_rejection() { + let store = ironclaw_turns::InMemoryTurnStateStore::default(); + let run_a = TurnRunId::new(); + let run_b = TurnRunId::new(); + let record = make_checkpoint(run_a, 1, LoopCheckpointKind::BeforeBlock); + let checkpoint_id = record.checkpoint_id; + + store.put_loop_checkpoint(record).await.unwrap(); + let fetched = store + .get_loop_checkpoint(checkpoint_id, run_b) + .await + .unwrap(); + assert!(fetched.is_none(), "cross-run lookup should return None"); +} + +#[tokio::test] +async fn inmemory_multiple_checkpoints_per_run() { + let store = ironclaw_turns::InMemoryTurnStateStore::default(); + let run_id = TurnRunId::new(); + + let r1 = make_checkpoint(run_id, 1, LoopCheckpointKind::BeforeModel); + let r2 = make_checkpoint(run_id, 2, LoopCheckpointKind::BeforeSideEffect); + let r3 = make_checkpoint(run_id, 3, LoopCheckpointKind::Final); + let id1 = r1.checkpoint_id; + let id2 = r2.checkpoint_id; + let id3 = r3.checkpoint_id; + + store.put_loop_checkpoint(r1).await.unwrap(); + store.put_loop_checkpoint(r2).await.unwrap(); + store.put_loop_checkpoint(r3).await.unwrap(); + + assert!(store.get_loop_checkpoint(id1, run_id).await.unwrap().is_some()); + assert!(store.get_loop_checkpoint(id2, run_id).await.unwrap().is_some()); + assert!(store.get_loop_checkpoint(id3, run_id).await.unwrap().is_some()); +} + +#[tokio::test] +async fn inmemory_idempotent_put() { + let store = ironclaw_turns::InMemoryTurnStateStore::default(); + let run_id = TurnRunId::new(); + let record = make_checkpoint(run_id, 1, LoopCheckpointKind::BeforeBlock); + + store.put_loop_checkpoint(record.clone()).await.unwrap(); + // Second put with same checkpoint_id should succeed without error. + store.put_loop_checkpoint(record).await.unwrap(); +} + +#[tokio::test] +async fn serde_backward_compat_missing_fields() { + // Simulate old persisted JSON that lacks kind, state_ref, and scope. + let json = r#"{ + "checkpoint_id": "00000000-0000-0000-0000-000000000001", + "run_id": "00000000-0000-0000-0000-000000000002", + "sequence": 1, + "status": "BlockedApproval", + "gate_ref": "gate:test", + "created_at": "2026-05-11T12:00:00Z" + }"#; + let record: TurnCheckpointRecord = serde_json::from_str(json).unwrap(); + assert_eq!(record.kind, LoopCheckpointKind::BeforeBlock); + assert_eq!( + record.state_ref, + LoopCheckpointStateRef::new("checkpoint:unknown").unwrap() + ); + assert_eq!(record.scope, None); +} + +// ── Tests against libSQL backend ───────────────────────────────────────────── + +#[cfg(feature = "libsql")] +async fn libsql_store() -> (Arc, tempfile::TempDir) { + let dir = tempfile::tempdir().unwrap(); + let db_path = dir.path().join("turns.db"); + let db = Arc::new(libsql::Builder::new_local(db_path).build().await.unwrap()); + let store = Arc::new(LibSqlTurnStateStore::new(db)); + store.run_migrations().await.unwrap(); + (store, dir) +} + +#[cfg(feature = "libsql")] +#[tokio::test] +async fn libsql_put_get_roundtrip() { + let (store, _dir) = libsql_store().await; + let run_id = TurnRunId::new(); + let record = make_checkpoint(run_id, 1, LoopCheckpointKind::BeforeModel); + let checkpoint_id = record.checkpoint_id; + + store.put_loop_checkpoint(record.clone()).await.unwrap(); + let fetched = store + .get_loop_checkpoint(checkpoint_id, run_id) + .await + .unwrap() + .expect("should find checkpoint"); + + assert_eq!(fetched.checkpoint_id, checkpoint_id); + assert_eq!(fetched.run_id, run_id); + assert_eq!(fetched.kind, LoopCheckpointKind::BeforeModel); + assert_eq!( + fetched.state_ref, + LoopCheckpointStateRef::new("checkpoint:test-state").unwrap() + ); + assert_eq!(fetched.scope, Some(test_scope("thread-a"))); +} + +#[cfg(feature = "libsql")] +#[tokio::test] +async fn libsql_cross_run_rejection() { + let (store, _dir) = libsql_store().await; + let run_a = TurnRunId::new(); + let run_b = TurnRunId::new(); + let record = make_checkpoint(run_a, 1, LoopCheckpointKind::BeforeBlock); + let checkpoint_id = record.checkpoint_id; + + store.put_loop_checkpoint(record).await.unwrap(); + let fetched = store + .get_loop_checkpoint(checkpoint_id, run_b) + .await + .unwrap(); + assert!(fetched.is_none(), "cross-run lookup should return None"); +} + +#[cfg(feature = "libsql")] +#[tokio::test] +async fn libsql_multiple_checkpoints_per_run() { + let (store, _dir) = libsql_store().await; + let run_id = TurnRunId::new(); + + let r1 = make_checkpoint(run_id, 1, LoopCheckpointKind::BeforeModel); + let r2 = make_checkpoint(run_id, 2, LoopCheckpointKind::BeforeSideEffect); + let r3 = make_checkpoint(run_id, 3, LoopCheckpointKind::Final); + let id1 = r1.checkpoint_id; + let id2 = r2.checkpoint_id; + let id3 = r3.checkpoint_id; + + store.put_loop_checkpoint(r1).await.unwrap(); + store.put_loop_checkpoint(r2).await.unwrap(); + store.put_loop_checkpoint(r3).await.unwrap(); + + assert!(store.get_loop_checkpoint(id1, run_id).await.unwrap().is_some()); + assert!(store.get_loop_checkpoint(id2, run_id).await.unwrap().is_some()); + assert!(store.get_loop_checkpoint(id3, run_id).await.unwrap().is_some()); +} + +#[cfg(feature = "libsql")] +#[tokio::test] +async fn libsql_idempotent_put() { + let (store, _dir) = libsql_store().await; + let run_id = TurnRunId::new(); + let record = make_checkpoint(run_id, 1, LoopCheckpointKind::BeforeBlock); + + store.put_loop_checkpoint(record.clone()).await.unwrap(); + store.put_loop_checkpoint(record).await.unwrap(); +} From fe26dbfec0d736a74213d09059d6e717dbd2f415 Mon Sep 17 00:00:00 2001 From: serrrfirat Date: Mon, 11 May 2026 12:53:38 +0200 Subject: [PATCH 5/8] fix(KB-002): remove unused UserId import in test --- crates/ironclaw_turns/tests/loop_checkpoint_store_contract.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/crates/ironclaw_turns/tests/loop_checkpoint_store_contract.rs b/crates/ironclaw_turns/tests/loop_checkpoint_store_contract.rs index f1656960ef3..7378235273b 100644 --- a/crates/ironclaw_turns/tests/loop_checkpoint_store_contract.rs +++ b/crates/ironclaw_turns/tests/loop_checkpoint_store_contract.rs @@ -1,7 +1,7 @@ #![cfg(any(feature = "libsql", feature = "postgres"))] use chrono::{TimeZone, Utc}; -use ironclaw_host_api::{AgentId, ProjectId, TenantId, ThreadId, UserId}; +use ironclaw_host_api::{AgentId, ProjectId, TenantId, ThreadId}; use ironclaw_turns::{ GateRef, LoopCheckpointStore, TurnCheckpointId, TurnCheckpointRecord, TurnRunId, TurnScope, TurnStatus, From 4fb85b4eb4065be2beb30cecc0dd9ae25cfef5a5 Mon Sep 17 00:00:00 2001 From: serrrfirat Date: Mon, 11 May 2026 15:45:18 +0200 Subject: [PATCH 6/8] fix(KB-002): use loop checkpoint mapping store --- crates/ironclaw_turns/src/db.rs | 450 ++++++++++++------ crates/ironclaw_turns/src/lib.rs | 8 +- crates/ironclaw_turns/src/memory.rs | 34 -- crates/ironclaw_turns/src/store.rs | 28 -- .../tests/checkpoint_state_store_contract.rs | 5 +- .../tests/loop_checkpoint_store_contract.rs | 305 ++++++------ 6 files changed, 461 insertions(+), 369 deletions(-) diff --git a/crates/ironclaw_turns/src/db.rs b/crates/ironclaw_turns/src/db.rs index cf5709ce569..b0624ca223d 100644 --- a/crates/ironclaw_turns/src/db.rs +++ b/crates/ironclaw_turns/src/db.rs @@ -1,5 +1,7 @@ use async_trait::async_trait; #[cfg(any(feature = "libsql", feature = "postgres"))] +use chrono::Utc; +#[cfg(any(feature = "libsql", feature = "postgres"))] use std::sync::Arc; use crate::{ @@ -9,9 +11,9 @@ use crate::{ PutLoopCheckpointRequest, ResolvedRunProfile, ResumeTurnRequest, ResumeTurnResponse, RunProfileResolutionError, RunProfileResolutionRequest, RunProfileResolver, SubmitTurnRequest, SubmitTurnResponse, TurnActiveLockRecord, TurnAdmissionLimitProvider, TurnAdmissionPolicy, - TurnAdmissionReservationRecord, TurnCheckpointId, TurnCheckpointRecord, TurnError, - TurnIdempotencyRecord, TurnLifecycleEvent, TurnPersistenceSnapshot, TurnRecord, TurnRunId, - TurnRunRecord, TurnRunState, TurnScope, TurnStateStore, + TurnAdmissionReservationRecord, TurnCheckpointRecord, TurnError, TurnIdempotencyRecord, + TurnLifecycleEvent, TurnPersistenceSnapshot, TurnRecord, TurnRunRecord, TurnRunState, + TurnScope, TurnStateStore, events::{EventCursor, TurnEventPage, TurnEventProjectionSource, project_turn_events}, runner::{ ApplyValidatedLoopExitRequest, BlockRunRequest, CancelRunCompletionRequest, @@ -402,27 +404,51 @@ impl LoopCheckpointStore for LibSqlTurnStateStore { request: PutLoopCheckpointRequest, ) -> Result { let conn = self.begin_immediate().await?; - let result = async { - let store = self.load_store_from_conn(&conn).await?; - let result = store.put_loop_checkpoint(request).await; - libsql_replace_snapshot(&conn, &store.persistence_snapshot()).await?; - Ok(result) - } - .await; - finish_libsql_transaction(&conn, result).await? + let record = LoopCheckpointRecord { + checkpoint_id: crate::TurnCheckpointId::new(), + scope: request.scope, + turn_id: request.turn_id, + run_id: request.run_id, + state_ref: request.state_ref, + schema_id: request.schema_id, + schema_version: request.schema_version, + kind: request.kind, + created_at: Utc::now(), + }; + let result = libsql_insert_loop_checkpoint_record(&conn, &record) + .await + .map(|()| record.clone()); + finish_libsql_transaction(&conn, result).await } async fn get_loop_checkpoint( &self, request: GetLoopCheckpointRequest, ) -> Result, TurnError> { - self.load_snapshot() - .await - .and_then(|snapshot| { - InMemoryTurnStateStore::from_persistence_snapshot(snapshot, self.limits) - })? - .get_loop_checkpoint(request) + let conn = self.connect().await?; + let scope_key = scope_key(&request.scope)?; + let mut rows = conn + .query( + "SELECT payload FROM turn_loop_checkpoints WHERE checkpoint_id = ?1 AND scope_key = ?2 AND turn_id = ?3 AND run_id = ?4", + libsql::params![ + request.checkpoint_id.as_uuid().to_string(), + scope_key, + request.turn_id.to_string(), + request.run_id.to_string(), + ], + ) .await + .map_err(db_error)?; + let Some(row) = rows.next().await.map_err(db_error)? else { + return Ok(None); + }; + let payload: String = row.get(0).map_err(db_error)?; + let record = serde_json::from_str(&payload).map_err(db_error)?; + if loop_checkpoint_record_matches_request(&record, &request) { + Ok(Some(record)) + } else { + Ok(None) + } } } @@ -556,56 +582,6 @@ impl TurnRunTransitionPort for LibSqlTurnStateStore { } } -#[cfg(feature = "libsql")] -#[async_trait] -impl LoopCheckpointStore for LibSqlTurnStateStore { - async fn put_loop_checkpoint( - &self, - record: TurnCheckpointRecord, - ) -> Result<(), TurnError> { - let conn = self.connect().await?; - let checkpoint_scope_key = record.scope.as_ref().map(scope_key).transpose()?.unwrap_or_default(); - conn.execute( - "INSERT OR REPLACE INTO turn_checkpoints (checkpoint_id, run_id, sequence, scope_key, kind, payload) VALUES (?1, ?2, ?3, ?4, ?5, ?6)", - libsql::params![ - record.checkpoint_id.as_uuid().to_string(), - record.run_id.to_string(), - record.sequence as i64, - checkpoint_scope_key, - record.kind.as_str(), - to_json(&record)?, - ], - ) - .await - .map_err(db_error)?; - Ok(()) - } - - async fn get_loop_checkpoint( - &self, - checkpoint_id: TurnCheckpointId, - run_id: TurnRunId, - ) -> Result, TurnError> { - let conn = self.connect().await?; - let mut rows = conn - .query( - "SELECT payload FROM turn_checkpoints WHERE checkpoint_id = ?1 AND run_id = ?2", - libsql::params![ - checkpoint_id.as_uuid().to_string(), - run_id.to_string(), - ], - ) - .await - .map_err(db_error)?; - let Some(row) = rows.next().await.map_err(db_error)? else { - return Ok(None); - }; - let payload: String = row.get(0).map_err(db_error)?; - let record: TurnCheckpointRecord = serde_json::from_str(&payload).map_err(db_error)?; - Ok(Some(record)) - } -} - #[cfg(feature = "postgres")] pub struct PostgresTurnStateStore { pool: deadpool_postgres::Pool, @@ -774,27 +750,50 @@ impl LoopCheckpointStore for PostgresTurnStateStore { &self, request: PutLoopCheckpointRequest, ) -> Result { - let mut client = self.client().await?; - let txn = client.transaction().await.map_err(db_error)?; - lock_postgres_turn_tables(&txn, "SHARE ROW EXCLUSIVE MODE").await?; - let store = self.load_store_from_txn(&txn).await?; - let result = store.put_loop_checkpoint(request).await; - postgres_replace_snapshot(&txn, &store.persistence_snapshot()).await?; - txn.commit().await.map_err(db_error)?; - result + let client = self.client().await?; + let record = LoopCheckpointRecord { + checkpoint_id: crate::TurnCheckpointId::new(), + scope: request.scope, + turn_id: request.turn_id, + run_id: request.run_id, + state_ref: request.state_ref, + schema_id: request.schema_id, + schema_version: request.schema_version, + kind: request.kind, + created_at: Utc::now(), + }; + postgres_insert_loop_checkpoint_record(&client, &record).await?; + Ok(record) } async fn get_loop_checkpoint( &self, request: GetLoopCheckpointRequest, ) -> Result, TurnError> { - self.load_snapshot() - .await - .and_then(|snapshot| { - InMemoryTurnStateStore::from_persistence_snapshot(snapshot, self.limits) - })? - .get_loop_checkpoint(request) + let client = self.client().await?; + let scope_key = scope_key(&request.scope)?; + let row = client + .query_opt( + "SELECT payload::text FROM turn_loop_checkpoints WHERE checkpoint_id = $1 AND scope_key = $2 AND turn_id = $3 AND run_id = $4", + &[ + &request.checkpoint_id.as_uuid().to_string(), + &scope_key, + &request.turn_id.to_string(), + &request.run_id.to_string(), + ], + ) .await + .map_err(db_error)?; + let Some(row) = row else { + return Ok(None); + }; + let payload: String = row.get(0); + let record = serde_json::from_str(&payload).map_err(db_error)?; + if loop_checkpoint_record_matches_request(&record, &request) { + Ok(Some(record)) + } else { + Ok(None) + } } } @@ -919,65 +918,6 @@ impl TurnRunTransitionPort for PostgresTurnStateStore { } } -#[cfg(feature = "postgres")] -#[async_trait] -impl LoopCheckpointStore for PostgresTurnStateStore { - async fn put_loop_checkpoint( - &self, - record: TurnCheckpointRecord, - ) -> Result<(), TurnError> { - let client = self.client().await?; - let payload = to_json(&record)?; - let checkpoint_scope_key = record.scope.as_ref().map(scope_key).transpose()?.unwrap_or_default(); - client - .execute( - "INSERT INTO turn_checkpoints (checkpoint_id, run_id, sequence, scope_key, kind, payload) - VALUES ($1, $2, $3, $4, $5, $6::jsonb) - ON CONFLICT (checkpoint_id) DO UPDATE SET - run_id = EXCLUDED.run_id, - sequence = EXCLUDED.sequence, - scope_key = EXCLUDED.scope_key, - kind = EXCLUDED.kind, - payload = EXCLUDED.payload", - &[ - &record.checkpoint_id.as_uuid().to_string(), - &record.run_id.to_string(), - &(record.sequence as i64), - &checkpoint_scope_key, - &record.kind.as_str(), - &payload, - ], - ) - .await - .map_err(db_error)?; - Ok(()) - } - - async fn get_loop_checkpoint( - &self, - checkpoint_id: TurnCheckpointId, - run_id: TurnRunId, - ) -> Result, TurnError> { - let client = self.client().await?; - let row = client - .query_opt( - "SELECT payload::text FROM turn_checkpoints WHERE checkpoint_id = $1 AND run_id = $2", - &[ - &checkpoint_id.as_uuid().to_string(), - &run_id.to_string(), - ], - ) - .await - .map_err(db_error)?; - let Some(row) = row else { - return Ok(None); - }; - let payload: String = row.get(0); - let record: TurnCheckpointRecord = serde_json::from_str(&payload).map_err(db_error)?; - Ok(Some(record)) - } -} - #[cfg(feature = "libsql")] async fn libsql_load_payloads(conn: &libsql::Connection, sql: &str) -> Result, TurnError> where @@ -992,6 +932,47 @@ where Ok(payloads) } +#[cfg(feature = "libsql")] +async fn libsql_insert_loop_checkpoint_record( + conn: &libsql::Connection, + record: &LoopCheckpointRecord, +) -> Result<(), TurnError> { + let rows = conn + .execute( + "INSERT OR IGNORE INTO turn_loop_checkpoints (checkpoint_id, scope_key, turn_id, run_id, created_at, payload) VALUES (?1, ?2, ?3, ?4, ?5, ?6)", + libsql::params![ + record.checkpoint_id.as_uuid().to_string(), + scope_key(&record.scope)?, + record.turn_id.to_string(), + record.run_id.to_string(), + record.created_at.to_rfc3339(), + to_json(record)?, + ], + ) + .await + .map_err(db_error)?; + if rows == 1 { + return Ok(()); + } + + let mut rows = conn + .query( + "SELECT payload FROM turn_loop_checkpoints WHERE checkpoint_id = ?1", + libsql::params![record.checkpoint_id.as_uuid().to_string()], + ) + .await + .map_err(db_error)?; + let Some(row) = rows.next().await.map_err(db_error)? else { + return Err(TurnError::Conflict { + reason: "loop checkpoint id insert conflicted but existing row was not readable" + .to_string(), + }); + }; + let payload: String = row.get(0).map_err(db_error)?; + let existing: LoopCheckpointRecord = serde_json::from_str(&payload).map_err(db_error)?; + ensure_loop_checkpoint_insert_is_idempotent(&existing, record) +} + #[cfg(feature = "libsql")] async fn libsql_load_event_retention_floor( conn: &libsql::Connection, @@ -1256,6 +1237,50 @@ where .collect() } +#[cfg(feature = "postgres")] +async fn postgres_insert_loop_checkpoint_record( + client: &impl deadpool_postgres::GenericClient, + record: &LoopCheckpointRecord, +) -> Result<(), TurnError> { + let payload = to_json(record)?; + let rows = client + .execute( + "INSERT INTO turn_loop_checkpoints (checkpoint_id, scope_key, turn_id, run_id, created_at, payload) + VALUES ($1, $2, $3, $4, $5::timestamptz, $6::jsonb) + ON CONFLICT (checkpoint_id) DO NOTHING", + &[ + &record.checkpoint_id.as_uuid().to_string(), + &scope_key(&record.scope)?, + &record.turn_id.to_string(), + &record.run_id.to_string(), + &record.created_at.to_rfc3339(), + &payload, + ], + ) + .await + .map_err(db_error)?; + if rows == 1 { + return Ok(()); + } + + let row = client + .query_opt( + "SELECT payload::text FROM turn_loop_checkpoints WHERE checkpoint_id = $1", + &[&record.checkpoint_id.as_uuid().to_string()], + ) + .await + .map_err(db_error)?; + let Some(row) = row else { + return Err(TurnError::Conflict { + reason: "loop checkpoint id insert conflicted but existing row was not readable" + .to_string(), + }); + }; + let payload: String = row.get(0); + let existing: LoopCheckpointRecord = serde_json::from_str(&payload).map_err(db_error)?; + ensure_loop_checkpoint_insert_is_idempotent(&existing, record) +} + #[cfg(feature = "postgres")] async fn postgres_load_event_retention_floor( client: &impl deadpool_postgres::GenericClient, @@ -1398,7 +1423,12 @@ async fn postgres_replace_snapshot( } for record in &snapshot.checkpoints { let payload = to_json(record)?; - let checkpoint_scope_key = record.scope.as_ref().map(scope_key).transpose()?.unwrap_or_default(); + let checkpoint_scope_key = record + .scope + .as_ref() + .map(scope_key) + .transpose()? + .unwrap_or_default(); txn.execute( "INSERT INTO turn_checkpoints (checkpoint_id, run_id, sequence, scope_key, kind, payload) VALUES ($1, $2, $3, $4, $5, $6::jsonb)", &[ @@ -1537,9 +1567,147 @@ fn turn_event_kind_key(event: &TurnLifecycleEvent) -> Result to_json(&event.kind) } +#[cfg(any(feature = "libsql", feature = "postgres"))] +fn loop_checkpoint_record_matches_request( + record: &LoopCheckpointRecord, + request: &GetLoopCheckpointRequest, +) -> bool { + record.scope == request.scope + && record.turn_id == request.turn_id + && record.run_id == request.run_id + && record.checkpoint_id == request.checkpoint_id +} + +#[cfg(any(feature = "libsql", feature = "postgres"))] +fn ensure_loop_checkpoint_insert_is_idempotent( + existing: &LoopCheckpointRecord, + incoming: &LoopCheckpointRecord, +) -> Result<(), TurnError> { + if existing == incoming { + Ok(()) + } else { + Err(TurnError::Conflict { + reason: "loop checkpoint id already belongs to a different checkpoint mapping" + .to_string(), + }) + } +} + fn db_error(error: impl std::fmt::Display) -> TurnError { tracing::debug!(%error, "turn state persistence operation failed"); TurnError::Unavailable { reason: "turn state persistence temporarily unavailable".to_string(), } } + +#[cfg(test)] +mod tests { + use super::*; + use ironclaw_host_api::{AgentId, ProjectId, TenantId, ThreadId}; + + fn test_scope(thread: &str) -> TurnScope { + TurnScope::new( + TenantId::new("tenant-db-checkpoint").unwrap(), + Some(AgentId::new("agent-db-checkpoint").unwrap()), + Some(ProjectId::new("project-db-checkpoint").unwrap()), + ThreadId::new(thread).unwrap(), + ) + } + + fn loop_checkpoint_record(thread: &str) -> LoopCheckpointRecord { + LoopCheckpointRecord { + checkpoint_id: crate::TurnCheckpointId::new(), + scope: test_scope(thread), + turn_id: crate::TurnId::new(), + run_id: crate::TurnRunId::new(), + state_ref: crate::LoopCheckpointStateRef::new("checkpoint:db-conflict").unwrap(), + schema_id: crate::CheckpointSchemaId::new("interactive_checkpoint_v1").unwrap(), + schema_version: crate::RunProfileVersion::new(1), + kind: crate::LoopCheckpointKind::BeforeBlock, + created_at: Utc::now(), + } + } + + #[cfg(feature = "libsql")] + #[tokio::test] + async fn libsql_loop_checkpoint_insert_conflicts_on_same_id_different_scope_or_run() { + let dir = tempfile::tempdir().unwrap(); + let db_path = dir.path().join("turns.db"); + let db = Arc::new(libsql::Builder::new_local(db_path).build().await.unwrap()); + let store = LibSqlTurnStateStore::new(Arc::clone(&db)); + store.run_migrations().await.unwrap(); + let conn = db.connect().unwrap(); + + let record = loop_checkpoint_record("libsql-conflict-a"); + libsql_insert_loop_checkpoint_record(&conn, &record) + .await + .unwrap(); + libsql_insert_loop_checkpoint_record(&conn, &record) + .await + .unwrap(); + + let mut conflicting = record.clone(); + conflicting.scope = test_scope("libsql-conflict-b"); + conflicting.run_id = crate::TurnRunId::new(); + let error = libsql_insert_loop_checkpoint_record(&conn, &conflicting) + .await + .unwrap_err(); + assert!(matches!(error, TurnError::Conflict { .. })); + } + + #[cfg(feature = "postgres")] + #[tokio::test] + async fn postgres_loop_checkpoint_insert_conflicts_on_same_id_different_scope_or_run() { + let Some(pool) = postgres_pool().await else { + return; + }; + let store = PostgresTurnStateStore::new(pool.clone()); + store.run_migrations().await.unwrap(); + let client = pool.get().await.unwrap(); + + let record = loop_checkpoint_record("postgres-conflict-a"); + postgres_insert_loop_checkpoint_record(&client, &record) + .await + .unwrap(); + postgres_insert_loop_checkpoint_record(&client, &record) + .await + .unwrap(); + + let mut conflicting = record.clone(); + conflicting.scope = test_scope("postgres-conflict-b"); + conflicting.run_id = crate::TurnRunId::new(); + let error = postgres_insert_loop_checkpoint_record(&client, &conflicting) + .await + .unwrap_err(); + assert!(matches!(error, TurnError::Conflict { .. })); + } + + #[cfg(feature = "postgres")] + async fn postgres_pool() -> Option { + let Ok(url) = std::env::var("IRONCLAW_TURNS_POSTGRES_URL") else { + eprintln!( + "skipping postgres loop checkpoint conflict test: IRONCLAW_TURNS_POSTGRES_URL not set" + ); + return None; + }; + let config: tokio_postgres::Config = match url.parse() { + Ok(config) => config, + Err(error) => { + eprintln!("skipping postgres loop checkpoint conflict test: invalid url ({error})"); + return None; + } + }; + let manager = deadpool_postgres::Manager::new(config, tokio_postgres::NoTls); + let pool = deadpool_postgres::Pool::builder(manager) + .max_size(4) + .build() + .unwrap(); + if let Err(error) = pool.get().await { + eprintln!( + "skipping postgres loop checkpoint conflict test: database unavailable ({error})" + ); + return None; + } + Some(pool) + } +} diff --git a/crates/ironclaw_turns/src/lib.rs b/crates/ironclaw_turns/src/lib.rs index b2d3ef4b9d0..c1e671772b7 100644 --- a/crates/ironclaw_turns/src/lib.rs +++ b/crates/ironclaw_turns/src/lib.rs @@ -81,8 +81,8 @@ pub use status::{ SanitizedFailure, TurnError, TurnErrorCategory, TurnRunProfile, TurnRunState, TurnStatus, }; pub use store::{ - LoopCheckpointStore, TurnActiveLockKey, TurnActiveLockRecord, TurnCheckpointRecord, - TurnIdempotencyErrorReplay, TurnIdempotencyOperationKind, TurnIdempotencyOutcomeKind, - TurnIdempotencyRecord, TurnIdempotencyReplay, TurnLockVersion, TurnPersistenceSnapshot, - TurnRecord, TurnRunRecord, TurnStateStore, + TurnActiveLockKey, TurnActiveLockRecord, TurnCheckpointRecord, TurnIdempotencyErrorReplay, + TurnIdempotencyOperationKind, TurnIdempotencyOutcomeKind, TurnIdempotencyRecord, + TurnIdempotencyReplay, TurnLockVersion, TurnPersistenceSnapshot, TurnRecord, TurnRunRecord, + TurnStateStore, }; diff --git a/crates/ironclaw_turns/src/memory.rs b/crates/ironclaw_turns/src/memory.rs index 802892a32c1..a98564c00ff 100644 --- a/crates/ironclaw_turns/src/memory.rs +++ b/crates/ironclaw_turns/src/memory.rs @@ -725,40 +725,6 @@ impl TurnRunTransitionPort for InMemoryTurnStateStore { } } -#[async_trait] -impl crate::LoopCheckpointStore for InMemoryTurnStateStore { - async fn put_loop_checkpoint( - &self, - record: TurnCheckpointRecord, - ) -> Result<(), crate::TurnError> { - let mut inner = self.lock_inner()?; - // Idempotent: update if same checkpoint_id exists, else push. - if let Some(existing) = inner - .checkpoints - .iter_mut() - .find(|c| c.checkpoint_id == record.checkpoint_id) - { - *existing = record; - } else { - inner.checkpoints.push(record); - } - Ok(()) - } - - async fn get_loop_checkpoint( - &self, - checkpoint_id: crate::TurnCheckpointId, - run_id: crate::TurnRunId, - ) -> Result, crate::TurnError> { - let inner = self.lock_inner()?; - Ok(inner - .checkpoints - .iter() - .find(|c| c.checkpoint_id == checkpoint_id && c.run_id == run_id) - .cloned()) - } -} - impl Inner { fn from_persistence_snapshot( snapshot: TurnPersistenceSnapshot, diff --git a/crates/ironclaw_turns/src/store.rs b/crates/ironclaw_turns/src/store.rs index 32ed0bcfc81..93bc5d85ecc 100644 --- a/crates/ironclaw_turns/src/store.rs +++ b/crates/ironclaw_turns/src/store.rs @@ -142,34 +142,6 @@ pub struct TurnCheckpointRecord { pub created_at: TurnTimestamp, } -/// Direct single-row checkpoint read/write operations. -/// -/// Unlike the snapshot-based path (`TurnRunTransitionPort::block_run` → -/// `libsql_replace_snapshot`/`postgres_replace_snapshot`) which loads the full -/// `TurnPersistenceSnapshot`, `LoopCheckpointStore` provides targeted -/// `INSERT`/`SELECT` operations scoped by checkpoint ID and run. -/// -/// Implementations must preserve idempotent `put` semantics (re-inserting the -/// same `checkpoint_id` is not an error) and cross-run isolation (`get` returns -/// `None` when the checkpoint belongs to a different `run_id`). -#[async_trait] -pub trait LoopCheckpointStore: Send + Sync { - /// Insert a single checkpoint record scoped to the run. - /// Idempotent: re-inserting the same `checkpoint_id` is not an error. - async fn put_loop_checkpoint( - &self, - record: TurnCheckpointRecord, - ) -> Result<(), TurnError>; - - /// Direct lookup by `checkpoint_id`. Returns `None` if not found or if the - /// checkpoint belongs to a different `run_id` (cross-run rejection). - async fn get_loop_checkpoint( - &self, - checkpoint_id: TurnCheckpointId, - run_id: TurnRunId, - ) -> Result, TurnError>; -} - #[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)] #[serde(rename_all = "snake_case")] pub enum TurnIdempotencyOperationKind { diff --git a/crates/ironclaw_turns/tests/checkpoint_state_store_contract.rs b/crates/ironclaw_turns/tests/checkpoint_state_store_contract.rs index bfa8dd6a5e7..929eea0c4a1 100644 --- a/crates/ironclaw_turns/tests/checkpoint_state_store_contract.rs +++ b/crates/ironclaw_turns/tests/checkpoint_state_store_contract.rs @@ -435,7 +435,7 @@ fn turn_checkpoint_public_status_does_not_expose_checkpoint_payload() { }; let event = TurnLifecycleEvent { cursor: EventCursor(2), - scope, + scope: scope.clone(), run_id, status: TurnStatus::BlockedApproval, kind: TurnEventKind::Blocked, @@ -445,9 +445,12 @@ fn turn_checkpoint_public_status_does_not_expose_checkpoint_payload() { checkpoints: vec![TurnCheckpointRecord { checkpoint_id, run_id, + scope: Some(scope.clone()), sequence: 1, status: TurnStatus::BlockedApproval, gate_ref: GateRef::new("gate-checkpoint-public").unwrap(), + kind: LoopCheckpointKind::BeforeBlock, + state_ref: LoopCheckpointStateRef::new("checkpoint:public-status").unwrap(), created_at: fixed_time(), }], events: vec![event.clone()], diff --git a/crates/ironclaw_turns/tests/loop_checkpoint_store_contract.rs b/crates/ironclaw_turns/tests/loop_checkpoint_store_contract.rs index 7378235273b..54d02748a99 100644 --- a/crates/ironclaw_turns/tests/loop_checkpoint_store_contract.rs +++ b/crates/ironclaw_turns/tests/loop_checkpoint_store_contract.rs @@ -1,16 +1,17 @@ -#![cfg(any(feature = "libsql", feature = "postgres"))] +#[cfg(any(feature = "libsql", feature = "postgres"))] +use std::sync::Arc; -use chrono::{TimeZone, Utc}; use ironclaw_host_api::{AgentId, ProjectId, TenantId, ThreadId}; use ironclaw_turns::{ - GateRef, LoopCheckpointStore, TurnCheckpointId, TurnCheckpointRecord, TurnRunId, TurnScope, - TurnStatus, - run_profile::{LoopCheckpointKind, LoopCheckpointStateRef}, + CheckpointSchemaId, GetLoopCheckpointRequest, InMemoryLoopCheckpointStore, + InMemoryTurnStateStore, LoopCheckpointStateRef, LoopCheckpointStore, PutLoopCheckpointRequest, + RunProfileVersion, TurnId, TurnRunId, TurnScope, run_profile::LoopCheckpointKind, }; -use std::sync::Arc; #[cfg(feature = "libsql")] use ironclaw_turns::LibSqlTurnStateStore; +#[cfg(feature = "postgres")] +use ironclaw_turns::PostgresTurnStateStore; fn test_scope(thread: &str) -> TurnScope { TurnScope::new( @@ -21,119 +22,106 @@ fn test_scope(thread: &str) -> TurnScope { ) } -fn make_checkpoint( - run_id: TurnRunId, - sequence: u64, - kind: LoopCheckpointKind, -) -> TurnCheckpointRecord { - TurnCheckpointRecord { - checkpoint_id: TurnCheckpointId::new(), +fn put_request(scope: TurnScope, turn_id: TurnId, run_id: TurnRunId) -> PutLoopCheckpointRequest { + PutLoopCheckpointRequest { + scope, + turn_id, run_id, - scope: Some(test_scope("thread-a")), - sequence, - status: TurnStatus::BlockedApproval, - gate_ref: GateRef::new("gate:test-gate").unwrap(), - kind, state_ref: LoopCheckpointStateRef::new("checkpoint:test-state").unwrap(), - created_at: Utc.with_ymd_and_hms(2026, 5, 11, 12, 0, 0).unwrap(), + schema_id: CheckpointSchemaId::new("interactive_checkpoint_v1").unwrap(), + schema_version: RunProfileVersion::new(1), + kind: LoopCheckpointKind::BeforeModel, } } -// ── Tests against in-memory backend ────────────────────────────────────────── - -#[tokio::test] -async fn inmemory_put_get_roundtrip() { - let store = ironclaw_turns::InMemoryTurnStateStore::default(); +async fn assert_loop_checkpoint_store_roundtrip(store: &(impl LoopCheckpointStore + ?Sized)) { + let scope = test_scope("thread-loop-checkpoint-roundtrip"); + let turn_id = TurnId::new(); let run_id = TurnRunId::new(); - let record = make_checkpoint(run_id, 1, LoopCheckpointKind::BeforeModel); - let checkpoint_id = record.checkpoint_id; - - store.put_loop_checkpoint(record.clone()).await.unwrap(); - let fetched = store - .get_loop_checkpoint(checkpoint_id, run_id) + let checkpoint = store + .put_loop_checkpoint(put_request(scope.clone(), turn_id, run_id)) .await .unwrap(); - let fetched = fetched.expect("should find checkpoint"); - assert_eq!(fetched.checkpoint_id, checkpoint_id); - assert_eq!(fetched.run_id, run_id); - assert_eq!(fetched.kind, LoopCheckpointKind::BeforeModel); - assert_eq!( - fetched.state_ref, - LoopCheckpointStateRef::new("checkpoint:test-state").unwrap() - ); - assert_eq!(fetched.scope, Some(test_scope("thread-a"))); -} - -#[tokio::test] -async fn inmemory_cross_run_rejection() { - let store = ironclaw_turns::InMemoryTurnStateStore::default(); - let run_a = TurnRunId::new(); - let run_b = TurnRunId::new(); - let record = make_checkpoint(run_a, 1, LoopCheckpointKind::BeforeBlock); - let checkpoint_id = record.checkpoint_id; - - store.put_loop_checkpoint(record).await.unwrap(); - let fetched = store - .get_loop_checkpoint(checkpoint_id, run_b) + let loaded = store + .get_loop_checkpoint(GetLoopCheckpointRequest { + scope: scope.clone(), + turn_id, + run_id, + checkpoint_id: checkpoint.checkpoint_id, + }) .await - .unwrap(); - assert!(fetched.is_none(), "cross-run lookup should return None"); + .unwrap() + .expect("checkpoint id should resolve to state ref"); + + assert_eq!(loaded, checkpoint); + assert_eq!(loaded.scope, scope); + assert_eq!(loaded.turn_id, turn_id); + assert_eq!(loaded.run_id, run_id); + assert_eq!(loaded.kind, LoopCheckpointKind::BeforeModel); } -#[tokio::test] -async fn inmemory_multiple_checkpoints_per_run() { - let store = ironclaw_turns::InMemoryTurnStateStore::default(); +async fn assert_loop_checkpoint_store_cross_scope_and_run_miss( + store: &(impl LoopCheckpointStore + ?Sized), +) { + let scope = test_scope("thread-loop-checkpoint-scope-a"); + let turn_id = TurnId::new(); let run_id = TurnRunId::new(); + let checkpoint = store + .put_loop_checkpoint(put_request(scope.clone(), turn_id, run_id)) + .await + .unwrap(); - let r1 = make_checkpoint(run_id, 1, LoopCheckpointKind::BeforeModel); - let r2 = make_checkpoint(run_id, 2, LoopCheckpointKind::BeforeSideEffect); - let r3 = make_checkpoint(run_id, 3, LoopCheckpointKind::Final); - let id1 = r1.checkpoint_id; - let id2 = r2.checkpoint_id; - let id3 = r3.checkpoint_id; - - store.put_loop_checkpoint(r1).await.unwrap(); - store.put_loop_checkpoint(r2).await.unwrap(); - store.put_loop_checkpoint(r3).await.unwrap(); - - assert!(store.get_loop_checkpoint(id1, run_id).await.unwrap().is_some()); - assert!(store.get_loop_checkpoint(id2, run_id).await.unwrap().is_some()); - assert!(store.get_loop_checkpoint(id3, run_id).await.unwrap().is_some()); + let cross_scope = store + .get_loop_checkpoint(GetLoopCheckpointRequest { + scope: test_scope("thread-loop-checkpoint-scope-b"), + turn_id, + run_id, + checkpoint_id: checkpoint.checkpoint_id, + }) + .await + .unwrap(); + assert!(cross_scope.is_none(), "cross-scope lookup must fail closed"); + + let cross_run = store + .get_loop_checkpoint(GetLoopCheckpointRequest { + scope, + turn_id, + run_id: TurnRunId::new(), + checkpoint_id: checkpoint.checkpoint_id, + }) + .await + .unwrap(); + assert!(cross_run.is_none(), "cross-run lookup must fail closed"); } #[tokio::test] -async fn inmemory_idempotent_put() { - let store = ironclaw_turns::InMemoryTurnStateStore::default(); - let run_id = TurnRunId::new(); - let record = make_checkpoint(run_id, 1, LoopCheckpointKind::BeforeBlock); - - store.put_loop_checkpoint(record.clone()).await.unwrap(); - // Second put with same checkpoint_id should succeed without error. - store.put_loop_checkpoint(record).await.unwrap(); +async fn inmemory_standalone_loop_checkpoint_roundtrip() { + let store = InMemoryLoopCheckpointStore::default(); + assert_loop_checkpoint_store_roundtrip(&store).await; + assert_loop_checkpoint_store_cross_scope_and_run_miss(&store).await; } #[tokio::test] -async fn serde_backward_compat_missing_fields() { - // Simulate old persisted JSON that lacks kind, state_ref, and scope. - let json = r#"{ - "checkpoint_id": "00000000-0000-0000-0000-000000000001", - "run_id": "00000000-0000-0000-0000-000000000002", - "sequence": 1, - "status": "BlockedApproval", - "gate_ref": "gate:test", - "created_at": "2026-05-11T12:00:00Z" - }"#; - let record: TurnCheckpointRecord = serde_json::from_str(json).unwrap(); - assert_eq!(record.kind, LoopCheckpointKind::BeforeBlock); - assert_eq!( - record.state_ref, - LoopCheckpointStateRef::new("checkpoint:unknown").unwrap() +async fn inmemory_turn_state_loop_checkpoint_roundtrip_and_snapshot() { + let store = InMemoryTurnStateStore::default(); + assert_loop_checkpoint_store_roundtrip(&store).await; + assert_loop_checkpoint_store_cross_scope_and_run_miss(&store).await; + + let snapshot = store.persistence_snapshot(); + assert_eq!(snapshot.loop_checkpoints.len(), 2); + assert!( + snapshot.checkpoints.is_empty(), + "loop checkpoint mappings must not use turn_checkpoints" ); - assert_eq!(record.scope, None); -} -// ── Tests against libSQL backend ───────────────────────────────────────────── + let reopened = InMemoryTurnStateStore::from_persistence_snapshot( + snapshot, + ironclaw_turns::InMemoryTurnStateStoreLimits::default(), + ) + .unwrap(); + assert_loop_checkpoint_store_cross_scope_and_run_miss(&reopened).await; +} #[cfg(feature = "libsql")] async fn libsql_store() -> (Arc, tempfile::TempDir) { @@ -147,75 +135,70 @@ async fn libsql_store() -> (Arc, tempfile::TempDir) { #[cfg(feature = "libsql")] #[tokio::test] -async fn libsql_put_get_roundtrip() { +async fn libsql_loop_checkpoint_roundtrip_uses_loop_mapping_table() { let (store, _dir) = libsql_store().await; - let run_id = TurnRunId::new(); - let record = make_checkpoint(run_id, 1, LoopCheckpointKind::BeforeModel); - let checkpoint_id = record.checkpoint_id; - - store.put_loop_checkpoint(record.clone()).await.unwrap(); - let fetched = store - .get_loop_checkpoint(checkpoint_id, run_id) - .await - .unwrap() - .expect("should find checkpoint"); - - assert_eq!(fetched.checkpoint_id, checkpoint_id); - assert_eq!(fetched.run_id, run_id); - assert_eq!(fetched.kind, LoopCheckpointKind::BeforeModel); - assert_eq!( - fetched.state_ref, - LoopCheckpointStateRef::new("checkpoint:test-state").unwrap() + assert_loop_checkpoint_store_roundtrip(store.as_ref()).await; + assert_loop_checkpoint_store_cross_scope_and_run_miss(store.as_ref()).await; + + let snapshot = store.persistence_snapshot().await.unwrap(); + assert_eq!(snapshot.loop_checkpoints.len(), 2); + assert!( + snapshot.checkpoints.is_empty(), + "libSQL loop mappings must not be written to turn_checkpoints" ); - assert_eq!(fetched.scope, Some(test_scope("thread-a"))); } -#[cfg(feature = "libsql")] +#[cfg(feature = "postgres")] #[tokio::test] -async fn libsql_cross_run_rejection() { - let (store, _dir) = libsql_store().await; - let run_a = TurnRunId::new(); - let run_b = TurnRunId::new(); - let record = make_checkpoint(run_a, 1, LoopCheckpointKind::BeforeBlock); - let checkpoint_id = record.checkpoint_id; - - store.put_loop_checkpoint(record).await.unwrap(); - let fetched = store - .get_loop_checkpoint(checkpoint_id, run_b) - .await - .unwrap(); - assert!(fetched.is_none(), "cross-run lookup should return None"); -} - -#[cfg(feature = "libsql")] -#[tokio::test] -async fn libsql_multiple_checkpoints_per_run() { - let (store, _dir) = libsql_store().await; - let run_id = TurnRunId::new(); - - let r1 = make_checkpoint(run_id, 1, LoopCheckpointKind::BeforeModel); - let r2 = make_checkpoint(run_id, 2, LoopCheckpointKind::BeforeSideEffect); - let r3 = make_checkpoint(run_id, 3, LoopCheckpointKind::Final); - let id1 = r1.checkpoint_id; - let id2 = r2.checkpoint_id; - let id3 = r3.checkpoint_id; - - store.put_loop_checkpoint(r1).await.unwrap(); - store.put_loop_checkpoint(r2).await.unwrap(); - store.put_loop_checkpoint(r3).await.unwrap(); - - assert!(store.get_loop_checkpoint(id1, run_id).await.unwrap().is_some()); - assert!(store.get_loop_checkpoint(id2, run_id).await.unwrap().is_some()); - assert!(store.get_loop_checkpoint(id3, run_id).await.unwrap().is_some()); +async fn postgres_loop_checkpoint_roundtrip_uses_loop_mapping_table() { + let Some(pool) = postgres_pool().await else { + return; + }; + let store = Arc::new(PostgresTurnStateStore::new(pool)); + store.run_migrations().await.unwrap(); + assert_loop_checkpoint_store_roundtrip(store.as_ref()).await; + assert_loop_checkpoint_store_cross_scope_and_run_miss(store.as_ref()).await; + + let snapshot = store.persistence_snapshot().await.unwrap(); + assert!( + snapshot + .loop_checkpoints + .iter() + .any(|record| record.state_ref.as_str() == "checkpoint:test-state"), + "Postgres loop mappings must be written to turn_loop_checkpoints" + ); + assert!( + snapshot + .checkpoints + .iter() + .all(|record| record.state_ref.as_str() != "checkpoint:test-state"), + "Postgres loop mappings must not be written to turn_checkpoints" + ); } -#[cfg(feature = "libsql")] -#[tokio::test] -async fn libsql_idempotent_put() { - let (store, _dir) = libsql_store().await; - let run_id = TurnRunId::new(); - let record = make_checkpoint(run_id, 1, LoopCheckpointKind::BeforeBlock); - - store.put_loop_checkpoint(record.clone()).await.unwrap(); - store.put_loop_checkpoint(record).await.unwrap(); +#[cfg(feature = "postgres")] +async fn postgres_pool() -> Option { + let Ok(url) = std::env::var("IRONCLAW_TURNS_POSTGRES_URL") else { + eprintln!( + "skipping postgres loop checkpoint contract: IRONCLAW_TURNS_POSTGRES_URL not set" + ); + return None; + }; + let config: tokio_postgres::Config = match url.parse() { + Ok(config) => config, + Err(error) => { + eprintln!("skipping postgres loop checkpoint contract: invalid url ({error})"); + return None; + } + }; + let manager = deadpool_postgres::Manager::new(config, tokio_postgres::NoTls); + let pool = deadpool_postgres::Pool::builder(manager) + .max_size(4) + .build() + .unwrap(); + if let Err(error) = pool.get().await { + eprintln!("skipping postgres loop checkpoint contract: database unavailable ({error})"); + return None; + } + Some(pool) } From 245ef21cb1480dbdb100e39e8a3d54e95170f75e Mon Sep 17 00:00:00 2001 From: serrrfirat Date: Mon, 11 May 2026 15:49:16 +0200 Subject: [PATCH 7/8] fix(KB-002): annotate checkpoint sentinel unwraps - Add inline safety comments required by no-panics gate for static checkpoint sentinel refs --- crates/ironclaw_turns/src/memory.rs | 2 +- crates/ironclaw_turns/src/store.rs | 3 +-- 2 files changed, 2 insertions(+), 3 deletions(-) diff --git a/crates/ironclaw_turns/src/memory.rs b/crates/ironclaw_turns/src/memory.rs index a98564c00ff..06153e911ab 100644 --- a/crates/ironclaw_turns/src/memory.rs +++ b/crates/ironclaw_turns/src/memory.rs @@ -1730,7 +1730,7 @@ impl Inner { // Placeholder — callers in block_run don't have the loop's actual state_ref. // Real values will be threaded in a follow-up task. state_ref: crate::run_profile::LoopCheckpointStateRef::new("checkpoint:block-state") - .unwrap(), + .unwrap(), // safety: literal satisfies the "checkpoint:" prefix and length constraints. created_at, }); } diff --git a/crates/ironclaw_turns/src/store.rs b/crates/ironclaw_turns/src/store.rs index 93bc5d85ecc..510fa06ee3d 100644 --- a/crates/ironclaw_turns/src/store.rs +++ b/crates/ironclaw_turns/src/store.rs @@ -118,8 +118,7 @@ fn default_checkpoint_kind() -> LoopCheckpointKind { /// deserializing old persisted data that predates the `state_ref` field. /// Real values will be threaded from loop callers in a follow-up task. fn default_checkpoint_state_ref() -> LoopCheckpointStateRef { - // Safety: literal satisfies the "checkpoint:" prefix and length constraints. - LoopCheckpointStateRef::new("checkpoint:unknown").unwrap() + LoopCheckpointStateRef::new("checkpoint:unknown").unwrap() // safety: literal satisfies the "checkpoint:" prefix and length constraints. } #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] From a37f794dea329f3e0473aafdd2a58cb0be37f56d Mon Sep 17 00:00:00 2001 From: serrrfirat Date: Mon, 11 May 2026 22:00:02 +0200 Subject: [PATCH 8/8] =?UTF-8?q?fix(turns):=20address=20zmanian=20review=20?= =?UTF-8?q?=E2=80=94=20checkpoint=20durability=20(#3468)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- crates/ironclaw_turns/Cargo.toml | 2 +- crates/ironclaw_turns/src/db.rs | 247 ++++++++++++------ crates/ironclaw_turns/src/loop_exit.rs | 6 +- crates/ironclaw_turns/src/memory.rs | 19 +- crates/ironclaw_turns/src/run_profile/host.rs | 4 + crates/ironclaw_turns/src/runner.rs | 8 +- crates/ironclaw_turns/src/store.rs | 5 +- .../tests/agent_loop_host_contract.rs | 4 +- .../tests/loop_checkpoint_store_contract.rs | 10 +- .../tests/loop_exit_contract.rs | 19 +- .../tests/turn_coordinator_contract.rs | 35 ++- docs/reborn/contracts/loop-exit.md | 2 +- docs/reborn/contracts/turn-persistence.md | 2 + 13 files changed, 253 insertions(+), 110 deletions(-) diff --git a/crates/ironclaw_turns/Cargo.toml b/crates/ironclaw_turns/Cargo.toml index a3811caf387..ff759a777f8 100644 --- a/crates/ironclaw_turns/Cargo.toml +++ b/crates/ironclaw_turns/Cargo.toml @@ -13,7 +13,7 @@ publish = false [features] default = [] -postgres = ["dep:deadpool-postgres", "dep:tokio-postgres"] +postgres = ["dep:deadpool-postgres", "dep:tokio-postgres", "tokio-postgres/with-chrono-0_4"] libsql = ["dep:libsql"] [dependencies] diff --git a/crates/ironclaw_turns/src/db.rs b/crates/ironclaw_turns/src/db.rs index b0624ca223d..bd330441925 100644 --- a/crates/ironclaw_turns/src/db.rs +++ b/crates/ironclaw_turns/src/db.rs @@ -255,17 +255,22 @@ impl LibSqlTurnStateStore { .await .map_err(db_error)?; - // Migration: add new columns to existing turn_checkpoints tables. - // For libSQL, ALTER TABLE ADD COLUMN fails if the column already exists, - // so we ignore "duplicate column name" errors. - for alter in [ - "ALTER TABLE turn_checkpoints ADD COLUMN scope_key TEXT NOT NULL DEFAULT ''", - "ALTER TABLE turn_checkpoints ADD COLUMN kind TEXT NOT NULL DEFAULT 'before_block'", + // Migration: add metadata columns to existing turn_checkpoints tables. + // Legacy rows predate scoped checkpoint metadata, so they keep an empty + // scope_key sentinel and the serialized payload remains the migration + // source of truth until a scoped backfill can prove each row's owner. + for (column, alter) in [ + ( + "scope_key", + "ALTER TABLE turn_checkpoints ADD COLUMN scope_key TEXT NOT NULL DEFAULT ''", + ), + ( + "kind", + "ALTER TABLE turn_checkpoints ADD COLUMN kind TEXT NOT NULL DEFAULT 'before_block'", + ), ] { - match conn.execute(alter, ()).await { - Ok(_) => {} - Err(e) if e.to_string().contains("duplicate column name") => {} - Err(e) => return Err(db_error(e)), + if !libsql_column_exists(&conn, "turn_checkpoints", column).await? { + conn.execute(alter, ()).await.map_err(db_error)?; } } Ok(()) @@ -444,11 +449,8 @@ impl LoopCheckpointStore for LibSqlTurnStateStore { }; let payload: String = row.get(0).map_err(db_error)?; let record = serde_json::from_str(&payload).map_err(db_error)?; - if loop_checkpoint_record_matches_request(&record, &request) { - Ok(Some(record)) - } else { - Ok(None) - } + ensure_loop_checkpoint_record_matches_request(&record, &request)?; + Ok(Some(record)) } } @@ -750,7 +752,8 @@ impl LoopCheckpointStore for PostgresTurnStateStore { &self, request: PutLoopCheckpointRequest, ) -> Result { - let client = self.client().await?; + let mut client = self.client().await?; + let txn = client.transaction().await.map_err(db_error)?; let record = LoopCheckpointRecord { checkpoint_id: crate::TurnCheckpointId::new(), scope: request.scope, @@ -762,8 +765,19 @@ impl LoopCheckpointStore for PostgresTurnStateStore { kind: request.kind, created_at: Utc::now(), }; - postgres_insert_loop_checkpoint_record(&client, &record).await?; - Ok(record) + let result = postgres_insert_loop_checkpoint_record(&txn, &record) + .await + .map(|()| record.clone()); + match result { + Ok(record) => { + txn.commit().await.map_err(db_error)?; + Ok(record) + } + Err(error) => { + let _ = txn.rollback().await; + Err(error) + } + } } async fn get_loop_checkpoint( @@ -789,11 +803,8 @@ impl LoopCheckpointStore for PostgresTurnStateStore { }; let payload: String = row.get(0); let record = serde_json::from_str(&payload).map_err(db_error)?; - if loop_checkpoint_record_matches_request(&record, &request) { - Ok(Some(record)) - } else { - Ok(None) - } + ensure_loop_checkpoint_record_matches_request(&record, &request)?; + Ok(Some(record)) } } @@ -918,6 +929,31 @@ impl TurnRunTransitionPort for PostgresTurnStateStore { } } +#[cfg(feature = "libsql")] +async fn libsql_column_exists( + conn: &libsql::Connection, + table: &str, + column: &str, +) -> Result { + if !table + .chars() + .all(|character| character.is_ascii_alphanumeric() || character == '_') + { + return Err(TurnError::Unavailable { + reason: "turn state persistence temporarily unavailable".to_string(), + }); + } + let sql = format!("PRAGMA table_info({table})"); + let mut rows = conn.query(sql.as_str(), ()).await.map_err(db_error)?; + while let Some(row) = rows.next().await.map_err(db_error)? { + let existing: String = row.get(1).map_err(db_error)?; + if existing == column { + return Ok(true); + } + } + Ok(false) +} + #[cfg(feature = "libsql")] async fn libsql_load_payloads(conn: &libsql::Connection, sql: &str) -> Result, TurnError> where @@ -937,16 +973,21 @@ async fn libsql_insert_loop_checkpoint_record( conn: &libsql::Connection, record: &LoopCheckpointRecord, ) -> Result<(), TurnError> { + let payload = to_json(record)?; let rows = conn .execute( - "INSERT OR IGNORE INTO turn_loop_checkpoints (checkpoint_id, scope_key, turn_id, run_id, created_at, payload) VALUES (?1, ?2, ?3, ?4, ?5, ?6)", + "INSERT INTO turn_loop_checkpoints (checkpoint_id, scope_key, turn_id, run_id, created_at, payload) + VALUES (?1, ?2, ?3, ?4, ?5, ?6) + ON CONFLICT(checkpoint_id) DO UPDATE + SET checkpoint_id = turn_loop_checkpoints.checkpoint_id + WHERE turn_loop_checkpoints.payload = excluded.payload", libsql::params![ record.checkpoint_id.as_uuid().to_string(), scope_key(&record.scope)?, record.turn_id.to_string(), record.run_id.to_string(), record.created_at.to_rfc3339(), - to_json(record)?, + payload, ], ) .await @@ -955,22 +996,9 @@ async fn libsql_insert_loop_checkpoint_record( return Ok(()); } - let mut rows = conn - .query( - "SELECT payload FROM turn_loop_checkpoints WHERE checkpoint_id = ?1", - libsql::params![record.checkpoint_id.as_uuid().to_string()], - ) - .await - .map_err(db_error)?; - let Some(row) = rows.next().await.map_err(db_error)? else { - return Err(TurnError::Conflict { - reason: "loop checkpoint id insert conflicted but existing row was not readable" - .to_string(), - }); - }; - let payload: String = row.get(0).map_err(db_error)?; - let existing: LoopCheckpointRecord = serde_json::from_str(&payload).map_err(db_error)?; - ensure_loop_checkpoint_insert_is_idempotent(&existing, record) + Err(TurnError::Conflict { + reason: "loop checkpoint id already belongs to a different checkpoint mapping".to_string(), + }) } #[cfg(feature = "libsql")] @@ -1244,41 +1272,31 @@ async fn postgres_insert_loop_checkpoint_record( ) -> Result<(), TurnError> { let payload = to_json(record)?; let rows = client - .execute( + .query( "INSERT INTO turn_loop_checkpoints (checkpoint_id, scope_key, turn_id, run_id, created_at, payload) - VALUES ($1, $2, $3, $4, $5::timestamptz, $6::jsonb) - ON CONFLICT (checkpoint_id) DO NOTHING", + VALUES ($1, $2, $3, $4, $5, $6::jsonb) + ON CONFLICT (checkpoint_id) DO UPDATE + SET checkpoint_id = turn_loop_checkpoints.checkpoint_id + WHERE turn_loop_checkpoints.payload = EXCLUDED.payload + RETURNING payload::text", &[ &record.checkpoint_id.as_uuid().to_string(), &scope_key(&record.scope)?, &record.turn_id.to_string(), &record.run_id.to_string(), - &record.created_at.to_rfc3339(), + &record.created_at, &payload, ], ) .await .map_err(db_error)?; - if rows == 1 { + if rows.len() == 1 { return Ok(()); } - let row = client - .query_opt( - "SELECT payload::text FROM turn_loop_checkpoints WHERE checkpoint_id = $1", - &[&record.checkpoint_id.as_uuid().to_string()], - ) - .await - .map_err(db_error)?; - let Some(row) = row else { - return Err(TurnError::Conflict { - reason: "loop checkpoint id insert conflicted but existing row was not readable" - .to_string(), - }); - }; - let payload: String = row.get(0); - let existing: LoopCheckpointRecord = serde_json::from_str(&payload).map_err(db_error)?; - ensure_loop_checkpoint_insert_is_idempotent(&existing, record) + Err(TurnError::Conflict { + reason: "loop checkpoint id already belongs to a different checkpoint mapping".to_string(), + }) } #[cfg(feature = "postgres")] @@ -1446,13 +1464,13 @@ async fn postgres_replace_snapshot( for record in &snapshot.loop_checkpoints { let payload = to_json(record)?; txn.execute( - "INSERT INTO turn_loop_checkpoints (checkpoint_id, scope_key, turn_id, run_id, created_at, payload) VALUES ($1, $2, $3, $4, $5::timestamptz, $6::jsonb)", + "INSERT INTO turn_loop_checkpoints (checkpoint_id, scope_key, turn_id, run_id, created_at, payload) VALUES ($1, $2, $3, $4, $5, $6::jsonb)", &[ &record.checkpoint_id.as_uuid().to_string(), &scope_key(&record.scope)?, &record.turn_id.to_string(), &record.run_id.to_string(), - &record.created_at.to_rfc3339(), + &record.created_at, &payload, ], ) @@ -1462,14 +1480,14 @@ async fn postgres_replace_snapshot( for record in &snapshot.idempotency_records { let payload = to_json(record)?; txn.execute( - "INSERT INTO turn_idempotency_records (record_key, scope_key, operation, run_id, idempotency_key, created_at, payload) VALUES ($1, $2, $3, $4, $5, $6::timestamptz, $7::jsonb)", + "INSERT INTO turn_idempotency_records (record_key, scope_key, operation, run_id, idempotency_key, created_at, payload) VALUES ($1, $2, $3, $4, $5, $6, $7::jsonb)", &[ &idempotency_record_key(record)?, &scope_key(&record.scope)?, &operation_key(record)?, &record.run_id.map(|run_id| run_id.to_string()), &record.key.as_str(), - &record.created_at.to_rfc3339(), + &record.created_at, &payload, ], ) @@ -1568,27 +1586,19 @@ fn turn_event_kind_key(event: &TurnLifecycleEvent) -> Result } #[cfg(any(feature = "libsql", feature = "postgres"))] -fn loop_checkpoint_record_matches_request( +fn ensure_loop_checkpoint_record_matches_request( record: &LoopCheckpointRecord, request: &GetLoopCheckpointRequest, -) -> bool { - record.scope == request.scope +) -> Result<(), TurnError> { + if record.scope == request.scope && record.turn_id == request.turn_id && record.run_id == request.run_id && record.checkpoint_id == request.checkpoint_id -} - -#[cfg(any(feature = "libsql", feature = "postgres"))] -fn ensure_loop_checkpoint_insert_is_idempotent( - existing: &LoopCheckpointRecord, - incoming: &LoopCheckpointRecord, -) -> Result<(), TurnError> { - if existing == incoming { + { Ok(()) } else { Err(TurnError::Conflict { - reason: "loop checkpoint id already belongs to a different checkpoint mapping" - .to_string(), + reason: "loop checkpoint row metadata conflicts with persisted payload".to_string(), }) } } @@ -1655,6 +1665,45 @@ mod tests { assert!(matches!(error, TurnError::Conflict { .. })); } + #[cfg(feature = "libsql")] + #[tokio::test] + async fn libsql_loop_checkpoint_get_errors_when_payload_identity_drifts() { + let dir = tempfile::tempdir().unwrap(); + let db_path = dir.path().join("turns.db"); + let db = Arc::new(libsql::Builder::new_local(db_path).build().await.unwrap()); + let store = LibSqlTurnStateStore::new(Arc::clone(&db)); + store.run_migrations().await.unwrap(); + let conn = db.connect().unwrap(); + + let record = loop_checkpoint_record("libsql-drift"); + libsql_insert_loop_checkpoint_record(&conn, &record) + .await + .unwrap(); + + let mut drifted = record.clone(); + drifted.run_id = crate::TurnRunId::new(); + conn.execute( + "UPDATE turn_loop_checkpoints SET payload = ?1 WHERE checkpoint_id = ?2", + libsql::params![ + to_json(&drifted).unwrap(), + record.checkpoint_id.as_uuid().to_string(), + ], + ) + .await + .unwrap(); + + let error = store + .get_loop_checkpoint(GetLoopCheckpointRequest { + scope: record.scope.clone(), + turn_id: record.turn_id, + run_id: record.run_id, + checkpoint_id: record.checkpoint_id, + }) + .await + .unwrap_err(); + assert!(matches!(error, TurnError::Conflict { .. })); + } + #[cfg(feature = "postgres")] #[tokio::test] async fn postgres_loop_checkpoint_insert_conflicts_on_same_id_different_scope_or_run() { @@ -1682,6 +1731,46 @@ mod tests { assert!(matches!(error, TurnError::Conflict { .. })); } + #[cfg(feature = "postgres")] + #[tokio::test] + async fn postgres_loop_checkpoint_get_errors_when_payload_identity_drifts() { + let Some(pool) = postgres_pool().await else { + return; + }; + let store = PostgresTurnStateStore::new(pool.clone()); + store.run_migrations().await.unwrap(); + let client = pool.get().await.unwrap(); + + let record = loop_checkpoint_record("postgres-drift"); + postgres_insert_loop_checkpoint_record(&client, &record) + .await + .unwrap(); + + let mut drifted = record.clone(); + drifted.run_id = crate::TurnRunId::new(); + client + .execute( + "UPDATE turn_loop_checkpoints SET payload = $1::jsonb WHERE checkpoint_id = $2", + &[ + &to_json(&drifted).unwrap(), + &record.checkpoint_id.as_uuid().to_string(), + ], + ) + .await + .unwrap(); + + let error = store + .get_loop_checkpoint(GetLoopCheckpointRequest { + scope: record.scope.clone(), + turn_id: record.turn_id, + run_id: record.run_id, + checkpoint_id: record.checkpoint_id, + }) + .await + .unwrap_err(); + assert!(matches!(error, TurnError::Conflict { .. })); + } + #[cfg(feature = "postgres")] async fn postgres_pool() -> Option { let Ok(url) = std::env::var("IRONCLAW_TURNS_POSTGRES_URL") else { diff --git a/crates/ironclaw_turns/src/loop_exit.rs b/crates/ironclaw_turns/src/loop_exit.rs index 76fbe3cb579..c43678e7cee 100644 --- a/crates/ironclaw_turns/src/loop_exit.rs +++ b/crates/ironclaw_turns/src/loop_exit.rs @@ -3,8 +3,8 @@ use std::{collections::HashSet, hash::Hash}; use serde::{Deserialize, Serialize, de}; use crate::{ - BlockedReason, GateRef, LoopDiagnosticRef, LoopExitId, LoopGateRef, LoopMessageRef, - LoopResultRef, LoopUsageSummaryRef, SanitizedFailure, TurnCheckpointId, + BlockedReason, GateRef, LoopCheckpointStateRef, LoopDiagnosticRef, LoopExitId, LoopGateRef, + LoopMessageRef, LoopResultRef, LoopUsageSummaryRef, SanitizedFailure, TurnCheckpointId, runner::TurnRunnerOutcome, }; @@ -37,6 +37,7 @@ impl LoopExit { exit_id, TurnRunnerOutcome::Blocked { checkpoint_id: exit.checkpoint_id, + state_ref: exit.state_ref, reason, }, ), @@ -130,6 +131,7 @@ pub struct LoopBlocked { pub kind: LoopBlockedKind, pub gate_ref: LoopGateRef, pub checkpoint_id: TurnCheckpointId, + pub state_ref: LoopCheckpointStateRef, pub exit_id: LoopExitId, } diff --git a/crates/ironclaw_turns/src/memory.rs b/crates/ironclaw_turns/src/memory.rs index 06153e911ab..7eab01b7ba1 100644 --- a/crates/ironclaw_turns/src/memory.rs +++ b/crates/ironclaw_turns/src/memory.rs @@ -654,6 +654,7 @@ impl TurnRunTransitionPort for InMemoryTurnStateStore { inner.record_checkpoint( &record, request.checkpoint_id, + request.state_ref, request.reason.gate_ref().clone(), now, ); @@ -1372,8 +1373,9 @@ impl Inner { } LoopExitMapping::RunnerOutcome(TurnRunnerOutcome::Blocked { checkpoint_id, + state_ref, reason, - }) => self.block_claimed_record(record, checkpoint_id, reason), + }) => self.block_claimed_record(record, checkpoint_id, state_ref, reason), LoopExitMapping::RunnerOutcome(TurnRunnerOutcome::Failed { failure }) => { self.fail_claimed_record(record, failure) } @@ -1463,6 +1465,7 @@ impl Inner { &mut self, mut record: RunRecord, checkpoint_id: TurnCheckpointId, + state_ref: crate::run_profile::LoopCheckpointStateRef, reason: BlockedReason, ) -> AppliedLoopTransition { if record.status != TurnStatus::Running { @@ -1483,7 +1486,13 @@ impl Inner { record.lease_token = None; record.lease_expires_at = None; record.event_cursor = self.next_cursor(); - self.record_checkpoint(&record, checkpoint_id, reason.gate_ref().clone(), now); + self.record_checkpoint( + &record, + checkpoint_id, + state_ref, + reason.gate_ref().clone(), + now, + ); self.update_active_lock(&record, now); let state = record.state(); self.push_event(&record, TurnEventKind::Blocked, None); @@ -1710,6 +1719,7 @@ impl Inner { &mut self, record: &RunRecord, checkpoint_id: TurnCheckpointId, + state_ref: crate::run_profile::LoopCheckpointStateRef, gate_ref: crate::GateRef, created_at: crate::TurnTimestamp, ) { @@ -1727,10 +1737,7 @@ impl Inner { status: record.status, gate_ref, kind: crate::run_profile::LoopCheckpointKind::BeforeBlock, - // Placeholder — callers in block_run don't have the loop's actual state_ref. - // Real values will be threaded in a follow-up task. - state_ref: crate::run_profile::LoopCheckpointStateRef::new("checkpoint:block-state") - .unwrap(), // safety: literal satisfies the "checkpoint:" prefix and length constraints. + state_ref, created_at, }); } diff --git a/crates/ironclaw_turns/src/run_profile/host.rs b/crates/ironclaw_turns/src/run_profile/host.rs index a376ebfadc2..89548c70d83 100644 --- a/crates/ironclaw_turns/src/run_profile/host.rs +++ b/crates/ironclaw_turns/src/run_profile/host.rs @@ -209,6 +209,10 @@ bounded_loop_ref!( bounded_loop_ref!(LoopProcessRef, "loop process ref", "process:", 256); impl LoopCheckpointStateRef { + pub(crate) fn legacy_unknown() -> Self { + Self("checkpoint:unknown".to_string()) + } + pub fn for_run(context: &LoopRunContext, token: impl Into) -> Result { let token = validate_loop_opaque_token(token.into(), "loop checkpoint state token", 96)?; Self::new(format!("checkpoint:{}:{token}", context.run_id)) diff --git a/crates/ironclaw_turns/src/runner.rs b/crates/ironclaw_turns/src/runner.rs index 45abb53560c..9236bfa11ee 100644 --- a/crates/ironclaw_turns/src/runner.rs +++ b/crates/ironclaw_turns/src/runner.rs @@ -2,9 +2,9 @@ use async_trait::async_trait; use serde::{Deserialize, Serialize}; use crate::{ - BlockedReason, LoopExit, LoopExitMapping, LoopExitValidationPolicy, ResolvedRunProfile, - SanitizedFailure, TurnCheckpointId, TurnError, TurnLeaseToken, TurnRunId, TurnRunState, - TurnRunnerId, TurnScope, TurnTimestamp, events::EventCursor, + BlockedReason, LoopCheckpointStateRef, LoopExit, LoopExitMapping, LoopExitValidationPolicy, + ResolvedRunProfile, SanitizedFailure, TurnCheckpointId, TurnError, TurnLeaseToken, TurnRunId, + TurnRunState, TurnRunnerId, TurnScope, TurnTimestamp, events::EventCursor, }; #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] @@ -46,6 +46,7 @@ pub struct BlockRunRequest { pub runner_id: TurnRunnerId, pub lease_token: TurnLeaseToken, pub checkpoint_id: TurnCheckpointId, + pub state_ref: LoopCheckpointStateRef, pub reason: BlockedReason, } @@ -102,6 +103,7 @@ pub enum TurnRunnerOutcome { Cancelled, Blocked { checkpoint_id: TurnCheckpointId, + state_ref: LoopCheckpointStateRef, reason: BlockedReason, }, Failed { diff --git a/crates/ironclaw_turns/src/store.rs b/crates/ironclaw_turns/src/store.rs index 510fa06ee3d..96df6b6228e 100644 --- a/crates/ironclaw_turns/src/store.rs +++ b/crates/ironclaw_turns/src/store.rs @@ -114,11 +114,10 @@ fn default_checkpoint_kind() -> LoopCheckpointKind { LoopCheckpointKind::BeforeBlock } -/// Serde default for `LoopCheckpointStateRef` — sentinel value used when +/// Serde default for `LoopCheckpointStateRef` — legacy sentinel used only when /// deserializing old persisted data that predates the `state_ref` field. -/// Real values will be threaded from loop callers in a follow-up task. fn default_checkpoint_state_ref() -> LoopCheckpointStateRef { - LoopCheckpointStateRef::new("checkpoint:unknown").unwrap() // safety: literal satisfies the "checkpoint:" prefix and length constraints. + LoopCheckpointStateRef::legacy_unknown() } #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] diff --git a/crates/ironclaw_turns/tests/agent_loop_host_contract.rs b/crates/ironclaw_turns/tests/agent_loop_host_contract.rs index da77727e705..33d9313e442 100644 --- a/crates/ironclaw_turns/tests/agent_loop_host_contract.rs +++ b/crates/ironclaw_turns/tests/agent_loop_host_contract.rs @@ -945,10 +945,11 @@ impl AgentLoopDriver for CapabilityDriver { reason_kind: "expected_approval".to_string(), }); }; + let state_ref = LoopCheckpointStateRef::new("checkpoint:approval-state").unwrap(); let checkpoint_id = host .checkpoint(LoopCheckpointRequest { kind: LoopCheckpointKind::BeforeBlock, - state_ref: LoopCheckpointStateRef::new("checkpoint:approval-state").unwrap(), + state_ref: state_ref.clone(), }) .await .map_err(driver_error)?; @@ -962,6 +963,7 @@ impl AgentLoopDriver for CapabilityDriver { kind: LoopBlockedKind::Approval, gate_ref, checkpoint_id, + state_ref, exit_id: LoopExitId::new("exit:capability-driver").unwrap(), })) } diff --git a/crates/ironclaw_turns/tests/loop_checkpoint_store_contract.rs b/crates/ironclaw_turns/tests/loop_checkpoint_store_contract.rs index 54d02748a99..6f1be011eb0 100644 --- a/crates/ironclaw_turns/tests/loop_checkpoint_store_contract.rs +++ b/crates/ironclaw_turns/tests/loop_checkpoint_store_contract.rs @@ -111,7 +111,10 @@ async fn inmemory_turn_state_loop_checkpoint_roundtrip_and_snapshot() { let snapshot = store.persistence_snapshot(); assert_eq!(snapshot.loop_checkpoints.len(), 2); assert!( - snapshot.checkpoints.is_empty(), + snapshot + .checkpoints + .iter() + .all(|record| record.state_ref.as_str() != "checkpoint:test-state"), "loop checkpoint mappings must not use turn_checkpoints" ); @@ -143,7 +146,10 @@ async fn libsql_loop_checkpoint_roundtrip_uses_loop_mapping_table() { let snapshot = store.persistence_snapshot().await.unwrap(); assert_eq!(snapshot.loop_checkpoints.len(), 2); assert!( - snapshot.checkpoints.is_empty(), + snapshot + .checkpoints + .iter() + .all(|record| record.state_ref.as_str() != "checkpoint:test-state"), "libSQL loop mappings must not be written to turn_checkpoints" ); } diff --git a/crates/ironclaw_turns/tests/loop_exit_contract.rs b/crates/ironclaw_turns/tests/loop_exit_contract.rs index 7d86a89d816..ea70360fbbb 100644 --- a/crates/ironclaw_turns/tests/loop_exit_contract.rs +++ b/crates/ironclaw_turns/tests/loop_exit_contract.rs @@ -1,9 +1,9 @@ use ironclaw_turns::{ BlockedReason, GateRef, LoopBlocked, LoopBlockedKind, LoopCancelled, LoopCancelledReasonKind, - LoopCompleted, LoopCompletionKind, LoopExit, LoopExitId, LoopExitInvalidHandling, - LoopExitValidationDecision, LoopExitValidationPolicy, LoopExitViolationKind, LoopFailureKind, - LoopGateRef, LoopMessageRef, LoopResultRef, SanitizedFailure, TurnCheckpointId, TurnStatus, - runner::TurnRunnerOutcome, + LoopCheckpointStateRef, LoopCompleted, LoopCompletionKind, LoopExit, LoopExitId, + LoopExitInvalidHandling, LoopExitValidationDecision, LoopExitValidationPolicy, + LoopExitViolationKind, LoopFailureKind, LoopGateRef, LoopMessageRef, LoopResultRef, + SanitizedFailure, TurnCheckpointId, TurnStatus, runner::TurnRunnerOutcome, }; use serde_json::json; @@ -197,10 +197,12 @@ fn blocked_exit_maps_to_block_run_outcome_with_verified_checkpoint_and_gate_ref( let checkpoint_id = TurnCheckpointId::new(); let loop_gate_ref = loop_gate_ref("gate:approval-gate"); let gate_ref = GateRef::new(loop_gate_ref.as_str()).unwrap(); + let state_ref = checkpoint_state_ref(); let decision = LoopExit::Blocked(LoopBlocked { kind: LoopBlockedKind::Approval, gate_ref: loop_gate_ref, checkpoint_id, + state_ref: state_ref.clone(), exit_id: exit_id("exit:blocked"), }) .validate(LoopExitValidationPolicy { @@ -217,6 +219,7 @@ fn blocked_exit_maps_to_block_run_outcome_with_verified_checkpoint_and_gate_ref( decision.mapping, TurnRunnerOutcome::Blocked { checkpoint_id, + state_ref, reason: BlockedReason::Approval { gate_ref }, } .into() @@ -229,6 +232,7 @@ fn blocked_exit_requires_host_verified_gate_and_checkpoint_before_trusted_mappin kind: LoopBlockedKind::Approval, gate_ref: loop_gate_ref("gate:approval-gate"), checkpoint_id: TurnCheckpointId::new(), + state_ref: checkpoint_state_ref(), exit_id: exit_id("exit:unverified-blocked"), }) .validate(LoopExitValidationPolicy { @@ -473,11 +477,13 @@ fn blocked_variants_map_to_correct_blocked_reason() { let checkpoint_id = TurnCheckpointId::new(); let lg = loop_gate_ref("gate:test-gate"); let gate_ref = GateRef::new(lg.as_str()).unwrap(); + let state_ref = checkpoint_state_ref(); let decision = LoopExit::Blocked(LoopBlocked { kind, gate_ref: lg, checkpoint_id, + state_ref: state_ref.clone(), exit_id: exit_id("exit:blocked-variant"), }) .validate(LoopExitValidationPolicy { @@ -496,6 +502,7 @@ fn blocked_variants_map_to_correct_blocked_reason() { decision.mapping, TurnRunnerOutcome::Blocked { checkpoint_id, + state_ref, reason: expected_reason, } .into() @@ -638,6 +645,10 @@ fn loop_gate_ref(value: &str) -> LoopGateRef { LoopGateRef::new(value).unwrap() } +fn checkpoint_state_ref() -> LoopCheckpointStateRef { + LoopCheckpointStateRef::new("checkpoint:blocked-state").unwrap() +} + fn result_ref(value: &str) -> LoopResultRef { LoopResultRef::new(value).unwrap() } diff --git a/crates/ironclaw_turns/tests/turn_coordinator_contract.rs b/crates/ironclaw_turns/tests/turn_coordinator_contract.rs index 23e60c658b2..cfb509292ac 100644 --- a/crates/ironclaw_turns/tests/turn_coordinator_contract.rs +++ b/crates/ironclaw_turns/tests/turn_coordinator_contract.rs @@ -14,14 +14,14 @@ use ironclaw_turns::{ AcceptedMessageRef, AdmissionRejection, AdmissionRejectionReason, AllowAllTurnAdmissionPolicy, BlockedReason, CancelRunRequest, DefaultTurnCoordinator, GateRef, GetRunStateRequest, IdempotencyKey, InMemoryRunProfileResolver, InMemoryTurnEventSink, InMemoryTurnStateStore, - InMemoryTurnStateStoreLimits, LoopCancelled, LoopCancelledReasonKind, LoopCompleted, - LoopCompletionKind, LoopDiagnosticRef, LoopExit, LoopExitId, LoopExitInvalidHandling, - LoopExitValidationPolicy, LoopFailed, LoopFailureKind, LoopGateRef, LoopMessageRef, - LoopUsageSummaryRef, ReplyTargetBindingRef, ResolvedRunProfile, ResumeTurnRequest, - RunProfileId, RunProfileRequest, RunProfileResolutionError, RunProfileResolutionRequest, - RunProfileResolver, RunProfileVersion, SanitizedCancelReason, SanitizedFailure, - SourceBindingRef, StaticTurnAdmissionLimitProvider, SubmitTurnRequest, SubmitTurnResponse, - ThreadBusy, TurnActor, TurnAdmissionAxisKind, TurnAdmissionBucketKind, + InMemoryTurnStateStoreLimits, LoopCancelled, LoopCancelledReasonKind, LoopCheckpointStateRef, + LoopCompleted, LoopCompletionKind, LoopDiagnosticRef, LoopExit, LoopExitId, + LoopExitInvalidHandling, LoopExitValidationPolicy, LoopFailed, LoopFailureKind, LoopGateRef, + LoopMessageRef, LoopUsageSummaryRef, ReplyTargetBindingRef, ResolvedRunProfile, + ResumeTurnRequest, RunProfileId, RunProfileRequest, RunProfileResolutionError, + RunProfileResolutionRequest, RunProfileResolver, RunProfileVersion, SanitizedCancelReason, + SanitizedFailure, SourceBindingRef, StaticTurnAdmissionLimitProvider, SubmitTurnRequest, + SubmitTurnResponse, ThreadBusy, TurnActor, TurnAdmissionAxisKind, TurnAdmissionBucketKind, TurnAdmissionBucketScope, TurnAdmissionCapacityDenial, TurnAdmissionClass, TurnAdmissionPolicy, TurnCheckpointId, TurnCoordinator, TurnError, TurnErrorCategory, TurnEventKind, TurnEventProjectionCursor, TurnEventProjectionError, TurnEventProjectionRequest, @@ -102,6 +102,7 @@ async fn turn_lifecycle_projection_replays_submit_block_resume_complete_without_ runner_id, lease_token, checkpoint_id: TurnCheckpointId::new(), + state_ref: block_state_ref(), reason: BlockedReason::Approval { gate_ref: gate_ref.clone(), }, @@ -762,6 +763,7 @@ async fn resume_turn_wakes_runner_for_same_run_after_requeue() { runner_id, lease_token, checkpoint_id: TurnCheckpointId::new(), + state_ref: block_state_ref(), reason: BlockedReason::Approval { gate_ref: gate_ref.clone(), }, @@ -847,6 +849,7 @@ async fn resume_turn_ignores_wake_notification_panic_after_requeue() { runner_id, lease_token, checkpoint_id: TurnCheckpointId::new(), + state_ref: block_state_ref(), reason: BlockedReason::Approval { gate_ref: gate_ref.clone(), }, @@ -1671,6 +1674,7 @@ async fn blocked_resume_and_recovery_required_keep_existing_admission_reservatio gate_ref: gate_ref.clone(), }, checkpoint_id: TurnCheckpointId::new(), + state_ref: block_state_ref(), }) .await .unwrap(); @@ -1807,12 +1811,14 @@ async fn runner_claim_and_block_update_persistent_run_lock_and_checkpoint_record let checkpoint_id = TurnCheckpointId::new(); let gate_ref = GateRef::new("approval-gate").unwrap(); + let state_ref = LoopCheckpointStateRef::new("checkpoint:requested-block-state").unwrap(); store .block_run(BlockRunRequest { run_id, runner_id, lease_token, checkpoint_id, + state_ref: state_ref.clone(), reason: BlockedReason::Approval { gate_ref: gate_ref.clone(), }, @@ -1844,6 +1850,7 @@ async fn runner_claim_and_block_update_persistent_run_lock_and_checkpoint_record assert_eq!(checkpoint.run_id, run_id); assert_eq!(checkpoint.sequence, 1); assert_eq!(checkpoint.gate_ref, gate_ref); + assert_eq!(checkpoint.state_ref, state_ref); } #[tokio::test] @@ -1873,6 +1880,7 @@ async fn resume_updates_persisted_run_binding_refs_and_replay_envelope() { runner_id, lease_token, checkpoint_id: TurnCheckpointId::new(), + state_ref: block_state_ref(), reason: BlockedReason::Approval { gate_ref: gate_ref.clone(), }, @@ -2190,6 +2198,7 @@ async fn idempotency_persistence_snapshot_retains_each_operation_kind_capacity() runner_id, lease_token, checkpoint_id: TurnCheckpointId::new(), + state_ref: block_state_ref(), reason: BlockedReason::Approval { gate_ref: gate_ref.clone(), }, @@ -2286,6 +2295,7 @@ async fn idempotency_replay_helpers_require_matching_operation_kind() { runner_id, lease_token, checkpoint_id: TurnCheckpointId::new(), + state_ref: block_state_ref(), reason: BlockedReason::Approval { gate_ref: gate_ref.clone(), }, @@ -3119,6 +3129,7 @@ async fn blocked_run_persists_checkpoint_and_keeps_same_thread_lock_until_resume runner_id, lease_token, checkpoint_id, + state_ref: block_state_ref(), reason: BlockedReason::Approval { gate_ref: gate_ref.clone(), }, @@ -3181,6 +3192,7 @@ async fn resume_turn_with_wrong_gate_resolution_ref_is_invalid_request() { runner_id, lease_token, checkpoint_id: TurnCheckpointId::new(), + state_ref: block_state_ref(), reason: BlockedReason::Approval { gate_ref: GateRef::new("approval-gate").unwrap(), }, @@ -3376,6 +3388,7 @@ async fn cancelled_running_run_cannot_be_reopened_as_blocked() { runner_id, lease_token, checkpoint_id: TurnCheckpointId::new(), + state_ref: block_state_ref(), reason: BlockedReason::Approval { gate_ref: GateRef::new("approval-gate").unwrap(), }, @@ -3559,6 +3572,10 @@ fn accepted_run_id(response: &SubmitTurnResponse) -> TurnRunId { *run_id } +fn block_state_ref() -> LoopCheckpointStateRef { + LoopCheckpointStateRef::new("checkpoint:block-state").unwrap() +} + fn scope(thread: &str) -> TurnScope { TurnScope::new( TenantId::new("tenant1").unwrap(), @@ -3843,6 +3860,7 @@ async fn loop_exit_application_blocks_with_checkpoint_and_keeps_lock() { .unwrap(); let checkpoint_id = TurnCheckpointId::new(); let gate_ref = LoopGateRef::new("gate:approval-gate").unwrap(); + let state_ref = block_state_ref(); let blocked = apply_loop_exit( store.as_ref(), @@ -3854,6 +3872,7 @@ async fn loop_exit_application_blocks_with_checkpoint_and_keeps_lock() { kind: ironclaw_turns::LoopBlockedKind::Approval, gate_ref: gate_ref.clone(), checkpoint_id, + state_ref, exit_id: ironclaw_turns::LoopExitId::new("exit:blocked").unwrap(), }), validation_policy: LoopExitValidationPolicy { diff --git a/docs/reborn/contracts/loop-exit.md b/docs/reborn/contracts/loop-exit.md index 18f289f951c..812f4725a7f 100644 --- a/docs/reborn/contracts/loop-exit.md +++ b/docs/reborn/contracts/loop-exit.md @@ -49,7 +49,7 @@ The driver-facing variants are fixed for the MVP: - `Completed` requires at least one durable reply-message ref or result ref, and the host/runner must verify those refs exist before mapping to a trusted completed outcome. Raw reply text is rejected by the wire shape and by strict loop-ref grammar. - `Completed` requires `final_checkpoint_id` only when the resolved run profile/checkpoint policy requires a terminal checkpoint. -- `Blocked` requires all of: blocked kind, durable `gate_ref`, and `checkpoint_id`, and the host/runner must verify the gate/checkpoint evidence before mapping to a trusted blocked outcome. The blocked kind is limited to approval, auth, and resource for MVP. +- `Blocked` requires all of: blocked kind, durable `gate_ref`, `checkpoint_id`, and opaque `state_ref`, and the host/runner must verify the gate/checkpoint evidence before mapping to a trusted blocked outcome. The blocked kind is limited to approval, auth, and resource for MVP. - `Cancelled` is accepted only when the host cancellation/interrupt input was observed by the runner/host policy. During application, terminal cancellation is still gated by durable run state in one transition-port operation: if the run is already `CancelRequested`, it becomes `Cancelled`; if an interrupt is observed before that durable state exists, the exit maps to recovery instead of terminal cancellation. - `Failed` uses stable sanitized failure kinds such as `iteration_limit`, `model_error`, `context_build_failed`, or `driver_bug`, and the host/runner must verify the failure evidence is safe to terminalize before mapping to a trusted failed outcome. - Ref lists are bounded and duplicate-free so a driver cannot force unbounded evidence verification work. diff --git a/docs/reborn/contracts/turn-persistence.md b/docs/reborn/contracts/turn-persistence.md index 2112c052e8b..3850aa18ed5 100644 --- a/docs/reborn/contracts/turn-persistence.md +++ b/docs/reborn/contracts/turn-persistence.md @@ -37,6 +37,8 @@ The `ironclaw_turns` contract models persistence with these record families: The initial PostgreSQL/libSQL adapter slice stores each logical record family in its own table with indexed metadata columns plus a serialized contract payload. Mutations hold a backend transaction/write lock across snapshot load, in-memory contract mutation, and snapshot replacement so active-lock and idempotency semantics remain atomic. Backends must preserve the same semantics as the in-memory contract tests while later slices add incremental row-level updates, targeted read paths, and service-graph wiring. +Legacy `turn_checkpoints` rows created before scoped checkpoint metadata may carry empty indexed `scope_key` values after migration. Those rows remain readable through the serialized payload; any future targeted `turn_checkpoints` read path must first add a scoped backfill plan or explicitly reject unbackfilled legacy rows instead of treating empty scope as a real owner. + --- ## 3. Active-lock rules