diff --git a/backend/.sqlx/query-fe1539db7384c8edc6d8ec672495fe3964efb8551ad6d74f141fc457034cc5b9.json b/backend/.sqlx/query-4f11760f5d283728ded5533e6a0b51f49fdb7c11bf7d47b8536d607e646bd265.json similarity index 58% rename from backend/.sqlx/query-fe1539db7384c8edc6d8ec672495fe3964efb8551ad6d74f141fc457034cc5b9.json rename to backend/.sqlx/query-4f11760f5d283728ded5533e6a0b51f49fdb7c11bf7d47b8536d607e646bd265.json index 7b7e250ef7..40c82a50aa 100644 --- a/backend/.sqlx/query-fe1539db7384c8edc6d8ec672495fe3964efb8551ad6d74f141fc457034cc5b9.json +++ b/backend/.sqlx/query-4f11760f5d283728ded5533e6a0b51f49fdb7c11bf7d47b8536d607e646bd265.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "SELECT job_id FROM debounce_key WHERE key = $1 FOR UPDATE", + "query": "SELECT job_id FROM debounce_key WHERE key = $1 AND job_id IN (SELECT id FROM v2_job_queue) FOR UPDATE", "describe": { "columns": [ { @@ -18,5 +18,5 @@ false ] }, - "hash": "fe1539db7384c8edc6d8ec672495fe3964efb8551ad6d74f141fc457034cc5b9" + "hash": "4f11760f5d283728ded5533e6a0b51f49fdb7c11bf7d47b8536d607e646bd265" } diff --git a/backend/.sqlx/query-529a52823913f4154786a2ada93f6112c39575a52a1bfc02f27d9b1185b0578e.json b/backend/.sqlx/query-529a52823913f4154786a2ada93f6112c39575a52a1bfc02f27d9b1185b0578e.json new file mode 100644 index 0000000000..732ae64ad6 --- /dev/null +++ b/backend/.sqlx/query-529a52823913f4154786a2ada93f6112c39575a52a1bfc02f27d9b1185b0578e.json @@ -0,0 +1,26 @@ +{ + "db_name": "PostgreSQL", + "query": "\nDELETE FROM debounce_key\nWHERE job_id NOT IN (SELECT id FROM v2_job_queue)\nRETURNING key,job_id\n ", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "key", + "type_info": "Varchar" + }, + { + "ordinal": 1, + "name": "job_id", + "type_info": "Uuid" + } + ], + "parameters": { + "Left": [] + }, + "nullable": [ + false, + false + ] + }, + "hash": "529a52823913f4154786a2ada93f6112c39575a52a1bfc02f27d9b1185b0578e" +} diff --git a/backend/src/monitor.rs b/backend/src/monitor.rs index 9d4a8a5fae..13cecca007 100644 --- a/backend/src/monitor.rs +++ b/backend/src/monitor.rs @@ -1589,6 +1589,16 @@ pub async fn monitor_db( } }; + let cleanup_debounce_keys_f = async { + if server_mode && iteration.is_some() && iteration.as_ref().unwrap().should_run(20) { + if let Some(db) = conn.as_sql() { + if let Err(e) = cleanup_debounce_orphaned_keys(&db).await { + tracing::error!("Error cleaning up debounce keys: {:?}", e); + } + } + } + }; + // run every hour (60 minutes / 30 seconds = 120) let cleanup_worker_group_stats_f = async { if server_mode && iteration.is_some() && iteration.as_ref().unwrap().should_run(120) { @@ -1707,6 +1717,7 @@ pub async fn monitor_db( update_min_worker_version_f, cleanup_concurrency_counters_f, cleanup_concurrency_counters_empty_keys_f, + cleanup_debounce_keys_f, cleanup_worker_group_stats_f, ); } @@ -2754,3 +2765,29 @@ pub async fn reload_jwt_secret_setting(db: &DB) -> error::Result<()> { Ok(()) } + +async fn cleanup_debounce_orphaned_keys(db: &DB) -> error::Result<()> { + let result = sqlx::query!( + " +DELETE FROM debounce_key +WHERE job_id NOT IN (SELECT id FROM v2_job_queue) +RETURNING key,job_id + ", + ) + .fetch_all(db) + .await?; + + tracing::debug!("Cleaning up debounce keys"); + + if result.len() > 0 { + tracing::info!("Cleaned up {} debounce keys", result.len()); + for row in result { + tracing::info!( + "Debounce key cleaned up: key: {}, job_id: {:?}", + row.key, + row.job_id + ); + } + } + Ok(()) +} diff --git a/backend/windmill-common/src/jobs.rs b/backend/windmill-common/src/jobs.rs index edb9f1f0cf..20107cc5ee 100644 --- a/backend/windmill-common/src/jobs.rs +++ b/backend/windmill-common/src/jobs.rs @@ -814,7 +814,7 @@ pub async fn lock_debounce_key<'c>( ); sqlx::query_scalar!( - "SELECT job_id FROM debounce_key WHERE key = $1 FOR UPDATE", + "SELECT job_id FROM debounce_key WHERE key = $1 AND job_id IN (SELECT id FROM v2_job_queue) FOR UPDATE", &key ) .fetch_optional(&mut **tx)