From 4071010299578947f676a36d87d480da924444ed Mon Sep 17 00:00:00 2001 From: serrrfirat Date: Mon, 18 May 2026 23:13:34 +0200 Subject: [PATCH 1/6] Add durable product workflow ledger --- Cargo.lock | 4 + crates/ironclaw_product_workflow/Cargo.toml | 6 + .../src/durable_ledger.rs | 605 ++++++++++++++++++ crates/ironclaw_product_workflow/src/lib.rs | 6 + .../tests/durable_ledger_contract.rs | 123 ++++ 5 files changed, 744 insertions(+) create mode 100644 crates/ironclaw_product_workflow/src/durable_ledger.rs create mode 100644 crates/ironclaw_product_workflow/tests/durable_ledger_contract.rs diff --git a/Cargo.lock b/Cargo.lock index 51b8e5f0c33..6160cf8736b 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -4565,6 +4565,7 @@ version = "0.1.0" dependencies = [ "async-trait", "chrono", + "deadpool-postgres", "ironclaw_host_api", "ironclaw_loop_support", "ironclaw_product_adapters", @@ -4572,10 +4573,13 @@ dependencies = [ "ironclaw_reborn", "ironclaw_threads", "ironclaw_turns", + "libsql", "serde", "serde_json", + "tempfile", "thiserror 2.0.18", "tokio", + "tokio-postgres", "tokio-util", "tracing", "uuid", diff --git a/crates/ironclaw_product_workflow/Cargo.toml b/crates/ironclaw_product_workflow/Cargo.toml index 167c6b72ca7..1b52d91d933 100644 --- a/crates/ironclaw_product_workflow/Cargo.toml +++ b/crates/ironclaw_product_workflow/Cargo.toml @@ -16,17 +16,22 @@ default = [] # FakeIdempotencyLedger, etc.). Off by default so production binaries don't carry # the fake state machinery; downstream crates enable it from `[dev-dependencies]`. test-support = ["ironclaw_product_adapters/test-support"] +libsql = ["dep:libsql"] +postgres = ["dep:deadpool-postgres", "dep:tokio-postgres"] [dependencies] async-trait = "0.1" chrono = { version = "0.4", features = ["serde"] } +deadpool-postgres = { version = "0.14", optional = true } ironclaw_host_api = { path = "../ironclaw_host_api", version = "0.1.0" } ironclaw_product_adapters = { path = "../ironclaw_product_adapters", version = "0.1.0" } ironclaw_threads = { path = "../ironclaw_threads", version = "0.1.0" } ironclaw_turns = { path = "../ironclaw_turns", version = "0.1.0" } +libsql = { version = "0.6", optional = true, default-features = false, features = ["core", "replication", "remote", "tls"] } serde = { version = "1", features = ["derive"] } serde_json = "1" thiserror = "2" +tokio-postgres = { version = "0.7", optional = true } tracing = "0.1" uuid = { version = "1", features = ["v4", "v5", "serde"] } @@ -35,5 +40,6 @@ ironclaw_loop_support = { path = "../ironclaw_loop_support" } ironclaw_reborn = { path = "../ironclaw_reborn" } ironclaw_product_workflow = { path = ".", features = ["test-support"] } ironclaw_product_adapters = { path = "../ironclaw_product_adapters", features = ["test-support", "host-auth-mint"] } +tempfile = "3" tokio = { version = "1", features = ["macros", "rt", "sync", "time"] } tokio-util = { version = "0.7", features = ["rt"] } diff --git a/crates/ironclaw_product_workflow/src/durable_ledger.rs b/crates/ironclaw_product_workflow/src/durable_ledger.rs new file mode 100644 index 00000000000..ae0e6792f9c --- /dev/null +++ b/crates/ironclaw_product_workflow/src/durable_ledger.rs @@ -0,0 +1,605 @@ +//! Durable [`IdempotencyLedger`] implementations. + +use chrono::{DateTime, Duration, Utc}; + +use crate::{ + ActionFingerprintKey, ActionPhase, IdempotencyDecision, ProductInboundAction, + ProductWorkflowError, +}; + +const DEFAULT_IN_FLIGHT_LEASE: Duration = Duration::seconds(60); + +const SCHEMA: &str = r#" +CREATE TABLE IF NOT EXISTS reborn_product_workflow_actions ( + adapter_id TEXT NOT NULL, + installation_id TEXT NOT NULL, + source_binding_key TEXT NOT NULL, + external_event_id TEXT NOT NULL, + action_id TEXT NOT NULL, + phase TEXT NOT NULL, + received_at TEXT NOT NULL, + settled_at TEXT, + payload TEXT NOT NULL, + PRIMARY KEY (adapter_id, installation_id, source_binding_key, external_event_id) +); + +CREATE INDEX IF NOT EXISTS idx_reborn_product_workflow_actions_phase + ON reborn_product_workflow_actions(phase, received_at); +"#; + +fn transient(reason: impl Into) -> ProductWorkflowError { + ProductWorkflowError::Transient { + reason: reason.into(), + } +} + +fn durable_error(operation: &'static str, error: impl std::fmt::Display) -> ProductWorkflowError { + transient(format!("idempotency ledger failed to {operation}: {error}")) +} + +fn to_json(action: &ProductInboundAction) -> Result { + serde_json::to_string(action).map_err(|error| durable_error("serialize action", error)) +} + +fn from_json(payload: &str) -> Result { + serde_json::from_str(payload).map_err(|error| durable_error("deserialize action", error)) +} + +fn phase_label(phase: ActionPhase) -> &'static str { + match phase { + ActionPhase::Received => "received", + ActionPhase::Dispatched => "dispatched", + ActionPhase::Settled => "settled", + ActionPhase::DeduplicatedReplay => "deduplicated_replay", + } +} + +fn fresh_in_flight( + action: &ProductInboundAction, + received_at: DateTime, + lease: Duration, +) -> bool { + !action.is_terminal() && action.received_at + lease > received_at +} + +fn in_flight_error() -> ProductWorkflowError { + transient("idempotency fingerprint already in flight; retry after recovery lease") +} + +#[cfg(feature = "libsql")] +mod libsql_impl { + use std::sync::Arc; + + use async_trait::async_trait; + use libsql::params; + + use super::*; + use crate::IdempotencyLedger; + + /// libSQL-backed product workflow idempotency ledger. + pub struct RebornLibSqlIdempotencyLedger { + db: Arc, + in_flight_lease: Duration, + } + + impl RebornLibSqlIdempotencyLedger { + pub fn new(db: Arc) -> Self { + Self::with_in_flight_lease(db, DEFAULT_IN_FLIGHT_LEASE) + } + + pub fn with_in_flight_lease(db: Arc, in_flight_lease: Duration) -> Self { + Self { + db, + in_flight_lease, + } + } + + pub async fn run_migrations(&self) -> Result<(), ProductWorkflowError> { + let conn = self.connect().await?; + conn.execute_batch(SCHEMA) + .await + .map(|_| ()) + .map_err(|error| durable_error("run migrations", error)) + } + + async fn connect(&self) -> Result { + let conn = self + .db + .connect() + .map_err(|error| durable_error("connect", error))?; + conn.query("PRAGMA busy_timeout = 5000", ()) + .await + .map_err(|error| durable_error("configure busy timeout", error))?; + Ok(conn) + } + + async fn begin_immediate(&self) -> Result { + let conn = self.connect().await?; + conn.execute("BEGIN IMMEDIATE", ()) + .await + .map_err(|error| durable_error("begin transaction", error))?; + Ok(conn) + } + } + + #[async_trait] + impl IdempotencyLedger for RebornLibSqlIdempotencyLedger { + async fn begin_or_replay( + &self, + fingerprint: ActionFingerprintKey, + received_at: DateTime, + ) -> Result { + self.run_migrations().await?; + let conn = self.begin_immediate().await?; + let result = + begin_or_replay_in_conn(&conn, fingerprint, received_at, self.in_flight_lease) + .await; + finish_transaction(&conn, result).await + } + + async fn settle(&self, action: ProductInboundAction) -> Result<(), ProductWorkflowError> { + self.run_migrations().await?; + let conn = self.begin_immediate().await?; + let result = settle_in_conn(&conn, action).await; + finish_transaction(&conn, result).await + } + + async fn release(&self, action: ProductInboundAction) -> Result<(), ProductWorkflowError> { + self.run_migrations().await?; + let conn = self.begin_immediate().await?; + let result = release_in_conn(&conn, action).await; + finish_transaction(&conn, result).await + } + } + + async fn begin_or_replay_in_conn( + conn: &libsql::Connection, + fingerprint: ActionFingerprintKey, + received_at: DateTime, + in_flight_lease: Duration, + ) -> Result { + let action = ProductInboundAction::begin(fingerprint.clone(), received_at); + let inserted = insert_action(conn, &action, "INSERT OR IGNORE").await?; + if inserted == 1 { + return Ok(IdempotencyDecision::New(action)); + } + + let Some(prior) = load_action(conn, &fingerprint).await? else { + return Err(transient("idempotency ledger conflict row disappeared")); + }; + if prior.is_terminal() { + return Ok(IdempotencyDecision::Replay(prior)); + } + if fresh_in_flight(&prior, received_at, in_flight_lease) { + return Err(in_flight_error()); + } + + update_action(conn, &action).await?; + Ok(IdempotencyDecision::New(action)) + } + + async fn settle_in_conn( + conn: &libsql::Connection, + action: ProductInboundAction, + ) -> Result<(), ProductWorkflowError> { + let Some(current) = load_action(conn, &action.fingerprint).await? else { + return Err(transient( + "idempotency reservation missing before terminal settle", + )); + }; + if current.is_terminal() { + if current.action_id == action.action_id { + return Ok(()); + } + return Err(transient( + "idempotency reservation was superseded before terminal settle", + )); + } + if current.action_id != action.action_id { + return Err(transient( + "idempotency reservation was superseded before terminal settle", + )); + } + update_action(conn, &action).await + } + + async fn release_in_conn( + conn: &libsql::Connection, + action: ProductInboundAction, + ) -> Result<(), ProductWorkflowError> { + conn.execute( + "DELETE FROM reborn_product_workflow_actions + WHERE adapter_id = ?1 + AND installation_id = ?2 + AND source_binding_key = ?3 + AND external_event_id = ?4 + AND action_id = ?5 + AND phase NOT IN ('settled', 'deduplicated_replay')", + params![ + action.fingerprint.adapter_id.as_str(), + action.fingerprint.installation_id.as_str(), + action.fingerprint.source_binding_key.as_str(), + action.fingerprint.external_event_id.as_str(), + action.action_id.to_string(), + ], + ) + .await + .map_err(|error| durable_error("release action", error))?; + Ok(()) + } + + async fn load_action( + conn: &libsql::Connection, + fingerprint: &ActionFingerprintKey, + ) -> Result, ProductWorkflowError> { + let mut rows = conn + .query( + "SELECT payload FROM reborn_product_workflow_actions + WHERE adapter_id = ?1 + AND installation_id = ?2 + AND source_binding_key = ?3 + AND external_event_id = ?4", + params![ + fingerprint.adapter_id.as_str(), + fingerprint.installation_id.as_str(), + fingerprint.source_binding_key.as_str(), + fingerprint.external_event_id.as_str(), + ], + ) + .await + .map_err(|error| durable_error("load action", error))?; + let Some(row) = rows + .next() + .await + .map_err(|error| durable_error("read action", error))? + else { + return Ok(None); + }; + let payload = row + .get::(0) + .map_err(|error| durable_error("read action payload", error))?; + Ok(Some(from_json(&payload)?)) + } + + async fn insert_action( + conn: &libsql::Connection, + action: &ProductInboundAction, + insert_prefix: &str, + ) -> Result { + let payload = to_json(action)?; + conn.execute( + &format!( + "{insert_prefix} INTO reborn_product_workflow_actions + (adapter_id, installation_id, source_binding_key, external_event_id, + action_id, phase, received_at, settled_at, payload) + VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9)" + ), + params![ + action.fingerprint.adapter_id.as_str(), + action.fingerprint.installation_id.as_str(), + action.fingerprint.source_binding_key.as_str(), + action.fingerprint.external_event_id.as_str(), + action.action_id.to_string(), + phase_label(action.phase), + action.received_at.to_rfc3339(), + action.settled_at.map(|value| value.to_rfc3339()), + payload, + ], + ) + .await + .map_err(|error| durable_error("insert action", error)) + } + + async fn update_action( + conn: &libsql::Connection, + action: &ProductInboundAction, + ) -> Result<(), ProductWorkflowError> { + let payload = to_json(action)?; + conn.execute( + "UPDATE reborn_product_workflow_actions + SET action_id = ?5, phase = ?6, received_at = ?7, settled_at = ?8, payload = ?9 + WHERE adapter_id = ?1 + AND installation_id = ?2 + AND source_binding_key = ?3 + AND external_event_id = ?4", + params![ + action.fingerprint.adapter_id.as_str(), + action.fingerprint.installation_id.as_str(), + action.fingerprint.source_binding_key.as_str(), + action.fingerprint.external_event_id.as_str(), + action.action_id.to_string(), + phase_label(action.phase), + action.received_at.to_rfc3339(), + action.settled_at.map(|value| value.to_rfc3339()), + payload, + ], + ) + .await + .map_err(|error| durable_error("update action", error))?; + Ok(()) + } + + async fn finish_transaction( + conn: &libsql::Connection, + result: Result, + ) -> Result { + match result { + Ok(value) => { + conn.execute("COMMIT", ()) + .await + .map_err(|error| durable_error("commit transaction", error))?; + Ok(value) + } + Err(error) => { + let _ = conn.execute("ROLLBACK", ()).await; + Err(error) + } + } + } +} + +#[cfg(feature = "postgres")] +mod postgres_impl { + use async_trait::async_trait; + + use super::*; + use crate::IdempotencyLedger; + + /// PostgreSQL-backed product workflow idempotency ledger. + pub struct RebornPostgresIdempotencyLedger { + pool: deadpool_postgres::Pool, + in_flight_lease: Duration, + } + + impl RebornPostgresIdempotencyLedger { + pub fn new(pool: deadpool_postgres::Pool) -> Self { + Self::with_in_flight_lease(pool, DEFAULT_IN_FLIGHT_LEASE) + } + + pub fn with_in_flight_lease( + pool: deadpool_postgres::Pool, + in_flight_lease: Duration, + ) -> Self { + Self { + pool, + in_flight_lease, + } + } + + pub async fn run_migrations(&self) -> Result<(), ProductWorkflowError> { + let client = self.client().await?; + client + .batch_execute(SCHEMA) + .await + .map_err(|error| durable_error("run migrations", error)) + } + + async fn client(&self) -> Result { + self.pool + .get() + .await + .map_err(|error| durable_error("connect", error)) + } + } + + #[async_trait] + impl IdempotencyLedger for RebornPostgresIdempotencyLedger { + async fn begin_or_replay( + &self, + fingerprint: ActionFingerprintKey, + received_at: DateTime, + ) -> Result { + self.run_migrations().await?; + let mut client = self.client().await?; + let txn = client + .transaction() + .await + .map_err(|error| durable_error("begin transaction", error))?; + let result = + begin_or_replay_in_txn(&txn, fingerprint, received_at, self.in_flight_lease).await; + finish_transaction(txn, result).await + } + + async fn settle(&self, action: ProductInboundAction) -> Result<(), ProductWorkflowError> { + self.run_migrations().await?; + let mut client = self.client().await?; + let txn = client + .transaction() + .await + .map_err(|error| durable_error("begin transaction", error))?; + let result = settle_in_txn(&txn, action).await; + finish_transaction(txn, result).await + } + + async fn release(&self, action: ProductInboundAction) -> Result<(), ProductWorkflowError> { + self.run_migrations().await?; + let mut client = self.client().await?; + let txn = client + .transaction() + .await + .map_err(|error| durable_error("begin transaction", error))?; + let result = release_in_txn(&txn, action).await; + finish_transaction(txn, result).await + } + } + + async fn begin_or_replay_in_txn( + txn: &deadpool_postgres::Transaction<'_>, + fingerprint: ActionFingerprintKey, + received_at: DateTime, + in_flight_lease: Duration, + ) -> Result { + let action = ProductInboundAction::begin(fingerprint.clone(), received_at); + let inserted = insert_action(txn, &action).await?; + if inserted == 1 { + return Ok(IdempotencyDecision::New(action)); + } + + let Some(prior) = load_action_for_update(txn, &fingerprint).await? else { + return Err(transient("idempotency ledger conflict row disappeared")); + }; + if prior.is_terminal() { + return Ok(IdempotencyDecision::Replay(prior)); + } + if fresh_in_flight(&prior, received_at, in_flight_lease) { + return Err(in_flight_error()); + } + + update_action(txn, &action).await?; + Ok(IdempotencyDecision::New(action)) + } + + async fn settle_in_txn( + txn: &deadpool_postgres::Transaction<'_>, + action: ProductInboundAction, + ) -> Result<(), ProductWorkflowError> { + let Some(current) = load_action_for_update(txn, &action.fingerprint).await? else { + return Err(transient( + "idempotency reservation missing before terminal settle", + )); + }; + if current.is_terminal() { + if current.action_id == action.action_id { + return Ok(()); + } + return Err(transient( + "idempotency reservation was superseded before terminal settle", + )); + } + if current.action_id != action.action_id { + return Err(transient( + "idempotency reservation was superseded before terminal settle", + )); + } + update_action(txn, &action).await + } + + async fn release_in_txn( + txn: &deadpool_postgres::Transaction<'_>, + action: ProductInboundAction, + ) -> Result<(), ProductWorkflowError> { + txn.execute( + "DELETE FROM reborn_product_workflow_actions + WHERE adapter_id = $1 + AND installation_id = $2 + AND source_binding_key = $3 + AND external_event_id = $4 + AND action_id = $5 + AND phase NOT IN ('settled', 'deduplicated_replay')", + &[ + &action.fingerprint.adapter_id.as_str(), + &action.fingerprint.installation_id.as_str(), + &action.fingerprint.source_binding_key.as_str(), + &action.fingerprint.external_event_id.as_str(), + &action.action_id.to_string(), + ], + ) + .await + .map_err(|error| durable_error("release action", error))?; + Ok(()) + } + + async fn load_action_for_update( + txn: &deadpool_postgres::Transaction<'_>, + fingerprint: &ActionFingerprintKey, + ) -> Result, ProductWorkflowError> { + let row = txn + .query_opt( + "SELECT payload FROM reborn_product_workflow_actions + WHERE adapter_id = $1 + AND installation_id = $2 + AND source_binding_key = $3 + AND external_event_id = $4 + FOR UPDATE", + &[ + &fingerprint.adapter_id.as_str(), + &fingerprint.installation_id.as_str(), + &fingerprint.source_binding_key.as_str(), + &fingerprint.external_event_id.as_str(), + ], + ) + .await + .map_err(|error| durable_error("load action", error))?; + row.map(|row| from_json(row.get::<_, &str>(0))).transpose() + } + + async fn insert_action( + txn: &deadpool_postgres::Transaction<'_>, + action: &ProductInboundAction, + ) -> Result { + let payload = to_json(action)?; + txn.execute( + "INSERT INTO reborn_product_workflow_actions + (adapter_id, installation_id, source_binding_key, external_event_id, + action_id, phase, received_at, settled_at, payload) + VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9) + ON CONFLICT (adapter_id, installation_id, source_binding_key, external_event_id) + DO NOTHING", + &[ + &action.fingerprint.adapter_id.as_str(), + &action.fingerprint.installation_id.as_str(), + &action.fingerprint.source_binding_key.as_str(), + &action.fingerprint.external_event_id.as_str(), + &action.action_id.to_string(), + &phase_label(action.phase), + &action.received_at.to_rfc3339(), + &action.settled_at.map(|value| value.to_rfc3339()), + &payload, + ], + ) + .await + .map_err(|error| durable_error("insert action", error)) + } + + async fn update_action( + txn: &deadpool_postgres::Transaction<'_>, + action: &ProductInboundAction, + ) -> Result<(), ProductWorkflowError> { + let payload = to_json(action)?; + txn.execute( + "UPDATE reborn_product_workflow_actions + SET action_id = $5, phase = $6, received_at = $7, settled_at = $8, payload = $9 + WHERE adapter_id = $1 + AND installation_id = $2 + AND source_binding_key = $3 + AND external_event_id = $4", + &[ + &action.fingerprint.adapter_id.as_str(), + &action.fingerprint.installation_id.as_str(), + &action.fingerprint.source_binding_key.as_str(), + &action.fingerprint.external_event_id.as_str(), + &action.action_id.to_string(), + &phase_label(action.phase), + &action.received_at.to_rfc3339(), + &action.settled_at.map(|value| value.to_rfc3339()), + &payload, + ], + ) + .await + .map_err(|error| durable_error("update action", error))?; + Ok(()) + } + + async fn finish_transaction( + txn: deadpool_postgres::Transaction<'_>, + result: Result, + ) -> Result { + match result { + Ok(value) => { + txn.commit() + .await + .map_err(|error| durable_error("commit transaction", error))?; + Ok(value) + } + Err(error) => { + let _ = txn.rollback().await; + Err(error) + } + } + } +} + +#[cfg(feature = "libsql")] +pub use libsql_impl::RebornLibSqlIdempotencyLedger; +#[cfg(feature = "postgres")] +pub use postgres_impl::RebornPostgresIdempotencyLedger; diff --git a/crates/ironclaw_product_workflow/src/lib.rs b/crates/ironclaw_product_workflow/src/lib.rs index 2ebae89ea6d..67ce7b82345 100644 --- a/crates/ironclaw_product_workflow/src/lib.rs +++ b/crates/ironclaw_product_workflow/src/lib.rs @@ -21,6 +21,8 @@ mod action; mod binding; +#[cfg(any(feature = "libsql", feature = "postgres"))] +mod durable_ledger; mod error; #[cfg(any(test, feature = "test-support"))] mod fakes; @@ -35,6 +37,10 @@ pub use action::{ ProductActionId, ProductCommandName, ProductInboundAction, SourceBindingKey, }; pub use binding::{ConversationBindingService, ResolveBindingRequest, ResolvedBinding}; +#[cfg(feature = "libsql")] +pub use durable_ledger::RebornLibSqlIdempotencyLedger; +#[cfg(feature = "postgres")] +pub use durable_ledger::RebornPostgresIdempotencyLedger; pub use error::ProductWorkflowError; #[cfg(any(test, feature = "test-support"))] pub use fakes::{FakeConversationBindingService, FakeIdempotencyLedger, FakeInboundTurnService}; diff --git a/crates/ironclaw_product_workflow/tests/durable_ledger_contract.rs b/crates/ironclaw_product_workflow/tests/durable_ledger_contract.rs new file mode 100644 index 00000000000..ed173c9289e --- /dev/null +++ b/crates/ironclaw_product_workflow/tests/durable_ledger_contract.rs @@ -0,0 +1,123 @@ +#![cfg(feature = "libsql")] + +use std::sync::Arc; + +use chrono::{Duration, Utc}; +use ironclaw_product_adapters::{ + AdapterInstallationId, ExternalEventId, ProductAdapterId, ProductInboundAck, +}; +use ironclaw_product_workflow::{ + ActionFingerprintKey, IdempotencyDecision, IdempotencyLedger, RebornLibSqlIdempotencyLedger, + SourceBindingKey, +}; + +fn fingerprint(suffix: &str) -> ActionFingerprintKey { + ActionFingerprintKey::new( + ProductAdapterId::new("test_adapter").expect("valid adapter"), + AdapterInstallationId::new("install_alpha").expect("valid installation"), + SourceBindingKey::new("space:0:;conversation:5:conv1;topic:0:;") + .expect("valid source binding key"), + ExternalEventId::new(format!("evt:{suffix}")).expect("valid event"), + ) +} + +async fn libsql_db(path: &str) -> Arc { + Arc::new( + libsql::Builder::new_local(path) + .build() + .await + .expect("build libsql db"), + ) +} + +#[tokio::test] +async fn libsql_settled_action_survives_reopen_and_replays() { + let dir = tempfile::tempdir().expect("tempdir"); + let db_path = dir.path().join("workflow-ledger.db"); + let db_path = db_path.display().to_string(); + let received_at = Utc::now(); + let fingerprint = fingerprint("settled-replay"); + + let ledger = RebornLibSqlIdempotencyLedger::new(libsql_db(&db_path).await); + let decision = ledger + .begin_or_replay(fingerprint.clone(), received_at) + .await + .expect("begin"); + let IdempotencyDecision::New(mut action) = decision else { + panic!("expected new action"); + }; + action.settle(ProductInboundAck::NoOp); + ledger.settle(action).await.expect("settle"); + + let reopened = RebornLibSqlIdempotencyLedger::new(libsql_db(&db_path).await); + let replay = reopened + .begin_or_replay(fingerprint, received_at + Duration::seconds(1)) + .await + .expect("replay"); + + let IdempotencyDecision::Replay(action) = replay else { + panic!("expected replay"); + }; + assert_eq!(action.outcome, Some(ProductInboundAck::NoOp)); +} + +#[tokio::test] +async fn libsql_in_flight_action_blocks_until_lease_expires() { + let dir = tempfile::tempdir().expect("tempdir"); + let db_path = dir.path().join("workflow-ledger.db"); + let ledger = RebornLibSqlIdempotencyLedger::with_in_flight_lease( + libsql_db(&db_path.display().to_string()).await, + Duration::seconds(10), + ); + let received_at = Utc::now(); + let fingerprint = fingerprint("lease"); + + assert!(matches!( + ledger + .begin_or_replay(fingerprint.clone(), received_at) + .await + .expect("begin"), + IdempotencyDecision::New(_) + )); + let blocked = ledger + .begin_or_replay(fingerprint.clone(), received_at + Duration::seconds(5)) + .await + .expect_err("fresh reservation should block"); + assert!(matches!( + blocked, + ironclaw_product_workflow::ProductWorkflowError::Transient { .. } + )); + + let reclaimed = ledger + .begin_or_replay(fingerprint, received_at + Duration::seconds(11)) + .await + .expect("expired reservation should be reclaimed"); + assert!(matches!(reclaimed, IdempotencyDecision::New(_))); +} + +#[tokio::test] +async fn libsql_release_allows_retry_without_waiting_for_lease() { + let dir = tempfile::tempdir().expect("tempdir"); + let db_path = dir.path().join("workflow-ledger.db"); + let ledger = RebornLibSqlIdempotencyLedger::with_in_flight_lease( + libsql_db(&db_path.display().to_string()).await, + Duration::seconds(60), + ); + let received_at = Utc::now(); + let fingerprint = fingerprint("release"); + + let decision = ledger + .begin_or_replay(fingerprint.clone(), received_at) + .await + .expect("begin"); + let IdempotencyDecision::New(action) = decision else { + panic!("expected new action"); + }; + ledger.release(action).await.expect("release"); + + let retry = ledger + .begin_or_replay(fingerprint, received_at + Duration::seconds(1)) + .await + .expect("retry after release"); + assert!(matches!(retry, IdempotencyDecision::New(_))); +} From c9332469f4aa839ba5333bd3e2e5a482fe7cf85c Mon Sep 17 00:00:00 2001 From: serrrfirat Date: Mon, 18 May 2026 23:38:16 +0200 Subject: [PATCH 2/6] fix(product-workflow): address review durable ledger storage (#3759) --- Cargo.lock | 20 +- Cargo.toml | 2 +- crates/ironclaw_product_workflow/Cargo.toml | 5 - crates/ironclaw_product_workflow/src/lib.rs | 6 - .../tests/durable_ledger_contract.rs | 123 ------ .../Cargo.toml | 32 ++ .../src/lib.rs} | 33 +- .../tests/durable_ledger_contract.rs | 364 ++++++++++++++++++ 8 files changed, 440 insertions(+), 145 deletions(-) delete mode 100644 crates/ironclaw_product_workflow/tests/durable_ledger_contract.rs create mode 100644 crates/ironclaw_product_workflow_storage/Cargo.toml rename crates/{ironclaw_product_workflow/src/durable_ledger.rs => ironclaw_product_workflow_storage/src/lib.rs} (95%) create mode 100644 crates/ironclaw_product_workflow_storage/tests/durable_ledger_contract.rs diff --git a/Cargo.lock b/Cargo.lock index 6160cf8736b..0511d19a03c 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -4565,7 +4565,6 @@ version = "0.1.0" dependencies = [ "async-trait", "chrono", - "deadpool-postgres", "ironclaw_host_api", "ironclaw_loop_support", "ironclaw_product_adapters", @@ -4573,18 +4572,33 @@ dependencies = [ "ironclaw_reborn", "ironclaw_threads", "ironclaw_turns", - "libsql", "serde", "serde_json", "tempfile", "thiserror 2.0.18", "tokio", - "tokio-postgres", "tokio-util", "tracing", "uuid", ] +[[package]] +name = "ironclaw_product_workflow_storage" +version = "0.1.0" +dependencies = [ + "async-trait", + "chrono", + "deadpool-postgres", + "ironclaw_product_adapters", + "ironclaw_product_workflow", + "libsql", + "serde_json", + "tempfile", + "tokio", + "tokio-postgres", + "tracing", +] + [[package]] name = "ironclaw_reborn" version = "0.1.0" diff --git a/Cargo.toml b/Cargo.toml index a2a1ebc2478..33bfd1654aa 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,5 +1,5 @@ [workspace] -members = [".", "crates/ironclaw_common", "crates/ironclaw_host_api", "crates/ironclaw_storage", "crates/ironclaw_filesystem", "crates/ironclaw_memory", "crates/ironclaw_events", "crates/ironclaw_event_projections", "crates/ironclaw_reborn_event_store", "crates/ironclaw_extensions", "crates/ironclaw_processes", "crates/ironclaw_dispatcher", "crates/ironclaw_scripts", "crates/ironclaw_mcp", "crates/ironclaw_wasm", "crates/ironclaw_wasm_sandbox_core", "crates/ironclaw_capabilities", "crates/ironclaw_secrets", "crates/ironclaw_network", "crates/ironclaw_host_runtime", "crates/ironclaw_runtime_policy", "crates/ironclaw_authorization", "crates/ironclaw_run_state", "crates/ironclaw_approvals", "crates/ironclaw_resources", "crates/ironclaw_trust", "crates/ironclaw_turns", "crates/ironclaw_agent_loop", "crates/ironclaw_threads", "crates/ironclaw_loop_support", "crates/ironclaw_reborn", "crates/ironclaw_reborn_config", "crates/ironclaw_reborn_composition", "crates/ironclaw_reborn_cli", "crates/ironclaw_conversations", "crates/ironclaw_product_adapters", "crates/ironclaw_product_workflow", "crates/ironclaw_product_adapter_registry", "crates/ironclaw_wasm_product_adapters", "crates/ironclaw_telegram_v2_adapter", "crates/ironclaw_outbound", "crates/ironclaw_architecture", "crates/ironclaw_safety", "crates/ironclaw_skills", "crates/ironclaw_llm", "crates/ironclaw_engine", "crates/ironclaw_gateway", "crates/ironclaw_tui"] +members = [".", "crates/ironclaw_common", "crates/ironclaw_host_api", "crates/ironclaw_storage", "crates/ironclaw_filesystem", "crates/ironclaw_memory", "crates/ironclaw_events", "crates/ironclaw_event_projections", "crates/ironclaw_reborn_event_store", "crates/ironclaw_extensions", "crates/ironclaw_processes", "crates/ironclaw_dispatcher", "crates/ironclaw_scripts", "crates/ironclaw_mcp", "crates/ironclaw_wasm", "crates/ironclaw_wasm_sandbox_core", "crates/ironclaw_capabilities", "crates/ironclaw_secrets", "crates/ironclaw_network", "crates/ironclaw_host_runtime", "crates/ironclaw_runtime_policy", "crates/ironclaw_authorization", "crates/ironclaw_run_state", "crates/ironclaw_approvals", "crates/ironclaw_resources", "crates/ironclaw_trust", "crates/ironclaw_turns", "crates/ironclaw_agent_loop", "crates/ironclaw_threads", "crates/ironclaw_loop_support", "crates/ironclaw_reborn", "crates/ironclaw_reborn_config", "crates/ironclaw_reborn_composition", "crates/ironclaw_reborn_cli", "crates/ironclaw_conversations", "crates/ironclaw_product_adapters", "crates/ironclaw_product_workflow", "crates/ironclaw_product_workflow_storage", "crates/ironclaw_product_adapter_registry", "crates/ironclaw_wasm_product_adapters", "crates/ironclaw_telegram_v2_adapter", "crates/ironclaw_outbound", "crates/ironclaw_architecture", "crates/ironclaw_safety", "crates/ironclaw_skills", "crates/ironclaw_llm", "crates/ironclaw_engine", "crates/ironclaw_gateway", "crates/ironclaw_tui"] exclude = [ "channels-src/discord", "channels-src/feishu", diff --git a/crates/ironclaw_product_workflow/Cargo.toml b/crates/ironclaw_product_workflow/Cargo.toml index 1b52d91d933..15eeb6cd8a8 100644 --- a/crates/ironclaw_product_workflow/Cargo.toml +++ b/crates/ironclaw_product_workflow/Cargo.toml @@ -16,22 +16,17 @@ default = [] # FakeIdempotencyLedger, etc.). Off by default so production binaries don't carry # the fake state machinery; downstream crates enable it from `[dev-dependencies]`. test-support = ["ironclaw_product_adapters/test-support"] -libsql = ["dep:libsql"] -postgres = ["dep:deadpool-postgres", "dep:tokio-postgres"] [dependencies] async-trait = "0.1" chrono = { version = "0.4", features = ["serde"] } -deadpool-postgres = { version = "0.14", optional = true } ironclaw_host_api = { path = "../ironclaw_host_api", version = "0.1.0" } ironclaw_product_adapters = { path = "../ironclaw_product_adapters", version = "0.1.0" } ironclaw_threads = { path = "../ironclaw_threads", version = "0.1.0" } ironclaw_turns = { path = "../ironclaw_turns", version = "0.1.0" } -libsql = { version = "0.6", optional = true, default-features = false, features = ["core", "replication", "remote", "tls"] } serde = { version = "1", features = ["derive"] } serde_json = "1" thiserror = "2" -tokio-postgres = { version = "0.7", optional = true } tracing = "0.1" uuid = { version = "1", features = ["v4", "v5", "serde"] } diff --git a/crates/ironclaw_product_workflow/src/lib.rs b/crates/ironclaw_product_workflow/src/lib.rs index 67ce7b82345..2ebae89ea6d 100644 --- a/crates/ironclaw_product_workflow/src/lib.rs +++ b/crates/ironclaw_product_workflow/src/lib.rs @@ -21,8 +21,6 @@ mod action; mod binding; -#[cfg(any(feature = "libsql", feature = "postgres"))] -mod durable_ledger; mod error; #[cfg(any(test, feature = "test-support"))] mod fakes; @@ -37,10 +35,6 @@ pub use action::{ ProductActionId, ProductCommandName, ProductInboundAction, SourceBindingKey, }; pub use binding::{ConversationBindingService, ResolveBindingRequest, ResolvedBinding}; -#[cfg(feature = "libsql")] -pub use durable_ledger::RebornLibSqlIdempotencyLedger; -#[cfg(feature = "postgres")] -pub use durable_ledger::RebornPostgresIdempotencyLedger; pub use error::ProductWorkflowError; #[cfg(any(test, feature = "test-support"))] pub use fakes::{FakeConversationBindingService, FakeIdempotencyLedger, FakeInboundTurnService}; diff --git a/crates/ironclaw_product_workflow/tests/durable_ledger_contract.rs b/crates/ironclaw_product_workflow/tests/durable_ledger_contract.rs deleted file mode 100644 index ed173c9289e..00000000000 --- a/crates/ironclaw_product_workflow/tests/durable_ledger_contract.rs +++ /dev/null @@ -1,123 +0,0 @@ -#![cfg(feature = "libsql")] - -use std::sync::Arc; - -use chrono::{Duration, Utc}; -use ironclaw_product_adapters::{ - AdapterInstallationId, ExternalEventId, ProductAdapterId, ProductInboundAck, -}; -use ironclaw_product_workflow::{ - ActionFingerprintKey, IdempotencyDecision, IdempotencyLedger, RebornLibSqlIdempotencyLedger, - SourceBindingKey, -}; - -fn fingerprint(suffix: &str) -> ActionFingerprintKey { - ActionFingerprintKey::new( - ProductAdapterId::new("test_adapter").expect("valid adapter"), - AdapterInstallationId::new("install_alpha").expect("valid installation"), - SourceBindingKey::new("space:0:;conversation:5:conv1;topic:0:;") - .expect("valid source binding key"), - ExternalEventId::new(format!("evt:{suffix}")).expect("valid event"), - ) -} - -async fn libsql_db(path: &str) -> Arc { - Arc::new( - libsql::Builder::new_local(path) - .build() - .await - .expect("build libsql db"), - ) -} - -#[tokio::test] -async fn libsql_settled_action_survives_reopen_and_replays() { - let dir = tempfile::tempdir().expect("tempdir"); - let db_path = dir.path().join("workflow-ledger.db"); - let db_path = db_path.display().to_string(); - let received_at = Utc::now(); - let fingerprint = fingerprint("settled-replay"); - - let ledger = RebornLibSqlIdempotencyLedger::new(libsql_db(&db_path).await); - let decision = ledger - .begin_or_replay(fingerprint.clone(), received_at) - .await - .expect("begin"); - let IdempotencyDecision::New(mut action) = decision else { - panic!("expected new action"); - }; - action.settle(ProductInboundAck::NoOp); - ledger.settle(action).await.expect("settle"); - - let reopened = RebornLibSqlIdempotencyLedger::new(libsql_db(&db_path).await); - let replay = reopened - .begin_or_replay(fingerprint, received_at + Duration::seconds(1)) - .await - .expect("replay"); - - let IdempotencyDecision::Replay(action) = replay else { - panic!("expected replay"); - }; - assert_eq!(action.outcome, Some(ProductInboundAck::NoOp)); -} - -#[tokio::test] -async fn libsql_in_flight_action_blocks_until_lease_expires() { - let dir = tempfile::tempdir().expect("tempdir"); - let db_path = dir.path().join("workflow-ledger.db"); - let ledger = RebornLibSqlIdempotencyLedger::with_in_flight_lease( - libsql_db(&db_path.display().to_string()).await, - Duration::seconds(10), - ); - let received_at = Utc::now(); - let fingerprint = fingerprint("lease"); - - assert!(matches!( - ledger - .begin_or_replay(fingerprint.clone(), received_at) - .await - .expect("begin"), - IdempotencyDecision::New(_) - )); - let blocked = ledger - .begin_or_replay(fingerprint.clone(), received_at + Duration::seconds(5)) - .await - .expect_err("fresh reservation should block"); - assert!(matches!( - blocked, - ironclaw_product_workflow::ProductWorkflowError::Transient { .. } - )); - - let reclaimed = ledger - .begin_or_replay(fingerprint, received_at + Duration::seconds(11)) - .await - .expect("expired reservation should be reclaimed"); - assert!(matches!(reclaimed, IdempotencyDecision::New(_))); -} - -#[tokio::test] -async fn libsql_release_allows_retry_without_waiting_for_lease() { - let dir = tempfile::tempdir().expect("tempdir"); - let db_path = dir.path().join("workflow-ledger.db"); - let ledger = RebornLibSqlIdempotencyLedger::with_in_flight_lease( - libsql_db(&db_path.display().to_string()).await, - Duration::seconds(60), - ); - let received_at = Utc::now(); - let fingerprint = fingerprint("release"); - - let decision = ledger - .begin_or_replay(fingerprint.clone(), received_at) - .await - .expect("begin"); - let IdempotencyDecision::New(action) = decision else { - panic!("expected new action"); - }; - ledger.release(action).await.expect("release"); - - let retry = ledger - .begin_or_replay(fingerprint, received_at + Duration::seconds(1)) - .await - .expect("retry after release"); - assert!(matches!(retry, IdempotencyDecision::New(_))); -} diff --git a/crates/ironclaw_product_workflow_storage/Cargo.toml b/crates/ironclaw_product_workflow_storage/Cargo.toml new file mode 100644 index 00000000000..3a3ceb4cacf --- /dev/null +++ b/crates/ironclaw_product_workflow_storage/Cargo.toml @@ -0,0 +1,32 @@ +[package] +name = "ironclaw_product_workflow_storage" +version = "0.1.0" +edition = "2024" +rust-version = "1.92" +description = "Durable storage adapters for IronClaw Reborn product workflow" +authors = ["NEAR AI "] +license = "MIT OR Apache-2.0" +homepage = "https://github.com/nearai/ironclaw" +repository = "https://github.com/nearai/ironclaw" +publish = false + +[features] +default = [] +libsql = ["dep:libsql"] +postgres = ["dep:deadpool-postgres", "dep:tokio-postgres"] + +[dependencies] +async-trait = "0.1" +chrono = { version = "0.4", features = ["serde"] } +deadpool-postgres = { version = "0.14", optional = true } +ironclaw_product_workflow = { path = "../ironclaw_product_workflow", version = "0.1.0" } +libsql = { version = "0.6", optional = true, default-features = false, features = ["core", "replication", "remote", "tls"] } +serde_json = "1" +tokio = { version = "1", features = ["sync"] } +tokio-postgres = { version = "0.7", optional = true } +tracing = "0.1" + +[dev-dependencies] +ironclaw_product_adapters = { path = "../ironclaw_product_adapters" } +tempfile = "3" +tokio = { version = "1", features = ["macros", "rt", "sync", "time"] } diff --git a/crates/ironclaw_product_workflow/src/durable_ledger.rs b/crates/ironclaw_product_workflow_storage/src/lib.rs similarity index 95% rename from crates/ironclaw_product_workflow/src/durable_ledger.rs rename to crates/ironclaw_product_workflow_storage/src/lib.rs index ae0e6792f9c..22bf55736ab 100644 --- a/crates/ironclaw_product_workflow/src/durable_ledger.rs +++ b/crates/ironclaw_product_workflow_storage/src/lib.rs @@ -1,10 +1,10 @@ -//! Durable [`IdempotencyLedger`] implementations. +//! Durable product workflow [`IdempotencyLedger`] storage adapters. use chrono::{DateTime, Duration, Utc}; -use crate::{ - ActionFingerprintKey, ActionPhase, IdempotencyDecision, ProductInboundAction, - ProductWorkflowError, +use ironclaw_product_workflow::{ + ActionFingerprintKey, ActionPhase, IdempotencyDecision, IdempotencyLedger, + ProductInboundAction, ProductWorkflowError, }; const DEFAULT_IN_FLIGHT_LEASE: Duration = Duration::seconds(60); @@ -34,7 +34,8 @@ fn transient(reason: impl Into) -> ProductWorkflowError { } fn durable_error(operation: &'static str, error: impl std::fmt::Display) -> ProductWorkflowError { - transient(format!("idempotency ledger failed to {operation}: {error}")) + tracing::error!(%error, operation, "product workflow idempotency ledger backend failed"); + transient(format!("idempotency ledger failed to {operation}")) } fn to_json(action: &ProductInboundAction) -> Result { @@ -74,12 +75,13 @@ mod libsql_impl { use libsql::params; use super::*; - use crate::IdempotencyLedger; + use tokio::sync::OnceCell; /// libSQL-backed product workflow idempotency ledger. pub struct RebornLibSqlIdempotencyLedger { db: Arc, in_flight_lease: Duration, + migrations: OnceCell<()>, } impl RebornLibSqlIdempotencyLedger { @@ -91,10 +93,18 @@ mod libsql_impl { Self { db, in_flight_lease, + migrations: OnceCell::new(), } } pub async fn run_migrations(&self) -> Result<(), ProductWorkflowError> { + self.migrations + .get_or_try_init(|| async { self.run_migrations_uncached().await }) + .await + .copied() + } + + async fn run_migrations_uncached(&self) -> Result<(), ProductWorkflowError> { let conn = self.connect().await?; conn.execute_batch(SCHEMA) .await @@ -343,12 +353,13 @@ mod postgres_impl { use async_trait::async_trait; use super::*; - use crate::IdempotencyLedger; + use tokio::sync::OnceCell; /// PostgreSQL-backed product workflow idempotency ledger. pub struct RebornPostgresIdempotencyLedger { pool: deadpool_postgres::Pool, in_flight_lease: Duration, + migrations: OnceCell<()>, } impl RebornPostgresIdempotencyLedger { @@ -363,10 +374,18 @@ mod postgres_impl { Self { pool, in_flight_lease, + migrations: OnceCell::new(), } } pub async fn run_migrations(&self) -> Result<(), ProductWorkflowError> { + self.migrations + .get_or_try_init(|| async { self.run_migrations_uncached().await }) + .await + .copied() + } + + async fn run_migrations_uncached(&self) -> Result<(), ProductWorkflowError> { let client = self.client().await?; client .batch_execute(SCHEMA) diff --git a/crates/ironclaw_product_workflow_storage/tests/durable_ledger_contract.rs b/crates/ironclaw_product_workflow_storage/tests/durable_ledger_contract.rs new file mode 100644 index 00000000000..ca20e454f2f --- /dev/null +++ b/crates/ironclaw_product_workflow_storage/tests/durable_ledger_contract.rs @@ -0,0 +1,364 @@ +#![cfg(any(feature = "libsql", feature = "postgres"))] + +#[cfg(feature = "libsql")] +use std::sync::Arc; +#[cfg(feature = "postgres")] +use std::time::{SystemTime, UNIX_EPOCH}; + +use chrono::{Duration, Utc}; +use ironclaw_product_adapters::{ + AdapterInstallationId, ExternalEventId, ProductAdapterId, ProductInboundAck, +}; +use ironclaw_product_workflow::{ + ActionFingerprintKey, IdempotencyDecision, IdempotencyLedger, ProductWorkflowError, + SourceBindingKey, +}; +#[cfg(feature = "libsql")] +use ironclaw_product_workflow_storage::RebornLibSqlIdempotencyLedger; +#[cfg(feature = "postgres")] +use ironclaw_product_workflow_storage::RebornPostgresIdempotencyLedger; + +fn fingerprint(suffix: &str) -> ActionFingerprintKey { + ActionFingerprintKey::new( + ProductAdapterId::new("test_adapter").expect("valid adapter"), + AdapterInstallationId::new("install_alpha").expect("valid installation"), + SourceBindingKey::new("space:0:;conversation:5:conv1;topic:0:;") + .expect("valid source binding key"), + ExternalEventId::new(format!("evt:{suffix}")).expect("valid event"), + ) +} + +#[cfg(feature = "postgres")] +fn unique_suffix(name: &str) -> String { + let nanos = SystemTime::now() + .duration_since(UNIX_EPOCH) + .expect("system clock after unix epoch") + .as_nanos(); + format!("{name}-{nanos}") +} + +#[cfg(feature = "libsql")] +async fn libsql_db(path: &str) -> Arc { + Arc::new( + libsql::Builder::new_local(path) + .build() + .await + .expect("build libsql db"), + ) +} + +async fn assert_settled_action_survives_reopen_and_replays( + ledger: &dyn IdempotencyLedger, + reopened: &dyn IdempotencyLedger, + suffix: &str, +) { + let received_at = Utc::now(); + let fingerprint = fingerprint(suffix); + + let decision = ledger + .begin_or_replay(fingerprint.clone(), received_at) + .await + .expect("begin"); + let IdempotencyDecision::New(mut action) = decision else { + panic!("expected new action"); + }; + action.settle(ProductInboundAck::NoOp); + ledger.settle(action).await.expect("settle"); + + let replay = reopened + .begin_or_replay(fingerprint, received_at + Duration::seconds(1)) + .await + .expect("replay"); + + let IdempotencyDecision::Replay(action) = replay else { + panic!("expected replay"); + }; + assert_eq!(action.outcome, Some(ProductInboundAck::NoOp)); +} + +async fn assert_in_flight_action_blocks_until_lease_expires( + ledger: &dyn IdempotencyLedger, + suffix: &str, +) { + let received_at = Utc::now(); + let fingerprint = fingerprint(suffix); + + assert!(matches!( + ledger + .begin_or_replay(fingerprint.clone(), received_at) + .await + .expect("begin"), + IdempotencyDecision::New(_) + )); + let blocked = ledger + .begin_or_replay(fingerprint.clone(), received_at + Duration::seconds(5)) + .await + .expect_err("fresh reservation should block"); + assert!(matches!(blocked, ProductWorkflowError::Transient { .. })); + + let reclaimed = ledger + .begin_or_replay(fingerprint, received_at + Duration::seconds(11)) + .await + .expect("expired reservation should be reclaimed"); + assert!(matches!(reclaimed, IdempotencyDecision::New(_))); +} + +async fn assert_release_allows_retry_without_waiting_for_lease( + ledger: &dyn IdempotencyLedger, + suffix: &str, +) { + let received_at = Utc::now(); + let fingerprint = fingerprint(suffix); + + let decision = ledger + .begin_or_replay(fingerprint.clone(), received_at) + .await + .expect("begin"); + let IdempotencyDecision::New(action) = decision else { + panic!("expected new action"); + }; + ledger.release(action).await.expect("release"); + + let retry = ledger + .begin_or_replay(fingerprint, received_at + Duration::seconds(1)) + .await + .expect("retry after release"); + assert!(matches!(retry, IdempotencyDecision::New(_))); +} + +async fn assert_duplicate_reservation_contention_serializes( + first: &dyn IdempotencyLedger, + second: &dyn IdempotencyLedger, + suffix: &str, +) { + let received_at = Utc::now(); + let fingerprint = fingerprint(suffix); + + let (left, right) = tokio::join!( + first.begin_or_replay(fingerprint.clone(), received_at), + second.begin_or_replay(fingerprint, received_at), + ); + let results = [left, right]; + let new_count = results + .iter() + .filter(|result| matches!(result, Ok(IdempotencyDecision::New(_)))) + .count(); + let blocked_count = results + .iter() + .filter(|result| matches!(result, Err(ProductWorkflowError::Transient { .. }))) + .count(); + + assert_eq!(new_count, 1); + assert_eq!(blocked_count, 1); +} + +async fn assert_superseded_reservation_cannot_settle(ledger: &dyn IdempotencyLedger, suffix: &str) { + let received_at = Utc::now(); + let fingerprint = fingerprint(suffix); + + let IdempotencyDecision::New(mut stale_action) = ledger + .begin_or_replay(fingerprint.clone(), received_at) + .await + .expect("begin") + else { + panic!("expected new action"); + }; + + let IdempotencyDecision::New(mut replacement) = ledger + .begin_or_replay(fingerprint, received_at + Duration::seconds(11)) + .await + .expect("expired reservation should be reclaimed") + else { + panic!("expected reclaimed action"); + }; + + stale_action.settle(ProductInboundAck::NoOp); + let stale_error = ledger + .settle(stale_action) + .await + .expect_err("superseded action must not settle"); + assert!(matches!( + stale_error, + ProductWorkflowError::Transient { .. } + )); + + replacement.settle(ProductInboundAck::NoOp); + ledger + .settle(replacement) + .await + .expect("replacement settle"); +} + +#[cfg(feature = "libsql")] +#[tokio::test] +async fn libsql_settled_action_survives_reopen_and_replays() { + let dir = tempfile::tempdir().expect("tempdir"); + let db_path = dir.path().join("workflow-ledger.db"); + let db_path = db_path.display().to_string(); + let ledger = RebornLibSqlIdempotencyLedger::new(libsql_db(&db_path).await); + let reopened = RebornLibSqlIdempotencyLedger::new(libsql_db(&db_path).await); + + assert_settled_action_survives_reopen_and_replays(&ledger, &reopened, "libsql-settled-replay") + .await; +} + +#[cfg(feature = "libsql")] +#[tokio::test] +async fn libsql_in_flight_action_blocks_until_lease_expires() { + let dir = tempfile::tempdir().expect("tempdir"); + let db_path = dir.path().join("workflow-ledger.db"); + let ledger = RebornLibSqlIdempotencyLedger::with_in_flight_lease( + libsql_db(&db_path.display().to_string()).await, + Duration::seconds(10), + ); + assert_in_flight_action_blocks_until_lease_expires(&ledger, "libsql-lease").await; +} + +#[cfg(feature = "libsql")] +#[tokio::test] +async fn libsql_release_allows_retry_without_waiting_for_lease() { + let dir = tempfile::tempdir().expect("tempdir"); + let db_path = dir.path().join("workflow-ledger.db"); + let ledger = RebornLibSqlIdempotencyLedger::with_in_flight_lease( + libsql_db(&db_path.display().to_string()).await, + Duration::seconds(60), + ); + assert_release_allows_retry_without_waiting_for_lease(&ledger, "libsql-release").await; +} + +#[cfg(feature = "libsql")] +#[tokio::test] +async fn libsql_duplicate_reservation_contention_serializes() { + let dir = tempfile::tempdir().expect("tempdir"); + let db_path = dir.path().join("workflow-ledger.db"); + let db_path = db_path.display().to_string(); + let first = RebornLibSqlIdempotencyLedger::with_in_flight_lease( + libsql_db(&db_path).await, + Duration::seconds(10), + ); + let second = RebornLibSqlIdempotencyLedger::with_in_flight_lease( + libsql_db(&db_path).await, + Duration::seconds(10), + ); + + assert_duplicate_reservation_contention_serializes(&first, &second, "libsql-contention").await; +} + +#[cfg(feature = "libsql")] +#[tokio::test] +async fn libsql_superseded_reservation_cannot_settle() { + let dir = tempfile::tempdir().expect("tempdir"); + let db_path = dir.path().join("workflow-ledger.db"); + let ledger = RebornLibSqlIdempotencyLedger::with_in_flight_lease( + libsql_db(&db_path.display().to_string()).await, + Duration::seconds(10), + ); + + assert_superseded_reservation_cannot_settle(&ledger, "libsql-superseded").await; +} + +#[cfg(feature = "postgres")] +#[tokio::test] +async fn postgres_settled_action_survives_reopen_and_replays_when_configured() { + let Some(pool) = postgres_pool().await else { + return; + }; + let ledger = RebornPostgresIdempotencyLedger::new(pool.clone()); + let reopened = RebornPostgresIdempotencyLedger::new(pool); + + assert_settled_action_survives_reopen_and_replays( + &ledger, + &reopened, + &unique_suffix("postgres-settled-replay"), + ) + .await; +} + +#[cfg(feature = "postgres")] +#[tokio::test] +async fn postgres_in_flight_action_blocks_until_lease_expires_when_configured() { + let Some(pool) = postgres_pool().await else { + return; + }; + let ledger = RebornPostgresIdempotencyLedger::with_in_flight_lease(pool, Duration::seconds(10)); + + assert_in_flight_action_blocks_until_lease_expires(&ledger, &unique_suffix("postgres-lease")) + .await; +} + +#[cfg(feature = "postgres")] +#[tokio::test] +async fn postgres_release_allows_retry_without_waiting_for_lease_when_configured() { + let Some(pool) = postgres_pool().await else { + return; + }; + let ledger = RebornPostgresIdempotencyLedger::with_in_flight_lease(pool, Duration::seconds(60)); + + assert_release_allows_retry_without_waiting_for_lease( + &ledger, + &unique_suffix("postgres-release"), + ) + .await; +} + +#[cfg(feature = "postgres")] +#[tokio::test] +async fn postgres_duplicate_reservation_contention_serializes_when_configured() { + let Some(pool) = postgres_pool().await else { + return; + }; + let first = + RebornPostgresIdempotencyLedger::with_in_flight_lease(pool.clone(), Duration::seconds(10)); + let second = RebornPostgresIdempotencyLedger::with_in_flight_lease(pool, Duration::seconds(10)); + + assert_duplicate_reservation_contention_serializes( + &first, + &second, + &unique_suffix("postgres-contention"), + ) + .await; +} + +#[cfg(feature = "postgres")] +#[tokio::test] +async fn postgres_superseded_reservation_cannot_settle_when_configured() { + let Some(pool) = postgres_pool().await else { + return; + }; + let ledger = RebornPostgresIdempotencyLedger::with_in_flight_lease(pool, Duration::seconds(10)); + + assert_superseded_reservation_cannot_settle(&ledger, &unique_suffix("postgres-superseded")) + .await; +} + +#[cfg(feature = "postgres")] +async fn postgres_pool() -> Option { + let url = match std::env::var("IRONCLAW_PRODUCT_WORKFLOW_POSTGRES_URL") { + Ok(url) => url, + Err(_) => { + eprintln!( + "skipping postgres product workflow ledger contract: IRONCLAW_PRODUCT_WORKFLOW_POSTGRES_URL not set" + ); + return None; + } + }; + let config = match url.parse::() { + Ok(config) => config, + Err(error) => { + eprintln!("skipping postgres product workflow ledger 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() + .expect("postgres pool builds"); + if let Err(error) = pool.get().await { + eprintln!( + "skipping postgres product workflow ledger contract: database unavailable ({error})" + ); + return None; + } + Some(pool) +} From 43c92c5a9cb989c0b39e2c20448a605e015ca4d4 Mon Sep 17 00:00:00 2001 From: serrrfirat Date: Mon, 18 May 2026 23:49:10 +0200 Subject: [PATCH 3/6] fix(product-workflow): address gemini ledger review (#3759) --- .../src/lib.rs | 22 +++++++++++++++---- 1 file changed, 18 insertions(+), 4 deletions(-) diff --git a/crates/ironclaw_product_workflow_storage/src/lib.rs b/crates/ironclaw_product_workflow_storage/src/lib.rs index 22bf55736ab..bbcf650034c 100644 --- a/crates/ironclaw_product_workflow_storage/src/lib.rs +++ b/crates/ironclaw_product_workflow_storage/src/lib.rs @@ -224,13 +224,15 @@ mod libsql_impl { AND source_binding_key = ?3 AND external_event_id = ?4 AND action_id = ?5 - AND phase NOT IN ('settled', 'deduplicated_replay')", + AND phase NOT IN (?6, ?7)", params![ action.fingerprint.adapter_id.as_str(), action.fingerprint.installation_id.as_str(), action.fingerprint.source_binding_key.as_str(), action.fingerprint.external_event_id.as_str(), action.action_id.to_string(), + phase_label(ActionPhase::Settled), + phase_label(ActionPhase::DeduplicatedReplay), ], ) .await @@ -341,7 +343,12 @@ mod libsql_impl { Ok(value) } Err(error) => { - let _ = conn.execute("ROLLBACK", ()).await; + if let Err(rollback_error) = conn.execute("ROLLBACK", ()).await { + tracing::warn!( + %rollback_error, + "product workflow idempotency ledger failed to rollback libSQL transaction" + ); + } Err(error) } } @@ -504,13 +511,15 @@ mod postgres_impl { AND source_binding_key = $3 AND external_event_id = $4 AND action_id = $5 - AND phase NOT IN ('settled', 'deduplicated_replay')", + AND phase NOT IN ($6, $7)", &[ &action.fingerprint.adapter_id.as_str(), &action.fingerprint.installation_id.as_str(), &action.fingerprint.source_binding_key.as_str(), &action.fingerprint.external_event_id.as_str(), &action.action_id.to_string(), + &phase_label(ActionPhase::Settled), + &phase_label(ActionPhase::DeduplicatedReplay), ], ) .await @@ -611,7 +620,12 @@ mod postgres_impl { Ok(value) } Err(error) => { - let _ = txn.rollback().await; + if let Err(rollback_error) = txn.rollback().await { + tracing::warn!( + %rollback_error, + "product workflow idempotency ledger failed to rollback PostgreSQL transaction" + ); + } Err(error) } } From 1715a9710888b447e13101fe0a5bb4c9c443a125 Mon Sep 17 00:00:00 2001 From: serrrfirat Date: Tue, 19 May 2026 00:37:51 +0200 Subject: [PATCH 4/6] refactor(product-workflow): use filesystem SQL ledger storage (#3759) --- Cargo.lock | 2 + .../Cargo.toml | 14 +- .../src/lib.rs | 895 +++++++----------- .../tests/durable_ledger_contract.rs | 73 +- 4 files changed, 392 insertions(+), 592 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 0511d19a03c..cac184b1e1b 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -4589,6 +4589,8 @@ dependencies = [ "async-trait", "chrono", "deadpool-postgres", + "ironclaw_filesystem", + "ironclaw_host_api", "ironclaw_product_adapters", "ironclaw_product_workflow", "libsql", diff --git a/crates/ironclaw_product_workflow_storage/Cargo.toml b/crates/ironclaw_product_workflow_storage/Cargo.toml index 3a3ceb4cacf..cf23838e71b 100644 --- a/crates/ironclaw_product_workflow_storage/Cargo.toml +++ b/crates/ironclaw_product_workflow_storage/Cargo.toml @@ -12,21 +12,23 @@ publish = false [features] default = [] -libsql = ["dep:libsql"] -postgres = ["dep:deadpool-postgres", "dep:tokio-postgres"] +libsql = ["ironclaw_filesystem/libsql"] +postgres = ["ironclaw_filesystem/postgres"] [dependencies] async-trait = "0.1" chrono = { version = "0.4", features = ["serde"] } -deadpool-postgres = { version = "0.14", optional = true } +ironclaw_filesystem = { path = "../ironclaw_filesystem", version = "0.1.0" } +ironclaw_host_api = { path = "../ironclaw_host_api", version = "0.1.0" } ironclaw_product_workflow = { path = "../ironclaw_product_workflow", version = "0.1.0" } -libsql = { version = "0.6", optional = true, default-features = false, features = ["core", "replication", "remote", "tls"] } serde_json = "1" -tokio = { version = "1", features = ["sync"] } -tokio-postgres = { version = "0.7", optional = true } tracing = "0.1" [dev-dependencies] +deadpool-postgres = "0.14" +ironclaw_filesystem = { path = "../ironclaw_filesystem", features = ["libsql", "postgres"] } ironclaw_product_adapters = { path = "../ironclaw_product_adapters" } +libsql = { version = "0.6", default-features = false, features = ["core", "replication", "remote", "tls"] } tempfile = "3" tokio = { version = "1", features = ["macros", "rt", "sync", "time"] } +tokio-postgres = "0.7" diff --git a/crates/ironclaw_product_workflow_storage/src/lib.rs b/crates/ironclaw_product_workflow_storage/src/lib.rs index bbcf650034c..85718dda6f0 100644 --- a/crates/ironclaw_product_workflow_storage/src/lib.rs +++ b/crates/ironclaw_product_workflow_storage/src/lib.rs @@ -1,638 +1,411 @@ //! Durable product workflow [`IdempotencyLedger`] storage adapters. -use chrono::{DateTime, Duration, Utc}; +#![cfg_attr( + not(any(feature = "libsql", feature = "postgres")), + allow(dead_code, unused_imports) +)] + +use std::sync::Arc; +#[cfg(any(feature = "libsql", feature = "postgres"))] +use async_trait::async_trait; +use chrono::{DateTime, Duration, Utc}; +#[cfg(feature = "libsql")] +use ironclaw_filesystem::LibSqlRootFilesystem; +#[cfg(feature = "postgres")] +use ironclaw_filesystem::PostgresRootFilesystem; +use ironclaw_filesystem::{ + CasExpectation, Entry, FilesystemError, IndexKey, IndexValue, RecordKind, RecordVersion, + RootFilesystem, +}; +use ironclaw_host_api::VirtualPath; +#[cfg(any(feature = "libsql", feature = "postgres"))] +use ironclaw_product_workflow::IdempotencyLedger; use ironclaw_product_workflow::{ - ActionFingerprintKey, ActionPhase, IdempotencyDecision, IdempotencyLedger, - ProductInboundAction, ProductWorkflowError, + ActionFingerprintKey, ActionPhase, IdempotencyDecision, ProductInboundAction, + ProductWorkflowError, }; const DEFAULT_IN_FLIGHT_LEASE: Duration = Duration::seconds(60); +const DEFAULT_LEDGER_ROOT: &str = "/engine/product_workflow/idempotency/actions"; +const ACTION_RECORD_KIND: &str = "product_workflow_action"; -const SCHEMA: &str = r#" -CREATE TABLE IF NOT EXISTS reborn_product_workflow_actions ( - adapter_id TEXT NOT NULL, - installation_id TEXT NOT NULL, - source_binding_key TEXT NOT NULL, - external_event_id TEXT NOT NULL, - action_id TEXT NOT NULL, - phase TEXT NOT NULL, - received_at TEXT NOT NULL, - settled_at TEXT, - payload TEXT NOT NULL, - PRIMARY KEY (adapter_id, installation_id, source_binding_key, external_event_id) -); - -CREATE INDEX IF NOT EXISTS idx_reborn_product_workflow_actions_phase - ON reborn_product_workflow_actions(phase, received_at); -"#; - -fn transient(reason: impl Into) -> ProductWorkflowError { - ProductWorkflowError::Transient { - reason: reason.into(), - } +struct FilesystemIdempotencyLedger { + filesystem: Arc, + root: String, + in_flight_lease: Duration, } -fn durable_error(operation: &'static str, error: impl std::fmt::Display) -> ProductWorkflowError { - tracing::error!(%error, operation, "product workflow idempotency ledger backend failed"); - transient(format!("idempotency ledger failed to {operation}")) -} - -fn to_json(action: &ProductInboundAction) -> Result { - serde_json::to_string(action).map_err(|error| durable_error("serialize action", error)) -} - -fn from_json(payload: &str) -> Result { - serde_json::from_str(payload).map_err(|error| durable_error("deserialize action", error)) -} - -fn phase_label(phase: ActionPhase) -> &'static str { - match phase { - ActionPhase::Received => "received", - ActionPhase::Dispatched => "dispatched", - ActionPhase::Settled => "settled", - ActionPhase::DeduplicatedReplay => "deduplicated_replay", +impl FilesystemIdempotencyLedger { + fn new(filesystem: Arc) -> Self { + Self::with_in_flight_lease(filesystem, DEFAULT_IN_FLIGHT_LEASE) } -} - -fn fresh_in_flight( - action: &ProductInboundAction, - received_at: DateTime, - lease: Duration, -) -> bool { - !action.is_terminal() && action.received_at + lease > received_at -} - -fn in_flight_error() -> ProductWorkflowError { - transient("idempotency fingerprint already in flight; retry after recovery lease") -} -#[cfg(feature = "libsql")] -mod libsql_impl { - use std::sync::Arc; - - use async_trait::async_trait; - use libsql::params; - - use super::*; - use tokio::sync::OnceCell; + fn with_in_flight_lease( + filesystem: Arc, + in_flight_lease: Duration, + ) -> Self { + Self { + filesystem, + root: DEFAULT_LEDGER_ROOT.to_string(), + in_flight_lease, + } + } - /// libSQL-backed product workflow idempotency ledger. - pub struct RebornLibSqlIdempotencyLedger { - db: Arc, + fn with_root( + filesystem: Arc, + root: VirtualPath, in_flight_lease: Duration, - migrations: OnceCell<()>, + ) -> Self { + Self { + filesystem, + root: root.as_str().to_string(), + in_flight_lease, + } } - impl RebornLibSqlIdempotencyLedger { - pub fn new(db: Arc) -> Self { - Self::with_in_flight_lease(db, DEFAULT_IN_FLIGHT_LEASE) + async fn begin_or_replay( + &self, + fingerprint: ActionFingerprintKey, + received_at: DateTime, + ) -> Result { + let path = action_path(&self.root, &fingerprint)?; + let action = ProductInboundAction::begin(fingerprint, received_at); + match self + .filesystem + .put(&path, entry_for_action(&action)?, CasExpectation::Absent) + .await + { + Ok(_) => return Ok(IdempotencyDecision::New(action)), + Err(FilesystemError::VersionMismatch { .. }) => {} + Err(error) => return Err(filesystem_error("reserve action", error)), } - pub fn with_in_flight_lease(db: Arc, in_flight_lease: Duration) -> Self { - Self { - db, - in_flight_lease, - migrations: OnceCell::new(), + loop { + let Some((prior, version)) = load_action(self.filesystem.as_ref(), &path).await? else { + return Err(transient("idempotency ledger conflict row disappeared")); + }; + if prior.is_terminal() { + return Ok(IdempotencyDecision::Replay(prior)); + } + if fresh_in_flight(&prior, received_at, self.in_flight_lease) { + return Err(in_flight_error()); } - } - pub async fn run_migrations(&self) -> Result<(), ProductWorkflowError> { - self.migrations - .get_or_try_init(|| async { self.run_migrations_uncached().await }) + let replacement = ProductInboundAction::begin(prior.fingerprint.clone(), received_at); + match self + .filesystem + .put( + &path, + entry_for_action(&replacement)?, + CasExpectation::Version(version), + ) .await - .copied() + { + Ok(_) => return Ok(IdempotencyDecision::New(replacement)), + Err(FilesystemError::VersionMismatch { .. }) => continue, + Err(error) => return Err(filesystem_error("reclaim action", error)), + } } + } - async fn run_migrations_uncached(&self) -> Result<(), ProductWorkflowError> { - let conn = self.connect().await?; - conn.execute_batch(SCHEMA) - .await - .map(|_| ()) - .map_err(|error| durable_error("run migrations", error)) - } + async fn settle(&self, action: ProductInboundAction) -> Result<(), ProductWorkflowError> { + let path = action_path(&self.root, &action.fingerprint)?; + loop { + let Some((current, version)) = load_action(self.filesystem.as_ref(), &path).await? + else { + return Err(transient( + "idempotency reservation missing before terminal settle", + )); + }; + if current.is_terminal() { + if current.action_id == action.action_id { + return Ok(()); + } + return Err(transient( + "idempotency reservation was superseded before terminal settle", + )); + } + if current.action_id != action.action_id { + return Err(transient( + "idempotency reservation was superseded before terminal settle", + )); + } - async fn connect(&self) -> Result { - let conn = self - .db - .connect() - .map_err(|error| durable_error("connect", error))?; - conn.query("PRAGMA busy_timeout = 5000", ()) + match self + .filesystem + .put( + &path, + entry_for_action(&action)?, + CasExpectation::Version(version), + ) .await - .map_err(|error| durable_error("configure busy timeout", error))?; - Ok(conn) + { + Ok(_) => return Ok(()), + Err(FilesystemError::VersionMismatch { .. }) => continue, + Err(error) => return Err(filesystem_error("settle action", error)), + } } + } - async fn begin_immediate(&self) -> Result { - let conn = self.connect().await?; - conn.execute("BEGIN IMMEDIATE", ()) + async fn release(&self, action: ProductInboundAction) -> Result<(), ProductWorkflowError> { + let path = action_path(&self.root, &action.fingerprint)?; + loop { + let Some((current, version)) = load_action(self.filesystem.as_ref(), &path).await? + else { + return Ok(()); + }; + if current.is_terminal() || current.action_id != action.action_id { + return Ok(()); + } + + let mut released = current; + released.received_at = expired_received_at(released.received_at, self.in_flight_lease); + match self + .filesystem + .put( + &path, + entry_for_action(&released)?, + CasExpectation::Version(version), + ) .await - .map_err(|error| durable_error("begin transaction", error))?; - Ok(conn) + { + Ok(_) => return Ok(()), + Err(FilesystemError::VersionMismatch { .. }) => continue, + Err(error) => return Err(filesystem_error("release action", error)), + } } } +} - #[async_trait] - impl IdempotencyLedger for RebornLibSqlIdempotencyLedger { - async fn begin_or_replay( - &self, - fingerprint: ActionFingerprintKey, - received_at: DateTime, - ) -> Result { - self.run_migrations().await?; - let conn = self.begin_immediate().await?; - let result = - begin_or_replay_in_conn(&conn, fingerprint, received_at, self.in_flight_lease) - .await; - finish_transaction(&conn, result).await - } - - async fn settle(&self, action: ProductInboundAction) -> Result<(), ProductWorkflowError> { - self.run_migrations().await?; - let conn = self.begin_immediate().await?; - let result = settle_in_conn(&conn, action).await; - finish_transaction(&conn, result).await - } +/// libSQL-backed product workflow idempotency ledger using the shared +/// SQL filesystem backend for persistence. +#[cfg(feature = "libsql")] +pub struct RebornLibSqlIdempotencyLedger { + inner: FilesystemIdempotencyLedger, +} - async fn release(&self, action: ProductInboundAction) -> Result<(), ProductWorkflowError> { - self.run_migrations().await?; - let conn = self.begin_immediate().await?; - let result = release_in_conn(&conn, action).await; - finish_transaction(&conn, result).await +#[cfg(feature = "libsql")] +impl RebornLibSqlIdempotencyLedger { + pub fn new(filesystem: Arc) -> Self { + Self { + inner: FilesystemIdempotencyLedger::new(filesystem), } } - async fn begin_or_replay_in_conn( - conn: &libsql::Connection, - fingerprint: ActionFingerprintKey, - received_at: DateTime, + pub fn with_in_flight_lease( + filesystem: Arc, in_flight_lease: Duration, - ) -> Result { - let action = ProductInboundAction::begin(fingerprint.clone(), received_at); - let inserted = insert_action(conn, &action, "INSERT OR IGNORE").await?; - if inserted == 1 { - return Ok(IdempotencyDecision::New(action)); + ) -> Self { + Self { + inner: FilesystemIdempotencyLedger::with_in_flight_lease(filesystem, in_flight_lease), } - - let Some(prior) = load_action(conn, &fingerprint).await? else { - return Err(transient("idempotency ledger conflict row disappeared")); - }; - if prior.is_terminal() { - return Ok(IdempotencyDecision::Replay(prior)); - } - if fresh_in_flight(&prior, received_at, in_flight_lease) { - return Err(in_flight_error()); - } - - update_action(conn, &action).await?; - Ok(IdempotencyDecision::New(action)) } - async fn settle_in_conn( - conn: &libsql::Connection, - action: ProductInboundAction, - ) -> Result<(), ProductWorkflowError> { - let Some(current) = load_action(conn, &action.fingerprint).await? else { - return Err(transient( - "idempotency reservation missing before terminal settle", - )); - }; - if current.is_terminal() { - if current.action_id == action.action_id { - return Ok(()); - } - return Err(transient( - "idempotency reservation was superseded before terminal settle", - )); - } - if current.action_id != action.action_id { - return Err(transient( - "idempotency reservation was superseded before terminal settle", - )); + pub fn with_root( + filesystem: Arc, + root: VirtualPath, + in_flight_lease: Duration, + ) -> Self { + Self { + inner: FilesystemIdempotencyLedger::with_root(filesystem, root, in_flight_lease), } - update_action(conn, &action).await - } - - async fn release_in_conn( - conn: &libsql::Connection, - action: ProductInboundAction, - ) -> Result<(), ProductWorkflowError> { - conn.execute( - "DELETE FROM reborn_product_workflow_actions - WHERE adapter_id = ?1 - AND installation_id = ?2 - AND source_binding_key = ?3 - AND external_event_id = ?4 - AND action_id = ?5 - AND phase NOT IN (?6, ?7)", - params![ - action.fingerprint.adapter_id.as_str(), - action.fingerprint.installation_id.as_str(), - action.fingerprint.source_binding_key.as_str(), - action.fingerprint.external_event_id.as_str(), - action.action_id.to_string(), - phase_label(ActionPhase::Settled), - phase_label(ActionPhase::DeduplicatedReplay), - ], - ) - .await - .map_err(|error| durable_error("release action", error))?; - Ok(()) } +} - async fn load_action( - conn: &libsql::Connection, - fingerprint: &ActionFingerprintKey, - ) -> Result, ProductWorkflowError> { - let mut rows = conn - .query( - "SELECT payload FROM reborn_product_workflow_actions - WHERE adapter_id = ?1 - AND installation_id = ?2 - AND source_binding_key = ?3 - AND external_event_id = ?4", - params![ - fingerprint.adapter_id.as_str(), - fingerprint.installation_id.as_str(), - fingerprint.source_binding_key.as_str(), - fingerprint.external_event_id.as_str(), - ], - ) - .await - .map_err(|error| durable_error("load action", error))?; - let Some(row) = rows - .next() - .await - .map_err(|error| durable_error("read action", error))? - else { - return Ok(None); - }; - let payload = row - .get::(0) - .map_err(|error| durable_error("read action payload", error))?; - Ok(Some(from_json(&payload)?)) +#[cfg(feature = "libsql")] +#[async_trait] +impl IdempotencyLedger for RebornLibSqlIdempotencyLedger { + async fn begin_or_replay( + &self, + fingerprint: ActionFingerprintKey, + received_at: DateTime, + ) -> Result { + self.inner.begin_or_replay(fingerprint, received_at).await } - async fn insert_action( - conn: &libsql::Connection, - action: &ProductInboundAction, - insert_prefix: &str, - ) -> Result { - let payload = to_json(action)?; - conn.execute( - &format!( - "{insert_prefix} INTO reborn_product_workflow_actions - (adapter_id, installation_id, source_binding_key, external_event_id, - action_id, phase, received_at, settled_at, payload) - VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9)" - ), - params![ - action.fingerprint.adapter_id.as_str(), - action.fingerprint.installation_id.as_str(), - action.fingerprint.source_binding_key.as_str(), - action.fingerprint.external_event_id.as_str(), - action.action_id.to_string(), - phase_label(action.phase), - action.received_at.to_rfc3339(), - action.settled_at.map(|value| value.to_rfc3339()), - payload, - ], - ) - .await - .map_err(|error| durable_error("insert action", error)) + async fn settle(&self, action: ProductInboundAction) -> Result<(), ProductWorkflowError> { + self.inner.settle(action).await } - async fn update_action( - conn: &libsql::Connection, - action: &ProductInboundAction, - ) -> Result<(), ProductWorkflowError> { - let payload = to_json(action)?; - conn.execute( - "UPDATE reborn_product_workflow_actions - SET action_id = ?5, phase = ?6, received_at = ?7, settled_at = ?8, payload = ?9 - WHERE adapter_id = ?1 - AND installation_id = ?2 - AND source_binding_key = ?3 - AND external_event_id = ?4", - params![ - action.fingerprint.adapter_id.as_str(), - action.fingerprint.installation_id.as_str(), - action.fingerprint.source_binding_key.as_str(), - action.fingerprint.external_event_id.as_str(), - action.action_id.to_string(), - phase_label(action.phase), - action.received_at.to_rfc3339(), - action.settled_at.map(|value| value.to_rfc3339()), - payload, - ], - ) - .await - .map_err(|error| durable_error("update action", error))?; - Ok(()) - } - - async fn finish_transaction( - conn: &libsql::Connection, - result: Result, - ) -> Result { - match result { - Ok(value) => { - conn.execute("COMMIT", ()) - .await - .map_err(|error| durable_error("commit transaction", error))?; - Ok(value) - } - Err(error) => { - if let Err(rollback_error) = conn.execute("ROLLBACK", ()).await { - tracing::warn!( - %rollback_error, - "product workflow idempotency ledger failed to rollback libSQL transaction" - ); - } - Err(error) - } - } + async fn release(&self, action: ProductInboundAction) -> Result<(), ProductWorkflowError> { + self.inner.release(action).await } } +/// PostgreSQL-backed product workflow idempotency ledger using the shared +/// SQL filesystem backend for persistence. #[cfg(feature = "postgres")] -mod postgres_impl { - use async_trait::async_trait; - - use super::*; - use tokio::sync::OnceCell; - - /// PostgreSQL-backed product workflow idempotency ledger. - pub struct RebornPostgresIdempotencyLedger { - pool: deadpool_postgres::Pool, - in_flight_lease: Duration, - migrations: OnceCell<()>, - } - - impl RebornPostgresIdempotencyLedger { - pub fn new(pool: deadpool_postgres::Pool) -> Self { - Self::with_in_flight_lease(pool, DEFAULT_IN_FLIGHT_LEASE) - } - - pub fn with_in_flight_lease( - pool: deadpool_postgres::Pool, - in_flight_lease: Duration, - ) -> Self { - Self { - pool, - in_flight_lease, - migrations: OnceCell::new(), - } - } - - pub async fn run_migrations(&self) -> Result<(), ProductWorkflowError> { - self.migrations - .get_or_try_init(|| async { self.run_migrations_uncached().await }) - .await - .copied() - } - - async fn run_migrations_uncached(&self) -> Result<(), ProductWorkflowError> { - let client = self.client().await?; - client - .batch_execute(SCHEMA) - .await - .map_err(|error| durable_error("run migrations", error)) - } +pub struct RebornPostgresIdempotencyLedger { + inner: FilesystemIdempotencyLedger, +} - async fn client(&self) -> Result { - self.pool - .get() - .await - .map_err(|error| durable_error("connect", error)) +#[cfg(feature = "postgres")] +impl RebornPostgresIdempotencyLedger { + pub fn new(filesystem: Arc) -> Self { + Self { + inner: FilesystemIdempotencyLedger::new(filesystem), } } - #[async_trait] - impl IdempotencyLedger for RebornPostgresIdempotencyLedger { - async fn begin_or_replay( - &self, - fingerprint: ActionFingerprintKey, - received_at: DateTime, - ) -> Result { - self.run_migrations().await?; - let mut client = self.client().await?; - let txn = client - .transaction() - .await - .map_err(|error| durable_error("begin transaction", error))?; - let result = - begin_or_replay_in_txn(&txn, fingerprint, received_at, self.in_flight_lease).await; - finish_transaction(txn, result).await - } - - async fn settle(&self, action: ProductInboundAction) -> Result<(), ProductWorkflowError> { - self.run_migrations().await?; - let mut client = self.client().await?; - let txn = client - .transaction() - .await - .map_err(|error| durable_error("begin transaction", error))?; - let result = settle_in_txn(&txn, action).await; - finish_transaction(txn, result).await + pub fn with_in_flight_lease( + filesystem: Arc, + in_flight_lease: Duration, + ) -> Self { + Self { + inner: FilesystemIdempotencyLedger::with_in_flight_lease(filesystem, in_flight_lease), } + } - async fn release(&self, action: ProductInboundAction) -> Result<(), ProductWorkflowError> { - self.run_migrations().await?; - let mut client = self.client().await?; - let txn = client - .transaction() - .await - .map_err(|error| durable_error("begin transaction", error))?; - let result = release_in_txn(&txn, action).await; - finish_transaction(txn, result).await + pub fn with_root( + filesystem: Arc, + root: VirtualPath, + in_flight_lease: Duration, + ) -> Self { + Self { + inner: FilesystemIdempotencyLedger::with_root(filesystem, root, in_flight_lease), } } +} - async fn begin_or_replay_in_txn( - txn: &deadpool_postgres::Transaction<'_>, +#[cfg(feature = "postgres")] +#[async_trait] +impl IdempotencyLedger for RebornPostgresIdempotencyLedger { + async fn begin_or_replay( + &self, fingerprint: ActionFingerprintKey, received_at: DateTime, - in_flight_lease: Duration, ) -> Result { - let action = ProductInboundAction::begin(fingerprint.clone(), received_at); - let inserted = insert_action(txn, &action).await?; - if inserted == 1 { - return Ok(IdempotencyDecision::New(action)); - } - - let Some(prior) = load_action_for_update(txn, &fingerprint).await? else { - return Err(transient("idempotency ledger conflict row disappeared")); - }; - if prior.is_terminal() { - return Ok(IdempotencyDecision::Replay(prior)); - } - if fresh_in_flight(&prior, received_at, in_flight_lease) { - return Err(in_flight_error()); - } - - update_action(txn, &action).await?; - Ok(IdempotencyDecision::New(action)) + self.inner.begin_or_replay(fingerprint, received_at).await } - async fn settle_in_txn( - txn: &deadpool_postgres::Transaction<'_>, - action: ProductInboundAction, - ) -> Result<(), ProductWorkflowError> { - let Some(current) = load_action_for_update(txn, &action.fingerprint).await? else { - return Err(transient( - "idempotency reservation missing before terminal settle", - )); - }; - if current.is_terminal() { - if current.action_id == action.action_id { - return Ok(()); - } - return Err(transient( - "idempotency reservation was superseded before terminal settle", - )); - } - if current.action_id != action.action_id { - return Err(transient( - "idempotency reservation was superseded before terminal settle", - )); - } - update_action(txn, &action).await + async fn settle(&self, action: ProductInboundAction) -> Result<(), ProductWorkflowError> { + self.inner.settle(action).await } - async fn release_in_txn( - txn: &deadpool_postgres::Transaction<'_>, - action: ProductInboundAction, - ) -> Result<(), ProductWorkflowError> { - txn.execute( - "DELETE FROM reborn_product_workflow_actions - WHERE adapter_id = $1 - AND installation_id = $2 - AND source_binding_key = $3 - AND external_event_id = $4 - AND action_id = $5 - AND phase NOT IN ($6, $7)", - &[ - &action.fingerprint.adapter_id.as_str(), - &action.fingerprint.installation_id.as_str(), - &action.fingerprint.source_binding_key.as_str(), - &action.fingerprint.external_event_id.as_str(), - &action.action_id.to_string(), - &phase_label(ActionPhase::Settled), - &phase_label(ActionPhase::DeduplicatedReplay), - ], - ) - .await - .map_err(|error| durable_error("release action", error))?; - Ok(()) + async fn release(&self, action: ProductInboundAction) -> Result<(), ProductWorkflowError> { + self.inner.release(action).await } +} - async fn load_action_for_update( - txn: &deadpool_postgres::Transaction<'_>, - fingerprint: &ActionFingerprintKey, - ) -> Result, ProductWorkflowError> { - let row = txn - .query_opt( - "SELECT payload FROM reborn_product_workflow_actions - WHERE adapter_id = $1 - AND installation_id = $2 - AND source_binding_key = $3 - AND external_event_id = $4 - FOR UPDATE", - &[ - &fingerprint.adapter_id.as_str(), - &fingerprint.installation_id.as_str(), - &fingerprint.source_binding_key.as_str(), - &fingerprint.external_event_id.as_str(), - ], - ) - .await - .map_err(|error| durable_error("load action", error))?; - row.map(|row| from_json(row.get::<_, &str>(0))).transpose() +fn transient(reason: impl Into) -> ProductWorkflowError { + ProductWorkflowError::Transient { + reason: reason.into(), } +} - async fn insert_action( - txn: &deadpool_postgres::Transaction<'_>, - action: &ProductInboundAction, - ) -> Result { - let payload = to_json(action)?; - txn.execute( - "INSERT INTO reborn_product_workflow_actions - (adapter_id, installation_id, source_binding_key, external_event_id, - action_id, phase, received_at, settled_at, payload) - VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9) - ON CONFLICT (adapter_id, installation_id, source_binding_key, external_event_id) - DO NOTHING", - &[ - &action.fingerprint.adapter_id.as_str(), - &action.fingerprint.installation_id.as_str(), - &action.fingerprint.source_binding_key.as_str(), - &action.fingerprint.external_event_id.as_str(), - &action.action_id.to_string(), - &phase_label(action.phase), - &action.received_at.to_rfc3339(), - &action.settled_at.map(|value| value.to_rfc3339()), - &payload, - ], - ) +fn durable_error(operation: &'static str, error: impl std::fmt::Display) -> ProductWorkflowError { + tracing::error!(%error, operation, "product workflow idempotency ledger failed"); + transient(format!("idempotency ledger failed to {operation}")) +} + +fn filesystem_error(operation: &'static str, error: FilesystemError) -> ProductWorkflowError { + durable_error(operation, error) +} + +fn fresh_in_flight( + action: &ProductInboundAction, + received_at: DateTime, + lease: Duration, +) -> bool { + !action.is_terminal() && action.received_at + lease > received_at +} + +fn in_flight_error() -> ProductWorkflowError { + transient("idempotency fingerprint already in flight; retry after recovery lease") +} + +fn expired_received_at(received_at: DateTime, lease: Duration) -> DateTime { + received_at - lease - Duration::seconds(1) +} + +async fn load_action( + filesystem: &dyn RootFilesystem, + path: &VirtualPath, +) -> Result, ProductWorkflowError> { + let Some(entry) = filesystem + .get(path) .await - .map_err(|error| durable_error("insert action", error)) - } + .map_err(|error| filesystem_error("load action", error))? + else { + return Ok(None); + }; + let action = entry + .entry + .parse_json() + .map_err(|error| durable_error("deserialize action", error))?; + Ok(Some((action, entry.version))) +} - async fn update_action( - txn: &deadpool_postgres::Transaction<'_>, - action: &ProductInboundAction, - ) -> Result<(), ProductWorkflowError> { - let payload = to_json(action)?; - txn.execute( - "UPDATE reborn_product_workflow_actions - SET action_id = $5, phase = $6, received_at = $7, settled_at = $8, payload = $9 - WHERE adapter_id = $1 - AND installation_id = $2 - AND source_binding_key = $3 - AND external_event_id = $4", - &[ - &action.fingerprint.adapter_id.as_str(), - &action.fingerprint.installation_id.as_str(), - &action.fingerprint.source_binding_key.as_str(), - &action.fingerprint.external_event_id.as_str(), - &action.action_id.to_string(), - &phase_label(action.phase), - &action.received_at.to_rfc3339(), - &action.settled_at.map(|value| value.to_rfc3339()), - &payload, - ], +fn entry_for_action(action: &ProductInboundAction) -> Result { + let payload = + serde_json::to_value(action).map_err(|error| durable_error("serialize action", error))?; + let kind = RecordKind::new(ACTION_RECORD_KIND) + .map_err(|error| durable_error("construct action record kind", error))?; + let entry = Entry::record(kind, &payload) + .map_err(|error| durable_error("serialize action entry", error))? + .with_indexed( + index_key("adapter_id")?, + text(action.fingerprint.adapter_id.as_str()), ) - .await - .map_err(|error| durable_error("update action", error))?; - Ok(()) - } + .with_indexed( + index_key("installation_id")?, + text(action.fingerprint.installation_id.as_str()), + ) + .with_indexed( + index_key("source_binding_key")?, + text(action.fingerprint.source_binding_key.as_str()), + ) + .with_indexed( + index_key("external_event_id")?, + text(action.fingerprint.external_event_id.as_str()), + ) + .with_indexed(index_key("phase")?, text(phase_label(action.phase))) + .with_indexed( + index_key("received_at_ms")?, + IndexValue::I64(action.received_at.timestamp_millis()), + ); + Ok(entry) +} - async fn finish_transaction( - txn: deadpool_postgres::Transaction<'_>, - result: Result, - ) -> Result { - match result { - Ok(value) => { - txn.commit() - .await - .map_err(|error| durable_error("commit transaction", error))?; - Ok(value) - } - Err(error) => { - if let Err(rollback_error) = txn.rollback().await { - tracing::warn!( - %rollback_error, - "product workflow idempotency ledger failed to rollback PostgreSQL transaction" - ); - } - Err(error) - } - } +fn index_key(value: &'static str) -> Result { + IndexKey::new(value).map_err(|error| durable_error("construct action index key", error)) +} + +fn text(value: &str) -> IndexValue { + IndexValue::Text(value.to_string()) +} + +fn phase_label(phase: ActionPhase) -> &'static str { + match phase { + ActionPhase::Received => "received", + ActionPhase::Dispatched => "dispatched", + ActionPhase::Settled => "settled", + ActionPhase::DeduplicatedReplay => "deduplicated_replay", } } -#[cfg(feature = "libsql")] -pub use libsql_impl::RebornLibSqlIdempotencyLedger; -#[cfg(feature = "postgres")] -pub use postgres_impl::RebornPostgresIdempotencyLedger; +fn action_path( + root: &str, + fingerprint: &ActionFingerprintKey, +) -> Result { + let path = format!( + "{}/{}/{}/{}/{}.json", + root.trim_end_matches('/'), + hex_component(fingerprint.adapter_id.as_str()), + hex_component(fingerprint.installation_id.as_str()), + hex_component(fingerprint.source_binding_key.as_str()), + hex_component(fingerprint.external_event_id.as_str()) + ); + VirtualPath::new(path).map_err(|error| durable_error("construct action path", error)) +} + +fn hex_component(value: &str) -> String { + const HEX: &[u8; 16] = b"0123456789abcdef"; + let mut encoded = String::with_capacity(value.len() * 2); + for byte in value.as_bytes() { + encoded.push(HEX[(byte >> 4) as usize] as char); + encoded.push(HEX[(byte & 0x0f) as usize] as char); + } + encoded +} diff --git a/crates/ironclaw_product_workflow_storage/tests/durable_ledger_contract.rs b/crates/ironclaw_product_workflow_storage/tests/durable_ledger_contract.rs index ca20e454f2f..5f58fe2cc8a 100644 --- a/crates/ironclaw_product_workflow_storage/tests/durable_ledger_contract.rs +++ b/crates/ironclaw_product_workflow_storage/tests/durable_ledger_contract.rs @@ -6,6 +6,10 @@ use std::sync::Arc; use std::time::{SystemTime, UNIX_EPOCH}; use chrono::{Duration, Utc}; +#[cfg(feature = "libsql")] +use ironclaw_filesystem::LibSqlRootFilesystem; +#[cfg(feature = "postgres")] +use ironclaw_filesystem::PostgresRootFilesystem; use ironclaw_product_adapters::{ AdapterInstallationId, ExternalEventId, ProductAdapterId, ProductInboundAck, }; @@ -38,13 +42,19 @@ fn unique_suffix(name: &str) -> String { } #[cfg(feature = "libsql")] -async fn libsql_db(path: &str) -> Arc { - Arc::new( +async fn libsql_filesystem(path: &str) -> Arc { + let db = Arc::new( libsql::Builder::new_local(path) .build() .await .expect("build libsql db"), - ) + ); + let filesystem = Arc::new(LibSqlRootFilesystem::new(db)); + filesystem + .run_migrations() + .await + .expect("run libsql filesystem migrations"); + filesystem } async fn assert_settled_action_survives_reopen_and_replays( @@ -195,8 +205,8 @@ async fn libsql_settled_action_survives_reopen_and_replays() { let dir = tempfile::tempdir().expect("tempdir"); let db_path = dir.path().join("workflow-ledger.db"); let db_path = db_path.display().to_string(); - let ledger = RebornLibSqlIdempotencyLedger::new(libsql_db(&db_path).await); - let reopened = RebornLibSqlIdempotencyLedger::new(libsql_db(&db_path).await); + let ledger = RebornLibSqlIdempotencyLedger::new(libsql_filesystem(&db_path).await); + let reopened = RebornLibSqlIdempotencyLedger::new(libsql_filesystem(&db_path).await); assert_settled_action_survives_reopen_and_replays(&ledger, &reopened, "libsql-settled-replay") .await; @@ -208,7 +218,7 @@ async fn libsql_in_flight_action_blocks_until_lease_expires() { let dir = tempfile::tempdir().expect("tempdir"); let db_path = dir.path().join("workflow-ledger.db"); let ledger = RebornLibSqlIdempotencyLedger::with_in_flight_lease( - libsql_db(&db_path.display().to_string()).await, + libsql_filesystem(&db_path.display().to_string()).await, Duration::seconds(10), ); assert_in_flight_action_blocks_until_lease_expires(&ledger, "libsql-lease").await; @@ -220,7 +230,7 @@ async fn libsql_release_allows_retry_without_waiting_for_lease() { let dir = tempfile::tempdir().expect("tempdir"); let db_path = dir.path().join("workflow-ledger.db"); let ledger = RebornLibSqlIdempotencyLedger::with_in_flight_lease( - libsql_db(&db_path.display().to_string()).await, + libsql_filesystem(&db_path.display().to_string()).await, Duration::seconds(60), ); assert_release_allows_retry_without_waiting_for_lease(&ledger, "libsql-release").await; @@ -233,11 +243,11 @@ async fn libsql_duplicate_reservation_contention_serializes() { let db_path = dir.path().join("workflow-ledger.db"); let db_path = db_path.display().to_string(); let first = RebornLibSqlIdempotencyLedger::with_in_flight_lease( - libsql_db(&db_path).await, + libsql_filesystem(&db_path).await, Duration::seconds(10), ); let second = RebornLibSqlIdempotencyLedger::with_in_flight_lease( - libsql_db(&db_path).await, + libsql_filesystem(&db_path).await, Duration::seconds(10), ); @@ -250,7 +260,7 @@ async fn libsql_superseded_reservation_cannot_settle() { let dir = tempfile::tempdir().expect("tempdir"); let db_path = dir.path().join("workflow-ledger.db"); let ledger = RebornLibSqlIdempotencyLedger::with_in_flight_lease( - libsql_db(&db_path.display().to_string()).await, + libsql_filesystem(&db_path.display().to_string()).await, Duration::seconds(10), ); @@ -260,11 +270,11 @@ async fn libsql_superseded_reservation_cannot_settle() { #[cfg(feature = "postgres")] #[tokio::test] async fn postgres_settled_action_survives_reopen_and_replays_when_configured() { - let Some(pool) = postgres_pool().await else { + let Some(filesystem) = postgres_filesystem().await else { return; }; - let ledger = RebornPostgresIdempotencyLedger::new(pool.clone()); - let reopened = RebornPostgresIdempotencyLedger::new(pool); + let ledger = RebornPostgresIdempotencyLedger::new(Arc::clone(&filesystem)); + let reopened = RebornPostgresIdempotencyLedger::new(filesystem); assert_settled_action_survives_reopen_and_replays( &ledger, @@ -277,10 +287,11 @@ async fn postgres_settled_action_survives_reopen_and_replays_when_configured() { #[cfg(feature = "postgres")] #[tokio::test] async fn postgres_in_flight_action_blocks_until_lease_expires_when_configured() { - let Some(pool) = postgres_pool().await else { + let Some(filesystem) = postgres_filesystem().await else { return; }; - let ledger = RebornPostgresIdempotencyLedger::with_in_flight_lease(pool, Duration::seconds(10)); + let ledger = + RebornPostgresIdempotencyLedger::with_in_flight_lease(filesystem, Duration::seconds(10)); assert_in_flight_action_blocks_until_lease_expires(&ledger, &unique_suffix("postgres-lease")) .await; @@ -289,10 +300,11 @@ async fn postgres_in_flight_action_blocks_until_lease_expires_when_configured() #[cfg(feature = "postgres")] #[tokio::test] async fn postgres_release_allows_retry_without_waiting_for_lease_when_configured() { - let Some(pool) = postgres_pool().await else { + let Some(filesystem) = postgres_filesystem().await else { return; }; - let ledger = RebornPostgresIdempotencyLedger::with_in_flight_lease(pool, Duration::seconds(60)); + let ledger = + RebornPostgresIdempotencyLedger::with_in_flight_lease(filesystem, Duration::seconds(60)); assert_release_allows_retry_without_waiting_for_lease( &ledger, @@ -304,12 +316,15 @@ async fn postgres_release_allows_retry_without_waiting_for_lease_when_configured #[cfg(feature = "postgres")] #[tokio::test] async fn postgres_duplicate_reservation_contention_serializes_when_configured() { - let Some(pool) = postgres_pool().await else { + let Some(filesystem) = postgres_filesystem().await else { return; }; - let first = - RebornPostgresIdempotencyLedger::with_in_flight_lease(pool.clone(), Duration::seconds(10)); - let second = RebornPostgresIdempotencyLedger::with_in_flight_lease(pool, Duration::seconds(10)); + let first = RebornPostgresIdempotencyLedger::with_in_flight_lease( + Arc::clone(&filesystem), + Duration::seconds(10), + ); + let second = + RebornPostgresIdempotencyLedger::with_in_flight_lease(filesystem, Duration::seconds(10)); assert_duplicate_reservation_contention_serializes( &first, @@ -322,17 +337,18 @@ async fn postgres_duplicate_reservation_contention_serializes_when_configured() #[cfg(feature = "postgres")] #[tokio::test] async fn postgres_superseded_reservation_cannot_settle_when_configured() { - let Some(pool) = postgres_pool().await else { + let Some(filesystem) = postgres_filesystem().await else { return; }; - let ledger = RebornPostgresIdempotencyLedger::with_in_flight_lease(pool, Duration::seconds(10)); + let ledger = + RebornPostgresIdempotencyLedger::with_in_flight_lease(filesystem, Duration::seconds(10)); assert_superseded_reservation_cannot_settle(&ledger, &unique_suffix("postgres-superseded")) .await; } #[cfg(feature = "postgres")] -async fn postgres_pool() -> Option { +async fn postgres_filesystem() -> Option> { let url = match std::env::var("IRONCLAW_PRODUCT_WORKFLOW_POSTGRES_URL") { Ok(url) => url, Err(_) => { @@ -360,5 +376,12 @@ async fn postgres_pool() -> Option { ); return None; } - Some(pool) + let filesystem = Arc::new(PostgresRootFilesystem::new(pool)); + if let Err(error) = filesystem.run_migrations().await { + eprintln!( + "skipping postgres product workflow ledger contract: filesystem migrations failed ({error})" + ); + return None; + } + Some(filesystem) } From b8c2736cccda7c93d8d8793a01942efa3dc10b6d Mon Sep 17 00:00:00 2001 From: serrrfirat Date: Tue, 19 May 2026 10:40:26 +0200 Subject: [PATCH 5/6] =?UTF-8?q?fix(product-workflow-storage):=20address=20?= =?UTF-8?q?serrrfirat=20review=20=E2=80=94=20durable=20ledger=20findings?= =?UTF-8?q?=20(#3759)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- crates/ironclaw_filesystem/src/libsql.rs | 72 +++++++---- crates/ironclaw_filesystem/src/postgres.rs | 78 +++++++----- crates/ironclaw_filesystem/src/root.rs | 16 +++ .../AGENTS.md | 36 ++++++ .../src/lib.rs | 29 +++-- .../tests/durable_ledger_contract.rs | 116 +++++++++++++++++- 6 files changed, 278 insertions(+), 69 deletions(-) create mode 100644 crates/ironclaw_product_workflow_storage/AGENTS.md diff --git a/crates/ironclaw_filesystem/src/libsql.rs b/crates/ironclaw_filesystem/src/libsql.rs index be1e80c3ea0..a397e310800 100644 --- a/crates/ironclaw_filesystem/src/libsql.rs +++ b/crates/ironclaw_filesystem/src/libsql.rs @@ -74,39 +74,13 @@ impl LibSqlRootFilesystem { .map_err(|error| infrastructure_libsql_error(FilesystemOperation::Stat, error))?; Ok(conn) } -} -#[cfg(feature = "libsql")] -#[async_trait] -impl RootFilesystem for LibSqlRootFilesystem { - fn capabilities(&self) -> BackendCapabilities { - // sql_typical covers read/write/append/list/stat/delete/records/query - // /IndexExact/IndexPrefix/CAS. The append/tail backing table is in - // place so Events is on; FTS5 is built into libSQL and a brute-force - // cosine ranker for vectors is implemented in Rust, so IndexFts and - // IndexVector are advertised here too. - BackendCapabilities::sql_typical() - .with(Capability::Events) - .with(Capability::IndexFts) - .with(Capability::IndexVector) - } - - async fn put( + async fn put_without_child_precheck( &self, path: &VirtualPath, entry: Entry, cas: CasExpectation, ) -> Result { - // Reject writes that would clobber a directory or a path that has - // children (mirrors `write_file` semantics so legacy and new ops - // stay consistent). - if matches!( - self.exact_entry(path).await?, - Some((_, FileType::Directory, _)) - ) || self.has_child_entry(path).await? - { - return Err(directory_write_error(path.clone())); - } let indexed_json = serde_json::to_string(&entry.indexed).map_err(|_| { FilesystemError::SerializeIndexed { path: path.clone(), @@ -232,6 +206,50 @@ impl RootFilesystem for LibSqlRootFilesystem { } } } +} + +#[cfg(feature = "libsql")] +#[async_trait] +impl RootFilesystem for LibSqlRootFilesystem { + fn capabilities(&self) -> BackendCapabilities { + // sql_typical covers read/write/append/list/stat/delete/records/query + // /IndexExact/IndexPrefix/CAS. The append/tail backing table is in + // place so Events is on; FTS5 is built into libSQL and a brute-force + // cosine ranker for vectors is implemented in Rust, so IndexFts and + // IndexVector are advertised here too. + BackendCapabilities::sql_typical() + .with(Capability::Events) + .with(Capability::IndexFts) + .with(Capability::IndexVector) + } + + async fn put( + &self, + path: &VirtualPath, + entry: Entry, + cas: CasExpectation, + ) -> Result { + // Reject writes that would clobber a directory or a path that has + // children (mirrors `write_file` semantics so legacy and new ops + // stay consistent). + if matches!( + self.exact_entry(path).await?, + Some((_, FileType::Directory, _)) + ) || self.has_child_entry(path).await? + { + return Err(directory_write_error(path.clone())); + } + self.put_without_child_precheck(path, entry, cas).await + } + + async fn put_known_leaf( + &self, + path: &VirtualPath, + entry: Entry, + cas: CasExpectation, + ) -> Result { + self.put_without_child_precheck(path, entry, cas).await + } async fn get(&self, path: &VirtualPath) -> Result, FilesystemError> { let conn = self.connect().await?; diff --git a/crates/ironclaw_filesystem/src/postgres.rs b/crates/ironclaw_filesystem/src/postgres.rs index 9112ba6c164..825210ed4e1 100644 --- a/crates/ironclaw_filesystem/src/postgres.rs +++ b/crates/ironclaw_filesystem/src/postgres.rs @@ -44,39 +44,14 @@ impl PostgresRootFilesystem { .await .map_err(|error| infrastructure_error(FilesystemOperation::Stat, error.to_string())) } -} - -#[cfg(feature = "postgres")] -#[async_trait] -impl RootFilesystem for PostgresRootFilesystem { - fn capabilities(&self) -> BackendCapabilities { - // sql_typical: read/write/append/list/stat/delete/records/query/ - // IndexExact/IndexPrefix/CAS. Events join the set with the V30 - // append/tail backing table. Postgres has native `tsvector` / - // `plainto_tsquery` so we advertise IndexFts. Vector indexing is - // currently a brute-force cosine ranker against `indexed->>'key'` - // values stored as IndexValue::Bytes; we advertise IndexVector but - // do not require pgvector. - BackendCapabilities::sql_typical() - .with(Capability::Events) - .with(Capability::IndexFts) - .with(Capability::IndexVector) - } - async fn put( + async fn put_with_client( &self, + client: &tokio_postgres::Client, path: &VirtualPath, entry: Entry, cas: CasExpectation, ) -> Result { - let client = self.client().await?; - if matches!( - self.exact_entry_with_client(&client, path).await?, - Some((_, FileType::Directory, _)) - ) || self.has_child_entry_with_client(&client, path).await? - { - return Err(directory_write_error(path.clone())); - } let indexed_json = serde_json::to_value(&entry.indexed).map_err(|_| { FilesystemError::SerializeIndexed { path: path.clone(), @@ -111,7 +86,7 @@ impl RootFilesystem for PostgresRootFilesystem { db_error(path.clone(), FilesystemOperation::WriteFile, error) })?; if rows == 0 { - let found = self.current_version_with_client(&client, path).await?; + let found = self.current_version_with_client(client, path).await?; return Err(FilesystemError::VersionMismatch { path: path.clone(), expected: None, @@ -148,7 +123,7 @@ impl RootFilesystem for PostgresRootFilesystem { db_error(path.clone(), FilesystemOperation::WriteFile, error) })?; if rows == 0 { - let found = self.current_version_with_client(&client, path).await?; + let found = self.current_version_with_client(client, path).await?; return Err(FilesystemError::VersionMismatch { path: path.clone(), expected: Some(expected), @@ -194,6 +169,51 @@ impl RootFilesystem for PostgresRootFilesystem { } } } +} + +#[cfg(feature = "postgres")] +#[async_trait] +impl RootFilesystem for PostgresRootFilesystem { + fn capabilities(&self) -> BackendCapabilities { + // sql_typical: read/write/append/list/stat/delete/records/query/ + // IndexExact/IndexPrefix/CAS. Events join the set with the V30 + // append/tail backing table. Postgres has native `tsvector` / + // `plainto_tsquery` so we advertise IndexFts. Vector indexing is + // currently a brute-force cosine ranker against `indexed->>'key'` + // values stored as IndexValue::Bytes; we advertise IndexVector but + // do not require pgvector. + BackendCapabilities::sql_typical() + .with(Capability::Events) + .with(Capability::IndexFts) + .with(Capability::IndexVector) + } + + async fn put( + &self, + path: &VirtualPath, + entry: Entry, + cas: CasExpectation, + ) -> Result { + let client = self.client().await?; + if matches!( + self.exact_entry_with_client(&client, path).await?, + Some((_, FileType::Directory, _)) + ) || self.has_child_entry_with_client(&client, path).await? + { + return Err(directory_write_error(path.clone())); + } + self.put_with_client(&client, path, entry, cas).await + } + + async fn put_known_leaf( + &self, + path: &VirtualPath, + entry: Entry, + cas: CasExpectation, + ) -> Result { + let client = self.client().await?; + self.put_with_client(&client, path, entry, cas).await + } async fn get(&self, path: &VirtualPath) -> Result, FilesystemError> { let client = self.client().await?; diff --git a/crates/ironclaw_filesystem/src/root.rs b/crates/ironclaw_filesystem/src/root.rs index 4e47fd548c6..4931f76584a 100644 --- a/crates/ironclaw_filesystem/src/root.rs +++ b/crates/ironclaw_filesystem/src/root.rs @@ -66,6 +66,22 @@ pub trait RootFilesystem: Send + Sync { unsupported(path, FilesystemOperation::WriteFile) } + /// Write a caller-owned leaf record with a compare-and-swap precondition. + /// + /// This is an optimization hook for stores whose path grammar guarantees + /// that the target path is always a file-like leaf. The default preserves + /// full [`put`](Self::put) semantics. Backends may override it to skip + /// expensive child-directory probes while still enforcing the exact-path + /// CAS and never clobbering a directory entry at `path`. + async fn put_known_leaf( + &self, + path: &VirtualPath, + entry: Entry, + cas: CasExpectation, + ) -> Result { + self.put(path, entry, cas).await + } + /// Read the entry at `path`, returning `None` if no entry is present. /// /// Default impl is `Unsupported`. Same recursion concern as `put`: diff --git a/crates/ironclaw_product_workflow_storage/AGENTS.md b/crates/ironclaw_product_workflow_storage/AGENTS.md new file mode 100644 index 00000000000..691c0c19c78 --- /dev/null +++ b/crates/ironclaw_product_workflow_storage/AGENTS.md @@ -0,0 +1,36 @@ +# ironclaw_product_workflow_storage + +Durable storage adapters for the product workflow idempotency ledger. + +## Purpose + +- Provide libSQL and PostgreSQL-backed implementations of + `ironclaw_product_workflow::IdempotencyLedger`. +- Persist product inbound action reservations and terminal outcomes through + `ironclaw_filesystem::RootFilesystem`. +- Preserve recovery-lease behavior for non-terminal reservations so retries do + not dispatch the same side effect concurrently. + +## Boundaries + +- This crate owns storage adapters only. Product workflow orchestration remains + in `ironclaw_product_workflow`. +- Keep durable records behind the existing `IdempotencyLedger` port; do not add + product workflow call paths around that trait. +- Use typed host and workflow values internally. Convert strings at boundaries. +- Keep libSQL and PostgreSQL behavior in parity when changing persistence + semantics. + +## Validation + +Run targeted checks from the workspace root: + +```bash +cargo test -p ironclaw_product_workflow_storage --features libsql +cargo test -p ironclaw_product_workflow_storage --features postgres --no-run +cargo check -p ironclaw_product_workflow_storage --features "libsql postgres" +cargo clippy -p ironclaw_product_workflow_storage --all-targets --features "libsql postgres" -- -D warnings +``` + +PostgreSQL runtime tests require `IRONCLAW_PRODUCT_WORKFLOW_POSTGRES_URL`; when +it is unset, postgres contract tests compile and skip execution. diff --git a/crates/ironclaw_product_workflow_storage/src/lib.rs b/crates/ironclaw_product_workflow_storage/src/lib.rs index 85718dda6f0..aecdfd8fd40 100644 --- a/crates/ironclaw_product_workflow_storage/src/lib.rs +++ b/crates/ironclaw_product_workflow_storage/src/lib.rs @@ -32,7 +32,7 @@ const ACTION_RECORD_KIND: &str = "product_workflow_action"; struct FilesystemIdempotencyLedger { filesystem: Arc, - root: String, + root: VirtualPath, in_flight_lease: Duration, } @@ -47,7 +47,7 @@ impl FilesystemIdempotencyLedger { ) -> Self { Self { filesystem, - root: DEFAULT_LEDGER_ROOT.to_string(), + root: default_ledger_root(), in_flight_lease, } } @@ -59,7 +59,7 @@ impl FilesystemIdempotencyLedger { ) -> Self { Self { filesystem, - root: root.as_str().to_string(), + root, in_flight_lease, } } @@ -73,7 +73,7 @@ impl FilesystemIdempotencyLedger { let action = ProductInboundAction::begin(fingerprint, received_at); match self .filesystem - .put(&path, entry_for_action(&action)?, CasExpectation::Absent) + .put_known_leaf(&path, entry_for_action(&action)?, CasExpectation::Absent) .await { Ok(_) => return Ok(IdempotencyDecision::New(action)), @@ -95,7 +95,7 @@ impl FilesystemIdempotencyLedger { let replacement = ProductInboundAction::begin(prior.fingerprint.clone(), received_at); match self .filesystem - .put( + .put_known_leaf( &path, entry_for_action(&replacement)?, CasExpectation::Version(version), @@ -134,7 +134,7 @@ impl FilesystemIdempotencyLedger { match self .filesystem - .put( + .put_known_leaf( &path, entry_for_action(&action)?, CasExpectation::Version(version), @@ -163,7 +163,7 @@ impl FilesystemIdempotencyLedger { released.received_at = expired_received_at(released.received_at, self.in_flight_lease); match self .filesystem - .put( + .put_known_leaf( &path, entry_for_action(&released)?, CasExpectation::Version(version), @@ -295,7 +295,12 @@ fn transient(reason: impl Into) -> ProductWorkflowError { } fn durable_error(operation: &'static str, error: impl std::fmt::Display) -> ProductWorkflowError { - tracing::error!(%error, operation, "product workflow idempotency ledger failed"); + let error_type = std::any::type_name_of_val(&error); + tracing::error!( + operation, + error_type, + "product workflow idempotency ledger failed" + ); transient(format!("idempotency ledger failed to {operation}")) } @@ -386,12 +391,12 @@ fn phase_label(phase: ActionPhase) -> &'static str { } fn action_path( - root: &str, + root: &VirtualPath, fingerprint: &ActionFingerprintKey, ) -> Result { let path = format!( "{}/{}/{}/{}/{}.json", - root.trim_end_matches('/'), + root.as_str().trim_end_matches('/'), hex_component(fingerprint.adapter_id.as_str()), hex_component(fingerprint.installation_id.as_str()), hex_component(fingerprint.source_binding_key.as_str()), @@ -400,6 +405,10 @@ fn action_path( VirtualPath::new(path).map_err(|error| durable_error("construct action path", error)) } +fn default_ledger_root() -> VirtualPath { + VirtualPath::new(DEFAULT_LEDGER_ROOT).expect("DEFAULT_LEDGER_ROOT is valid") // safety: hard-coded /engine virtual path literal. +} + fn hex_component(value: &str) -> String { const HEX: &[u8; 16] = b"0123456789abcdef"; let mut encoded = String::with_capacity(value.len() * 2); diff --git a/crates/ironclaw_product_workflow_storage/tests/durable_ledger_contract.rs b/crates/ironclaw_product_workflow_storage/tests/durable_ledger_contract.rs index 5f58fe2cc8a..da7b5591abd 100644 --- a/crates/ironclaw_product_workflow_storage/tests/durable_ledger_contract.rs +++ b/crates/ironclaw_product_workflow_storage/tests/durable_ledger_contract.rs @@ -1,6 +1,6 @@ #![cfg(any(feature = "libsql", feature = "postgres"))] -#[cfg(feature = "libsql")] +#[cfg(any(feature = "libsql", feature = "postgres"))] use std::sync::Arc; #[cfg(feature = "postgres")] use std::time::{SystemTime, UNIX_EPOCH}; @@ -10,12 +10,13 @@ use chrono::{Duration, Utc}; use ironclaw_filesystem::LibSqlRootFilesystem; #[cfg(feature = "postgres")] use ironclaw_filesystem::PostgresRootFilesystem; +use ironclaw_host_api::VirtualPath; use ironclaw_product_adapters::{ AdapterInstallationId, ExternalEventId, ProductAdapterId, ProductInboundAck, }; use ironclaw_product_workflow::{ - ActionFingerprintKey, IdempotencyDecision, IdempotencyLedger, ProductWorkflowError, - SourceBindingKey, + ActionFingerprintKey, IdempotencyDecision, IdempotencyLedger, ProductInboundAction, + ProductWorkflowError, SourceBindingKey, }; #[cfg(feature = "libsql")] use ironclaw_product_workflow_storage::RebornLibSqlIdempotencyLedger; @@ -32,6 +33,13 @@ fn fingerprint(suffix: &str) -> ActionFingerprintKey { ) } +fn custom_root(suffix: &str) -> VirtualPath { + VirtualPath::new(format!( + "/engine/product_workflow/idempotency/test_roots/{suffix}" + )) + .expect("valid custom ledger root") +} + #[cfg(feature = "postgres")] fn unique_suffix(name: &str) -> String { let nanos = SystemTime::now() @@ -199,6 +207,45 @@ async fn assert_superseded_reservation_cannot_settle(ledger: &dyn IdempotencyLed .expect("replacement settle"); } +async fn assert_settle_missing_reservation_returns_transient( + ledger: &dyn IdempotencyLedger, + suffix: &str, +) { + let received_at = Utc::now(); + let mut action = ProductInboundAction::begin(fingerprint(suffix), received_at); + action.settle(ProductInboundAck::NoOp); + + let error = ledger + .settle(action) + .await + .expect_err("missing reservation must not settle"); + assert!(matches!(error, ProductWorkflowError::Transient { .. })); +} + +async fn assert_custom_root_isolated_from_default_root( + custom: &dyn IdempotencyLedger, + default: &dyn IdempotencyLedger, + suffix: &str, +) { + let received_at = Utc::now(); + let fingerprint = fingerprint(suffix); + let IdempotencyDecision::New(mut action) = custom + .begin_or_replay(fingerprint.clone(), received_at) + .await + .expect("begin in custom root") + else { + panic!("expected new custom-root action"); + }; + action.settle(ProductInboundAck::NoOp); + custom.settle(action).await.expect("settle custom root"); + + let default_decision = default + .begin_or_replay(fingerprint, received_at + Duration::seconds(1)) + .await + .expect("begin in default root"); + assert!(matches!(default_decision, IdempotencyDecision::New(_))); +} + #[cfg(feature = "libsql")] #[tokio::test] async fn libsql_settled_action_survives_reopen_and_replays() { @@ -267,6 +314,33 @@ async fn libsql_superseded_reservation_cannot_settle() { assert_superseded_reservation_cannot_settle(&ledger, "libsql-superseded").await; } +#[cfg(feature = "libsql")] +#[tokio::test] +async fn libsql_settle_missing_reservation_returns_transient() { + let dir = tempfile::tempdir().expect("tempdir"); + let db_path = dir.path().join("workflow-ledger.db"); + let ledger = + RebornLibSqlIdempotencyLedger::new(libsql_filesystem(&db_path.display().to_string()).await); + + assert_settle_missing_reservation_returns_transient(&ledger, "libsql-missing-settle").await; +} + +#[cfg(feature = "libsql")] +#[tokio::test] +async fn libsql_custom_root_isolated_from_default_root() { + let dir = tempfile::tempdir().expect("tempdir"); + let db_path = dir.path().join("workflow-ledger.db"); + let filesystem = libsql_filesystem(&db_path.display().to_string()).await; + let custom = RebornLibSqlIdempotencyLedger::with_root( + Arc::clone(&filesystem), + custom_root("libsql"), + Duration::seconds(60), + ); + let default = RebornLibSqlIdempotencyLedger::new(filesystem); + + assert_custom_root_isolated_from_default_root(&custom, &default, "libsql-custom-root").await; +} + #[cfg(feature = "postgres")] #[tokio::test] async fn postgres_settled_action_survives_reopen_and_replays_when_configured() { @@ -347,6 +421,42 @@ async fn postgres_superseded_reservation_cannot_settle_when_configured() { .await; } +#[cfg(feature = "postgres")] +#[tokio::test] +async fn postgres_settle_missing_reservation_returns_transient_when_configured() { + let Some(filesystem) = postgres_filesystem().await else { + return; + }; + let ledger = RebornPostgresIdempotencyLedger::new(filesystem); + + assert_settle_missing_reservation_returns_transient( + &ledger, + &unique_suffix("postgres-missing-settle"), + ) + .await; +} + +#[cfg(feature = "postgres")] +#[tokio::test] +async fn postgres_custom_root_isolated_from_default_root_when_configured() { + let Some(filesystem) = postgres_filesystem().await else { + return; + }; + let custom = RebornPostgresIdempotencyLedger::with_root( + Arc::clone(&filesystem), + custom_root("postgres"), + Duration::seconds(60), + ); + let default = RebornPostgresIdempotencyLedger::new(filesystem); + + assert_custom_root_isolated_from_default_root( + &custom, + &default, + &unique_suffix("postgres-custom-root"), + ) + .await; +} + #[cfg(feature = "postgres")] async fn postgres_filesystem() -> Option> { let url = match std::env::var("IRONCLAW_PRODUCT_WORKFLOW_POSTGRES_URL") { From 8c97598f3885b440cb7a6ebf06134448a40d400a Mon Sep 17 00:00:00 2001 From: serrrfirat Date: Tue, 19 May 2026 10:49:16 +0200 Subject: [PATCH 6/6] fix(product-workflow-storage): keep durable ledger fixes scoped (#3759) --- crates/ironclaw_filesystem/src/libsql.rs | 72 +++++++---------- crates/ironclaw_filesystem/src/postgres.rs | 78 +++++++------------ crates/ironclaw_filesystem/src/root.rs | 16 ---- .../src/lib.rs | 8 +- 4 files changed, 60 insertions(+), 114 deletions(-) diff --git a/crates/ironclaw_filesystem/src/libsql.rs b/crates/ironclaw_filesystem/src/libsql.rs index a397e310800..be1e80c3ea0 100644 --- a/crates/ironclaw_filesystem/src/libsql.rs +++ b/crates/ironclaw_filesystem/src/libsql.rs @@ -74,13 +74,39 @@ impl LibSqlRootFilesystem { .map_err(|error| infrastructure_libsql_error(FilesystemOperation::Stat, error))?; Ok(conn) } +} - async fn put_without_child_precheck( +#[cfg(feature = "libsql")] +#[async_trait] +impl RootFilesystem for LibSqlRootFilesystem { + fn capabilities(&self) -> BackendCapabilities { + // sql_typical covers read/write/append/list/stat/delete/records/query + // /IndexExact/IndexPrefix/CAS. The append/tail backing table is in + // place so Events is on; FTS5 is built into libSQL and a brute-force + // cosine ranker for vectors is implemented in Rust, so IndexFts and + // IndexVector are advertised here too. + BackendCapabilities::sql_typical() + .with(Capability::Events) + .with(Capability::IndexFts) + .with(Capability::IndexVector) + } + + async fn put( &self, path: &VirtualPath, entry: Entry, cas: CasExpectation, ) -> Result { + // Reject writes that would clobber a directory or a path that has + // children (mirrors `write_file` semantics so legacy and new ops + // stay consistent). + if matches!( + self.exact_entry(path).await?, + Some((_, FileType::Directory, _)) + ) || self.has_child_entry(path).await? + { + return Err(directory_write_error(path.clone())); + } let indexed_json = serde_json::to_string(&entry.indexed).map_err(|_| { FilesystemError::SerializeIndexed { path: path.clone(), @@ -206,50 +232,6 @@ impl LibSqlRootFilesystem { } } } -} - -#[cfg(feature = "libsql")] -#[async_trait] -impl RootFilesystem for LibSqlRootFilesystem { - fn capabilities(&self) -> BackendCapabilities { - // sql_typical covers read/write/append/list/stat/delete/records/query - // /IndexExact/IndexPrefix/CAS. The append/tail backing table is in - // place so Events is on; FTS5 is built into libSQL and a brute-force - // cosine ranker for vectors is implemented in Rust, so IndexFts and - // IndexVector are advertised here too. - BackendCapabilities::sql_typical() - .with(Capability::Events) - .with(Capability::IndexFts) - .with(Capability::IndexVector) - } - - async fn put( - &self, - path: &VirtualPath, - entry: Entry, - cas: CasExpectation, - ) -> Result { - // Reject writes that would clobber a directory or a path that has - // children (mirrors `write_file` semantics so legacy and new ops - // stay consistent). - if matches!( - self.exact_entry(path).await?, - Some((_, FileType::Directory, _)) - ) || self.has_child_entry(path).await? - { - return Err(directory_write_error(path.clone())); - } - self.put_without_child_precheck(path, entry, cas).await - } - - async fn put_known_leaf( - &self, - path: &VirtualPath, - entry: Entry, - cas: CasExpectation, - ) -> Result { - self.put_without_child_precheck(path, entry, cas).await - } async fn get(&self, path: &VirtualPath) -> Result, FilesystemError> { let conn = self.connect().await?; diff --git a/crates/ironclaw_filesystem/src/postgres.rs b/crates/ironclaw_filesystem/src/postgres.rs index 825210ed4e1..9112ba6c164 100644 --- a/crates/ironclaw_filesystem/src/postgres.rs +++ b/crates/ironclaw_filesystem/src/postgres.rs @@ -44,14 +44,39 @@ impl PostgresRootFilesystem { .await .map_err(|error| infrastructure_error(FilesystemOperation::Stat, error.to_string())) } +} + +#[cfg(feature = "postgres")] +#[async_trait] +impl RootFilesystem for PostgresRootFilesystem { + fn capabilities(&self) -> BackendCapabilities { + // sql_typical: read/write/append/list/stat/delete/records/query/ + // IndexExact/IndexPrefix/CAS. Events join the set with the V30 + // append/tail backing table. Postgres has native `tsvector` / + // `plainto_tsquery` so we advertise IndexFts. Vector indexing is + // currently a brute-force cosine ranker against `indexed->>'key'` + // values stored as IndexValue::Bytes; we advertise IndexVector but + // do not require pgvector. + BackendCapabilities::sql_typical() + .with(Capability::Events) + .with(Capability::IndexFts) + .with(Capability::IndexVector) + } - async fn put_with_client( + async fn put( &self, - client: &tokio_postgres::Client, path: &VirtualPath, entry: Entry, cas: CasExpectation, ) -> Result { + let client = self.client().await?; + if matches!( + self.exact_entry_with_client(&client, path).await?, + Some((_, FileType::Directory, _)) + ) || self.has_child_entry_with_client(&client, path).await? + { + return Err(directory_write_error(path.clone())); + } let indexed_json = serde_json::to_value(&entry.indexed).map_err(|_| { FilesystemError::SerializeIndexed { path: path.clone(), @@ -86,7 +111,7 @@ impl PostgresRootFilesystem { db_error(path.clone(), FilesystemOperation::WriteFile, error) })?; if rows == 0 { - let found = self.current_version_with_client(client, path).await?; + let found = self.current_version_with_client(&client, path).await?; return Err(FilesystemError::VersionMismatch { path: path.clone(), expected: None, @@ -123,7 +148,7 @@ impl PostgresRootFilesystem { db_error(path.clone(), FilesystemOperation::WriteFile, error) })?; if rows == 0 { - let found = self.current_version_with_client(client, path).await?; + let found = self.current_version_with_client(&client, path).await?; return Err(FilesystemError::VersionMismatch { path: path.clone(), expected: Some(expected), @@ -169,51 +194,6 @@ impl PostgresRootFilesystem { } } } -} - -#[cfg(feature = "postgres")] -#[async_trait] -impl RootFilesystem for PostgresRootFilesystem { - fn capabilities(&self) -> BackendCapabilities { - // sql_typical: read/write/append/list/stat/delete/records/query/ - // IndexExact/IndexPrefix/CAS. Events join the set with the V30 - // append/tail backing table. Postgres has native `tsvector` / - // `plainto_tsquery` so we advertise IndexFts. Vector indexing is - // currently a brute-force cosine ranker against `indexed->>'key'` - // values stored as IndexValue::Bytes; we advertise IndexVector but - // do not require pgvector. - BackendCapabilities::sql_typical() - .with(Capability::Events) - .with(Capability::IndexFts) - .with(Capability::IndexVector) - } - - async fn put( - &self, - path: &VirtualPath, - entry: Entry, - cas: CasExpectation, - ) -> Result { - let client = self.client().await?; - if matches!( - self.exact_entry_with_client(&client, path).await?, - Some((_, FileType::Directory, _)) - ) || self.has_child_entry_with_client(&client, path).await? - { - return Err(directory_write_error(path.clone())); - } - self.put_with_client(&client, path, entry, cas).await - } - - async fn put_known_leaf( - &self, - path: &VirtualPath, - entry: Entry, - cas: CasExpectation, - ) -> Result { - let client = self.client().await?; - self.put_with_client(&client, path, entry, cas).await - } async fn get(&self, path: &VirtualPath) -> Result, FilesystemError> { let client = self.client().await?; diff --git a/crates/ironclaw_filesystem/src/root.rs b/crates/ironclaw_filesystem/src/root.rs index 4931f76584a..4e47fd548c6 100644 --- a/crates/ironclaw_filesystem/src/root.rs +++ b/crates/ironclaw_filesystem/src/root.rs @@ -66,22 +66,6 @@ pub trait RootFilesystem: Send + Sync { unsupported(path, FilesystemOperation::WriteFile) } - /// Write a caller-owned leaf record with a compare-and-swap precondition. - /// - /// This is an optimization hook for stores whose path grammar guarantees - /// that the target path is always a file-like leaf. The default preserves - /// full [`put`](Self::put) semantics. Backends may override it to skip - /// expensive child-directory probes while still enforcing the exact-path - /// CAS and never clobbering a directory entry at `path`. - async fn put_known_leaf( - &self, - path: &VirtualPath, - entry: Entry, - cas: CasExpectation, - ) -> Result { - self.put(path, entry, cas).await - } - /// Read the entry at `path`, returning `None` if no entry is present. /// /// Default impl is `Unsupported`. Same recursion concern as `put`: diff --git a/crates/ironclaw_product_workflow_storage/src/lib.rs b/crates/ironclaw_product_workflow_storage/src/lib.rs index aecdfd8fd40..2307a920b4e 100644 --- a/crates/ironclaw_product_workflow_storage/src/lib.rs +++ b/crates/ironclaw_product_workflow_storage/src/lib.rs @@ -73,7 +73,7 @@ impl FilesystemIdempotencyLedger { let action = ProductInboundAction::begin(fingerprint, received_at); match self .filesystem - .put_known_leaf(&path, entry_for_action(&action)?, CasExpectation::Absent) + .put(&path, entry_for_action(&action)?, CasExpectation::Absent) .await { Ok(_) => return Ok(IdempotencyDecision::New(action)), @@ -95,7 +95,7 @@ impl FilesystemIdempotencyLedger { let replacement = ProductInboundAction::begin(prior.fingerprint.clone(), received_at); match self .filesystem - .put_known_leaf( + .put( &path, entry_for_action(&replacement)?, CasExpectation::Version(version), @@ -134,7 +134,7 @@ impl FilesystemIdempotencyLedger { match self .filesystem - .put_known_leaf( + .put( &path, entry_for_action(&action)?, CasExpectation::Version(version), @@ -163,7 +163,7 @@ impl FilesystemIdempotencyLedger { released.received_at = expired_received_at(released.received_at, self.in_flight_lease); match self .filesystem - .put_known_leaf( + .put( &path, entry_for_action(&released)?, CasExpectation::Version(version),