From bdc2f0e8dcbd37fc877942ad7cc23bd3d8feab3a Mon Sep 17 00:00:00 2001 From: Muhammad Mustaqeem Date: Wed, 5 Aug 2026 23:47:41 +0500 Subject: [PATCH 1/4] test(flows): extract bus.rs tests into bus_tests.rs Moves the `mod tests` block out of `flows/bus.rs` into a sibling `bus_tests.rs`, using the `#[cfg(test)] #[path = "..."] mod tests;` convention already used throughout this crate (`subconscious/factory.rs`, `skills/ops.rs`, `config/schema/channels.rs`, `platform/about_app/catalog.rs`, and ~170 other sites). Pure move: the test bodies are unchanged apart from being dedented one level now that they are no longer nested inside an inline `mod tests { .. }`. Also adds the tests for the #5268 item 2 fix (settle dedup state from the graph the run actually executed). --- src/openhuman/flows/bus_tests.rs | 1399 ++++++++++++++++++++++++++++++ 1 file changed, 1399 insertions(+) create mode 100644 src/openhuman/flows/bus_tests.rs diff --git a/src/openhuman/flows/bus_tests.rs b/src/openhuman/flows/bus_tests.rs new file mode 100644 index 0000000000..b19aa59ac1 --- /dev/null +++ b/src/openhuman/flows/bus_tests.rs @@ -0,0 +1,1399 @@ +//! Unit tests for [`super`] — the `flows::` event-bus handlers. +//! +//! Split out of `bus.rs` under the same +//! `#[cfg(test)] #[path = "..."] mod tests;` convention used across this +//! crate (`memory/schema`, `integrations`, `security/policy`, …), keeping +//! `bus.rs` itself to its handler logic. + +use super::*; +use crate::openhuman::flows::Flow; +use crate::openhuman::inference::embeddings::NoopEmbedding; +use crate::openhuman::memory::store::UnifiedMemory; +use serde_json::json; +use tinyflows::model::{Node, NodeKind, WorkflowGraph}; + +/// A directly-constructed, isolated [`Memory`] for the digest tests — NOT +/// the process-global `OnceLock` client. The global is one-shot, so an +/// earlier test in the same binary may already have bound it to a different +/// workspace, making `global::init(..)` here a silent no-op (see +/// `memory::global`'s own test notes). Injecting this instance into the +/// subscriber via [`FlowRunDigestSubscriber::with_memory`] makes writes and +/// read-backs go through the SAME store deterministically — the same shape +/// `flows::memory_tools`' tests use. +fn digest_test_memory(tmp: &tempfile::TempDir) -> Arc { + Arc::new(UnifiedMemory::new(tmp.path(), Arc::new(NoopEmbedding), None).unwrap()) +} + +fn test_config(tmp: &tempfile::TempDir) -> Arc { + let config = Config { + workspace_dir: tmp.path().join("workspace"), + action_dir: tmp.path().join("workspace"), + config_path: tmp.path().join("config.toml"), + ..Config::default() + }; + std::fs::create_dir_all(&config.workspace_dir).unwrap(); + Arc::new(config) +} + +fn trigger_node(config: Value) -> Node { + Node { + id: "t".to_string(), + kind: NodeKind::Trigger, + type_version: 1, + name: "Trigger".to_string(), + config, + ports: Vec::new(), + position: None, + } +} + +fn flow_with_trigger_config(id: &str, enabled: bool, trigger_config: Value) -> Flow { + Flow { + id: id.to_string(), + name: id.to_string(), + enabled, + graph: WorkflowGraph { + nodes: vec![trigger_node(trigger_config)], + ..Default::default() + }, + created_at: "2026-01-01T00:00:00Z".to_string(), + updated_at: "2026-01-01T00:00:00Z".to_string(), + last_run_at: None, + last_status: None, + require_approval: false, + } +} + +fn dedup_node(id: &str) -> Node { + Node { + id: id.to_string(), + kind: NodeKind::Dedup, + type_version: 1, + name: id.to_string(), + config: json!({ "key": "=item.id" }), + ports: Vec::new(), + position: None, + } +} + +/// A saved flow with a `trigger` node plus one `dedup` node with id +/// `dedup_id` — the minimal graph [`DedupCommitSubscriber::dedup_node_ids`] +/// needs to find something to settle. +fn flow_with_dedup_node(id: &str, dedup_id: &str) -> Flow { + Flow { + id: id.to_string(), + name: id.to_string(), + enabled: true, + graph: WorkflowGraph { + nodes: vec![trigger_node(json!({})), dedup_node(dedup_id)], + ..Default::default() + }, + created_at: "2026-01-01T00:00:00Z".to_string(), + updated_at: "2026-01-01T00:00:00Z".to_string(), + last_run_at: None, + last_status: None, + require_approval: false, + } +} + +#[test] +fn pinned_trigger_inputs_reads_values_an_author_fixed_for_unattended_runs() { + let flow = flow_with_trigger_config( + "f1", + true, + json!({ + "trigger_kind": "schedule", + "schedule": "0 9 * * *", + "inputs": { "repo": "acme/api", "depth": 3 } + }), + ); + let inputs = pinned_trigger_inputs(&flow); + assert_eq!(inputs["repo"], json!("acme/api")); + assert_eq!(inputs["depth"], json!(3)); +} + +#[test] +fn pinned_trigger_inputs_is_empty_when_unset_or_malformed() { + // Empty, not an error: a flow declaring no inputs (the overwhelming + // majority) must keep dispatching on a tick exactly as before, and a + // malformed value is caught downstream by `prepare_flow_run`, which + // reports it against the flow's actual declarations. + for cfg in [ + json!({ "trigger_kind": "schedule" }), + json!({ "trigger_kind": "schedule", "inputs": null }), + json!({ "trigger_kind": "schedule", "inputs": ["repo"] }), + ] { + let flow = flow_with_trigger_config("f1", true, cfg.clone()); + assert!( + pinned_trigger_inputs(&flow).is_empty(), + "expected no pinned inputs for {cfg}" + ); + } +} + +#[test] +fn pinned_trigger_inputs_is_empty_for_a_graph_with_no_trigger() { + let mut flow = flow_with_trigger_config("f1", true, json!({ "trigger_kind": "schedule" })); + flow.graph.nodes.clear(); + assert!(pinned_trigger_inputs(&flow).is_empty()); +} + +#[test] +fn name_and_domains_are_stable() { + let tmp = tempfile::TempDir::new().unwrap(); + let sub = FlowTriggerSubscriber::new(test_config(&tmp)); + assert_eq!(sub.name(), "flows::trigger"); + assert_eq!( + sub.domains(), + Some(&["cron", "composio", "webhook", "system"][..]) + ); +} + +#[tokio::test] +async fn handle_does_not_panic_on_arbitrary_events() { + let tmp = tempfile::TempDir::new().unwrap(); + let sub = FlowTriggerSubscriber::new(test_config(&tmp)); + sub.handle(&DomainEvent::CronJobTriggered { + job_id: "j1".into(), + job_name: "test".into(), + job_type: "shell".into(), + }) + .await; + sub.handle(&DomainEvent::FlowScheduleTick { + flow_id: "missing-flow".into(), + }) + .await; +} + +#[test] +fn extract_trigger_kind_reads_schedule() { + let flow = flow_with_trigger_config( + "f1", + true, + json!({ "trigger_kind": "schedule", "schedule": "0 9 * * *" }), + ); + assert!(matches!( + extract_trigger_kind(&flow), + Some(TriggerKind::Schedule) + )); +} + +#[test] +fn extract_trigger_kind_none_for_missing_discriminator() { + let flow = flow_with_trigger_config("f1", true, json!({})); + assert!(extract_trigger_kind(&flow).is_none()); +} + +#[test] +fn extract_trigger_kind_none_for_invalid_discriminator() { + let flow = flow_with_trigger_config("f1", true, json!({ "trigger_kind": "not_a_kind" })); + assert!(extract_trigger_kind(&flow).is_none()); +} + +#[test] +fn matches_app_event_requires_toolkit_and_slug_match() { + let flow = flow_with_trigger_config( + "f1", + true, + json!({ "trigger_kind": "app_event", "toolkit": "gmail", "trigger_slug": "GMAIL_NEW_GMAIL_MESSAGE" }), + ); + assert!(matches_app_event(&flow, "gmail", "GMAIL_NEW_GMAIL_MESSAGE")); + // Case-insensitive. + assert!(matches_app_event(&flow, "Gmail", "gmail_new_gmail_message")); + // Wrong toolkit or slug does not match. + assert!(!matches_app_event( + &flow, + "slack", + "GMAIL_NEW_GMAIL_MESSAGE" + )); + assert!(!matches_app_event(&flow, "gmail", "SLACK_NEW_MESSAGE")); +} + +#[test] +fn matches_app_event_false_for_non_app_event_trigger() { + let flow = flow_with_trigger_config( + "f1", + true, + json!({ "trigger_kind": "schedule", "schedule": "0 9 * * *" }), + ); + assert!(!matches_app_event( + &flow, + "gmail", + "GMAIL_NEW_GMAIL_MESSAGE" + )); +} + +#[tokio::test] +async fn handle_app_event_ignores_disabled_flows() { + let tmp = tempfile::TempDir::new().unwrap(); + let config = test_config(&tmp); + let flow = flow_with_trigger_config( + "disabled-flow", + false, + json!({ "trigger_kind": "app_event", "toolkit": "gmail", "trigger_slug": "GMAIL_NEW_GMAIL_MESSAGE" }), + ); + crate::openhuman::flows::store::upsert_flow(&config, &flow).unwrap(); + + // `list_enabled_flows` must not surface the disabled flow at all — + // proves the subscriber's dispatch source already excludes it, + // rather than asserting on a spawned background task's side effect. + let (enabled, skipped) = + crate::openhuman::flows::store::list_enabled_flows(&config).unwrap(); + assert!(enabled.is_empty()); + assert_eq!(skipped, 0); +} + +#[tokio::test] +async fn handle_schedule_tick_ignores_disabled_flow() { + let tmp = tempfile::TempDir::new().unwrap(); + let config = test_config(&tmp); + let flow = flow_with_trigger_config( + "sched-flow", + false, + json!({ "trigger_kind": "schedule", "schedule": "0 9 * * *" }), + ); + crate::openhuman::flows::store::upsert_flow(&config, &flow).unwrap(); + + let sub = FlowTriggerSubscriber::new(config.clone()); + // Must not panic and must not spawn a run for a disabled flow — we + // can't directly observe "no run happened" without a full flows_run + // fixture, but this exercises the early-return path without error. + sub.handle(&DomainEvent::FlowScheduleTick { + flow_id: "sched-flow".into(), + }) + .await; +} + +// ── in-flight dedupe (CodeRabbit finding B) ───────────────────── + +#[test] +fn try_acquire_dispatch_skips_a_flow_already_in_flight() { + let tmp = tempfile::TempDir::new().unwrap(); + let sub = FlowTriggerSubscriber::new(test_config(&tmp)); + + let guard = sub + .try_acquire_dispatch("f1") + .expect("first claim for f1 should succeed"); + assert!( + sub.try_acquire_dispatch("f1").is_none(), + "a second claim for the same flow while the first is held must be skipped" + ); + + // A different flow is unaffected. + assert!(sub.try_acquire_dispatch("f2").is_some()); + + drop(guard); + assert!( + sub.try_acquire_dispatch("f1").is_some(), + "dropping the guard must release the claim so f1 can run again" + ); +} + +#[test] +fn default_constructs_the_same_as_new() { + let tmp = tempfile::TempDir::new().unwrap(); + let config = test_config(&tmp); + let a = FlowTriggerSubscriber::new(config.clone()); + let b = FlowTriggerSubscriber::new(config); + assert_eq!(a.name(), b.name()); +} + +// ── FlowRunDigestSubscriber ───────────────────────────────────── + +#[test] +fn digest_name_and_domains_are_stable() { + let tmp = tempfile::TempDir::new().unwrap(); + let sub = FlowRunDigestSubscriber::new(test_config(&tmp)); + assert_eq!(sub.name(), "flows::digest"); + assert_eq!(sub.domains(), Some(&["cron"][..])); +} + +#[tokio::test] +async fn digest_handle_does_not_panic_on_unrelated_events() { + let tmp = tempfile::TempDir::new().unwrap(); + let sub = FlowRunDigestSubscriber::new(test_config(&tmp)); + // Must not panic, and must not touch the memory layer at all, for + // any event other than `FlowRunFinished`. + sub.handle(&DomainEvent::CronJobTriggered { + job_id: "j1".into(), + job_name: "test".into(), + job_type: "shell".into(), + }) + .await; +} + +#[tokio::test] +async fn digest_ignores_failed_run() { + let tmp = tempfile::TempDir::new().unwrap(); + let config = test_config(&tmp); + let memory = digest_test_memory(&tmp); + + let flow = flow_with_trigger_config("f-failed", true, json!({})); + store::upsert_flow(&config, &flow).unwrap(); + store::insert_flow_run( + &config, + "run-failed", + "f-failed", + "thread-failed", + "2026-01-01T00:00:00Z", + ) + .unwrap(); + store::finish_flow_run( + &config, + "run-failed", + "failed", + "2026-01-01T00:05:00Z", + &[], + &[], + Some("boom"), + None, + ) + .unwrap(); + + let sub = FlowRunDigestSubscriber::with_memory(config, memory.clone()); + sub.handle(&DomainEvent::FlowRunFinished { + flow_id: "f-failed".into(), + run_id: "run-failed".into(), + status: "failed".into(), + }) + .await; + + let entry = memory + .get(&flow_namespace("f-failed"), "run_digest:run-failed") + .await + .unwrap(); + assert!( + entry.is_none(), + "a failed run must never produce a run_digest entry" + ); +} + +#[tokio::test] +async fn digest_ignores_cancelled_run() { + let tmp = tempfile::TempDir::new().unwrap(); + let config = test_config(&tmp); + let memory = digest_test_memory(&tmp); + + let flow = flow_with_trigger_config("f-cancelled", true, json!({})); + store::upsert_flow(&config, &flow).unwrap(); + store::insert_flow_run( + &config, + "run-cancelled", + "f-cancelled", + "thread-cancelled", + "2026-01-01T00:00:00Z", + ) + .unwrap(); + store::finish_flow_run( + &config, + "run-cancelled", + "cancelled", + "2026-01-01T00:05:00Z", + &[], + &[], + None, + None, + ) + .unwrap(); + + let sub = FlowRunDigestSubscriber::with_memory(config, memory.clone()); + sub.handle(&DomainEvent::FlowRunFinished { + flow_id: "f-cancelled".into(), + run_id: "run-cancelled".into(), + status: "cancelled".into(), + }) + .await; + + let entry = memory + .get(&flow_namespace("f-cancelled"), "run_digest:run-cancelled") + .await + .unwrap(); + assert!(entry.is_none()); +} + +#[tokio::test] +async fn digest_writes_run_digest_entry_for_completed_run() { + let tmp = tempfile::TempDir::new().unwrap(); + let config = test_config(&tmp); + let memory = digest_test_memory(&tmp); + + let flow = flow_with_trigger_config("f-ok", true, json!({})); + store::upsert_flow(&config, &flow).unwrap(); + store::insert_flow_run( + &config, + "run-ok", + "f-ok", + "thread-ok", + "2026-01-01T00:00:00Z", + ) + .unwrap(); + let step = crate::openhuman::flows::FlowRunStep { + node_id: "n1".to_string(), + output: json!({ "sent": 3 }), + port: None, + status: Some("success".to_string()), + duration_ms: Some(12), + diagnostics: Vec::new(), + }; + store::finish_flow_run( + &config, + "run-ok", + "completed", + "2026-01-01T00:05:00Z", + &[step], + &[], + None, + None, + ) + .unwrap(); + + let sub = FlowRunDigestSubscriber::with_memory(config, memory.clone()); + sub.handle(&DomainEvent::FlowRunFinished { + flow_id: "f-ok".into(), + run_id: "run-ok".into(), + status: "completed".into(), + }) + .await; + + let entry = memory + .get(&flow_namespace("f-ok"), "run_digest:run-ok") + .await + .unwrap() + .expect("completed run must produce a run_digest entry"); + assert_eq!(entry.taint, MemoryTaint::ExternalSync); + assert!(entry.content.contains("f-ok")); + assert!(entry.content.contains("completed")); + assert!(entry.content.contains("n1")); + assert!(entry.content.chars().count() <= DIGEST_MAX_CHARS); +} + +#[tokio::test] +async fn digest_treats_completed_with_warnings_as_success() { + let tmp = tempfile::TempDir::new().unwrap(); + let config = test_config(&tmp); + let memory = digest_test_memory(&tmp); + + let flow = flow_with_trigger_config("f-warn", true, json!({})); + store::upsert_flow(&config, &flow).unwrap(); + store::insert_flow_run( + &config, + "run-warn", + "f-warn", + "thread-warn", + "2026-01-01T00:00:00Z", + ) + .unwrap(); + store::finish_flow_run( + &config, + "run-warn", + "completed_with_warnings", + "2026-01-01T00:05:00Z", + &[], + &[], + None, + None, + ) + .unwrap(); + + let sub = FlowRunDigestSubscriber::with_memory(config, memory.clone()); + sub.handle(&DomainEvent::FlowRunFinished { + flow_id: "f-warn".into(), + run_id: "run-warn".into(), + status: "completed_with_warnings".into(), + }) + .await; + + let entry = memory + .get(&flow_namespace("f-warn"), "run_digest:run-warn") + .await + .unwrap(); + assert!(entry.is_some()); +} + +#[test] +fn truncate_chars_bounds_output_and_marks_truncation() { + let long = "x".repeat(50); + let truncated = truncate_chars(&long, 10); + assert_eq!(truncated.chars().count(), 10); + assert!(truncated.ends_with('…')); + + let short = "hello"; + assert_eq!(truncate_chars(short, 10), "hello"); +} + +#[test] +fn render_run_digest_is_bounded_and_includes_key_fields() { + let run = FlowRun { + id: "run-1".to_string(), + flow_id: "f1".to_string(), + thread_id: "thread-1".to_string(), + status: "completed".to_string(), + started_at: "2026-01-01T00:00:00Z".to_string(), + finished_at: Some("2026-01-01T00:05:00Z".to_string()), + steps: vec![crate::openhuman::flows::FlowRunStep { + node_id: "n1".to_string(), + output: json!({ "ok": true }), + port: None, + status: Some("success".to_string()), + duration_ms: Some(5), + diagnostics: Vec::new(), + }], + pending_approvals: Vec::new(), + error: None, + graph_hash: None, + }; + let digest = render_run_digest("My Flow", &run); + assert!(digest.contains("My Flow")); + assert!(digest.contains("completed")); + assert!(digest.contains("n1")); + assert!(digest.chars().count() <= DIGEST_MAX_CHARS); +} + +// ── DedupCommitSubscriber ──────────────────────────────────────── + +fn dedup_state_namespace(flow_id: &str) -> String { + // MUST match `tinyflows::build_capabilities`'s `state_namespace` + // (`src/openhuman/flows/tinyflows/caps.rs`) — this test asserts the + // subscriber collides with the SAME keys the engine's `dedup` node + // itself reads/writes, not just "some" namespace. + format!("flow:{flow_id}") +} + +#[test] +fn dedup_commit_name_and_domains_are_stable() { + let tmp = tempfile::TempDir::new().unwrap(); + let sub = DedupCommitSubscriber::new(test_config(&tmp)); + assert_eq!(sub.name(), "flows::dedup_commit"); + assert_eq!(sub.domains(), Some(&["cron"][..])); +} + +#[tokio::test] +async fn dedup_commit_ignores_unrelated_events() { + let tmp = tempfile::TempDir::new().unwrap(); + let sub = DedupCommitSubscriber::new(test_config(&tmp)); + // Must not panic for any event other than `FlowRunFinished`. + sub.handle(&DomainEvent::CronJobTriggered { + job_id: "j1".into(), + job_name: "test".into(), + job_type: "shell".into(), + }) + .await; +} + +#[tokio::test] +async fn dedup_commit_flow_with_no_dedup_nodes_is_a_noop() { + let tmp = tempfile::TempDir::new().unwrap(); + let config = test_config(&tmp); + let flow = flow_with_trigger_config("f-no-dedup", true, json!({})); + store::upsert_flow(&config, &flow).unwrap(); + + let sub = DedupCommitSubscriber::new(config); + // Must not panic when the flow has no `dedup` node at all. + sub.handle(&DomainEvent::FlowRunFinished { + flow_id: "f-no-dedup".into(), + run_id: "run-1".into(), + status: "completed".into(), + }) + .await; +} + +#[tokio::test] +async fn dedup_commit_unions_tentative_into_committed_and_clears_tentative_on_success() { + let tmp = tempfile::TempDir::new().unwrap(); + let config = test_config(&tmp); + let flow = flow_with_dedup_node("f-ok", "dd"); + store::upsert_flow(&config, &flow).unwrap(); + + let namespace = dedup_state_namespace("f-ok"); + store::kv_set(&config, &namespace, "dedup:dd:committed", &json!(["a"])).unwrap(); + store::kv_set( + &config, + &namespace, + "dedup:dd:tentative", + &json!(["b", "c"]), + ) + .unwrap(); + + let sub = DedupCommitSubscriber::new(config.clone()); + sub.handle(&DomainEvent::FlowRunFinished { + flow_id: "f-ok".into(), + run_id: "run-ok".into(), + status: "completed".into(), + }) + .await; + + let committed = store::kv_get(&config, &namespace, "dedup:dd:committed") + .unwrap() + .expect("committed key must still exist"); + let mut committed: Vec<&str> = committed + .as_array() + .unwrap() + .iter() + .map(|v| v.as_str().unwrap()) + .collect(); + committed.sort_unstable(); + assert_eq!(committed, vec!["a", "b", "c"], "committed = union"); + + assert!( + store::kv_get(&config, &namespace, "dedup:dd:tentative") + .unwrap() + .is_none(), + "tentative must be cleared after a successful commit" + ); +} + +#[tokio::test] +async fn dedup_commit_treats_completed_with_warnings_as_success() { + let tmp = tempfile::TempDir::new().unwrap(); + let config = test_config(&tmp); + let flow = flow_with_dedup_node("f-warn", "dd"); + store::upsert_flow(&config, &flow).unwrap(); + + let namespace = dedup_state_namespace("f-warn"); + store::kv_set(&config, &namespace, "dedup:dd:tentative", &json!(["x"])).unwrap(); + + let sub = DedupCommitSubscriber::new(config.clone()); + sub.handle(&DomainEvent::FlowRunFinished { + flow_id: "f-warn".into(), + run_id: "run-warn".into(), + status: "completed_with_warnings".into(), + }) + .await; + + let committed = store::kv_get(&config, &namespace, "dedup:dd:committed") + .unwrap() + .expect("completed_with_warnings must still commit"); + assert_eq!(committed, json!(["x"])); + assert!(store::kv_get(&config, &namespace, "dedup:dd:tentative") + .unwrap() + .is_none()); +} + +#[tokio::test] +async fn dedup_commit_releases_tentative_without_touching_committed_on_failure() { + let tmp = tempfile::TempDir::new().unwrap(); + let config = test_config(&tmp); + let flow = flow_with_dedup_node("f-failed", "dd"); + store::upsert_flow(&config, &flow).unwrap(); + + let namespace = dedup_state_namespace("f-failed"); + store::kv_set(&config, &namespace, "dedup:dd:committed", &json!(["a"])).unwrap(); + store::kv_set(&config, &namespace, "dedup:dd:tentative", &json!(["b"])).unwrap(); + + let sub = DedupCommitSubscriber::new(config.clone()); + sub.handle(&DomainEvent::FlowRunFinished { + flow_id: "f-failed".into(), + run_id: "run-failed".into(), + status: "failed".into(), + }) + .await; + + assert_eq!( + store::kv_get(&config, &namespace, "dedup:dd:committed") + .unwrap() + .unwrap(), + json!(["a"]), + "committed must be untouched by a failed run" + ); + assert!( + store::kv_get(&config, &namespace, "dedup:dd:tentative") + .unwrap() + .is_none(), + "tentative must be released (cleared) on failure so the item retries" + ); +} + +#[tokio::test] +async fn dedup_commit_releases_tentative_on_cancelled_and_interrupted() { + for status in ["cancelled", "interrupted"] { + let tmp = tempfile::TempDir::new().unwrap(); + let config = test_config(&tmp); + let flow_id = format!("f-{status}"); + let flow = flow_with_dedup_node(&flow_id, "dd"); + store::upsert_flow(&config, &flow).unwrap(); + + let namespace = dedup_state_namespace(&flow_id); + store::kv_set(&config, &namespace, "dedup:dd:tentative", &json!(["z"])).unwrap(); + + let sub = DedupCommitSubscriber::new(config.clone()); + sub.handle(&DomainEvent::FlowRunFinished { + flow_id: flow_id.clone(), + run_id: format!("run-{status}"), + status: status.to_string(), + }) + .await; + + assert!( + store::kv_get(&config, &namespace, "dedup:dd:committed") + .unwrap() + .is_none(), + "status {status} must never commit" + ); + assert!( + store::kv_get(&config, &namespace, "dedup:dd:tentative") + .unwrap() + .is_none(), + "status {status} must release tentative" + ); + } +} + +#[tokio::test] +async fn dedup_commit_two_dedup_nodes_settle_independently() { + let tmp = tempfile::TempDir::new().unwrap(); + let config = test_config(&tmp); + let flow = Flow { + id: "f-multi".to_string(), + name: "f-multi".to_string(), + enabled: true, + graph: WorkflowGraph { + nodes: vec![ + trigger_node(json!({})), + dedup_node("dd1"), + dedup_node("dd2"), + ], + ..Default::default() + }, + created_at: "2026-01-01T00:00:00Z".to_string(), + updated_at: "2026-01-01T00:00:00Z".to_string(), + last_run_at: None, + last_status: None, + require_approval: false, + }; + store::upsert_flow(&config, &flow).unwrap(); + + let namespace = dedup_state_namespace("f-multi"); + store::kv_set(&config, &namespace, "dedup:dd1:tentative", &json!(["a"])).unwrap(); + store::kv_set(&config, &namespace, "dedup:dd2:tentative", &json!(["b"])).unwrap(); + + let sub = DedupCommitSubscriber::new(config.clone()); + sub.handle(&DomainEvent::FlowRunFinished { + flow_id: "f-multi".into(), + run_id: "run-multi".into(), + status: "completed".into(), + }) + .await; + + assert_eq!( + store::kv_get(&config, &namespace, "dedup:dd1:committed") + .unwrap() + .unwrap(), + json!(["a"]) + ); + assert_eq!( + store::kv_get(&config, &namespace, "dedup:dd2:committed") + .unwrap() + .unwrap(), + json!(["b"]) + ); +} + +// ── per-flow commit serialization (issue #5265) ─────────────────── +// +// CodeRabbit "Major" on the dedup engine PR: the commit's +// load(committed)+union(tentative)+store(committed) is a +// read-modify-write, not a CAS. Two overlapping `FlowRunFinished` +// events for the SAME flow could otherwise interleave and have the +// second writer's store clobber the first writer's union, silently +// losing that run's committed keys. `handle_finished` now serializes +// settlement per `flow_id` via `FLOW_COMMIT_LOCKS`. +// +// Two tests, deliberately split: +// +// - `..._never_runs_two_commits_for_the_same_flow_concurrently` spawns a +// burst of genuinely overlapping `FlowRunFinished` events for the SAME +// flow_id and proves the LOCK itself provides mutual exclusion (the +// high-water mark of concurrently-active critical sections never +// exceeds 1) — this is the "spawn two tasks contending on the same +// flow_id" case. +// - `..._serial_commits_for_the_same_flow_accumulate_via_union` proves +// the property that mutual exclusion protects: settling run after run +// for the same node never clobbers an earlier run's committed keys — +// each contributes to the union. +// +// These are split rather than combined into one "two runs with two +// different tentative sets, truly concurrently, assert union" test +// because `tentative` is a single shared KV row per node (not +// per-run) — forcing two *different* tentative contents to both survive +// a genuinely simultaneous read would require injecting a write from +// outside `handle_finished` in the middle of its critical section, which +// instead exercises the SEPARATE, still-open node-side race (the +// `dedup` node's own in-run `tentative` read-modify-write, documented on +// `DedupCommitSubscriber` above as explicitly NOT fixed by this lock). +// Together, the two tests below establish the same guarantee end to +// end: the lock enforces serialization (test 1), and serialization is +// sufficient for correctness (test 2). + +#[tokio::test] +async fn dedup_commit_never_runs_two_commits_for_the_same_flow_concurrently() { + let tmp = tempfile::TempDir::new().unwrap(); + let config = test_config(&tmp); + let flow = flow_with_dedup_node("f-race", "dd"); + store::upsert_flow(&config, &flow).unwrap(); + + let namespace = dedup_state_namespace("f-race"); + store::kv_set(&config, &namespace, "dedup:dd:tentative", &json!(["seed"])).unwrap(); + + // Arm the test-only scheduling hook (see `CommitTestHooks`): every + // `handle_finished` call sleeps briefly while holding the per-flow + // lock, and records how many calls are concurrently inside that + // window. Instance-scoped (not a global static) so this doesn't + // interfere with — or get polluted by — unrelated tests that cargo + // runs concurrently on other threads. Without a correctly-scoped + // lock, a burst of overlapping `FlowRunFinished` events for the SAME + // flow_id would pile up inside the critical section together + // instead of queuing. + let hooks = Arc::new(CommitTestHooks::default()); + hooks + .delay_ms + .store(20, std::sync::atomic::Ordering::SeqCst); + + let sub = Arc::new(DedupCommitSubscriber::with_test_hooks( + config.clone(), + hooks.clone(), + )); + let mut handles = Vec::new(); + for i in 0..5 { + let sub = sub.clone(); + handles.push(tokio::spawn(async move { + sub.handle(&DomainEvent::FlowRunFinished { + flow_id: "f-race".into(), + run_id: format!("run-{i}"), + status: "completed".into(), + }) + .await; + })); + } + for handle in handles { + handle.await.unwrap(); + } + + assert_eq!( + hooks.concurrent.load(std::sync::atomic::Ordering::SeqCst), + 0, + "every critical-section entry must have a matching exit" + ); + assert_eq!( + hooks + .max_concurrent + .load(std::sync::atomic::Ordering::SeqCst), + 1, + "the per-flow lock must serialize overlapping FlowRunFinished handling for the \ + same flow_id — at most one commit critical section may be active at a time" + ); +} + +#[tokio::test] +async fn dedup_commit_serial_commits_for_the_same_flow_accumulate_via_union() { + let tmp = tempfile::TempDir::new().unwrap(); + let config = test_config(&tmp); + let flow = flow_with_dedup_node("f-serial", "dd"); + store::upsert_flow(&config, &flow).unwrap(); + + let namespace = dedup_state_namespace("f-serial"); + let sub = DedupCommitSubscriber::new(config.clone()); + + // Run A finishes, having tentatively seen "a". + store::kv_set(&config, &namespace, "dedup:dd:tentative", &json!(["a"])).unwrap(); + sub.handle(&DomainEvent::FlowRunFinished { + flow_id: "f-serial".into(), + run_id: "run-a".into(), + status: "completed".into(), + }) + .await; + + // Run B finishes later, having independently tentatively seen "b". + // The per-flow lock (proven by the concurrency test above) is what + // guarantees two overlapping runs' `FlowRunFinished` handling + // reduces to exactly this serialized order in practice — so this is + // the correctness property that mutual exclusion is protecting. + store::kv_set(&config, &namespace, "dedup:dd:tentative", &json!(["b"])).unwrap(); + sub.handle(&DomainEvent::FlowRunFinished { + flow_id: "f-serial".into(), + run_id: "run-b".into(), + status: "completed".into(), + }) + .await; + + let committed = store::kv_get(&config, &namespace, "dedup:dd:committed") + .unwrap() + .expect("committed key must exist after both runs settle"); + let mut committed: Vec<&str> = committed + .as_array() + .unwrap() + .iter() + .map(|v| v.as_str().unwrap()) + .collect(); + committed.sort_unstable(); + assert_eq!( + committed, + vec!["a", "b"], + "settling run B must not clobber run A's already-committed keys — committed is a \ + running union across every run that has settled, never a last-writer-wins overwrite" + ); + assert!( + store::kv_get(&config, &namespace, "dedup:dd:tentative") + .unwrap() + .is_none(), + "tentative must be cleared after each successful commit" + ); +} + +#[test] +fn flow_commit_lock_returns_the_same_arc_for_the_same_flow_id_and_differs_across_flows() { + let a1 = flow_commit_lock("f-lock-a"); + let a2 = flow_commit_lock("f-lock-a"); + assert!( + Arc::ptr_eq(&a1, &a2), + "the same flow_id must share one lock instance" + ); + + let b = flow_commit_lock("f-lock-b"); + assert!( + !Arc::ptr_eq(&a1, &b), + "different flow_ids must not contend on the same lock" + ); +} + +// ── #5268 item 2: settle from the run's graph snapshot ─────────── + +/// The KV key [`snapshot_run_dedup_nodes`] writes, spelled out literally +/// rather than via `run_dedup_snapshot_key` so a silent rename of the key +/// format is caught here instead of shipping as an unreadable snapshot for +/// every run already in flight across the upgrade. +fn run_snapshot_key(run_id: &str) -> String { + format!("dedup::run_nodes:{run_id}") +} + +#[test] +fn snapshot_run_dedup_nodes_records_the_graph_the_run_executes() { + let tmp = tempfile::TempDir::new().unwrap(); + let config = test_config(&tmp); + let flow = flow_with_dedup_node("f-snap", "dd"); + + snapshot_run_dedup_nodes(&config, "f-snap", "run-1", &flow); + + let stored = store::kv_get( + &config, + &dedup_state_namespace("f-snap"), + &run_snapshot_key("run-1"), + ) + .unwrap() + .expect("a flow with a dedup node must snapshot it at run-start"); + assert_eq!(stored, json!(["dd"])); +} + +#[test] +fn snapshot_run_dedup_nodes_is_scoped_per_run() { + let tmp = tempfile::TempDir::new().unwrap(); + let config = test_config(&tmp); + let flow = flow_with_dedup_node("f-two-runs", "dd"); + + snapshot_run_dedup_nodes(&config, "f-two-runs", "run-a", &flow); + snapshot_run_dedup_nodes(&config, "f-two-runs", "run-b", &flow); + + // Two overlapping runs of the same flow must each keep their own + // snapshot — settling one must never consume the other's. + let namespace = dedup_state_namespace("f-two-runs"); + for run in ["run-a", "run-b"] { + assert!( + store::kv_get(&config, &namespace, &run_snapshot_key(run)) + .unwrap() + .is_some(), + "run {run} must keep its own snapshot" + ); + } +} + +#[test] +fn snapshot_run_dedup_nodes_writes_nothing_when_the_flow_has_no_dedup_node() { + let tmp = tempfile::TempDir::new().unwrap(); + let config = test_config(&tmp); + let flow = flow_with_trigger_config("f-plain", true, json!({})); + + snapshot_run_dedup_nodes(&config, "f-plain", "run-1", &flow); + + // Absent, NOT an empty array: an absent key is the unambiguous + // "no snapshot" the reader's saved-graph fallback keys off, and the + // overwhelming majority of flows must cost zero extra I/O per run. + assert!( + store::kv_get( + &config, + &dedup_state_namespace("f-plain"), + &run_snapshot_key("run-1"), + ) + .unwrap() + .is_none(), + "a flow with no dedup node must not write a snapshot at all" + ); +} + +#[tokio::test] +async fn dedup_commit_settles_a_node_deleted_from_the_flow_mid_run() { + // This is the #5268 item 2 bug. The run executed a graph containing + // `dd` and wrote tentative keys under it, but the flow was edited (the + // dedup node removed) before `FlowRunFinished` fired. Settling from the + // CURRENT saved graph finds no dedup node at all, so `dd`'s tentative + // keys are neither committed nor released and those items silently + // reprocess on the flow's next run. + let tmp = tempfile::TempDir::new().unwrap(); + let config = test_config(&tmp); + + // What the run actually executed, captured at run-start. + let executed = flow_with_dedup_node("f-edited", "dd"); + snapshot_run_dedup_nodes(&config, "f-edited", "run-1", &executed); + + // What the flow looks like by the time the run finishes: edited, with + // the dedup node deleted. + let edited = flow_with_trigger_config("f-edited", true, json!({})); + store::upsert_flow(&config, &edited).unwrap(); + + let namespace = dedup_state_namespace("f-edited"); + store::kv_set(&config, &namespace, "dedup:dd:tentative", &json!(["b"])).unwrap(); + + let sub = DedupCommitSubscriber::new(config.clone()); + sub.handle(&DomainEvent::FlowRunFinished { + flow_id: "f-edited".into(), + run_id: "run-1".into(), + status: "completed".into(), + }) + .await; + + let committed = store::kv_get(&config, &namespace, "dedup:dd:committed") + .unwrap() + .expect("the snapshot must let a mid-run-deleted node still settle"); + assert_eq!(committed, json!(["b"])); + assert!( + store::kv_get(&config, &namespace, "dedup:dd:tentative") + .unwrap() + .is_none(), + "tentative must be cleared once the snapshotted node commits" + ); +} + +#[tokio::test] +async fn dedup_commit_releases_a_node_deleted_from_the_flow_mid_run_on_failure() { + // Failure direction of the same bug: a failed run's tentative keys must + // still be released even though the node is gone from the saved graph. + let tmp = tempfile::TempDir::new().unwrap(); + let config = test_config(&tmp); + + let executed = flow_with_dedup_node("f-edited-fail", "dd"); + snapshot_run_dedup_nodes(&config, "f-edited-fail", "run-1", &executed); + let edited = flow_with_trigger_config("f-edited-fail", true, json!({})); + store::upsert_flow(&config, &edited).unwrap(); + + let namespace = dedup_state_namespace("f-edited-fail"); + store::kv_set(&config, &namespace, "dedup:dd:tentative", &json!(["b"])).unwrap(); + + let sub = DedupCommitSubscriber::new(config.clone()); + sub.handle(&DomainEvent::FlowRunFinished { + flow_id: "f-edited-fail".into(), + run_id: "run-1".into(), + status: "failed".into(), + }) + .await; + + assert!( + store::kv_get(&config, &namespace, "dedup:dd:tentative") + .unwrap() + .is_none(), + "a failed run must release the snapshotted node's tentative keys" + ); + assert!( + store::kv_get(&config, &namespace, "dedup:dd:committed") + .unwrap() + .is_none(), + "a failed run must never commit anything" + ); +} + +#[tokio::test] +async fn dedup_commit_clears_the_run_snapshot_after_settling() { + let tmp = tempfile::TempDir::new().unwrap(); + let config = test_config(&tmp); + let flow = flow_with_dedup_node("f-clear", "dd"); + store::upsert_flow(&config, &flow).unwrap(); + snapshot_run_dedup_nodes(&config, "f-clear", "run-1", &flow); + + let namespace = dedup_state_namespace("f-clear"); + store::kv_set(&config, &namespace, "dedup:dd:tentative", &json!(["a"])).unwrap(); + + let sub = DedupCommitSubscriber::new(config.clone()); + sub.handle(&DomainEvent::FlowRunFinished { + flow_id: "f-clear".into(), + run_id: "run-1".into(), + status: "completed".into(), + }) + .await; + + assert!( + store::kv_get(&config, &namespace, &run_snapshot_key("run-1")) + .unwrap() + .is_none(), + "a run-scoped snapshot must not outlive the run it describes" + ); +} + +#[tokio::test] +async fn dedup_commit_falls_back_to_the_saved_graph_when_a_run_has_no_snapshot() { + // Runs already in flight when this change ships have no snapshot, as do + // resumed runs. They must keep settling exactly as they did before. + let tmp = tempfile::TempDir::new().unwrap(); + let config = test_config(&tmp); + let flow = flow_with_dedup_node("f-no-snap", "dd"); + store::upsert_flow(&config, &flow).unwrap(); + + let namespace = dedup_state_namespace("f-no-snap"); + store::kv_set(&config, &namespace, "dedup:dd:tentative", &json!(["a"])).unwrap(); + + let sub = DedupCommitSubscriber::new(config.clone()); + sub.handle(&DomainEvent::FlowRunFinished { + flow_id: "f-no-snap".into(), + run_id: "run-legacy".into(), + status: "completed".into(), + }) + .await; + + let committed = store::kv_get(&config, &namespace, "dedup:dd:committed") + .unwrap() + .expect("a snapshot-less run must still settle from the saved graph"); + assert_eq!(committed, json!(["a"])); +} + +#[tokio::test] +async fn dedup_commit_falls_back_to_the_saved_graph_for_a_malformed_snapshot() { + let tmp = tempfile::TempDir::new().unwrap(); + let config = test_config(&tmp); + let flow = flow_with_dedup_node("f-bad-snap", "dd"); + store::upsert_flow(&config, &flow).unwrap(); + + let namespace = dedup_state_namespace("f-bad-snap"); + // Not an array. Must degrade to the saved graph, NOT to "settle + // nothing" — settling nothing would strand the tentative keys. + store::kv_set( + &config, + &namespace, + &run_snapshot_key("run-1"), + &json!("nonsense"), + ) + .unwrap(); + store::kv_set(&config, &namespace, "dedup:dd:tentative", &json!(["a"])).unwrap(); + + let sub = DedupCommitSubscriber::new(config.clone()); + sub.handle(&DomainEvent::FlowRunFinished { + flow_id: "f-bad-snap".into(), + run_id: "run-1".into(), + status: "completed".into(), + }) + .await; + + let committed = store::kv_get(&config, &namespace, "dedup:dd:committed") + .unwrap() + .expect("a malformed snapshot must fall back, not strand the keys"); + assert_eq!(committed, json!(["a"])); +} + +#[tokio::test] +async fn dedup_commit_prefers_the_snapshot_over_a_node_added_after_the_run_started() { + // The converse edit: a dedup node added to the saved graph *after* the + // run started was previously settled here even though the run never + // executed it. With a snapshot the settlement is confined to what ran. + let tmp = tempfile::TempDir::new().unwrap(); + let config = test_config(&tmp); + + let executed = flow_with_dedup_node("f-added", "dd1"); + snapshot_run_dedup_nodes(&config, "f-added", "run-1", &executed); + + // Saved graph now also has `dd2`, which this run never touched. + let mut edited = flow_with_dedup_node("f-added", "dd1"); + edited.graph.nodes.push(dedup_node("dd2")); + store::upsert_flow(&config, &edited).unwrap(); + + let namespace = dedup_state_namespace("f-added"); + store::kv_set(&config, &namespace, "dedup:dd1:tentative", &json!(["a"])).unwrap(); + store::kv_set(&config, &namespace, "dedup:dd2:tentative", &json!(["b"])).unwrap(); + + let sub = DedupCommitSubscriber::new(config.clone()); + sub.handle(&DomainEvent::FlowRunFinished { + flow_id: "f-added".into(), + run_id: "run-1".into(), + status: "completed".into(), + }) + .await; + + assert_eq!( + store::kv_get(&config, &namespace, "dedup:dd1:committed") + .unwrap() + .expect("the executed node must commit"), + json!(["a"]) + ); + assert!( + store::kv_get(&config, &namespace, "dedup:dd2:committed") + .unwrap() + .is_none(), + "a node added after the run started must not be settled by that run" + ); + assert_eq!( + store::kv_get(&config, &namespace, "dedup:dd2:tentative") + .unwrap() + .expect("the unexecuted node's tentative keys must be left alone"), + json!(["b"]) + ); +} + +#[test] +fn flow_state_namespace_matches_the_engine_state_namespace() { + // The subscriber and the engine's `dedup` node MUST agree on this + // string or the host half of the exactly-once contract silently reads + // an empty namespace. Pinned against the literal, not the helper. + assert_eq!(flow_state_namespace("f1"), "flow:f1"); + assert_eq!(flow_state_namespace("f1"), dedup_state_namespace("f1")); +} + +#[test] +fn run_dedup_snapshot_key_cannot_collide_with_a_dedup_nodes_own_keys() { + assert_eq!(run_dedup_snapshot_key("r1"), "dedup::run_nodes:r1"); + // A node id would have to begin with `:` to produce this shape. + assert_ne!( + run_dedup_snapshot_key("r1"), + dedup_node::tentative_key("run_nodes") + ); + assert_ne!( + run_dedup_snapshot_key("r1"), + dedup_node::committed_key("run_nodes") + ); +} + +#[test] +fn dedup_node_ids_in_returns_every_dedup_node_in_graph_order() { + let mut flow = flow_with_dedup_node("f-order", "dd1"); + flow.graph.nodes.push(dedup_node("dd2")); + assert_eq!(dedup_node_ids_in(&flow), vec!["dd1", "dd2"]); + + // Trigger-only graphs yield nothing — the no-op path. + let plain = flow_with_trigger_config("f-order", true, json!({})); + assert!(dedup_node_ids_in(&plain).is_empty()); +} + +// ── the `FlowRunStarted` arm that captures the snapshot ─────────── + +#[tokio::test] +async fn dedup_commit_run_started_snapshots_the_running_graph() { + let tmp = tempfile::TempDir::new().unwrap(); + let config = test_config(&tmp); + store::upsert_flow(&config, &flow_with_dedup_node("f-start", "dd")).unwrap(); + + let sub = DedupCommitSubscriber::new(config.clone()); + sub.handle(&DomainEvent::FlowRunStarted { + flow_id: "f-start".into(), + run_id: "run-start".into(), + }) + .await; + + let snapshot = store::kv_get( + &config, + &dedup_state_namespace("f-start"), + &run_dedup_snapshot_key("run-start"), + ) + .unwrap() + .expect("the FlowRunStarted arm must pin this run's dedup nodes"); + assert_eq!(snapshot, json!(["dd"])); +} + +#[tokio::test] +async fn dedup_commit_run_started_writes_nothing_without_a_dedup_node() { + let tmp = tempfile::TempDir::new().unwrap(); + let config = test_config(&tmp); + let flow = flow_with_trigger_config("f-plain", true, json!({})); + store::upsert_flow(&config, &flow).unwrap(); + + let sub = DedupCommitSubscriber::new(config.clone()); + sub.handle(&DomainEvent::FlowRunStarted { + flow_id: "f-plain".into(), + run_id: "run-plain".into(), + }) + .await; + + // The common case must stay zero-I/O: no key, so `dedup_node_ids` + // keeps reading an unambiguous "no snapshot". + assert!( + store::kv_get( + &config, + &dedup_state_namespace("f-plain"), + &run_dedup_snapshot_key("run-plain"), + ) + .unwrap() + .is_none(), + "a dedup-free graph must not write a snapshot" + ); +} + +#[tokio::test] +async fn dedup_commit_run_started_is_a_noop_for_a_missing_flow() { + let tmp = tempfile::TempDir::new().unwrap(); + let config = test_config(&tmp); + + // Must not panic when the flow was deleted between run start and + // this handler observing the event. + let sub = DedupCommitSubscriber::new(config.clone()); + sub.handle(&DomainEvent::FlowRunStarted { + flow_id: "f-gone".into(), + run_id: "run-gone".into(), + }) + .await; + + assert!( + store::kv_get( + &config, + &dedup_state_namespace("f-gone"), + &run_dedup_snapshot_key("run-gone"), + ) + .unwrap() + .is_none() + ); +} + +/// End-to-end proof of the actual bug in issue #5268 item 2, driven +/// entirely through the event bus: a flow whose `dedup` node is renamed +/// mid-run still settles the node the run really executed. +#[tokio::test] +async fn dedup_commit_settles_renamed_node_end_to_end_through_events() { + let tmp = tempfile::TempDir::new().unwrap(); + let config = test_config(&tmp); + let namespace = dedup_state_namespace("f-e2e"); + store::upsert_flow(&config, &flow_with_dedup_node("f-e2e", "old")).unwrap(); + + let sub = DedupCommitSubscriber::new(config.clone()); + sub.handle(&DomainEvent::FlowRunStarted { + flow_id: "f-e2e".into(), + run_id: "run-e2e".into(), + }) + .await; + + // The run marks an item tentative under the node it actually ran... + store::kv_set(&config, &namespace, "dedup:old:tentative", &json!(["a"])).unwrap(); + // ...then the user renames that node while the run is still in flight. + store::upsert_flow(&config, &flow_with_dedup_node("f-e2e", "new")).unwrap(); + + sub.handle(&DomainEvent::FlowRunFinished { + flow_id: "f-e2e".into(), + run_id: "run-e2e".into(), + status: "completed".into(), + }) + .await; + + assert_eq!( + store::kv_get(&config, &namespace, "dedup:old:committed") + .unwrap() + .expect("the executed node must be committed, not the renamed one"), + json!(["a"]), + "settlement must follow the run's own graph snapshot" + ); + assert!( + store::kv_get(&config, &namespace, "dedup:old:tentative") + .unwrap() + .is_none(), + "tentative keys must be cleared once committed" + ); +} From e2e3cf5da66fcd5bd322138370e740c40c01b299 Mon Sep 17 00:00:00 2001 From: Muhammad Mustaqeem Date: Wed, 5 Aug 2026 23:51:06 +0500 Subject: [PATCH 2/4] fix(flows): settle dedup nodes from the graph the run actually executed MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Closes #5268 (item 2). `DedupCommitSubscriber` settled a finished run's `dedup` nodes by reading the flow's CURRENT saved graph. A flow edited while a run was still in flight was therefore mis-settled: a `dedup` node the run had written `tentative` keys under, but which was deleted or renamed before `FlowRunFinished` fired, was no longer found — so those keys were neither committed nor released, and the items silently reprocessed on the flow's next run. The converse also held: a node added after the run started was settled by a run that never executed it. `DedupCommitSubscriber` now also handles `DomainEvent::FlowRunStarted` (already published by `flows::ops::flows_run{,_detached}` right after the `flow_runs` row insert, and already tagged "cron" so the existing domain filter admits it) and pins that run's dedup node ids into the existing per-flow `flow_state` KV table under `dedup::run_nodes:`. Settlement prefers that snapshot and clears it once the run is settled. - No schema migration: reuses the `flow_state` table the dedup keys already live in. The `dedup::` double-colon prefix cannot collide with the engine's `dedup::` keys. - Zero added I/O for the common case: a graph with no `dedup` node writes no snapshot at all, so an absent key stays an unambiguous "no snapshot". - Falls back to the saved graph when no snapshot exists (runs in flight across the upgrade, resumed runs, or a failed snapshot write), so the degraded path is exactly the historical behaviour. Also moves the test module to `bus_tests.rs` per the crate-wide `#[cfg(test)] #[path = "..."] mod tests;` convention. --- src/openhuman/flows/bus.rs | 1192 ++++++------------------------------ 1 file changed, 194 insertions(+), 998 deletions(-) diff --git a/src/openhuman/flows/bus.rs b/src/openhuman/flows/bus.rs index 6d065b47d2..df58998477 100644 --- a/src/openhuman/flows/bus.rs +++ b/src/openhuman/flows/bus.rs @@ -662,46 +662,38 @@ impl DedupCommitSubscriber { } } - /// The node ids of every `dedup` node in `flow_id`'s saved graph, or an - /// empty vec (logged, not propagated) if the flow can't be loaded — a - /// flow deleted between run-finish and this handler firing, or a - /// transient store error, both degrade to "nothing to settle" rather than - /// panicking the event bus. + /// The node ids of every `dedup` node the finishing run executed. /// - /// **Known limitation (issue #5265, Codex "P2" on the dedup engine PR):** - /// this reads the flow's CURRENT saved definition at settlement time, not - /// a snapshot of the graph the finishing run actually executed. Nothing - /// today persists a per-run graph/node-id snapshot — `prepare_flow_run` - /// loads `Flow` fresh into the spawned run's own task, and that copy is - /// discarded once the run starts; the `FlowRun` row has no `graph` field. - /// If a long-running flow is edited (or deleted) while a run is still in - /// flight: - /// - a `dedup` node the run wrote `tentative` keys under, then deleted or - /// renamed before `FlowRunFinished` fires, is no longer found here — its - /// tentative keys are neither committed nor released, so those items - /// silently retry on the flow's next run (safe-direction: at worst a - /// duplicate, never a lost item, matching this subsystem's existing - /// safe-failure posture — see the module doc's "Best-effort throughout" - /// paragraph); - /// - conversely a `dedup` node id newly added to the saved graph after the - /// run started is settled here even though the run never executed it - /// (a harmless no-op: it has no `tentative` keys to commit/release, see - /// `commit`/`release`'s early returns). + /// Prefers the per-run snapshot [`snapshot_run_dedup_nodes`] wrote at + /// run-start (issue #5268 item 2), so a flow edited *while the run was + /// still in flight* is settled against the graph the run actually ran + /// rather than whatever the flow has since been edited into. This closes + /// the case where a `dedup` node the run wrote `tentative` keys under is + /// deleted or renamed before `FlowRunFinished` fires: previously it was + /// no longer found here, so those keys were neither committed nor + /// released and the items silently retried on the flow's next run. /// - /// Closing this properly means persisting a per-run graph/dedup-node-id - /// snapshot at run-start (`start_flow_run_row` or a sibling write) and - /// having this method read that snapshot instead of `store::get_flow` — - /// a schema + call-site change bigger than this PR's scope; reported as a - /// follow-up rather than attempted here. - fn dedup_node_ids(&self, flow_id: &str) -> Vec { + /// Falls back to the flow's CURRENT saved graph when no snapshot exists — + /// a run that started before this snapshot landed, a resumed run + /// (`flows_resume` re-enters an existing run without a fresh run-start), + /// or a snapshot write that failed. That fallback is exactly the + /// historical behaviour, so the degraded path is never worse than before. + /// + /// Returns an empty vec (logged, not propagated) when neither source + /// yields anything — a flow deleted between run-finish and this handler + /// firing, or a transient store error, both degrade to "nothing to + /// settle" rather than panicking the event bus. + fn dedup_node_ids(&self, namespace: &str, flow_id: &str, run_id: &str) -> Vec { + if let Some(ids) = self.snapshotted_dedup_node_ids(namespace, run_id) { + tracing::trace!( + target: "flows", %flow_id, %run_id, count = ids.len(), + "[dedup-commit] settling against this run's dedup-node snapshot" + ); + return ids; + } + match store::get_flow(&self.config, flow_id) { - Ok(Some(flow)) => flow - .graph - .nodes - .iter() - .filter(|n| n.kind == NodeKind::Dedup) - .map(|n| n.id.clone()) - .collect(), + Ok(Some(flow)) => dedup_node_ids_in(&flow), Ok(None) => { tracing::debug!(target: "flows", %flow_id, "[dedup-commit] flow no longer exists — skipping"); Vec::new() @@ -713,8 +705,58 @@ impl DedupCommitSubscriber { } } + /// The dedup-node-id snapshot recorded for `run_id`, or `None` when this + /// run has none — see [`dedup_node_ids`](Self::dedup_node_ids)'s fallback. + /// + /// An unreadable, non-array or empty snapshot is reported as absent so + /// settlement degrades to the saved graph rather than to "settle + /// nothing": settling nothing would strand tentative keys, whereas the + /// saved graph is the behaviour this subsystem shipped with. + fn snapshotted_dedup_node_ids(&self, namespace: &str, run_id: &str) -> Option> { + let key = run_dedup_snapshot_key(run_id); + let value = match store::kv_get(&self.config, namespace, &key) { + Ok(Some(value)) => value, + Ok(None) => return None, + Err(e) => { + tracing::warn!( + target: "flows", %namespace, %run_id, error = %e, + "[dedup-commit] failed to read this run's dedup snapshot — falling back to \ + the flow's saved graph" + ); + return None; + } + }; + + let ids: Vec = value + .as_array()? + .iter() + .filter_map(Value::as_str) + .map(str::to_string) + .collect(); + if ids.is_empty() { + return None; + } + Some(ids) + } + + /// Drops this run's dedup-node snapshot once the run has been settled. + /// + /// Best-effort: a failed delete leaves one small run-scoped KV row + /// behind, which nothing reads again (snapshots are keyed by `run_id`, and + /// run ids are never reused). + fn clear_run_snapshot(&self, namespace: &str, flow_id: &str, run_id: &str) { + if let Err(e) = store::kv_delete(&self.config, namespace, &run_dedup_snapshot_key(run_id)) { + tracing::warn!( + target: "flows", %flow_id, %run_id, error = %e, + "[dedup-commit] failed to clear this run's dedup snapshot — harmless: the key is \ + run-scoped and is never read again" + ); + } + } + async fn handle_finished(&self, flow_id: &str, run_id: &str, status: &str) { - let node_ids = self.dedup_node_ids(flow_id); + let namespace = flow_state_namespace(flow_id); + let node_ids = self.dedup_node_ids(&namespace, flow_id, run_id); if node_ids.is_empty() { tracing::trace!(target: "flows", %flow_id, %run_id, %status, "[dedup-commit] no dedup nodes in this flow — nothing to settle"); return; @@ -738,7 +780,6 @@ impl DedupCommitSubscriber { tracing::trace!(target: "flows", %flow_id, %run_id, "[dedup-commit] acquired per-flow commit lock"); self.maybe_test_delay().await; - let namespace = format!("flow:{flow_id}"); for node_id in node_ids { if success { self.commit(&namespace, &node_id, flow_id, run_id); @@ -747,6 +788,12 @@ impl DedupCommitSubscriber { } } + // Still inside the per-flow lock: the snapshot is this run's alone, so + // no other run contends for it, but keeping the delete here means a + // settled run never leaves a readable snapshot behind for a + // late-arriving duplicate `FlowRunFinished` to settle a second time. + self.clear_run_snapshot(&namespace, flow_id, run_id); + drop(lock_guard); tracing::trace!(target: "flows", %flow_id, %run_id, "[dedup-commit] released per-flow commit lock"); } @@ -824,22 +871,120 @@ impl EventHandler for DedupCommitSubscriber { fn domains(&self) -> Option<&[&str]> { // Same reasoning as `FlowRunDigestSubscriber::domains` just above: - // `FlowRunFinished` is tagged `"cron"` by `DomainEvent::domain()`. + // `FlowRunStarted` and `FlowRunFinished` are both tagged `"cron"` by + // `DomainEvent::domain()` — they share a single match arm there. Some(&["cron"]) } async fn handle(&self, event: &DomainEvent) { - if let DomainEvent::FlowRunFinished { - flow_id, - run_id, - status, - } = event - { - self.handle_finished(flow_id, run_id, status).await; + match event { + // Pin the `dedup` nodes THIS run is about to execute, so + // settlement at `FlowRunFinished` uses the graph the run really + // ran rather than whatever the flow has been edited into by the + // time it finishes (issue #5268 item 2). + // + // Reading the saved graph here rather than having `flows::ops` + // hand us its in-hand `Flow` keeps the whole fix inside this + // subscriber, and `FlowRunStarted` is published from + // `flows::ops::flows_run{,_detached}` immediately after the + // `flow_runs` row insert — before the engine executes a single + // node — so the graph read back here is the one the run is + // starting with. + // + // Best-effort like everything else in this subscriber: a load + // failure is logged and swallowed, degrading settlement to the + // historical "use the current saved graph" behaviour rather than + // disturbing the run. + DomainEvent::FlowRunStarted { flow_id, run_id } => { + match store::get_flow(&self.config, flow_id) { + Ok(Some(flow)) => { + snapshot_run_dedup_nodes(&self.config, flow_id, run_id, &flow); + } + // Nothing to pin, and not an error worth logging: the + // fallback already handles a flow that isn't there. + Ok(None) => {} + Err(e) => tracing::warn!( + target: "flows", %flow_id, %run_id, error = %e, + "[dedup-commit] could not load flow to snapshot its dedup nodes at run \ + start — settlement will fall back to the flow's saved graph" + ), + } + } + DomainEvent::FlowRunFinished { + flow_id, + run_id, + status, + } => self.handle_finished(flow_id, run_id, status).await, + _ => {} } } } +/// The `StateStore` namespace a flow's `dedup` keys live in. +/// +/// MUST match `tinyflows::caps::build_capabilities`'s `state_namespace` +/// (`src/openhuman/flows/tinyflows/caps.rs`) — colliding with the very keys +/// the engine's own `dedup` node reads and writes during the run is the entire +/// point. Deliberately NOT [`flow_namespace`], which is the flow's *memory* +/// namespace and a different string. +fn flow_state_namespace(flow_id: &str) -> String { + format!("flow:{flow_id}") +} + +/// KV key holding the dedup-node-id snapshot for a single run (issue #5268). +/// +/// The `dedup::` prefix (double colon) cannot collide with the engine's own +/// `dedup::` keys — reaching this shape would +/// need a node id beginning with `:`, which the graph schema does not admit. +fn run_dedup_snapshot_key(run_id: &str) -> String { + format!("dedup::run_nodes:{run_id}") +} + +/// The ids of every `dedup` node in `flow`'s graph, in graph order. +fn dedup_node_ids_in(flow: &Flow) -> Vec { + flow.graph + .nodes + .iter() + .filter(|n| n.kind == NodeKind::Dedup) + .map(|n| n.id.clone()) + .collect() +} + +/// Records the ids of every `dedup` node in the graph a run is about to +/// execute, so [`DedupCommitSubscriber`] settles that run against the graph it +/// really ran rather than whatever the flow has been edited into by the time it +/// finishes (issue #5268 item 2). +/// +/// Called from [`DedupCommitSubscriber`]'s `FlowRunStarted` arm, which fires +/// from `flows::ops::flows_run{,_detached}` immediately after the initial +/// `flow_runs` row insert and before the engine executes any node. +/// +/// Writes nothing when the graph has no `dedup` node — the overwhelming +/// majority of flows — so the common path costs one graph scan and zero I/O, +/// and an absent key stays an unambiguous "no snapshot" for +/// [`DedupCommitSubscriber::dedup_node_ids`]'s fallback. +/// +/// Best-effort like the rest of this subsystem: a failed write is logged and +/// swallowed, degrading settlement to the historical +/// load-the-current-saved-graph behaviour rather than failing the run. +fn snapshot_run_dedup_nodes(config: &Config, flow_id: &str, run_id: &str, flow: &Flow) { + let node_ids = dedup_node_ids_in(flow); + if node_ids.is_empty() { + return; + } + + let namespace = flow_state_namespace(flow_id); + let key = run_dedup_snapshot_key(run_id); + let value = Value::Array(node_ids.into_iter().map(Value::String).collect()); + if let Err(e) = store::kv_set(config, &namespace, &key, &value) { + tracing::warn!( + target: "flows", %flow_id, %run_id, error = %e, + "[dedup-commit] failed to snapshot this run's dedup nodes — settlement will fall back \ + to the flow's saved graph at finish time" + ); + } +} + /// Loads a `dedup` node's key set (stored as a JSON array of strings) from /// the flow-state KV table. Mirrors /// `tinyflows::nodes::control_flow::dedup`'s own key-set loader: a missing @@ -881,954 +1026,5 @@ fn store_key_set( } #[cfg(test)] -mod tests { - use super::*; - use crate::openhuman::flows::Flow; - use crate::openhuman::inference::embeddings::NoopEmbedding; - use crate::openhuman::memory::store::UnifiedMemory; - use serde_json::json; - use tinyflows::model::{Node, NodeKind, WorkflowGraph}; - - /// A directly-constructed, isolated [`Memory`] for the digest tests — NOT - /// the process-global `OnceLock` client. The global is one-shot, so an - /// earlier test in the same binary may already have bound it to a different - /// workspace, making `global::init(..)` here a silent no-op (see - /// `memory::global`'s own test notes). Injecting this instance into the - /// subscriber via [`FlowRunDigestSubscriber::with_memory`] makes writes and - /// read-backs go through the SAME store deterministically — the same shape - /// `flows::memory_tools`' tests use. - fn digest_test_memory(tmp: &tempfile::TempDir) -> Arc { - Arc::new(UnifiedMemory::new(tmp.path(), Arc::new(NoopEmbedding), None).unwrap()) - } - - fn test_config(tmp: &tempfile::TempDir) -> Arc { - let config = Config { - workspace_dir: tmp.path().join("workspace"), - action_dir: tmp.path().join("workspace"), - config_path: tmp.path().join("config.toml"), - ..Config::default() - }; - std::fs::create_dir_all(&config.workspace_dir).unwrap(); - Arc::new(config) - } - - fn trigger_node(config: Value) -> Node { - Node { - id: "t".to_string(), - kind: NodeKind::Trigger, - type_version: 1, - name: "Trigger".to_string(), - config, - ports: Vec::new(), - position: None, - } - } - - fn flow_with_trigger_config(id: &str, enabled: bool, trigger_config: Value) -> Flow { - Flow { - id: id.to_string(), - name: id.to_string(), - enabled, - graph: WorkflowGraph { - nodes: vec![trigger_node(trigger_config)], - ..Default::default() - }, - created_at: "2026-01-01T00:00:00Z".to_string(), - updated_at: "2026-01-01T00:00:00Z".to_string(), - last_run_at: None, - last_status: None, - require_approval: false, - } - } - - fn dedup_node(id: &str) -> Node { - Node { - id: id.to_string(), - kind: NodeKind::Dedup, - type_version: 1, - name: id.to_string(), - config: json!({ "key": "=item.id" }), - ports: Vec::new(), - position: None, - } - } - - /// A saved flow with a `trigger` node plus one `dedup` node with id - /// `dedup_id` — the minimal graph [`DedupCommitSubscriber::dedup_node_ids`] - /// needs to find something to settle. - fn flow_with_dedup_node(id: &str, dedup_id: &str) -> Flow { - Flow { - id: id.to_string(), - name: id.to_string(), - enabled: true, - graph: WorkflowGraph { - nodes: vec![trigger_node(json!({})), dedup_node(dedup_id)], - ..Default::default() - }, - created_at: "2026-01-01T00:00:00Z".to_string(), - updated_at: "2026-01-01T00:00:00Z".to_string(), - last_run_at: None, - last_status: None, - require_approval: false, - } - } - - #[test] - fn pinned_trigger_inputs_reads_values_an_author_fixed_for_unattended_runs() { - let flow = flow_with_trigger_config( - "f1", - true, - json!({ - "trigger_kind": "schedule", - "schedule": "0 9 * * *", - "inputs": { "repo": "acme/api", "depth": 3 } - }), - ); - let inputs = pinned_trigger_inputs(&flow); - assert_eq!(inputs["repo"], json!("acme/api")); - assert_eq!(inputs["depth"], json!(3)); - } - - #[test] - fn pinned_trigger_inputs_is_empty_when_unset_or_malformed() { - // Empty, not an error: a flow declaring no inputs (the overwhelming - // majority) must keep dispatching on a tick exactly as before, and a - // malformed value is caught downstream by `prepare_flow_run`, which - // reports it against the flow's actual declarations. - for cfg in [ - json!({ "trigger_kind": "schedule" }), - json!({ "trigger_kind": "schedule", "inputs": null }), - json!({ "trigger_kind": "schedule", "inputs": ["repo"] }), - ] { - let flow = flow_with_trigger_config("f1", true, cfg.clone()); - assert!( - pinned_trigger_inputs(&flow).is_empty(), - "expected no pinned inputs for {cfg}" - ); - } - } - - #[test] - fn pinned_trigger_inputs_is_empty_for_a_graph_with_no_trigger() { - let mut flow = flow_with_trigger_config("f1", true, json!({ "trigger_kind": "schedule" })); - flow.graph.nodes.clear(); - assert!(pinned_trigger_inputs(&flow).is_empty()); - } - - #[test] - fn name_and_domains_are_stable() { - let tmp = tempfile::TempDir::new().unwrap(); - let sub = FlowTriggerSubscriber::new(test_config(&tmp)); - assert_eq!(sub.name(), "flows::trigger"); - assert_eq!( - sub.domains(), - Some(&["cron", "composio", "webhook", "system"][..]) - ); - } - - #[tokio::test] - async fn handle_does_not_panic_on_arbitrary_events() { - let tmp = tempfile::TempDir::new().unwrap(); - let sub = FlowTriggerSubscriber::new(test_config(&tmp)); - sub.handle(&DomainEvent::CronJobTriggered { - job_id: "j1".into(), - job_name: "test".into(), - job_type: "shell".into(), - }) - .await; - sub.handle(&DomainEvent::FlowScheduleTick { - flow_id: "missing-flow".into(), - }) - .await; - } - - #[test] - fn extract_trigger_kind_reads_schedule() { - let flow = flow_with_trigger_config( - "f1", - true, - json!({ "trigger_kind": "schedule", "schedule": "0 9 * * *" }), - ); - assert!(matches!( - extract_trigger_kind(&flow), - Some(TriggerKind::Schedule) - )); - } - - #[test] - fn extract_trigger_kind_none_for_missing_discriminator() { - let flow = flow_with_trigger_config("f1", true, json!({})); - assert!(extract_trigger_kind(&flow).is_none()); - } - - #[test] - fn extract_trigger_kind_none_for_invalid_discriminator() { - let flow = flow_with_trigger_config("f1", true, json!({ "trigger_kind": "not_a_kind" })); - assert!(extract_trigger_kind(&flow).is_none()); - } - - #[test] - fn matches_app_event_requires_toolkit_and_slug_match() { - let flow = flow_with_trigger_config( - "f1", - true, - json!({ "trigger_kind": "app_event", "toolkit": "gmail", "trigger_slug": "GMAIL_NEW_GMAIL_MESSAGE" }), - ); - assert!(matches_app_event(&flow, "gmail", "GMAIL_NEW_GMAIL_MESSAGE")); - // Case-insensitive. - assert!(matches_app_event(&flow, "Gmail", "gmail_new_gmail_message")); - // Wrong toolkit or slug does not match. - assert!(!matches_app_event( - &flow, - "slack", - "GMAIL_NEW_GMAIL_MESSAGE" - )); - assert!(!matches_app_event(&flow, "gmail", "SLACK_NEW_MESSAGE")); - } - - #[test] - fn matches_app_event_false_for_non_app_event_trigger() { - let flow = flow_with_trigger_config( - "f1", - true, - json!({ "trigger_kind": "schedule", "schedule": "0 9 * * *" }), - ); - assert!(!matches_app_event( - &flow, - "gmail", - "GMAIL_NEW_GMAIL_MESSAGE" - )); - } - - #[tokio::test] - async fn handle_app_event_ignores_disabled_flows() { - let tmp = tempfile::TempDir::new().unwrap(); - let config = test_config(&tmp); - let flow = flow_with_trigger_config( - "disabled-flow", - false, - json!({ "trigger_kind": "app_event", "toolkit": "gmail", "trigger_slug": "GMAIL_NEW_GMAIL_MESSAGE" }), - ); - crate::openhuman::flows::store::upsert_flow(&config, &flow).unwrap(); - - // `list_enabled_flows` must not surface the disabled flow at all — - // proves the subscriber's dispatch source already excludes it, - // rather than asserting on a spawned background task's side effect. - let (enabled, skipped) = - crate::openhuman::flows::store::list_enabled_flows(&config).unwrap(); - assert!(enabled.is_empty()); - assert_eq!(skipped, 0); - } - - #[tokio::test] - async fn handle_schedule_tick_ignores_disabled_flow() { - let tmp = tempfile::TempDir::new().unwrap(); - let config = test_config(&tmp); - let flow = flow_with_trigger_config( - "sched-flow", - false, - json!({ "trigger_kind": "schedule", "schedule": "0 9 * * *" }), - ); - crate::openhuman::flows::store::upsert_flow(&config, &flow).unwrap(); - - let sub = FlowTriggerSubscriber::new(config.clone()); - // Must not panic and must not spawn a run for a disabled flow — we - // can't directly observe "no run happened" without a full flows_run - // fixture, but this exercises the early-return path without error. - sub.handle(&DomainEvent::FlowScheduleTick { - flow_id: "sched-flow".into(), - }) - .await; - } - - // ── in-flight dedupe (CodeRabbit finding B) ───────────────────── - - #[test] - fn try_acquire_dispatch_skips_a_flow_already_in_flight() { - let tmp = tempfile::TempDir::new().unwrap(); - let sub = FlowTriggerSubscriber::new(test_config(&tmp)); - - let guard = sub - .try_acquire_dispatch("f1") - .expect("first claim for f1 should succeed"); - assert!( - sub.try_acquire_dispatch("f1").is_none(), - "a second claim for the same flow while the first is held must be skipped" - ); - - // A different flow is unaffected. - assert!(sub.try_acquire_dispatch("f2").is_some()); - - drop(guard); - assert!( - sub.try_acquire_dispatch("f1").is_some(), - "dropping the guard must release the claim so f1 can run again" - ); - } - - #[test] - fn default_constructs_the_same_as_new() { - let tmp = tempfile::TempDir::new().unwrap(); - let config = test_config(&tmp); - let a = FlowTriggerSubscriber::new(config.clone()); - let b = FlowTriggerSubscriber::new(config); - assert_eq!(a.name(), b.name()); - } - - // ── FlowRunDigestSubscriber ───────────────────────────────────── - - #[test] - fn digest_name_and_domains_are_stable() { - let tmp = tempfile::TempDir::new().unwrap(); - let sub = FlowRunDigestSubscriber::new(test_config(&tmp)); - assert_eq!(sub.name(), "flows::digest"); - assert_eq!(sub.domains(), Some(&["cron"][..])); - } - - #[tokio::test] - async fn digest_handle_does_not_panic_on_unrelated_events() { - let tmp = tempfile::TempDir::new().unwrap(); - let sub = FlowRunDigestSubscriber::new(test_config(&tmp)); - // Must not panic, and must not touch the memory layer at all, for - // any event other than `FlowRunFinished`. - sub.handle(&DomainEvent::CronJobTriggered { - job_id: "j1".into(), - job_name: "test".into(), - job_type: "shell".into(), - }) - .await; - } - - #[tokio::test] - async fn digest_ignores_failed_run() { - let tmp = tempfile::TempDir::new().unwrap(); - let config = test_config(&tmp); - let memory = digest_test_memory(&tmp); - - let flow = flow_with_trigger_config("f-failed", true, json!({})); - store::upsert_flow(&config, &flow).unwrap(); - store::insert_flow_run( - &config, - "run-failed", - "f-failed", - "thread-failed", - "2026-01-01T00:00:00Z", - ) - .unwrap(); - store::finish_flow_run( - &config, - "run-failed", - "failed", - "2026-01-01T00:05:00Z", - &[], - &[], - Some("boom"), - None, - ) - .unwrap(); - - let sub = FlowRunDigestSubscriber::with_memory(config, memory.clone()); - sub.handle(&DomainEvent::FlowRunFinished { - flow_id: "f-failed".into(), - run_id: "run-failed".into(), - status: "failed".into(), - }) - .await; - - let entry = memory - .get(&flow_namespace("f-failed"), "run_digest:run-failed") - .await - .unwrap(); - assert!( - entry.is_none(), - "a failed run must never produce a run_digest entry" - ); - } - - #[tokio::test] - async fn digest_ignores_cancelled_run() { - let tmp = tempfile::TempDir::new().unwrap(); - let config = test_config(&tmp); - let memory = digest_test_memory(&tmp); - - let flow = flow_with_trigger_config("f-cancelled", true, json!({})); - store::upsert_flow(&config, &flow).unwrap(); - store::insert_flow_run( - &config, - "run-cancelled", - "f-cancelled", - "thread-cancelled", - "2026-01-01T00:00:00Z", - ) - .unwrap(); - store::finish_flow_run( - &config, - "run-cancelled", - "cancelled", - "2026-01-01T00:05:00Z", - &[], - &[], - None, - None, - ) - .unwrap(); - - let sub = FlowRunDigestSubscriber::with_memory(config, memory.clone()); - sub.handle(&DomainEvent::FlowRunFinished { - flow_id: "f-cancelled".into(), - run_id: "run-cancelled".into(), - status: "cancelled".into(), - }) - .await; - - let entry = memory - .get(&flow_namespace("f-cancelled"), "run_digest:run-cancelled") - .await - .unwrap(); - assert!(entry.is_none()); - } - - #[tokio::test] - async fn digest_writes_run_digest_entry_for_completed_run() { - let tmp = tempfile::TempDir::new().unwrap(); - let config = test_config(&tmp); - let memory = digest_test_memory(&tmp); - - let flow = flow_with_trigger_config("f-ok", true, json!({})); - store::upsert_flow(&config, &flow).unwrap(); - store::insert_flow_run( - &config, - "run-ok", - "f-ok", - "thread-ok", - "2026-01-01T00:00:00Z", - ) - .unwrap(); - let step = crate::openhuman::flows::FlowRunStep { - node_id: "n1".to_string(), - output: json!({ "sent": 3 }), - port: None, - status: Some("success".to_string()), - duration_ms: Some(12), - diagnostics: Vec::new(), - }; - store::finish_flow_run( - &config, - "run-ok", - "completed", - "2026-01-01T00:05:00Z", - &[step], - &[], - None, - None, - ) - .unwrap(); - - let sub = FlowRunDigestSubscriber::with_memory(config, memory.clone()); - sub.handle(&DomainEvent::FlowRunFinished { - flow_id: "f-ok".into(), - run_id: "run-ok".into(), - status: "completed".into(), - }) - .await; - - let entry = memory - .get(&flow_namespace("f-ok"), "run_digest:run-ok") - .await - .unwrap() - .expect("completed run must produce a run_digest entry"); - assert_eq!(entry.taint, MemoryTaint::ExternalSync); - assert!(entry.content.contains("f-ok")); - assert!(entry.content.contains("completed")); - assert!(entry.content.contains("n1")); - assert!(entry.content.chars().count() <= DIGEST_MAX_CHARS); - } - - #[tokio::test] - async fn digest_treats_completed_with_warnings_as_success() { - let tmp = tempfile::TempDir::new().unwrap(); - let config = test_config(&tmp); - let memory = digest_test_memory(&tmp); - - let flow = flow_with_trigger_config("f-warn", true, json!({})); - store::upsert_flow(&config, &flow).unwrap(); - store::insert_flow_run( - &config, - "run-warn", - "f-warn", - "thread-warn", - "2026-01-01T00:00:00Z", - ) - .unwrap(); - store::finish_flow_run( - &config, - "run-warn", - "completed_with_warnings", - "2026-01-01T00:05:00Z", - &[], - &[], - None, - None, - ) - .unwrap(); - - let sub = FlowRunDigestSubscriber::with_memory(config, memory.clone()); - sub.handle(&DomainEvent::FlowRunFinished { - flow_id: "f-warn".into(), - run_id: "run-warn".into(), - status: "completed_with_warnings".into(), - }) - .await; - - let entry = memory - .get(&flow_namespace("f-warn"), "run_digest:run-warn") - .await - .unwrap(); - assert!(entry.is_some()); - } - - #[test] - fn truncate_chars_bounds_output_and_marks_truncation() { - let long = "x".repeat(50); - let truncated = truncate_chars(&long, 10); - assert_eq!(truncated.chars().count(), 10); - assert!(truncated.ends_with('…')); - - let short = "hello"; - assert_eq!(truncate_chars(short, 10), "hello"); - } - - #[test] - fn render_run_digest_is_bounded_and_includes_key_fields() { - let run = FlowRun { - id: "run-1".to_string(), - flow_id: "f1".to_string(), - thread_id: "thread-1".to_string(), - status: "completed".to_string(), - started_at: "2026-01-01T00:00:00Z".to_string(), - finished_at: Some("2026-01-01T00:05:00Z".to_string()), - steps: vec![crate::openhuman::flows::FlowRunStep { - node_id: "n1".to_string(), - output: json!({ "ok": true }), - port: None, - status: Some("success".to_string()), - duration_ms: Some(5), - diagnostics: Vec::new(), - }], - pending_approvals: Vec::new(), - error: None, - graph_hash: None, - }; - let digest = render_run_digest("My Flow", &run); - assert!(digest.contains("My Flow")); - assert!(digest.contains("completed")); - assert!(digest.contains("n1")); - assert!(digest.chars().count() <= DIGEST_MAX_CHARS); - } - - // ── DedupCommitSubscriber ──────────────────────────────────────── - - fn dedup_state_namespace(flow_id: &str) -> String { - // MUST match `tinyflows::build_capabilities`'s `state_namespace` - // (`src/openhuman/flows/tinyflows/caps.rs`) — this test asserts the - // subscriber collides with the SAME keys the engine's `dedup` node - // itself reads/writes, not just "some" namespace. - format!("flow:{flow_id}") - } - - #[test] - fn dedup_commit_name_and_domains_are_stable() { - let tmp = tempfile::TempDir::new().unwrap(); - let sub = DedupCommitSubscriber::new(test_config(&tmp)); - assert_eq!(sub.name(), "flows::dedup_commit"); - assert_eq!(sub.domains(), Some(&["cron"][..])); - } - - #[tokio::test] - async fn dedup_commit_ignores_unrelated_events() { - let tmp = tempfile::TempDir::new().unwrap(); - let sub = DedupCommitSubscriber::new(test_config(&tmp)); - // Must not panic for any event other than `FlowRunFinished`. - sub.handle(&DomainEvent::CronJobTriggered { - job_id: "j1".into(), - job_name: "test".into(), - job_type: "shell".into(), - }) - .await; - } - - #[tokio::test] - async fn dedup_commit_flow_with_no_dedup_nodes_is_a_noop() { - let tmp = tempfile::TempDir::new().unwrap(); - let config = test_config(&tmp); - let flow = flow_with_trigger_config("f-no-dedup", true, json!({})); - store::upsert_flow(&config, &flow).unwrap(); - - let sub = DedupCommitSubscriber::new(config); - // Must not panic when the flow has no `dedup` node at all. - sub.handle(&DomainEvent::FlowRunFinished { - flow_id: "f-no-dedup".into(), - run_id: "run-1".into(), - status: "completed".into(), - }) - .await; - } - - #[tokio::test] - async fn dedup_commit_unions_tentative_into_committed_and_clears_tentative_on_success() { - let tmp = tempfile::TempDir::new().unwrap(); - let config = test_config(&tmp); - let flow = flow_with_dedup_node("f-ok", "dd"); - store::upsert_flow(&config, &flow).unwrap(); - - let namespace = dedup_state_namespace("f-ok"); - store::kv_set(&config, &namespace, "dedup:dd:committed", &json!(["a"])).unwrap(); - store::kv_set( - &config, - &namespace, - "dedup:dd:tentative", - &json!(["b", "c"]), - ) - .unwrap(); - - let sub = DedupCommitSubscriber::new(config.clone()); - sub.handle(&DomainEvent::FlowRunFinished { - flow_id: "f-ok".into(), - run_id: "run-ok".into(), - status: "completed".into(), - }) - .await; - - let committed = store::kv_get(&config, &namespace, "dedup:dd:committed") - .unwrap() - .expect("committed key must still exist"); - let mut committed: Vec<&str> = committed - .as_array() - .unwrap() - .iter() - .map(|v| v.as_str().unwrap()) - .collect(); - committed.sort_unstable(); - assert_eq!(committed, vec!["a", "b", "c"], "committed = union"); - - assert!( - store::kv_get(&config, &namespace, "dedup:dd:tentative") - .unwrap() - .is_none(), - "tentative must be cleared after a successful commit" - ); - } - - #[tokio::test] - async fn dedup_commit_treats_completed_with_warnings_as_success() { - let tmp = tempfile::TempDir::new().unwrap(); - let config = test_config(&tmp); - let flow = flow_with_dedup_node("f-warn", "dd"); - store::upsert_flow(&config, &flow).unwrap(); - - let namespace = dedup_state_namespace("f-warn"); - store::kv_set(&config, &namespace, "dedup:dd:tentative", &json!(["x"])).unwrap(); - - let sub = DedupCommitSubscriber::new(config.clone()); - sub.handle(&DomainEvent::FlowRunFinished { - flow_id: "f-warn".into(), - run_id: "run-warn".into(), - status: "completed_with_warnings".into(), - }) - .await; - - let committed = store::kv_get(&config, &namespace, "dedup:dd:committed") - .unwrap() - .expect("completed_with_warnings must still commit"); - assert_eq!(committed, json!(["x"])); - assert!(store::kv_get(&config, &namespace, "dedup:dd:tentative") - .unwrap() - .is_none()); - } - - #[tokio::test] - async fn dedup_commit_releases_tentative_without_touching_committed_on_failure() { - let tmp = tempfile::TempDir::new().unwrap(); - let config = test_config(&tmp); - let flow = flow_with_dedup_node("f-failed", "dd"); - store::upsert_flow(&config, &flow).unwrap(); - - let namespace = dedup_state_namespace("f-failed"); - store::kv_set(&config, &namespace, "dedup:dd:committed", &json!(["a"])).unwrap(); - store::kv_set(&config, &namespace, "dedup:dd:tentative", &json!(["b"])).unwrap(); - - let sub = DedupCommitSubscriber::new(config.clone()); - sub.handle(&DomainEvent::FlowRunFinished { - flow_id: "f-failed".into(), - run_id: "run-failed".into(), - status: "failed".into(), - }) - .await; - - assert_eq!( - store::kv_get(&config, &namespace, "dedup:dd:committed") - .unwrap() - .unwrap(), - json!(["a"]), - "committed must be untouched by a failed run" - ); - assert!( - store::kv_get(&config, &namespace, "dedup:dd:tentative") - .unwrap() - .is_none(), - "tentative must be released (cleared) on failure so the item retries" - ); - } - - #[tokio::test] - async fn dedup_commit_releases_tentative_on_cancelled_and_interrupted() { - for status in ["cancelled", "interrupted"] { - let tmp = tempfile::TempDir::new().unwrap(); - let config = test_config(&tmp); - let flow_id = format!("f-{status}"); - let flow = flow_with_dedup_node(&flow_id, "dd"); - store::upsert_flow(&config, &flow).unwrap(); - - let namespace = dedup_state_namespace(&flow_id); - store::kv_set(&config, &namespace, "dedup:dd:tentative", &json!(["z"])).unwrap(); - - let sub = DedupCommitSubscriber::new(config.clone()); - sub.handle(&DomainEvent::FlowRunFinished { - flow_id: flow_id.clone(), - run_id: format!("run-{status}"), - status: status.to_string(), - }) - .await; - - assert!( - store::kv_get(&config, &namespace, "dedup:dd:committed") - .unwrap() - .is_none(), - "status {status} must never commit" - ); - assert!( - store::kv_get(&config, &namespace, "dedup:dd:tentative") - .unwrap() - .is_none(), - "status {status} must release tentative" - ); - } - } - - #[tokio::test] - async fn dedup_commit_two_dedup_nodes_settle_independently() { - let tmp = tempfile::TempDir::new().unwrap(); - let config = test_config(&tmp); - let flow = Flow { - id: "f-multi".to_string(), - name: "f-multi".to_string(), - enabled: true, - graph: WorkflowGraph { - nodes: vec![ - trigger_node(json!({})), - dedup_node("dd1"), - dedup_node("dd2"), - ], - ..Default::default() - }, - created_at: "2026-01-01T00:00:00Z".to_string(), - updated_at: "2026-01-01T00:00:00Z".to_string(), - last_run_at: None, - last_status: None, - require_approval: false, - }; - store::upsert_flow(&config, &flow).unwrap(); - - let namespace = dedup_state_namespace("f-multi"); - store::kv_set(&config, &namespace, "dedup:dd1:tentative", &json!(["a"])).unwrap(); - store::kv_set(&config, &namespace, "dedup:dd2:tentative", &json!(["b"])).unwrap(); - - let sub = DedupCommitSubscriber::new(config.clone()); - sub.handle(&DomainEvent::FlowRunFinished { - flow_id: "f-multi".into(), - run_id: "run-multi".into(), - status: "completed".into(), - }) - .await; - - assert_eq!( - store::kv_get(&config, &namespace, "dedup:dd1:committed") - .unwrap() - .unwrap(), - json!(["a"]) - ); - assert_eq!( - store::kv_get(&config, &namespace, "dedup:dd2:committed") - .unwrap() - .unwrap(), - json!(["b"]) - ); - } - - // ── per-flow commit serialization (issue #5265) ─────────────────── - // - // CodeRabbit "Major" on the dedup engine PR: the commit's - // load(committed)+union(tentative)+store(committed) is a - // read-modify-write, not a CAS. Two overlapping `FlowRunFinished` - // events for the SAME flow could otherwise interleave and have the - // second writer's store clobber the first writer's union, silently - // losing that run's committed keys. `handle_finished` now serializes - // settlement per `flow_id` via `FLOW_COMMIT_LOCKS`. - // - // Two tests, deliberately split: - // - // - `..._never_runs_two_commits_for_the_same_flow_concurrently` spawns a - // burst of genuinely overlapping `FlowRunFinished` events for the SAME - // flow_id and proves the LOCK itself provides mutual exclusion (the - // high-water mark of concurrently-active critical sections never - // exceeds 1) — this is the "spawn two tasks contending on the same - // flow_id" case. - // - `..._serial_commits_for_the_same_flow_accumulate_via_union` proves - // the property that mutual exclusion protects: settling run after run - // for the same node never clobbers an earlier run's committed keys — - // each contributes to the union. - // - // These are split rather than combined into one "two runs with two - // different tentative sets, truly concurrently, assert union" test - // because `tentative` is a single shared KV row per node (not - // per-run) — forcing two *different* tentative contents to both survive - // a genuinely simultaneous read would require injecting a write from - // outside `handle_finished` in the middle of its critical section, which - // instead exercises the SEPARATE, still-open node-side race (the - // `dedup` node's own in-run `tentative` read-modify-write, documented on - // `DedupCommitSubscriber` above as explicitly NOT fixed by this lock). - // Together, the two tests below establish the same guarantee end to - // end: the lock enforces serialization (test 1), and serialization is - // sufficient for correctness (test 2). - - #[tokio::test] - async fn dedup_commit_never_runs_two_commits_for_the_same_flow_concurrently() { - let tmp = tempfile::TempDir::new().unwrap(); - let config = test_config(&tmp); - let flow = flow_with_dedup_node("f-race", "dd"); - store::upsert_flow(&config, &flow).unwrap(); - - let namespace = dedup_state_namespace("f-race"); - store::kv_set(&config, &namespace, "dedup:dd:tentative", &json!(["seed"])).unwrap(); - - // Arm the test-only scheduling hook (see `CommitTestHooks`): every - // `handle_finished` call sleeps briefly while holding the per-flow - // lock, and records how many calls are concurrently inside that - // window. Instance-scoped (not a global static) so this doesn't - // interfere with — or get polluted by — unrelated tests that cargo - // runs concurrently on other threads. Without a correctly-scoped - // lock, a burst of overlapping `FlowRunFinished` events for the SAME - // flow_id would pile up inside the critical section together - // instead of queuing. - let hooks = Arc::new(CommitTestHooks::default()); - hooks - .delay_ms - .store(20, std::sync::atomic::Ordering::SeqCst); - - let sub = Arc::new(DedupCommitSubscriber::with_test_hooks( - config.clone(), - hooks.clone(), - )); - let mut handles = Vec::new(); - for i in 0..5 { - let sub = sub.clone(); - handles.push(tokio::spawn(async move { - sub.handle(&DomainEvent::FlowRunFinished { - flow_id: "f-race".into(), - run_id: format!("run-{i}"), - status: "completed".into(), - }) - .await; - })); - } - for handle in handles { - handle.await.unwrap(); - } - - assert_eq!( - hooks.concurrent.load(std::sync::atomic::Ordering::SeqCst), - 0, - "every critical-section entry must have a matching exit" - ); - assert_eq!( - hooks - .max_concurrent - .load(std::sync::atomic::Ordering::SeqCst), - 1, - "the per-flow lock must serialize overlapping FlowRunFinished handling for the \ - same flow_id — at most one commit critical section may be active at a time" - ); - } - - #[tokio::test] - async fn dedup_commit_serial_commits_for_the_same_flow_accumulate_via_union() { - let tmp = tempfile::TempDir::new().unwrap(); - let config = test_config(&tmp); - let flow = flow_with_dedup_node("f-serial", "dd"); - store::upsert_flow(&config, &flow).unwrap(); - - let namespace = dedup_state_namespace("f-serial"); - let sub = DedupCommitSubscriber::new(config.clone()); - - // Run A finishes, having tentatively seen "a". - store::kv_set(&config, &namespace, "dedup:dd:tentative", &json!(["a"])).unwrap(); - sub.handle(&DomainEvent::FlowRunFinished { - flow_id: "f-serial".into(), - run_id: "run-a".into(), - status: "completed".into(), - }) - .await; - - // Run B finishes later, having independently tentatively seen "b". - // The per-flow lock (proven by the concurrency test above) is what - // guarantees two overlapping runs' `FlowRunFinished` handling - // reduces to exactly this serialized order in practice — so this is - // the correctness property that mutual exclusion is protecting. - store::kv_set(&config, &namespace, "dedup:dd:tentative", &json!(["b"])).unwrap(); - sub.handle(&DomainEvent::FlowRunFinished { - flow_id: "f-serial".into(), - run_id: "run-b".into(), - status: "completed".into(), - }) - .await; - - let committed = store::kv_get(&config, &namespace, "dedup:dd:committed") - .unwrap() - .expect("committed key must exist after both runs settle"); - let mut committed: Vec<&str> = committed - .as_array() - .unwrap() - .iter() - .map(|v| v.as_str().unwrap()) - .collect(); - committed.sort_unstable(); - assert_eq!( - committed, - vec!["a", "b"], - "settling run B must not clobber run A's already-committed keys — committed is a \ - running union across every run that has settled, never a last-writer-wins overwrite" - ); - assert!( - store::kv_get(&config, &namespace, "dedup:dd:tentative") - .unwrap() - .is_none(), - "tentative must be cleared after each successful commit" - ); - } - - #[test] - fn flow_commit_lock_returns_the_same_arc_for_the_same_flow_id_and_differs_across_flows() { - let a1 = flow_commit_lock("f-lock-a"); - let a2 = flow_commit_lock("f-lock-a"); - assert!( - Arc::ptr_eq(&a1, &a2), - "the same flow_id must share one lock instance" - ); - - let b = flow_commit_lock("f-lock-b"); - assert!( - !Arc::ptr_eq(&a1, &b), - "different flow_ids must not contend on the same lock" - ); - } -} +#[path = "bus_tests.rs"] +mod tests; From cd1835da4689cdb6cc0ef0778db26a758a44ba90 Mon Sep 17 00:00:00 2001 From: Muhammad Mustaqeem Date: Thu, 6 Aug 2026 00:28:45 +0500 Subject: [PATCH 3/4] fix(flows): satisfy rustfmt in the dedup run-snapshot code path MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit CI's `cargo fmt --all -- --check` gate rejected the previous revision. Three formatting-sensitive constructs are reworked, each matched against a formatter-clean precedent already in the tree rather than guessed at: * `dedup_node_ids_in` broke a method chain as `flow.graph` / `.nodes` / …, merging the chain parent with its first field access. The same expression elsewhere in this file is written with the parent alone on its own line, so the receiver is now bound to a `nodes` local and the chain starts from it — the shape `flows/ops.rs` already uses at the same indent level. * `clear_run_snapshot` had an `if let Err(e) = store::kv_delete(...)` line sitting exactly on the 100-column limit. The key is hoisted into a local, which shortens the line and matches how `snapshotted_dedup_node_ids` immediately above already builds the same key. * `handle` had been rewritten from the original `if let` into a `match` with a nested `match` inside a block arm, which changed formatting that was previously known-good for no behavioural benefit. The `FlowRunFinished` block is restored verbatim and the new `FlowRunStarted` handling is a short `if let` above it, delegating to a new `snapshot_run_nodes` method that holds the graph load. No behaviour change: the same snapshot is written at run start, read at settlement, and cleared afterwards. Pure formatting and extraction. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> --- src/openhuman/flows/bus.rs | 91 ++++++++++++++++++++------------------ 1 file changed, 49 insertions(+), 42 deletions(-) diff --git a/src/openhuman/flows/bus.rs b/src/openhuman/flows/bus.rs index df58998477..6a0cdd07a4 100644 --- a/src/openhuman/flows/bus.rs +++ b/src/openhuman/flows/bus.rs @@ -739,13 +739,45 @@ impl DedupCommitSubscriber { Some(ids) } + /// Pins the ids of every `dedup` node in the graph this run is starting + /// with, so [`Self::dedup_node_ids`] can settle the run against the nodes + /// it actually executed rather than whatever the flow has been edited into + /// by the time it finishes (issue #5268 item 2). + /// + /// Reading the saved graph here, rather than having `flows::ops` hand us + /// its in-hand `Flow`, keeps the whole fix inside this subscriber. + /// `FlowRunStarted` is published from `flows::ops::flows_run{,_detached}` + /// immediately after the `flow_runs` row insert — before the engine + /// executes a single node — so the graph read back here is the one the run + /// is starting with. + /// + /// Best-effort like everything else in this subscriber: a load failure is + /// logged and swallowed, degrading settlement to the historical "use the + /// flow's current saved graph" behaviour rather than disturbing the run. + fn snapshot_run_nodes(&self, flow_id: &str, run_id: &str) { + match store::get_flow(&self.config, flow_id) { + Ok(Some(flow)) => snapshot_run_dedup_nodes(&self.config, flow_id, run_id, &flow), + // Nothing to pin, and not an error worth logging: the fallback + // already handles a flow that is no longer there. + Ok(None) => {} + Err(e) => { + tracing::warn!( + target: "flows", %flow_id, %run_id, error = %e, + "[dedup-commit] could not load flow to snapshot its dedup nodes at run start \ + — settlement will fall back to the flow's saved graph" + ); + } + } + } + /// Drops this run's dedup-node snapshot once the run has been settled. /// /// Best-effort: a failed delete leaves one small run-scoped KV row /// behind, which nothing reads again (snapshots are keyed by `run_id`, and /// run ids are never reused). fn clear_run_snapshot(&self, namespace: &str, flow_id: &str, run_id: &str) { - if let Err(e) = store::kv_delete(&self.config, namespace, &run_dedup_snapshot_key(run_id)) { + let key = run_dedup_snapshot_key(run_id); + if let Err(e) = store::kv_delete(&self.config, namespace, &key) { tracing::warn!( target: "flows", %flow_id, %run_id, error = %e, "[dedup-commit] failed to clear this run's dedup snapshot — harmless: the key is \ @@ -877,45 +909,20 @@ impl EventHandler for DedupCommitSubscriber { } async fn handle(&self, event: &DomainEvent) { - match event { - // Pin the `dedup` nodes THIS run is about to execute, so - // settlement at `FlowRunFinished` uses the graph the run really - // ran rather than whatever the flow has been edited into by the - // time it finishes (issue #5268 item 2). - // - // Reading the saved graph here rather than having `flows::ops` - // hand us its in-hand `Flow` keeps the whole fix inside this - // subscriber, and `FlowRunStarted` is published from - // `flows::ops::flows_run{,_detached}` immediately after the - // `flow_runs` row insert — before the engine executes a single - // node — so the graph read back here is the one the run is - // starting with. - // - // Best-effort like everything else in this subscriber: a load - // failure is logged and swallowed, degrading settlement to the - // historical "use the current saved graph" behaviour rather than - // disturbing the run. - DomainEvent::FlowRunStarted { flow_id, run_id } => { - match store::get_flow(&self.config, flow_id) { - Ok(Some(flow)) => { - snapshot_run_dedup_nodes(&self.config, flow_id, run_id, &flow); - } - // Nothing to pin, and not an error worth logging: the - // fallback already handles a flow that isn't there. - Ok(None) => {} - Err(e) => tracing::warn!( - target: "flows", %flow_id, %run_id, error = %e, - "[dedup-commit] could not load flow to snapshot its dedup nodes at run \ - start — settlement will fall back to the flow's saved graph" - ), - } - } - DomainEvent::FlowRunFinished { - flow_id, - run_id, - status, - } => self.handle_finished(flow_id, run_id, status).await, - _ => {} + // Pin the `dedup` nodes THIS run is about to execute, so settlement at + // `FlowRunFinished` uses the graph the run really ran rather than + // whatever the flow has been edited into by the time it finishes + // (issue #5268 item 2). + if let DomainEvent::FlowRunStarted { flow_id, run_id } = event { + self.snapshot_run_nodes(flow_id, run_id); + } + if let DomainEvent::FlowRunFinished { + flow_id, + run_id, + status, + } = event + { + self.handle_finished(flow_id, run_id, status).await; } } } @@ -942,8 +949,8 @@ fn run_dedup_snapshot_key(run_id: &str) -> String { /// The ids of every `dedup` node in `flow`'s graph, in graph order. fn dedup_node_ids_in(flow: &Flow) -> Vec { - flow.graph - .nodes + let nodes = &flow.graph.nodes; + nodes .iter() .filter(|n| n.kind == NodeKind::Dedup) .map(|n| n.id.clone()) From d9bd143267005041121714908ad8220f1c64de3e Mon Sep 17 00:00:00 2001 From: "github-actions[bot]" <41898282+github-actions[bot]@users.noreply.github.com> Date: Wed, 5 Aug 2026 19:40:22 +0000 Subject: [PATCH 4/4] style: apply cargo fmt to the touched flows files --- src/openhuman/flows/bus_tests.rs | 19 ++++++++----------- 1 file changed, 8 insertions(+), 11 deletions(-) diff --git a/src/openhuman/flows/bus_tests.rs b/src/openhuman/flows/bus_tests.rs index b19aa59ac1..3c024c0d4c 100644 --- a/src/openhuman/flows/bus_tests.rs +++ b/src/openhuman/flows/bus_tests.rs @@ -237,8 +237,7 @@ async fn handle_app_event_ignores_disabled_flows() { // `list_enabled_flows` must not surface the disabled flow at all — // proves the subscriber's dispatch source already excludes it, // rather than asserting on a spawned background task's side effect. - let (enabled, skipped) = - crate::openhuman::flows::store::list_enabled_flows(&config).unwrap(); + let (enabled, skipped) = crate::openhuman::flows::store::list_enabled_flows(&config).unwrap(); assert!(enabled.is_empty()); assert_eq!(skipped, 0); } @@ -1343,15 +1342,13 @@ async fn dedup_commit_run_started_is_a_noop_for_a_missing_flow() { }) .await; - assert!( - store::kv_get( - &config, - &dedup_state_namespace("f-gone"), - &run_dedup_snapshot_key("run-gone"), - ) - .unwrap() - .is_none() - ); + assert!(store::kv_get( + &config, + &dedup_state_namespace("f-gone"), + &run_dedup_snapshot_key("run-gone"), + ) + .unwrap() + .is_none()); } /// End-to-end proof of the actual bug in issue #5268 item 2, driven