diff --git a/backend/.sqlx/query-8597cd40f80e69edbf1bc7d7402baca32e33e871be454acb5175c11361fe1b0a.json b/backend/.sqlx/query-8597cd40f80e69edbf1bc7d7402baca32e33e871be454acb5175c11361fe1b0a.json new file mode 100644 index 0000000000..fd82867507 --- /dev/null +++ b/backend/.sqlx/query-8597cd40f80e69edbf1bc7d7402baca32e33e871be454acb5175c11361fe1b0a.json @@ -0,0 +1,14 @@ +{ + "db_name": "PostgreSQL", + "query": "DELETE FROM background_task_state\n WHERE name LIKE $1\n AND updated_at < NOW() - INTERVAL '7 days'", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "Text" + ] + }, + "nullable": [] + }, + "hash": "8597cd40f80e69edbf1bc7d7402baca32e33e871be454acb5175c11361fe1b0a" +} diff --git a/backend/src/monitor.rs b/backend/src/monitor.rs index defd47489c..13f52037e7 100644 --- a/backend/src/monitor.rs +++ b/backend/src/monitor.rs @@ -2705,6 +2705,26 @@ pub async fn monitor_db( } }; + // run every hour (120 iterations * 30s = 3600s) + let cleanup_stale_server_heartbeats_f = async { + if server_mode && iteration.is_some() && iteration.as_ref().unwrap().should_run(120) { + if let Some(db) = conn.as_sql() { + match windmill_api::cleanup_stale_server_heartbeats(db).await { + Ok(count) if count > 0 => { + tracing::info!( + "Deleted {} stale server_heartbeat background_task_state rows", + count + ); + } + Err(e) => { + tracing::error!("Error cleaning up stale server_heartbeat rows: {:?}", e); + } + _ => {} + } + } + } + }; + // run every hour (120 iterations * 30s = 3600s) let manage_audit_partitions_f = async { if server_mode && iteration.is_some() && iteration.as_ref().unwrap().should_run(120) { @@ -2763,6 +2783,7 @@ pub async fn monitor_db( native_triggers_sync_f, cleanup_notify_events_f, check_expiring_tokens_f, + cleanup_stale_server_heartbeats_f, manage_audit_partitions_f, export_audit_logs_to_object_store_f, cleanup_scheduled_job_deletions_f, diff --git a/backend/windmill-api/src/lib.rs b/backend/windmill-api/src/lib.rs index 8b2b9d8128..41872d0353 100644 --- a/backend/windmill-api/src/lib.rs +++ b/backend/windmill-api/src/lib.rs @@ -1261,3 +1261,26 @@ pub async fn check_any_server_started(db: &DB, not_before: chrono::DateTime not_before` (the moment a restart was initiated), +/// so rows older than the cutoff cannot influence any restart decision and +/// are safe to delete. +pub async fn cleanup_stale_server_heartbeats(db: &DB) -> anyhow::Result { + let prefix = format!("{SERVER_HEARTBEAT_TASK}:"); + let res = sqlx::query!( + "DELETE FROM background_task_state + WHERE name LIKE $1 + AND updated_at < NOW() - INTERVAL '7 days'", + format!("{prefix}%"), + ) + .execute(db) + .await?; + Ok(res.rows_affected()) +}