diff --git a/backend/.sqlx/query-0035bf99ce6fc00c7338bebfeb7e79bb9e7bc3d216b84279dee0018603965941.json b/backend/.sqlx/query-0035bf99ce6fc00c7338bebfeb7e79bb9e7bc3d216b84279dee0018603965941.json new file mode 100644 index 0000000000..cde17acf76 --- /dev/null +++ b/backend/.sqlx/query-0035bf99ce6fc00c7338bebfeb7e79bb9e7bc3d216b84279dee0018603965941.json @@ -0,0 +1,20 @@ +{ + "db_name": "PostgreSQL", + "query": "WITH del AS (\n DELETE FROM v2_job_debounce_batch\n WHERE consumed_at IS NOT NULL AND consumed_at < now() - interval '10 minutes'\n RETURNING 1\n ) SELECT count(*) FROM del", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "count", + "type_info": "Int8" + } + ], + "parameters": { + "Left": [] + }, + "nullable": [ + null + ] + }, + "hash": "0035bf99ce6fc00c7338bebfeb7e79bb9e7bc3d216b84279dee0018603965941" +} diff --git a/backend/.sqlx/query-299b94a7972443267dd664c178a1704d195a7fc0d4e66e1014a18398e3a294f4.json b/backend/.sqlx/query-299b94a7972443267dd664c178a1704d195a7fc0d4e66e1014a18398e3a294f4.json new file mode 100644 index 0000000000..642723decf --- /dev/null +++ b/backend/.sqlx/query-299b94a7972443267dd664c178a1704d195a7fc0d4e66e1014a18398e3a294f4.json @@ -0,0 +1,20 @@ +{ + "db_name": "PostgreSQL", + "query": "WITH del AS (\n DELETE FROM v2_job_debounce_batch\n WHERE consumed_at IS NOT NULL AND consumed_at < now() - interval '10 minutes'\n RETURNING 1\n ) SELECT count(*) as \"c!\" FROM del", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "c!", + "type_info": "Int8" + } + ], + "parameters": { + "Left": [] + }, + "nullable": [ + null + ] + }, + "hash": "299b94a7972443267dd664c178a1704d195a7fc0d4e66e1014a18398e3a294f4" +} diff --git a/backend/.sqlx/query-3b06ecd4339e966bab32de0b85b7197b2f99174a0066e25d975c303b2e60a2e2.json b/backend/.sqlx/query-3b06ecd4339e966bab32de0b85b7197b2f99174a0066e25d975c303b2e60a2e2.json new file mode 100644 index 0000000000..8a79de3aa4 --- /dev/null +++ b/backend/.sqlx/query-3b06ecd4339e966bab32de0b85b7197b2f99174a0066e25d975c303b2e60a2e2.json @@ -0,0 +1,20 @@ +{ + "db_name": "PostgreSQL", + "query": "SELECT count(*) as \"c!\" FROM v2_job_debounce_batch", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "c!", + "type_info": "Int8" + } + ], + "parameters": { + "Left": [] + }, + "nullable": [ + null + ] + }, + "hash": "3b06ecd4339e966bab32de0b85b7197b2f99174a0066e25d975c303b2e60a2e2" +} diff --git a/backend/.sqlx/query-ae8c0d0397609cf275bcffb92a5014665d47fe95b9eee0e1785441b0b4497a3c.json b/backend/.sqlx/query-ae8c0d0397609cf275bcffb92a5014665d47fe95b9eee0e1785441b0b4497a3c.json new file mode 100644 index 0000000000..3413ddc87f --- /dev/null +++ b/backend/.sqlx/query-ae8c0d0397609cf275bcffb92a5014665d47fe95b9eee0e1785441b0b4497a3c.json @@ -0,0 +1,22 @@ +{ + "db_name": "PostgreSQL", + "query": "SELECT EXISTS(SELECT 1 FROM v2_job_debounce_batch WHERE id = $1) as \"e!\"", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "e!", + "type_info": "Bool" + } + ], + "parameters": { + "Left": [ + "Uuid" + ] + }, + "nullable": [ + null + ] + }, + "hash": "ae8c0d0397609cf275bcffb92a5014665d47fe95b9eee0e1785441b0b4497a3c" +} diff --git a/backend/.sqlx/query-c886e8af0fc8a3999a813371855c0053571e79960280f0714616d13a456d7bed.json b/backend/.sqlx/query-c886e8af0fc8a3999a813371855c0053571e79960280f0714616d13a456d7bed.json new file mode 100644 index 0000000000..7361b645b5 --- /dev/null +++ b/backend/.sqlx/query-c886e8af0fc8a3999a813371855c0053571e79960280f0714616d13a456d7bed.json @@ -0,0 +1,16 @@ +{ + "db_name": "PostgreSQL", + "query": "INSERT INTO v2_job_debounce_batch (id, debounce_batch, consumed_at) VALUES\n ($1, nextval('debounce_batch_seq'), now() - interval '20 minutes'),\n ($2, nextval('debounce_batch_seq'), now() - interval '1 minute'),\n ($3, nextval('debounce_batch_seq'), NULL)", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "Uuid", + "Uuid", + "Uuid" + ] + }, + "nullable": [] + }, + "hash": "c886e8af0fc8a3999a813371855c0053571e79960280f0714616d13a456d7bed" +} diff --git a/backend/src/monitor.rs b/backend/src/monitor.rs index df71c1710d..4b8e1b678b 100644 --- a/backend/src/monitor.rs +++ b/backend/src/monitor.rs @@ -4404,10 +4404,13 @@ RETURNING key,job_id /// own (already-accumulated, persisted) args, so the grace period is not correctness- /// critical. async fn cleanup_consumed_debounce_batches(db: &DB) -> error::Result<()> { + // 10 minutes is far longer than any debounce window, so any survivor that could + // still reference a consumed row has been pulled well before then; keeping it short + // bounds table growth under high-throughput debounce. let deleted = sqlx::query_scalar!( "WITH del AS ( DELETE FROM v2_job_debounce_batch - WHERE consumed_at IS NOT NULL AND consumed_at < now() - interval '1 hour' + WHERE consumed_at IS NOT NULL AND consumed_at < now() - interval '10 minutes' RETURNING 1 ) SELECT count(*) FROM del" ) diff --git a/backend/windmill-queue/tests/debounce_test.rs b/backend/windmill-queue/tests/debounce_test.rs index 996dbd51cd..f3c63db602 100644 --- a/backend/windmill-queue/tests/debounce_test.rs +++ b/backend/windmill-queue/tests/debounce_test.rs @@ -4706,6 +4706,294 @@ mod debounce { Ok(()) } + /// Helper: insert a script job and put it on the SAME debounce batch as `of_job` + /// (simulating a chained survivor). Returns its args JSON. + async fn add_survivor_to_batch_of( + db: &Pool, + id: Uuid, + items: Vec, + of_job: Uuid, + rs_handle: Option, + ) -> serde_json::Value { + let args = serde_json::json!({ "items": items }); + insert_script_job_with_args(db, id, "test-workspace", "f/test/script", &args).await; + sqlx::query!( + "UPDATE v2_job_queue SET runnable_settings_handle = $1 WHERE id = $2", + rs_handle, + id, + ) + .execute(db) + .await + .unwrap(); + sqlx::query!( + "INSERT INTO v2_job_debounce_batch (id, debounce_batch) + SELECT $1, debounce_batch FROM v2_job_debounce_batch WHERE id = $2", + id, + of_job, + ) + .execute(db) + .await + .unwrap(); + args + } + + /// Helper: read the accumulated `items` of a pulled job as a sorted Vec. + fn items_of(result: &windmill_queue::PulledJobResult) -> Vec { + let job = result.job.as_ref().expect("job present"); + let raw = job.job.args.as_ref().unwrap().get("items").unwrap(); + let mut v: Vec = serde_json::from_str::>(raw.get()) + .unwrap() + .iter() + .map(|x| x.as_i64().unwrap()) + .collect(); + v.sort(); + v + } + + /// Edge: an accumulate-debounced job that was NEVER batched (CE / workers behind v2: + /// no v2_job_debounce_batch row) must keep its own args, not be emptied. + #[sqlx::test(migrations = "../migrations", fixtures("base"))] + async fn test_debounce_never_batched_keeps_own_args(db: Pool) -> anyhow::Result<()> { + let settings = DebouncingSettings { + debounce_delay_s: Some(5), + debounce_key: Some("never_batched_key".to_string()), + debounce_args_to_accumulate: Some(vec!["items".to_string()]), + ..Default::default() + }; + let rs_handle = setup_debouncing_settings(&db, &settings).await; + + // Insert a job with the debounce handle but DO NOT push through maybe_debounce, + // so it has no batch row at all. + let j = Uuid::new_v4(); + let args = serde_json::json!({ "items": [7, 8] }); + insert_script_job_with_args(&db, j, "test-workspace", "f/test/script", &args).await; + sqlx::query!( + "UPDATE v2_job_queue SET runnable_settings_handle = $1 WHERE id = $2", + rs_handle, + j, + ) + .execute(&db) + .await?; + + let mut res = make_pulled_job_result( + j, + "test-workspace", + "f/test/script", + &args, + JobKind::Script, + "deno", + rs_handle, + ); + res.maybe_apply_debouncing(&db).await?; + assert!(res.job.is_some(), "never-batched job still runs"); + assert_accumulated_items(&res, &[7, 8], "items"); + Ok(()) + } + + /// Edge: two survivors of one batch pulled CONCURRENTLY. The atomic claim must + /// partition the batch disjointly — the union of what they each accumulate is the + /// full set, with NO item processed by both. + #[sqlx::test(migrations = "../migrations", fixtures("base"))] + async fn test_debounce_concurrent_claim_disjoint(db: Pool) -> anyhow::Result<()> { + let key = "concurrent_claim_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; // superseded + let j2_args = push_debounced_script(&db, j2, vec![2], &settings, rs_handle).await; // survivor 1 + let j3 = Uuid::new_v4(); + let j3_args = add_survivor_to_batch_of(&db, j3, vec![3], j2, rs_handle).await; // survivor 2 + + let pull = |id: Uuid, args: serde_json::Value| { + let db = db.clone(); + async move { + let mut res = make_pulled_job_result( + id, + "test-workspace", + "f/test/script", + &args, + JobKind::Script, + "deno", + rs_handle, + ); + res.maybe_apply_debouncing(&db).await.unwrap(); + res + } + }; + let (r2, r3) = tokio::join!(pull(j2, j2_args), pull(j3, j3_args)); + + let mut union = items_of(&r2); + union.extend(items_of(&r3)); + union.sort(); + assert_eq!( + union, + vec![1, 2, 3], + "every item accumulated exactly once across the two concurrent survivors" + ); + // disjoint: no overlap between the two survivors' items + let i2 = items_of(&r2); + let i3 = items_of(&r3); + assert!( + !i2.iter().any(|x| i3.contains(x)), + "no item processed by both survivors; got j2={i2:?} j3={i3:?}" + ); + Ok(()) + } + + /// Edge: three survivors on one batch pulled in sequence. The first claims the whole + /// batch; the rest find themselves consumed and run empty. + #[sqlx::test(migrations = "../migrations", fixtures("base"))] + async fn test_debounce_three_survivors_first_takes_all( + db: Pool, + ) -> anyhow::Result<()> { + let key = "three_survivors_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; + let j3 = Uuid::new_v4(); + let j3_args = add_survivor_to_batch_of(&db, j3, vec![3], j2, rs_handle).await; + let j4 = Uuid::new_v4(); + let j4_args = add_survivor_to_batch_of(&db, j4, vec![4], j2, rs_handle).await; + + let mk = |id, args: &serde_json::Value| { + make_pulled_job_result( + id, + "test-workspace", + "f/test/script", + args, + JobKind::Script, + "deno", + rs_handle, + ) + }; + let (mut r2, mut r3, mut r4) = (mk(j2, &j2_args), mk(j3, &j3_args), mk(j4, &j4_args)); + r2.maybe_apply_debouncing(&db).await?; + r3.maybe_apply_debouncing(&db).await?; + r4.maybe_apply_debouncing(&db).await?; + + assert_eq!( + items_of(&r2), + vec![1, 2, 3, 4], + "first survivor takes the whole batch" + ); + assert!(items_of(&r3).is_empty(), "second survivor runs empty"); + assert!(items_of(&r4).is_empty(), "third survivor runs empty"); + Ok(()) + } + + /// Edge: plain debounce (no accumulate args) must HARD-DELETE its batch rows on pull + /// (not leave consumed rows lingering), so the non-accumulate path doesn't leak. + #[sqlx::test(migrations = "../migrations", fixtures("base"))] + async fn test_debounce_non_accumulate_deletes_batch(db: Pool) -> anyhow::Result<()> { + let key = "non_accum_key"; + let settings = DebouncingSettings { + debounce_delay_s: Some(5), + debounce_key: Some(key.to_string()), + // no debounce_args_to_accumulate + ..Default::default() + }; + let rs_handle = setup_debouncing_settings(&db, &settings).await; + + let j1 = Uuid::new_v4(); + let j2 = Uuid::new_v4(); + // push_debounced_script sends {items:[...]} but with no accumulate arg configured, + // the batch is created yet never accumulated. + 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; + + let batch_rows_before: i64 = + sqlx::query_scalar!("SELECT count(*) as \"c!\" FROM v2_job_debounce_batch") + .fetch_one(&db) + .await?; + assert!(batch_rows_before >= 2, "batch rows exist before pull"); + + let mut res = make_pulled_job_result( + j2, + "test-workspace", + "f/test/script", + &j2_args, + JobKind::Script, + "deno", + rs_handle, + ); + res.maybe_apply_debouncing(&db).await?; + + let remaining: i64 = + sqlx::query_scalar!("SELECT count(*) as \"c!\" FROM v2_job_debounce_batch") + .fetch_one(&db) + .await?; + assert_eq!( + remaining, 0, + "non-accumulate pull hard-deletes the batch rows" + ); + Ok(()) + } + + /// Edge: GC sweep deletes consumed rows past the grace period but keeps recently + /// consumed and not-yet-consumed rows. + #[sqlx::test(migrations = "../migrations", fixtures("base"))] + async fn test_debounce_gc_consumed_batches(db: Pool) -> anyhow::Result<()> { + let old = Uuid::new_v4(); + let recent = Uuid::new_v4(); + let unconsumed = Uuid::new_v4(); + sqlx::query!( + "INSERT INTO v2_job_debounce_batch (id, debounce_batch, consumed_at) VALUES + ($1, nextval('debounce_batch_seq'), now() - interval '20 minutes'), + ($2, nextval('debounce_batch_seq'), now() - interval '1 minute'), + ($3, nextval('debounce_batch_seq'), NULL)", + old, + recent, + unconsumed, + ) + .execute(&db) + .await?; + + // Mirror the monitor GC sweep. + let deleted = sqlx::query_scalar!( + "WITH del AS ( + DELETE FROM v2_job_debounce_batch + WHERE consumed_at IS NOT NULL AND consumed_at < now() - interval '10 minutes' + RETURNING 1 + ) SELECT count(*) as \"c!\" FROM del" + ) + .fetch_one(&db) + .await?; + assert_eq!(deleted, 1, "only the old consumed row is GC'd"); + + let exists = |id: Uuid, db: Pool| async move { + sqlx::query_scalar!( + "SELECT EXISTS(SELECT 1 FROM v2_job_debounce_batch WHERE id = $1) as \"e!\"", + id + ) + .fetch_one(&db) + .await + .unwrap() + }; + assert!(!exists(old, db.clone()).await, "old consumed row gone"); + assert!( + exists(recent, db.clone()).await, + "recently consumed row kept" + ); + assert!(exists(unconsumed, db.clone()).await, "unconsumed row kept"); + 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)