Repository navigation
feat: (product-workflow) Add durable product workflow ledger - #3759
Conversation
There was a problem hiding this comment.
Code Review
This pull request introduces durable IdempotencyLedger implementations for libSQL and PostgreSQL to the ironclaw_product_workflow crate, along with associated integration tests and dependencies. The feedback highlights opportunities to optimize performance by ensuring migrations run only once per instance, improve maintainability by using helper functions instead of hardcoded strings in SQL queries, and enhance observability by logging failures during transaction rollbacks.
| } | ||
| } | ||
|
|
||
| pub async fn run_migrations(&self) -> Result<(), ProductWorkflowError> { |
There was a problem hiding this comment.
Running migrations on every ledger operation introduces significant overhead, as it requires a roundtrip to the database to check the schema for every begin_or_replay, settle, and release call. Consider running migrations once during application startup or using a synchronization primitive like tokio::sync::OnceCell within the struct to ensure they only run once per instance lifetime.
| 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(), | ||
| ], | ||
| ) |
There was a problem hiding this comment.
The list of terminal phases is hardcoded as strings in the SQL query. This is fragile and may lead to data loss or incorrect behavior if new terminal phases are added to the ActionPhase enum but not updated here. It is safer to construct the query using the phase_label function to ensure consistency with the rest of the implementation.
conn.execute(
&format!(
"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 ('{}', '{}')",
phase_label(ActionPhase::Settled),
phase_label(ActionPhase::DeduplicatedReplay)
),
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(),
],
)| let _ = conn.execute("ROLLBACK", ()).await; | ||
| Err(error) | ||
| } | ||
| } |
There was a problem hiding this comment.
When a transaction operation fails and a rollback is attempted, any failure during the rollback itself should be logged as a warning. This prevents secondary failures from being completely silent while still ensuring the original primary error is returned to the caller.
Err(error) => {
if let Err(rollback_error) = conn.execute("ROLLBACK", ()).await {
tracing::warn!("idempotency ledger failed to rollback transaction: {rollback_error}");
}
Err(error)
}References
- When a primary operation fails and a subsequent cleanup or secondary operation also fails, log the secondary failure as a warning and ensure the original primary error is returned to the caller to prevent it from being swallowed.
| } | ||
| } | ||
|
|
||
| pub async fn run_migrations(&self) -> 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(), | ||
| ], | ||
| ) |
There was a problem hiding this comment.
The terminal phases are hardcoded in the SQL query, which poses a maintenance risk. Construction of the query should ideally use the phase_label function to stay in sync with the enum definitions.
txn.execute(
&format!(
"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 ('{}', '{}')",
phase_label(ActionPhase::Settled),
phase_label(ActionPhase::DeduplicatedReplay)
),
&[
&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(),
],
)| let _ = txn.rollback().await; | ||
| Err(error) | ||
| } | ||
| } |
There was a problem hiding this comment.
Rollback failures should be logged as warnings to ensure visibility into secondary cleanup errors, while preserving the original error for the caller.
Err(error) => {
if let Err(rollback_error) = txn.rollback().await {
tracing::warn!("idempotency ledger failed to rollback transaction: {rollback_error}");
}
Err(error)
}References
- When a primary operation fails and a subsequent cleanup or secondary operation also fails, log the secondary failure as a warning and ensure the original primary error is returned to the caller to prevent it from being swallowed.
serrrfirat
left a comment
There was a problem hiding this comment.
Multi-agent code review for PR #3759 (--force used because the PR is draft).
Summary:
- Security: no findings
- Bugs: no findings
- Performance/Concurrency: 1 finding
- Tests: 3 findings
- Conventions: 1 finding
I would request changes for these findings, but GitHub does not allow requesting changes on my own PR, so this is posted as a review comment.
|
|
||
| const DEFAULT_IN_FLIGHT_LEASE: Duration = Duration::seconds(60); | ||
|
|
||
| const SCHEMA: &str = r#" |
There was a problem hiding this comment.
[Medium][conventions/boundary] ironclaw_product_workflow/AGENTS.md says storage backend details must not move into this crate. This new module embeds SQL schema and libSQL/Postgres implementations directly in product_workflow. Move the concrete storage implementation to the storage/DB-owning layer or a storage adapter crate, and keep this crate behind the existing IdempotencyLedger port.
| use super::*; | ||
| use crate::IdempotencyLedger; | ||
|
|
||
| /// PostgreSQL-backed product workflow idempotency ledger. |
There was a problem hiding this comment.
[High][tests] The PostgreSQL-backed ledger has its own SQL, transaction, FOR UPDATE, lease, settle, and release paths, but the contract test file is cfg(feature = "libsql") only. Add PostgreSQL parity tests covering settled replay across reopen, in-flight lease block/reclaim, and release-before-retry.
|
|
||
| #[async_trait] | ||
| impl IdempotencyLedger for RebornLibSqlIdempotencyLedger { | ||
| async fn begin_or_replay( |
There was a problem hiding this comment.
[High][tests/concurrency] begin_or_replay relies on database uniqueness/transactions for the core duplicate-reservation race, but the durable tests only make sequential duplicate calls. Add a contention test where two concurrent begin_or_replay calls for the same fingerprint produce exactly one New and one transient/in-flight result.
| "idempotency reservation was superseded before terminal settle", | ||
| )); | ||
| } | ||
| if current.action_id != action.action_id { |
There was a problem hiding this comment.
[Medium][tests] This stale-action mismatch path is important after lease reclaim: the old holder must not be able to settle over the newer reservation. Add a contract test that begins an action, lets the lease expire/reclaims it, then verifies settling the original action returns the superseded-reservation error and does not overwrite the current row.
| fingerprint: ActionFingerprintKey, | ||
| received_at: DateTime<Utc>, | ||
| ) -> Result<IdempotencyDecision, ProductWorkflowError> { | ||
| self.run_migrations().await?; |
There was a problem hiding this comment.
[Medium][performance/concurrency] Every ledger operation reruns migrations before doing DML. On the product inbound hot path this adds an extra connection and DDL/catalog work per action; for Postgres it can also take DDL-related locks under concurrent traffic. Run migrations at construction/startup or guard them with a per-instance OnceCell/OnceLock so normal ledger operations only perform the transaction they need.
…roduct-workflow-filesystem-port # Conflicts: # Cargo.toml # crates/ironclaw_product_workflow/Cargo.toml
serrrfirat
left a comment
There was a problem hiding this comment.
Multi-agent code review completed with --force.
Event: COMMENT
GitHub rejected REQUEST_CHANGES for this account because it owns the pull request, so this is posted as a comment review. Treat the High finding as blocking.
Findings kept after confidence filtering:
- High: 1
- Medium: 4
- Low: 3
Primary blocker:
cargo test -p ironclaw_product_workflow_storage --features postgres --no-runfails because the postgres-only contract test build usesArcwhile the import is gated behindfeature = "libsql".
Reviewer lanes: security, bugs, performance/concurrency, tests, conventions.
| #![cfg(any(feature = "libsql", feature = "postgres"))] | ||
|
|
||
| #[cfg(feature = "libsql")] | ||
| use std::sync::Arc; |
There was a problem hiding this comment.
[High][build] Arc is only imported when feature = "libsql", but the postgres-only test build uses it in the postgres tests and in postgres_filesystem. I verified cargo test -p ironclaw_product_workflow_storage --features postgres --no-run fails with undeclared Arc, so the stated postgres feature combination is broken. Gate this import on any(feature = "libsql", feature = "postgres") or import it unconditionally.
| } | ||
|
|
||
| fn durable_error(operation: &'static str, error: impl std::fmt::Display) -> ProductWorkflowError { | ||
| tracing::error!(%error, operation, "product workflow idempotency ledger failed"); |
There was a problem hiding this comment.
[Low][security] This logs the full backend FilesystemError with %error. Those errors include the virtual path, and the ledger path is built from reversible hex of adapter, installation, source binding, and external event identifiers. On backend failures, logs can expose conversation/event identifiers and backend details. Please log a sanitized category/status here, or make the path components non-reversible before they reach error display.
| 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? |
There was a problem hiding this comment.
[Medium][tests] The missing-reservation branch in settle is new behavior, but the contract tests only cover successful settle and superseded reservations. A regression that silently accepts an unreserved action would not be caught. Add libsql_settle_missing_reservation_returns_transient and the postgres *_when_configured equivalent that call settle(ProductInboundAction::begin(...)) before any reservation.
| } | ||
| } | ||
|
|
||
| pub fn with_root( |
There was a problem hiding this comment.
[Low][tests] RebornLibSqlIdempotencyLedger::with_root is public but never exercised by the adjacent contract tests. If this constructor accidentally ignored the supplied root, the existing reopen/replay and lease tests would still pass. Add a libsql_custom_root_isolated_from_default_root contract case.
| } | ||
| } | ||
|
|
||
| pub fn with_root( |
There was a problem hiding this comment.
[Low][tests] RebornPostgresIdempotencyLedger::with_root is public but not covered by the postgres contract path. If it falls back to the default root, configured postgres tests would still pass. Add postgres_custom_root_isolated_from_default_root_when_configured.
| ) -> Result<IdempotencyDecision, ProductWorkflowError> { | ||
| let path = action_path(&self.root, &fingerprint)?; | ||
| let action = ProductInboundAction::begin(fingerprint, received_at); | ||
| match self |
There was a problem hiding this comment.
[Medium][performance] begin_or_replay runs on every inbound product action, but this durable ledger goes through the generic filesystem put/get path. On the SQL backends each put performs directory/child prechecks before the actual write, and settle/release add a read plus another checked write, so a normal successful action needs several SQL round trips before ack completion. If this is the intended hot-path implementation, consider adding a RootFilesystem-level batch/transition helper or another architecture-compatible fast path before enabling it for high inbound volume.
| @@ -1,5 +1,5 @@ | |||
| [workspace] | |||
| members = [".", "crates/ironclaw_common", "crates/ironclaw_host_api", "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_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"] | |||
There was a problem hiding this comment.
[Medium][conventions] This adds a new workspace crate, but crates/ironclaw_product_workflow_storage has no crate-local AGENTS.md. crates/AGENTS.md says each crate with a Cargo.toml has local instructions and agents should read them first. Please add crate ownership, boundaries, and validation guidance for this new crate.
| const DEFAULT_LEDGER_ROOT: &str = "/engine/product_workflow/idempotency/actions"; | ||
| const ACTION_RECORD_KIND: &str = "product_workflow_action"; | ||
|
|
||
| struct FilesystemIdempotencyLedger { |
There was a problem hiding this comment.
[Medium][conventions] FilesystemIdempotencyLedger stores the durable root as a raw String, even though callers provide a validated VirtualPath and filesystem APIs operate on VirtualPath. The typed-internals rule says raw strings are boundary formats; keep this as a VirtualPath internally and format child paths only at the final construction point.
|
Addressed the review findings in b8c2736:\n\n- fixed the postgres-only contract test compile failure by enabling the Arc import for postgres\n- added libSQL/postgres contract coverage for missing-reservation settle and custom-root isolation\n- kept ledger roots typed as VirtualPath internally\n- sanitized durable ledger error logging so backend paths/details are not logged\n- added a RootFilesystem known-leaf CAS write hook and SQL overrides, then used it from the durable ledger hot path\n- added crate-local AGENTS.md guidance for ironclaw_product_workflow_storage\n\nLocal verification run:\n- cargo test -p ironclaw_product_workflow_storage --features libsql\n- cargo test -p ironclaw_product_workflow_storage --features postgres --no-run\n- cargo test -p ironclaw_product_workflow_storage --features postgres\n- cargo check -p ironclaw_product_workflow_storage --features "libsql postgres"\n- cargo clippy -p ironclaw_product_workflow_storage --all-targets --features "libsql postgres" -- -D warnings\n- cargo test -p ironclaw_filesystem --features "libsql postgres"\n- bash scripts/pre-commit-safety.sh |
…roduct-workflow-durable-ledger # Conflicts: # Cargo.toml
Resolves four conflicts surfaced by reborn-integration moving forward since the previous merge: * `Cargo.toml`: keep this branch's workspace `members` list, which includes `crates/ironclaw_reborn_telegram_v2_host` (added by this PR). The only conflict was that membership line. * `crates/ironclaw_product_workflow_storage/Cargo.toml` (add/add): keep this branch's dependency set (the storage crate also ships the outbound-state delivery sink, which needs `ironclaw_outbound`, `ironclaw_product_adapters`, `ironclaw_threads`, `ironclaw_turns`, plus `sha2`/`hex`/`uuid` for the filesystem-ledger path scheme). Picked up RI's package metadata polish (authors, homepage, repository, license fields). * `crates/ironclaw_product_workflow_storage/src/lib.rs` (add/add): keep this branch's module-index form. RI shipped a competing inline `FilesystemIdempotencyLedger` implementation (PR #3759) plus backend-specific `RebornLibSqlIdempotencyLedger` / `RebornPostgresIdempotencyLedger` wrappers. Those wrappers have no downstream consumers yet, while this branch's `FilesystemIdempotencyLedger` over `ScopedFilesystem` is what the composition/runtime in this PR already consumes. The two implementations have the same intent (durable ledger over the universal FS dispatch fabric) but incompatible APIs (`ScopedFilesystem` vs raw `RootFilesystem`), so they can't coexist; RI's design can either rebase against this PR or land later as a follow-up that replaces this one cleanly. * `Cargo.lock`: regenerated. Also removes `crates/ironclaw_product_workflow_storage/tests/durable_ledger_contract.rs` (came in cleanly from RI but references the dropped `RebornLibSqlIdempotencyLedger` / `RebornPostgresIdempotencyLedger` types); the existing `ledger_filesystem_contract.rs` in this branch provides equivalent CAS/concurrency coverage. `InboundTurnService` trait gained two methods upstream (`replay_accepted_user_message`, `accept_user_message_with_before_policy`) — implemented on `StubInboundTurnService` with the smallest honest behavior: replay probe returns `Ok(None)` (no session-thread store is wired in the tracer), and the before-policy variant applies the policy gate then delegates to `accept_user_message`. The `#[non_exhaustive]` `BeforeInboundPolicyOutcome` enum is handled with an explicit fallback arm so a new variant upstream surfaces as `TurnSubmissionRejected` rather than a silent allow/reject in this stub. Verification: `cargo fmt`, `cargo clippy --tests --all-features` (clean), full test suite for `ironclaw_reborn_telegram_v2_host` / `ironclaw_host_runtime` / `ironclaw_product_workflow_storage` — all green. Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
…3759) * Add durable product workflow ledger * fix(product-workflow): address review durable ledger storage (nearai#3759) * fix(product-workflow): address gemini ledger review (nearai#3759) * refactor(product-workflow): use filesystem SQL ledger storage (nearai#3759) * fix(product-workflow-storage): address serrrfirat review — durable ledger findings (nearai#3759) * fix(product-workflow-storage): keep durable ledger fixes scoped (nearai#3759)
Summary
Validation