From f4eac049550e2d45d10c73213e758678b9e6d714 Mon Sep 17 00:00:00 2001 From: Ruben Fiszel Date: Thu, 25 Jun 2026 09:01:40 +0000 Subject: [PATCH] fix(debounce): never supersede a running debounce survivor Companion to the windmill-ee-private change in upsert_debounce_key. With debounce_args_to_accumulate + a concurrent_limit, a message arriving while its debounce survivor is already running was marked completed/skipped ("Debounced Running by ...") and the running survivor deleted from the queue, silently dropping accumulated elements. A slow step + concurrent limit keeps the survivor running for a long window, so any arrival during it was lost. The fix leaves a running survivor untouched and starts a fresh debounce window for the late arrival. Adds regression coverage in windmill-queue/tests/debounce_test.rs (push, flow post-preprocessing, no-accumulation, committed-running, and max-count-window cases) and refreshes the SQLx cache for the changed upsert_debounce_key queries. Co-Authored-By: Claude Opus 4.8 (1M context) --- ...083c06a32eb573bc7fbb045fcd622abffad1e.json | 35 ++ ...40649ed1a6f6d4f1e1b636c71bdd36c31b618.json | 35 -- ...e77fb4255568b44b23868723f7d1d66a27af4.json | 35 ++ ...3fc104a2c83f698f7dd657d6e6c95c4ff1f3b.json | 35 -- backend/ee-repo-ref.txt | 2 +- backend/windmill-queue/tests/debounce_test.rs | 492 ++++++++++++++++++ 6 files changed, 563 insertions(+), 71 deletions(-) create mode 100644 backend/.sqlx/query-1763c95498aeb5536824a2ba4ee083c06a32eb573bc7fbb045fcd622abffad1e.json delete mode 100644 backend/.sqlx/query-45de6be332f4ec89482782e3b1640649ed1a6f6d4f1e1b636c71bdd36c31b618.json create mode 100644 backend/.sqlx/query-5a218315f36fd357ff69e599292e77fb4255568b44b23868723f7d1d66a27af4.json delete mode 100644 backend/.sqlx/query-8657c21ace89a9bafe4d184b30e3fc104a2c83f698f7dd657d6e6c95c4ff1f3b.json diff --git a/backend/.sqlx/query-1763c95498aeb5536824a2ba4ee083c06a32eb573bc7fbb045fcd622abffad1e.json b/backend/.sqlx/query-1763c95498aeb5536824a2ba4ee083c06a32eb573bc7fbb045fcd622abffad1e.json new file mode 100644 index 0000000000..9ea5545579 --- /dev/null +++ b/backend/.sqlx/query-1763c95498aeb5536824a2ba4ee083c06a32eb573bc7fbb045fcd622abffad1e.json @@ -0,0 +1,35 @@ +{ + "db_name": "PostgreSQL", + "query": "\n WITH prev_running AS (\n SELECT COALESCE(q.running, false) AS running\n FROM debounce_key dk_prev\n LEFT JOIN v2_job_queue q ON q.id = dk_prev.job_id\n WHERE dk_prev.key = $2\n )\n INSERT INTO debounce_key (job_id, key)\n VALUES ($1, $2)\n ON CONFLICT (key)\n DO UPDATE SET\n previous_job_id = CASE WHEN (SELECT running FROM prev_running)\n THEN NULL ELSE debounce_key.job_id END,\n job_id = EXCLUDED.job_id,\n debounced_times = CASE WHEN (SELECT running FROM prev_running)\n THEN 0 ELSE debounce_key.debounced_times + 1 END,\n first_started_at = CASE WHEN (SELECT running FROM prev_running)\n THEN now() ELSE debounce_key.first_started_at END\n RETURNING\n debounced_times,\n first_started_at,\n previous_job_id AS job_id_to_debounce\n ", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "debounced_times", + "type_info": "Int4" + }, + { + "ordinal": 1, + "name": "first_started_at", + "type_info": "Timestamptz" + }, + { + "ordinal": 2, + "name": "job_id_to_debounce", + "type_info": "Uuid" + } + ], + "parameters": { + "Left": [ + "Uuid", + "Text" + ] + }, + "nullable": [ + false, + false, + true + ] + }, + "hash": "1763c95498aeb5536824a2ba4ee083c06a32eb573bc7fbb045fcd622abffad1e" +} diff --git a/backend/.sqlx/query-45de6be332f4ec89482782e3b1640649ed1a6f6d4f1e1b636c71bdd36c31b618.json b/backend/.sqlx/query-45de6be332f4ec89482782e3b1640649ed1a6f6d4f1e1b636c71bdd36c31b618.json deleted file mode 100644 index 3995bffd64..0000000000 --- a/backend/.sqlx/query-45de6be332f4ec89482782e3b1640649ed1a6f6d4f1e1b636c71bdd36c31b618.json +++ /dev/null @@ -1,35 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "\n INSERT INTO debounce_key (job_id, key)\n VALUES ($1, $2)\n ON CONFLICT (key)\n DO UPDATE SET\n previous_job_id = debounce_key.job_id,\n job_id = EXCLUDED.job_id,\n debounced_times = debounce_key.debounced_times + 1\n RETURNING\n debounced_times,\n first_started_at,\n previous_job_id AS job_id_to_debounce\n ", - "describe": { - "columns": [ - { - "ordinal": 0, - "name": "debounced_times", - "type_info": "Int4" - }, - { - "ordinal": 1, - "name": "first_started_at", - "type_info": "Timestamptz" - }, - { - "ordinal": 2, - "name": "job_id_to_debounce", - "type_info": "Uuid" - } - ], - "parameters": { - "Left": [ - "Uuid", - "Varchar" - ] - }, - "nullable": [ - false, - false, - true - ] - }, - "hash": "45de6be332f4ec89482782e3b1640649ed1a6f6d4f1e1b636c71bdd36c31b618" -} diff --git a/backend/.sqlx/query-5a218315f36fd357ff69e599292e77fb4255568b44b23868723f7d1d66a27af4.json b/backend/.sqlx/query-5a218315f36fd357ff69e599292e77fb4255568b44b23868723f7d1d66a27af4.json new file mode 100644 index 0000000000..eac183b347 --- /dev/null +++ b/backend/.sqlx/query-5a218315f36fd357ff69e599292e77fb4255568b44b23868723f7d1d66a27af4.json @@ -0,0 +1,35 @@ +{ + "db_name": "PostgreSQL", + "query": "\n WITH prev_running AS (\n SELECT COALESCE(q.running, false) AS running\n FROM debounce_key dk_prev\n LEFT JOIN v2_job_queue q ON q.id = dk_prev.job_id\n WHERE dk_prev.key = $2\n ), dk AS (\n INSERT INTO debounce_key (job_id, key)\n VALUES ($1, $2)\n ON CONFLICT (key)\n DO UPDATE SET\n previous_job_id = CASE WHEN (SELECT running FROM prev_running)\n THEN NULL ELSE debounce_key.job_id END,\n job_id = EXCLUDED.job_id,\n debounced_times = CASE WHEN (SELECT running FROM prev_running)\n THEN 0 ELSE debounce_key.debounced_times + 1 END,\n first_started_at = CASE WHEN (SELECT running FROM prev_running)\n THEN now() ELSE debounce_key.first_started_at END\n RETURNING\n debounced_times,\n first_started_at,\n previous_job_id AS job_id_to_debounce\n ), _batch AS (\n INSERT INTO v2_job_debounce_batch (id, debounce_batch)\n SELECT\n $1,\n COALESCE(\n (SELECT debounce_batch FROM v2_job_debounce_batch WHERE id = dk.job_id_to_debounce LIMIT 1),\n nextval('debounce_batch_seq')\n )\n FROM dk\n )\n SELECT debounced_times, first_started_at, job_id_to_debounce FROM dk\n ", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "debounced_times", + "type_info": "Int4" + }, + { + "ordinal": 1, + "name": "first_started_at", + "type_info": "Timestamptz" + }, + { + "ordinal": 2, + "name": "job_id_to_debounce", + "type_info": "Uuid" + } + ], + "parameters": { + "Left": [ + "Uuid", + "Text" + ] + }, + "nullable": [ + false, + false, + true + ] + }, + "hash": "5a218315f36fd357ff69e599292e77fb4255568b44b23868723f7d1d66a27af4" +} diff --git a/backend/.sqlx/query-8657c21ace89a9bafe4d184b30e3fc104a2c83f698f7dd657d6e6c95c4ff1f3b.json b/backend/.sqlx/query-8657c21ace89a9bafe4d184b30e3fc104a2c83f698f7dd657d6e6c95c4ff1f3b.json deleted file mode 100644 index 96e3439612..0000000000 --- a/backend/.sqlx/query-8657c21ace89a9bafe4d184b30e3fc104a2c83f698f7dd657d6e6c95c4ff1f3b.json +++ /dev/null @@ -1,35 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "\n WITH dk AS (\n INSERT INTO debounce_key (job_id, key)\n VALUES ($1, $2)\n ON CONFLICT (key)\n DO UPDATE SET\n previous_job_id = debounce_key.job_id,\n job_id = EXCLUDED.job_id,\n debounced_times = debounce_key.debounced_times + 1\n RETURNING\n debounced_times,\n first_started_at,\n previous_job_id AS job_id_to_debounce\n ), _batch AS (\n INSERT INTO v2_job_debounce_batch (id, debounce_batch)\n SELECT\n $1,\n COALESCE(\n (SELECT debounce_batch FROM v2_job_debounce_batch WHERE id = dk.job_id_to_debounce LIMIT 1),\n nextval('debounce_batch_seq')\n )\n FROM dk\n )\n SELECT debounced_times, first_started_at, job_id_to_debounce FROM dk\n ", - "describe": { - "columns": [ - { - "ordinal": 0, - "name": "debounced_times", - "type_info": "Int4" - }, - { - "ordinal": 1, - "name": "first_started_at", - "type_info": "Timestamptz" - }, - { - "ordinal": 2, - "name": "job_id_to_debounce", - "type_info": "Uuid" - } - ], - "parameters": { - "Left": [ - "Uuid", - "Varchar" - ] - }, - "nullable": [ - false, - false, - true - ] - }, - "hash": "8657c21ace89a9bafe4d184b30e3fc104a2c83f698f7dd657d6e6c95c4ff1f3b" -} diff --git a/backend/ee-repo-ref.txt b/backend/ee-repo-ref.txt index c7439d751f..cb5d08fcc9 100644 --- a/backend/ee-repo-ref.txt +++ b/backend/ee-repo-ref.txt @@ -1 +1 @@ -ed89574be9117cda5e2d7d9de02cb5db066e93e3 +a70f822114103f15eb6b7047a61297c145466482 \ No newline at end of file diff --git a/backend/windmill-queue/tests/debounce_test.rs b/backend/windmill-queue/tests/debounce_test.rs index a2c6971c1c..db0e3f35fa 100644 --- a/backend/windmill-queue/tests/debounce_test.rs +++ b/backend/windmill-queue/tests/debounce_test.rs @@ -3878,6 +3878,498 @@ mod debounce { Ok(()) } + /// Helper: push a script job through push-time `maybe_debounce` with the given key. + /// Returns the args JSON it was pushed with. + async fn push_debounced_script( + db: &Pool, + id: Uuid, + items: Vec, + settings: &DebouncingSettings, + rs_handle: Option, + ) -> serde_json::Value { + let args_val = serde_json::json!({ "items": items }); + insert_script_job_with_args(db, id, "test-workspace", "f/test/script", &args_val).await; + sqlx::query!( + "UPDATE v2_job_queue SET runnable_settings_handle = $1 WHERE id = $2", + rs_handle, + id, + ) + .execute(db) + .await + .unwrap(); + let args_hm: HashMap> = + serde_json::from_value(args_val.clone()).unwrap(); + let push_args = PushArgs::from(&args_hm); + let mut scheduled_for = None; + let mut tx = db.begin().await.unwrap(); + windmill_queue::jobs_ee::maybe_debounce( + settings, + &mut scheduled_for, + &Some("f/test/script".to_string()), + "test-workspace", + JobKind::Script, + id, + &push_args, + &mut tx, + ) + .await + .unwrap(); + tx.commit().await.unwrap(); + args_val + } + + /// Regression test for the "running survivor" data-loss bug. + /// + /// When a debounce survivor has already been pulled and is executing (running), it + /// has committed its batch and can no longer accumulate later arrivals. The old + /// behavior superseded the running survivor anyway: it was completed/skipped + /// ("Debounced Running by ...") and deleted from the queue, silently dropping its + /// accumulated work, while the late arrival could not merge into it. + /// + /// Fix: a late arrival that finds the current survivor already running starts a + /// FRESH debounce window. The running survivor is left to finish with its own + /// accumulated batch; the late arrival accumulates only its own batch. No job is + /// killed and no item is dropped or double-run. + #[sqlx::test(migrations = "../migrations", fixtures("base"))] + async fn test_debounce_running_survivor_not_superseded( + db: Pool, + ) -> anyhow::Result<()> { + let key = "running_survivor_key"; + let settings = DebouncingSettings { + debounce_delay_s: Some(5), + debounce_key: Some(key.to_string()), + debounce_args_to_accumulate: Some(vec!["items".to_string()]), + ..Default::default() + }; + let rs_handle = setup_debouncing_settings(&db, &settings).await; + + // J1, J2 share a window; J2 is the survivor with batch {J1, J2}. + let j1 = Uuid::new_v4(); + let j2 = Uuid::new_v4(); + push_debounced_script(&db, j1, vec![1], &settings, rs_handle).await; + push_debounced_script(&db, j2, vec![2], &settings, rs_handle).await; + assert!(is_completed(&db, &j1).await, "J1 should be debounced by J2"); + assert!(is_queued(&db, &j2).await, "J2 should be the survivor"); + + // Worker pulls J2 and marks it running: the window where the key still points + // to J2 and its batch is intact, but J2 can no longer accumulate new arrivals. + sqlx::query!( + "UPDATE v2_job_queue SET running = true, started_at = now() WHERE id = $1", + j2 + ) + .execute(&db) + .await?; + + // Late arrival J3 while J2 is running. + let j3 = Uuid::new_v4(); + let j3_args = push_debounced_script(&db, j3, vec![3], &settings, rs_handle).await; + + // The running survivor J2 must NOT be superseded: still queued, not completed. + assert!( + is_queued(&db, &j2).await, + "running survivor J2 must stay in the queue" + ); + assert!( + !is_completed(&db, &j2).await, + "running survivor J2 must not be completed/skipped" + ); + + // J3 must own the debounce key as the head of a FRESH window (no previous job). + let (dk_job, dk_prev, dk_times) = get_debounce_key(&db, key) + .await + .expect("debounce key exists"); + assert_eq!(dk_job, j3, "J3 should hold the debounce key"); + assert!( + dk_prev.is_none(), + "J3 should start a fresh window with no previous job (got {dk_prev:?})" + ); + assert_eq!( + dk_times, 0, + "fresh window should reset debounced_times to 0" + ); + + // The running survivor J2 accumulates only its own committed batch: [1, 2]. + let mut j2_res = make_pulled_job_result( + j2, + "test-workspace", + "f/test/script", + &serde_json::json!({"items": [2]}), + JobKind::Script, + "deno", + rs_handle, + ); + j2_res.maybe_apply_debouncing(&db).await?; + assert!( + j2_res.job.is_some(), + "running survivor J2 must still execute (not nulled out)" + ); + assert_accumulated_items(&j2_res, &[1, 2], "items"); + + // The late arrival J3 accumulates only its own batch: [3]. No overlap with J2. + let mut j3_res = make_pulled_job_result( + j3, + "test-workspace", + "f/test/script", + &j3_args, + JobKind::Script, + "deno", + rs_handle, + ); + j3_res.maybe_apply_debouncing(&db).await?; + assert!(j3_res.job.is_some(), "J3 must execute"); + assert_accumulated_items(&j3_res, &[3], "items"); + + Ok(()) + } + + /// A second arrival debouncing a NON-running survivor must keep accumulating into + /// the same batch (the normal debounce behavior must be unchanged by the + /// running-survivor guard). + #[sqlx::test(migrations = "../migrations", fixtures("base"))] + async fn test_debounce_queued_survivor_still_accumulates( + db: Pool, + ) -> anyhow::Result<()> { + let key = "queued_survivor_key"; + let settings = DebouncingSettings { + debounce_delay_s: Some(5), + debounce_key: Some(key.to_string()), + debounce_args_to_accumulate: Some(vec!["items".to_string()]), + ..Default::default() + }; + let rs_handle = setup_debouncing_settings(&db, &settings).await; + + // Three arrivals, none running: classic debounce, all accumulate into J3. + let j1 = Uuid::new_v4(); + let j2 = Uuid::new_v4(); + let j3 = Uuid::new_v4(); + push_debounced_script(&db, j1, vec![1], &settings, rs_handle).await; + push_debounced_script(&db, j2, vec![2], &settings, rs_handle).await; + let j3_args = push_debounced_script(&db, j3, vec![3], &settings, rs_handle).await; + + assert!(is_completed(&db, &j1).await, "J1 debounced"); + assert!(is_completed(&db, &j2).await, "J2 debounced"); + assert!(is_queued(&db, &j3).await, "J3 is the survivor"); + + // Window keeps growing: J3 is the third arrival in the same batch. + let (dk_job, dk_prev, dk_times) = get_debounce_key(&db, key) + .await + .expect("debounce key exists"); + assert_eq!(dk_job, j3); + assert_eq!(dk_prev, Some(j2), "previous job should be J2"); + assert_eq!(dk_times, 2, "debounced_times should keep incrementing"); + + let mut j3_res = make_pulled_job_result( + j3, + "test-workspace", + "f/test/script", + &j3_args, + JobKind::Script, + "deno", + rs_handle, + ); + j3_res.maybe_apply_debouncing(&db).await?; + assert_accumulated_items(&j3_res, &[1, 2, 3], "items"); + + Ok(()) + } + + /// The running-survivor guard must also apply to plain debounce (delay only, no + /// argument accumulation): a running survivor must never be completed/skipped. + #[sqlx::test(migrations = "../migrations", fixtures("base"))] + async fn test_debounce_running_survivor_no_accumulation( + db: Pool, + ) -> anyhow::Result<()> { + let key = "running_no_accum_key"; + let settings = DebouncingSettings { + debounce_delay_s: Some(5), + debounce_key: Some(key.to_string()), + // No debounce_args_to_accumulate. + ..Default::default() + }; + + let j1 = Uuid::new_v4(); + let j2 = Uuid::new_v4(); + insert_noop_job(&db, j1, "test-workspace").await; + insert_noop_job(&db, j2, "test-workspace").await; + + let push = |id: Uuid| { + let settings = settings.clone(); + let db = db.clone(); + async move { + let args_hm = empty_args(); + let args = PushArgs::from(&args_hm); + let mut scheduled_for = None; + let mut tx = db.begin().await.unwrap(); + windmill_queue::jobs_ee::maybe_debounce( + &settings, + &mut scheduled_for, + &Some("f/test/script".to_string()), + "test-workspace", + JobKind::Noop, + id, + &args, + &mut tx, + ) + .await + .unwrap(); + tx.commit().await.unwrap(); + } + }; + + push(j1).await; + push(j2).await; + assert!(is_completed(&db, &j1).await, "J1 debounced by J2"); + + // J2 starts running. + sqlx::query!( + "UPDATE v2_job_queue SET running = true, started_at = now() WHERE id = $1", + j2 + ) + .execute(&db) + .await?; + + // Late arrival J3. + let j3 = Uuid::new_v4(); + insert_noop_job(&db, j3, "test-workspace").await; + push(j3).await; + + // Running survivor J2 is preserved; J3 takes over a fresh window. + assert!(is_queued(&db, &j2).await, "running J2 stays queued"); + assert!(!is_completed(&db, &j2).await, "running J2 not completed"); + let (dk_job, dk_prev, dk_times) = get_debounce_key(&db, key) + .await + .expect("debounce key exists"); + assert_eq!(dk_job, j3); + assert!(dk_prev.is_none(), "fresh window: no previous job"); + assert_eq!(dk_times, 0, "fresh window resets debounced_times"); + + Ok(()) + } + + /// Once a survivor has fully been pulled (batch + key consumed by + /// maybe_apply_debouncing) and is running, a later arrival naturally starts a new + /// window. This locks in that the committed-running case stays correct alongside + /// the in-flight-running guard. + #[sqlx::test(migrations = "../migrations", fixtures("base"))] + async fn test_debounce_committed_running_survivor_independent( + db: Pool, + ) -> anyhow::Result<()> { + let key = "committed_running_key"; + let settings = DebouncingSettings { + debounce_delay_s: Some(5), + debounce_key: Some(key.to_string()), + debounce_args_to_accumulate: Some(vec!["items".to_string()]), + ..Default::default() + }; + let rs_handle = setup_debouncing_settings(&db, &settings).await; + + let j1 = Uuid::new_v4(); + let j2 = Uuid::new_v4(); + push_debounced_script(&db, j1, vec![1], &settings, rs_handle).await; + let j2_args = push_debounced_script(&db, j2, vec![2], &settings, rs_handle).await; + + // J2 is pulled: accumulate its batch and consume key + batch. + let mut j2_res = make_pulled_job_result( + j2, + "test-workspace", + "f/test/script", + &j2_args, + JobKind::Script, + "deno", + rs_handle, + ); + j2_res.maybe_apply_debouncing(&db).await?; + assert_accumulated_items(&j2_res, &[1, 2], "items"); + assert!( + get_debounce_key(&db, key).await.is_none(), + "key consumed when survivor pulled" + ); + + // J2 now running. + sqlx::query!( + "UPDATE v2_job_queue SET running = true, started_at = now() WHERE id = $1", + j2 + ) + .execute(&db) + .await?; + + // Late arrival J3: fresh window, independent batch, J2 untouched. + let j3 = Uuid::new_v4(); + let j3_args = push_debounced_script(&db, j3, vec![3], &settings, rs_handle).await; + assert!(is_queued(&db, &j2).await, "running J2 untouched"); + assert!(!is_completed(&db, &j2).await, "running J2 not completed"); + let (dk_job, _, dk_times) = get_debounce_key(&db, key).await.expect("key exists"); + assert_eq!(dk_job, j3); + assert_eq!(dk_times, 0); + + let mut j3_res = make_pulled_job_result( + j3, + "test-workspace", + "f/test/script", + &j3_args, + JobKind::Script, + "deno", + rs_handle, + ); + j3_res.maybe_apply_debouncing(&db).await?; + assert_accumulated_items(&j3_res, &[3], "items"); + + Ok(()) + } + + /// Flow post-preprocessing debounce must apply the same running-survivor guard: + /// a running flow survivor must not be completed/skipped by a late flow arrival. + #[sqlx::test(migrations = "../migrations", fixtures("base"))] + async fn test_post_preprocessing_running_survivor_not_superseded( + db: Pool, + ) -> anyhow::Result<()> { + let key = "pp_running_survivor_key"; + let settings = DebouncingSettings { + debounce_delay_s: Some(5), + debounce_key: Some(key.to_string()), + ..Default::default() + }; + let args_hm = empty_args(); + + let flow1 = Uuid::new_v4(); + let flow2 = Uuid::new_v4(); + insert_flow_job(&db, flow1, "test-workspace", "f/test/flow").await; + insert_flow_job(&db, flow2, "test-workspace", "f/test/flow").await; + + let pp = |id: Uuid| { + let settings = settings.clone(); + let db = db.clone(); + let args_hm = args_hm.clone(); + async move { + let args = PushArgs::from(&args_hm); + windmill_queue::jobs_ee::maybe_debounce_post_preprocessing( + &settings, + &Some("f/test/flow".to_string()), + "test-workspace", + id, + &args, + &db, + ) + .await + .unwrap() + } + }; + + pp(flow1).await; + pp(flow2).await; + assert!(is_completed(&db, &flow1).await, "flow1 debounced by flow2"); + + // flow2 (the survivor) starts running. + sqlx::query!( + "UPDATE v2_job_queue SET running = true, started_at = now() WHERE id = $1", + flow2 + ) + .execute(&db) + .await?; + + // Late flow3 arrival while flow2 is running. + let flow3 = Uuid::new_v4(); + insert_flow_job(&db, flow3, "test-workspace", "f/test/flow").await; + let sched = pp(flow3).await; + assert!(sched.is_some(), "flow3 should be debounced (fresh window)"); + + // Running flow2 must be preserved; flow3 owns a fresh window. + assert!(is_queued(&db, &flow2).await, "running flow2 stays queued"); + assert!( + !is_completed(&db, &flow2).await, + "running flow2 must not be completed" + ); + let (dk_job, dk_prev, dk_times) = get_debounce_key(&db, key).await.expect("key exists"); + assert_eq!(dk_job, flow3, "flow3 holds the key"); + assert!(dk_prev.is_none(), "fresh window: no previous job"); + assert_eq!(dk_times, 0, "fresh window resets debounced_times"); + + Ok(()) + } + + /// A running survivor resets the debounce window for the late arrival, so an + /// inherited high `debounced_times` cannot push the new arrival over + /// max_total_debounces_amount and force an immediate (un-debounced) run. + #[sqlx::test(migrations = "../migrations", fixtures("base"))] + async fn test_debounce_running_survivor_resets_limit_window( + db: Pool, + ) -> anyhow::Result<()> { + let key = "running_limit_key"; + let settings = DebouncingSettings { + debounce_delay_s: Some(5), + debounce_key: Some(key.to_string()), + // Limit of 3: J1, J2 stay debounced; a third arrival in the SAME window + // would trip the limit (current_amount + 1 >= 3) and fire immediately. + max_total_debounces_amount: Some(3), + debounce_args_to_accumulate: Some(vec!["items".to_string()]), + ..Default::default() + }; + let rs_handle = setup_debouncing_settings(&db, &settings).await; + + // Build up the window close to the limit: J1, J2 (debounced_times = 1 on J2). + let j1 = Uuid::new_v4(); + let j2 = Uuid::new_v4(); + push_debounced_script(&db, j1, vec![1], &settings, rs_handle).await; + push_debounced_script(&db, j2, vec![2], &settings, rs_handle).await; + let (_, _, times_before) = get_debounce_key(&db, key).await.expect("key exists"); + assert_eq!(times_before, 1); + + // J2 starts running. + sqlx::query!( + "UPDATE v2_job_queue SET running = true, started_at = now() WHERE id = $1", + j2 + ) + .execute(&db) + .await?; + + // J3 arrives. Without the reset it would inherit debounced_times and could trip + // the max-count limit and fire immediately, killing running J2. With the guard + // it starts a fresh window (debounced_times = 0) and is debounced normally. + let j3 = Uuid::new_v4(); + let mut scheduled_for = None; + { + let args_val = serde_json::json!({ "items": [3] }); + insert_script_job_with_args(&db, j3, "test-workspace", "f/test/script", &args_val) + .await; + sqlx::query!( + "UPDATE v2_job_queue SET runnable_settings_handle = $1 WHERE id = $2", + rs_handle, + j3, + ) + .execute(&db) + .await?; + let args_hm: HashMap> = serde_json::from_value(args_val).unwrap(); + let push_args = PushArgs::from(&args_hm); + let mut tx = db.begin().await?; + windmill_queue::jobs_ee::maybe_debounce( + &settings, + &mut scheduled_for, + &Some("f/test/script".to_string()), + "test-workspace", + JobKind::Script, + j3, + &push_args, + &mut tx, + ) + .await?; + tx.commit().await?; + } + + // J3 is debounced (scheduled_for set, not fired immediately) and J2 survives. + assert!( + scheduled_for.is_some(), + "J3 should be debounced, not fired immediately" + ); + assert!(is_queued(&db, &j2).await, "running J2 stays queued"); + assert!(!is_completed(&db, &j2).await, "running J2 not completed"); + let (dk_job, dk_prev, dk_times) = get_debounce_key(&db, key).await.expect("key exists"); + assert_eq!(dk_job, j3); + assert!(dk_prev.is_none()); + assert_eq!(dk_times, 0, "fresh window resets the limit counter"); + + Ok(()) + } + /// Test: Push-time (script) debounce with max_total_debounces_amount=2. /// 5 calls, each sending {x: [i]}. Expected: /// Call 1: debounced (scheduled_for set)