diff --git a/backend/.sqlx/query-6eadb8d64cce490f82c1f8e328a0d9a5c1a0015bb82326a1efc4f194b6d4ba70.json b/backend/.sqlx/query-6eadb8d64cce490f82c1f8e328a0d9a5c1a0015bb82326a1efc4f194b6d4ba70.json deleted file mode 100644 index 1f11f3135c..0000000000 --- a/backend/.sqlx/query-6eadb8d64cce490f82c1f8e328a0d9a5c1a0015bb82326a1efc4f194b6d4ba70.json +++ /dev/null @@ -1,17 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "INSERT INTO jobs_pending_deletion (id, marked_by)\n SELECT jc.id, $4 FROM v2_job_completed jc\n LEFT JOIN v2_job j ON j.id = jc.id\n WHERE jc.completed_at <= now() - ($1::bigint::text || ' s')::interval\n AND COALESCE(j.root_job, j.flow_innermost_root_job, jc.id) != ALL($3)\n ORDER BY jc.completed_at ASC\n LIMIT $2\n ON CONFLICT DO NOTHING", - "describe": { - "columns": [], - "parameters": { - "Left": [ - "Int8", - "Int8", - "UuidArray", - "Varchar" - ] - }, - "nullable": [] - }, - "hash": "6eadb8d64cce490f82c1f8e328a0d9a5c1a0015bb82326a1efc4f194b6d4ba70" -} diff --git a/backend/.sqlx/query-7420b825e8e2ee74b7839dbffec5709b99faecb45044069bfa895262fb5ca940.json b/backend/.sqlx/query-7420b825e8e2ee74b7839dbffec5709b99faecb45044069bfa895262fb5ca940.json new file mode 100644 index 0000000000..c4043614f4 --- /dev/null +++ b/backend/.sqlx/query-7420b825e8e2ee74b7839dbffec5709b99faecb45044069bfa895262fb5ca940.json @@ -0,0 +1,14 @@ +{ + "db_name": "PostgreSQL", + "query": "UPDATE jobs_pending_deletion SET marked_by = $1, marked_at = now()\n WHERE marked_at < now() - interval '30 minutes'", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "Varchar" + ] + }, + "nullable": [] + }, + "hash": "7420b825e8e2ee74b7839dbffec5709b99faecb45044069bfa895262fb5ca940" +} diff --git a/backend/.sqlx/query-a9d80d005362d1748e354c6bb6f19b5205ecd4c92f16f2b756dddd3485b1b7d6.json b/backend/.sqlx/query-a9d80d005362d1748e354c6bb6f19b5205ecd4c92f16f2b756dddd3485b1b7d6.json new file mode 100644 index 0000000000..7894a920de --- /dev/null +++ b/backend/.sqlx/query-a9d80d005362d1748e354c6bb6f19b5205ecd4c92f16f2b756dddd3485b1b7d6.json @@ -0,0 +1,17 @@ +{ + "db_name": "PostgreSQL", + "query": "WITH marked AS (\n INSERT INTO jobs_pending_deletion (id, marked_by)\n SELECT jc.id, $4 FROM v2_job_completed jc\n LEFT JOIN v2_job j ON j.id = jc.id\n WHERE jc.completed_at <= now() - ($1::bigint::text || ' s')::interval\n AND COALESCE(j.root_job, j.flow_innermost_root_job, jc.id) != ALL($3)\n ORDER BY jc.completed_at ASC\n LIMIT $2\n ON CONFLICT DO NOTHING\n RETURNING id\n )\n DELETE FROM v2_job_completed WHERE id IN (SELECT id FROM marked)", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "Int8", + "Int8", + "UuidArray", + "Varchar" + ] + }, + "nullable": [] + }, + "hash": "a9d80d005362d1748e354c6bb6f19b5205ecd4c92f16f2b756dddd3485b1b7d6" +} diff --git a/backend/.sqlx/query-c25f6aa099050ded451fb6316e9c07630030d4c9cb09eb6bd5e46b043e798322.json b/backend/.sqlx/query-c25f6aa099050ded451fb6316e9c07630030d4c9cb09eb6bd5e46b043e798322.json deleted file mode 100644 index 5d0c33a2af..0000000000 --- a/backend/.sqlx/query-c25f6aa099050ded451fb6316e9c07630030d4c9cb09eb6bd5e46b043e798322.json +++ /dev/null @@ -1,18 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "INSERT INTO jobs_pending_deletion (id, marked_by)\n SELECT jc.id, $5 FROM v2_job_completed jc\n JOIN v2_job j ON j.id = jc.id\n WHERE jc.completed_at <= now() - ($1::bigint::text || ' s')::interval\n AND COALESCE(j.root_job, j.flow_innermost_root_job, jc.id) != ALL($3)\n AND j.tag = $4\n ORDER BY jc.completed_at ASC\n LIMIT $2\n ON CONFLICT DO NOTHING", - "describe": { - "columns": [], - "parameters": { - "Left": [ - "Int8", - "Int8", - "UuidArray", - "Text", - "Varchar" - ] - }, - "nullable": [] - }, - "hash": "c25f6aa099050ded451fb6316e9c07630030d4c9cb09eb6bd5e46b043e798322" -} diff --git a/backend/.sqlx/query-d5145506c6f835522547bb18facf1dea9447fe9e06ab7eb3f50bdf879029b95f.json b/backend/.sqlx/query-d5145506c6f835522547bb18facf1dea9447fe9e06ab7eb3f50bdf879029b95f.json deleted file mode 100644 index 8d084cf089..0000000000 --- a/backend/.sqlx/query-d5145506c6f835522547bb18facf1dea9447fe9e06ab7eb3f50bdf879029b95f.json +++ /dev/null @@ -1,14 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "DELETE FROM v2_job_completed USING jobs_pending_deletion d\n WHERE v2_job_completed.id = d.id AND d.marked_by = $1", - "describe": { - "columns": [], - "parameters": { - "Left": [ - "Text" - ] - }, - "nullable": [] - }, - "hash": "d5145506c6f835522547bb18facf1dea9447fe9e06ab7eb3f50bdf879029b95f" -} diff --git a/backend/.sqlx/query-db3b300319cd3a5dee5dd078dd9df4729e68170d85ab6988e19ac6c8fc753253.json b/backend/.sqlx/query-db3b300319cd3a5dee5dd078dd9df4729e68170d85ab6988e19ac6c8fc753253.json new file mode 100644 index 0000000000..e05bc357e7 --- /dev/null +++ b/backend/.sqlx/query-db3b300319cd3a5dee5dd078dd9df4729e68170d85ab6988e19ac6c8fc753253.json @@ -0,0 +1,18 @@ +{ + "db_name": "PostgreSQL", + "query": "WITH marked AS (\n INSERT INTO jobs_pending_deletion (id, marked_by)\n SELECT jc.id, $5 FROM v2_job_completed jc\n JOIN v2_job j ON j.id = jc.id\n WHERE jc.completed_at <= now() - ($1::bigint::text || ' s')::interval\n AND COALESCE(j.root_job, j.flow_innermost_root_job, jc.id) != ALL($3)\n AND j.tag = $4\n ORDER BY jc.completed_at ASC\n LIMIT $2\n ON CONFLICT DO NOTHING\n RETURNING id\n )\n DELETE FROM v2_job_completed WHERE id IN (SELECT id FROM marked)", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "Int8", + "Int8", + "UuidArray", + "Text", + "Varchar" + ] + }, + "nullable": [] + }, + "hash": "db3b300319cd3a5dee5dd078dd9df4729e68170d85ab6988e19ac6c8fc753253" +} diff --git a/backend/migrations/20260215000000_jobs_pending_deletion.up.sql b/backend/migrations/20260215000000_jobs_pending_deletion.up.sql index 1e528b670f..0b9bdde74d 100644 --- a/backend/migrations/20260215000000_jobs_pending_deletion.up.sql +++ b/backend/migrations/20260215000000_jobs_pending_deletion.up.sql @@ -1,6 +1,7 @@ CREATE TABLE IF NOT EXISTS jobs_pending_deletion ( id UUID PRIMARY KEY, - marked_by VARCHAR(10) NOT NULL + marked_by VARCHAR(10) NOT NULL, + marked_at TIMESTAMPTZ NOT NULL DEFAULT now() ); CREATE INDEX IF NOT EXISTS idx_jobs_pending_deletion_marked_by diff --git a/backend/src/monitor.rs b/backend/src/monitor.rs index ee91fdfbc2..a7cfd5495b 100644 --- a/backend/src/monitor.rs +++ b/backend/src/monitor.rs @@ -1064,7 +1064,8 @@ pub async fn delete_expired_items(db: &DB) -> () { /// No job UUIDs are held in application memory — all joins stay inside Postgres. /// /// Rows are scoped by `marked_by` (the server's `INSTANCE_NAME`) so multiple servers -/// can run cleanup concurrently without interfering with each other. +/// can run cleanup concurrently without interfering with each other. Stale rows from dead +/// servers are adopted via `marked_at` timestamp after 30 minutes. /// /// Returns the number of jobs marked for deletion in this batch. async fn delete_expired_jobs_batch( @@ -1074,8 +1075,16 @@ async fn delete_expired_jobs_batch( ) -> error::Result { let instance = &*INSTANCE_NAME; - // Step 0: Crash recovery — check for leftover rows from a previous incomplete run - // by this same server instance. + // Step 0: Adopt orphaned rows from dead servers (stale for >30 min), then check + // for leftover rows from a previous incomplete run by this same instance. + sqlx::query!( + "UPDATE jobs_pending_deletion SET marked_by = $1, marked_at = now() + WHERE marked_at < now() - interval '30 minutes'", + instance + ) + .execute(db) + .await?; + let pending_count: i64 = sqlx::query_scalar!( "SELECT COUNT(*) FROM jobs_pending_deletion WHERE marked_by = $1", instance @@ -1092,7 +1101,7 @@ async fn delete_expired_jobs_batch( ); pending_count as usize } else { - // Step 1: Mark expired jobs into the staging table. + // Step 1: Mark expired jobs and atomically remove them from v2_job_completed. let active_root_job_ids: Vec = sqlx::query_scalar!( "SELECT q.id FROM v2_job_queue q JOIN v2_job j ON j.id = q.id @@ -1104,14 +1113,18 @@ async fn delete_expired_jobs_batch( .await?; let result = sqlx::query!( - "INSERT INTO jobs_pending_deletion (id, marked_by) - SELECT jc.id, $4 FROM v2_job_completed jc - LEFT JOIN v2_job j ON j.id = jc.id - WHERE jc.completed_at <= now() - ($1::bigint::text || ' s')::interval - AND COALESCE(j.root_job, j.flow_innermost_root_job, jc.id) != ALL($3) - ORDER BY jc.completed_at ASC - LIMIT $2 - ON CONFLICT DO NOTHING", + "WITH marked AS ( + INSERT INTO jobs_pending_deletion (id, marked_by) + SELECT jc.id, $4 FROM v2_job_completed jc + LEFT JOIN v2_job j ON j.id = jc.id + WHERE jc.completed_at <= now() - ($1::bigint::text || ' s')::interval + AND COALESCE(j.root_job, j.flow_innermost_root_job, jc.id) != ALL($3) + ORDER BY jc.completed_at ASC + LIMIT $2 + ON CONFLICT DO NOTHING + RETURNING id + ) + DELETE FROM v2_job_completed WHERE id IN (SELECT id FROM marked)", job_retention_secs, batch_size, &active_root_job_ids, @@ -1167,7 +1180,8 @@ async fn delete_expired_jobs_batch( Err(e) => tracing::error!("Error deleting job logs: {:?}", e), } - // Step 4: Delete from v2_job_completed + // Step 4: Delete from v2_job_completed (idempotent — already done for fresh batches + // via the CTE in step 1, but needed for crash-recovery and adopted orphan rows). if let Err(e) = sqlx::query!( "DELETE FROM v2_job_completed USING jobs_pending_deletion d WHERE v2_job_completed.id = d.id AND d.marked_by = $1", diff --git a/backend/windmill-queue/tests/job_cleanup_test.rs b/backend/windmill-queue/tests/job_cleanup_test.rs index 40f841bc8d..2920378685 100644 --- a/backend/windmill-queue/tests/job_cleanup_test.rs +++ b/backend/windmill-queue/tests/job_cleanup_test.rs @@ -239,15 +239,19 @@ async fn run_new_method( .unwrap(); let result = sqlx::query!( - "INSERT INTO jobs_pending_deletion (id, marked_by) - SELECT jc.id, $5 FROM v2_job_completed jc - JOIN v2_job j ON j.id = jc.id - WHERE jc.completed_at <= now() - ($1::bigint::text || ' s')::interval - AND COALESCE(j.root_job, j.flow_innermost_root_job, jc.id) != ALL($3) - AND j.tag = $4 - ORDER BY jc.completed_at ASC - LIMIT $2 - ON CONFLICT DO NOTHING", + "WITH marked AS ( + INSERT INTO jobs_pending_deletion (id, marked_by) + SELECT jc.id, $5 FROM v2_job_completed jc + JOIN v2_job j ON j.id = jc.id + WHERE jc.completed_at <= now() - ($1::bigint::text || ' s')::interval + AND COALESCE(j.root_job, j.flow_innermost_root_job, jc.id) != ALL($3) + AND j.tag = $4 + ORDER BY jc.completed_at ASC + LIMIT $2 + ON CONFLICT DO NOTHING + RETURNING id + ) + DELETE FROM v2_job_completed WHERE id IN (SELECT id FROM marked)", retention_secs, batch_size, &active_roots, @@ -279,14 +283,6 @@ async fn run_new_method( .execute(db) .await .ok(); - sqlx::query!( - "DELETE FROM v2_job_completed USING jobs_pending_deletion d - WHERE v2_job_completed.id = d.id AND d.marked_by = $1", - instance - ) - .execute(db) - .await - .ok(); sqlx::query!( "DELETE FROM v2_job USING jobs_pending_deletion d WHERE v2_job.id = d.id AND d.marked_by = $1",