From f7bef8de11a1fae9ed5eb6ae44efd0ca9bdce785 Mon Sep 17 00:00:00 2001 From: npub1jmc9dt2lyvzu3h0kxlwxt5zg4fxp9476awyxw6gwxn72g6cw7exqs64whm <96f056ad5f2305c8ddf637dc65d048aa4c12d7daeb8867690e34fca46b0ef64c@sprout-oss.stage.blox.sqprod.co> Date: Fri, 26 Jun 2026 14:06:39 -0400 Subject: [PATCH] test(workflow): audit engine-layer cron-dedup invariants MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Two `#[ignore]` live-Postgres tests layered on Max's `process_scheduled_workflow_at` seam (df997d7cc), each pinning an invariant the existing F4 coverage intentionally does NOT pin: orphan_claim_blocks_refire_at_same_canonical_instant Pre-INSERTs a `scheduled_workflow_fires` claim row with `workflow_run_id = NULL` (simulating a prior pod that won the claim then died before `create_workflow_run`), then calls the seam for the same canonical `scheduled_for`. Asserts `Option == None`, zero rows in `workflow_runs`, single claim row unchanged. Proves the claim row alone — not claim + attached run — is the dedupe boundary. A refactor that ever made `attach_scheduled_workflow_run` (or the `workflow_runs` row) the gate instead of the claim row would still pass `scheduled_cron_claim_is_exactly_once_across_two_engines` (success-path only) but must fail this one. Uses a fixed daily cron (`0 12 * * *`) with `now = 12:00:30Z` so the canonical instant is trivially `12:00:00Z` — closes the alignment hazard Sami flagged where an off-by-one `scheduled_for` would silently exercise the wrong PK and produce a false failure (verified: shifting the orphan INSERT to `12:00:01Z` breaks the test loudly). adjacent_interval_boundaries_produce_distinct_runs Engine-layer companion to Sami's DB-layer `clock_skewed_adjacent_windows_each_claim_independently` (9652037ec). Feeds two `now` values straddling a 60s interval boundary through the seam; assert each returns a `Some(run_id)` and the two run ids differ. Asserts 2 claim rows + 2 workflow_runs rows. Exercises the seam's own canonical-boundary math (`interval_should_fire`) end-to- end — the property `71da65e51` was designed to provide is now observable through the engine surface. Adversarial verify: replacing the canonical-boundary computation with `Some(anchor)` (max coalescing) makes the test fail loudly at `expect("second canonical boundary must fire")`. Both tests reuse Max's `connect_test_db` helper and the unique-uuid- per-test convention (Sami's F1/F4) — no harness work owed. Validation (all hermit cargo): - fmt --check -p buzz-workflow ✅ - clippy -p buzz-workflow --tests --no-deps -- -D warnings ✅ - test -p buzz-workflow --lib ✅ — 145 passed, 3 ignored - test -p buzz-workflow --lib -- --ignored ✅ — 3 passed (orphan_claim, adjacent_interval, scheduled_cron_claim) - check -p buzz-relay ✅ Co-authored-by: Tyler Longwell Signed-off-by: Tyler Longwell --- crates/buzz-workflow/src/lib.rs | 241 ++++++++++++++++++++++++++++++++ 1 file changed, 241 insertions(+) diff --git a/crates/buzz-workflow/src/lib.rs b/crates/buzz-workflow/src/lib.rs index 0518d8394..1e53c5666 100644 --- a/crates/buzz-workflow/src/lib.rs +++ b/crates/buzz-workflow/src/lib.rs @@ -1044,6 +1044,247 @@ steps: assert_eq!(run_count, 1); } + /// Engine-layer crash-mid-run audit (b): proves the `scheduled_workflow_fires` + /// claim row alone — *with no attached run* — is the dedupe boundary the + /// scheduler relies on. Simulates: a prior pod won the claim, then died + /// before `create_workflow_run` (or before `attach_scheduled_workflow_run`), + /// leaving an orphan claim row. A subsequent tick for the same canonical + /// `scheduled_for` must no-op rather than create a duplicate run. + /// + /// A future refactor that ever made `attach_scheduled_workflow_run` (or the + /// `workflow_runs` row itself) the dedupe gate instead of the claim row + /// would still pass `scheduled_cron_claim_is_exactly_once_across_two_engines` + /// (which only exercises the success path) but must fail this one. + #[tokio::test] + #[ignore = "requires Postgres"] + async fn orphan_claim_blocks_refire_at_same_canonical_instant() { + let (db, pool) = connect_test_db().await; + let community = db + .ensure_configured_community(&format!("workflow-cron-orphan-{}.test", Uuid::new_v4())) + .await + .expect("create community"); + let owner_pubkey = [43_u8; 32]; + db.ensure_user(&owner_pubkey).await.expect("create owner"); + + let channel_id = Uuid::new_v4(); + sqlx::query( + r#" + INSERT INTO channels + (id, community_id, name, channel_type, visibility, created_by) + VALUES ($1, $2, $3, 'stream', 'open', $4) + "#, + ) + .bind(channel_id) + .bind(community.id.as_uuid()) + .bind("workflow cron orphan test") + .bind(owner_pubkey.as_slice()) + .execute(&pool) + .await + .expect("create channel"); + + // Fixed daily cron so the canonical instant is trivially known. + let (_def, definition_json) = WorkflowEngine::parse_yaml( + r#" +name: Cron orphan test +trigger: + on: schedule + cron: '0 12 * * *' +steps: + - id: wait + action: delay + duration: 0s +"#, + ) + .expect("valid workflow yaml"); + + let workflow_id = db + .create_workflow( + community.id, + Some(channel_id), + &owner_pubkey, + "Cron orphan test", + &definition_json, + b"test-definition-hash", + ) + .await + .expect("create workflow"); + let workflow = db.get_workflow(workflow_id).await.expect("load workflow"); + + // Cron `0 12 * * *` + `now = 12:00:30Z` → canonical scheduled_for = 12:00:00Z. + let canonical_scheduled_for = chrono::DateTime::parse_from_rfc3339("2026-06-15T12:00:00Z") + .unwrap() + .with_timezone(&Utc); + let now = chrono::DateTime::parse_from_rfc3339("2026-06-15T12:00:30Z") + .unwrap() + .with_timezone(&Utc); + + // Pre-insert an orphan claim row: claim exists, workflow_run_id is NULL + // (simulating prior pod crashed after claim, before create_workflow_run). + // Source community_id from the workflow row to keep the FK pair sound — + // the same invariant F4 made schema-impossible to forge. + sqlx::query( + r#" + INSERT INTO scheduled_workflow_fires + (community_id, workflow_id, scheduled_for, workflow_run_id) + SELECT w.community_id, w.id, $2, NULL + FROM workflows w + WHERE w.id = $1 + "#, + ) + .bind(workflow_id) + .bind(canonical_scheduled_for) + .execute(&pool) + .await + .expect("seed orphan claim row"); + + let engine = Arc::new(WorkflowEngine::new(db.clone(), WorkflowConfig::default())); + let result = engine.process_scheduled_workflow_at(&workflow, now).await; + + assert_eq!( + result, None, + "orphan claim row must block refire — the claim row, not attach/run, is the dedupe boundary" + ); + + let run_count = sqlx::query_scalar::<_, i64>( + "SELECT COUNT(*) FROM workflow_runs WHERE workflow_id = $1", + ) + .bind(workflow_id) + .fetch_one(&pool) + .await + .expect("count workflow runs"); + assert_eq!( + run_count, 0, + "no workflow_runs row should exist for this workflow — the orphan claim must short-circuit before create_workflow_run" + ); + + let fire_count = sqlx::query_scalar::<_, i64>( + "SELECT COUNT(*) FROM scheduled_workflow_fires WHERE workflow_id = $1", + ) + .bind(workflow_id) + .fetch_one(&pool) + .await + .expect("count scheduled fires"); + assert_eq!( + fire_count, 1, + "the single orphan claim row must remain — no duplicate claim was inserted" + ); + } + + /// Engine-layer canonical-boundary audit (c): adjacent interval windows + /// fed through the seam produce distinct run ids. + /// + /// Sami's `clock_skewed_adjacent_windows_each_claim_independently` proves + /// the DB layer keeps adjacent `scheduled_for` keys independent under the + /// composite PK. This test proves the engine-side canonicalization + /// (`interval_should_fire`) actually *produces* distinct keys for adjacent + /// `now` values when the anchor shifts between ticks — i.e. the property + /// `71da65e51` was designed to provide is observable end-to-end through + /// the seam. + #[tokio::test] + #[ignore = "requires Postgres"] + async fn adjacent_interval_boundaries_produce_distinct_runs() { + let (db, pool) = connect_test_db().await; + let community = db + .ensure_configured_community(&format!("workflow-cron-adjacent-{}.test", Uuid::new_v4())) + .await + .expect("create community"); + let owner_pubkey = [44_u8; 32]; + db.ensure_user(&owner_pubkey).await.expect("create owner"); + + let channel_id = Uuid::new_v4(); + sqlx::query( + r#" + INSERT INTO channels + (id, community_id, name, channel_type, visibility, created_by) + VALUES ($1, $2, $3, 'stream', 'open', $4) + "#, + ) + .bind(channel_id) + .bind(community.id.as_uuid()) + .bind("workflow cron adjacent test") + .bind(owner_pubkey.as_slice()) + .execute(&pool) + .await + .expect("create channel"); + + // Interval workflow exercises the seam's own canonical-boundary math + // through `interval_should_fire`. The DB-authoritative anchor + // (latest_scheduled_workflow_fire) shifts after the first call, which is + // the exact property under test. + let (_def, definition_json) = WorkflowEngine::parse_yaml( + r#" +name: Interval adjacent test +trigger: + on: schedule + interval: 60s +steps: + - id: wait + action: delay + duration: 0s +"#, + ) + .expect("valid workflow yaml"); + + let workflow_id = db + .create_workflow( + community.id, + Some(channel_id), + &owner_pubkey, + "Interval adjacent test", + &definition_json, + b"test-definition-hash", + ) + .await + .expect("create workflow"); + let workflow = db.get_workflow(workflow_id).await.expect("load workflow"); + + let engine = Arc::new(WorkflowEngine::new(db.clone(), WorkflowConfig::default())); + + // Anchor is workflow.created_at. With interval = 60s: + // now_1 = anchor + 90s → canonical scheduled_for_1 = anchor + 60s + // (after call 1, latest_scheduled_workflow_fire == anchor + 60s) + // now_2 = anchor + 150s → canonical scheduled_for_2 = anchor + 120s + // (= new_anchor + 60s) + let anchor = workflow.created_at; + let now_1 = anchor + chrono::Duration::seconds(90); + let now_2 = anchor + chrono::Duration::seconds(150); + + let run_1 = engine + .process_scheduled_workflow_at(&workflow, now_1) + .await + .expect("first canonical boundary must fire"); + let run_2 = engine + .process_scheduled_workflow_at(&workflow, now_2) + .await + .expect("second canonical boundary must fire"); + + assert_ne!( + run_1, run_2, + "adjacent canonical boundaries must produce distinct run ids" + ); + + let fire_count = sqlx::query_scalar::<_, i64>( + "SELECT COUNT(*) FROM scheduled_workflow_fires WHERE workflow_id = $1", + ) + .bind(workflow_id) + .fetch_one(&pool) + .await + .expect("count scheduled fires"); + assert_eq!( + fire_count, 2, + "two distinct canonical instants must produce two claim rows" + ); + + let run_count = sqlx::query_scalar::<_, i64>( + "SELECT COUNT(*) FROM workflow_runs WHERE workflow_id = $1", + ) + .bind(workflow_id) + .fetch_one(&pool) + .await + .expect("count workflow runs"); + assert_eq!(run_count, 2); + } + async fn connect_test_db() -> (Db, sqlx::PgPool) { let database_url = std::env::var("BUZZ_TEST_DATABASE_URL") .or_else(|_| std::env::var("DATABASE_URL"))