diff --git a/backend/.sqlx/query-07168aaf14cb6beff0ad4274b441f7f387f5055c47f493271d26731336257384.json b/backend/.sqlx/query-07168aaf14cb6beff0ad4274b441f7f387f5055c47f493271d26731336257384.json index d29a18c691..e7ed0aee65 100644 --- a/backend/.sqlx/query-07168aaf14cb6beff0ad4274b441f7f387f5055c47f493271d26731336257384.json +++ b/backend/.sqlx/query-07168aaf14cb6beff0ad4274b441f7f387f5055c47f493271d26731336257384.json @@ -46,11 +46,11 @@ ] }, "nullable": [ - true, - true, - true, - true, - true, + false, + false, + false, + false, + false, true, true ] diff --git a/backend/.sqlx/query-306e0156ee1541710c1c6512ecb4f61baeb3ae6f31ba3fd57a3ec485108a7f49.json b/backend/.sqlx/query-306e0156ee1541710c1c6512ecb4f61baeb3ae6f31ba3fd57a3ec485108a7f49.json deleted file mode 100644 index 1e05f032f0..0000000000 --- a/backend/.sqlx/query-306e0156ee1541710c1c6512ecb4f61baeb3ae6f31ba3fd57a3ec485108a7f49.json +++ /dev/null @@ -1,23 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "DELETE FROM v2_job_completed\n WHERE id IN (\n SELECT id FROM v2_job_completed\n WHERE completed_at <= now() - ($1::bigint::text || ' s')::interval\n ORDER BY completed_at ASC\n LIMIT $2\n FOR UPDATE SKIP LOCKED\n )\n RETURNING id", - "describe": { - "columns": [ - { - "ordinal": 0, - "name": "id", - "type_info": "Uuid" - } - ], - "parameters": { - "Left": [ - "Int8", - "Int8" - ] - }, - "nullable": [ - false - ] - }, - "hash": "306e0156ee1541710c1c6512ecb4f61baeb3ae6f31ba3fd57a3ec485108a7f49" -} diff --git a/backend/.sqlx/query-45997fcb4d9d62c7f7011966bf59bdb86e12ce1d0c8e925e738d2645121a5c1f.json b/backend/.sqlx/query-45997fcb4d9d62c7f7011966bf59bdb86e12ce1d0c8e925e738d2645121a5c1f.json new file mode 100644 index 0000000000..f63afbc986 --- /dev/null +++ b/backend/.sqlx/query-45997fcb4d9d62c7f7011966bf59bdb86e12ce1d0c8e925e738d2645121a5c1f.json @@ -0,0 +1,24 @@ +{ + "db_name": "PostgreSQL", + "query": "DELETE FROM v2_job_completed\n WHERE id IN (\n SELECT jc.id 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 FOR UPDATE OF jc SKIP LOCKED\n )\n RETURNING id", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "id", + "type_info": "Uuid" + } + ], + "parameters": { + "Left": [ + "Int8", + "Int8", + "UuidArray" + ] + }, + "nullable": [ + false + ] + }, + "hash": "45997fcb4d9d62c7f7011966bf59bdb86e12ce1d0c8e925e738d2645121a5c1f" +} diff --git a/backend/.sqlx/query-5a219a2532517869578c4504ff3153c43903f929ae5d62fbba12610f89c36d55.json b/backend/.sqlx/query-5a219a2532517869578c4504ff3153c43903f929ae5d62fbba12610f89c36d55.json index 713ccb9dd3..36ddb8ab9f 100644 --- a/backend/.sqlx/query-5a219a2532517869578c4504ff3153c43903f929ae5d62fbba12610f89c36d55.json +++ b/backend/.sqlx/query-5a219a2532517869578c4504ff3153c43903f929ae5d62fbba12610f89c36d55.json @@ -15,7 +15,7 @@ ] }, "nullable": [ - null + true ] }, "hash": "5a219a2532517869578c4504ff3153c43903f929ae5d62fbba12610f89c36d55" diff --git a/backend/.sqlx/query-c01a94ec75990c8fb4068485c91211af295a1d2630861bf53a33966abc4562ec.json b/backend/.sqlx/query-c01a94ec75990c8fb4068485c91211af295a1d2630861bf53a33966abc4562ec.json new file mode 100644 index 0000000000..6e190acda3 --- /dev/null +++ b/backend/.sqlx/query-c01a94ec75990c8fb4068485c91211af295a1d2630861bf53a33966abc4562ec.json @@ -0,0 +1,22 @@ +{ + "db_name": "PostgreSQL", + "query": "SELECT q.id FROM v2_job_queue q\n JOIN v2_job j ON j.id = q.id\n WHERE j.parent_job IS NULL\n AND j.created_at <= now() - ($1::bigint::text || ' s')::interval", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "id", + "type_info": "Uuid" + } + ], + "parameters": { + "Left": [ + "Int8" + ] + }, + "nullable": [ + false + ] + }, + "hash": "c01a94ec75990c8fb4068485c91211af295a1d2630861bf53a33966abc4562ec" +} diff --git a/backend/src/monitor.rs b/backend/src/monitor.rs index adc00a9098..3e4e7dce92 100644 --- a/backend/src/monitor.rs +++ b/backend/src/monitor.rs @@ -1059,20 +1059,36 @@ async fn delete_expired_jobs_batch( ) -> error::Result { let mut tx = db.begin().await?; + // Fetch active ROOT job IDs that started before the retention period. We only care about + // these because their child jobs could be old enough to be deletion candidates. + // Jobs started after the retention period can't have children old enough to delete. + 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 + WHERE j.parent_job IS NULL + AND j.created_at <= now() - ($1::bigint::text || ' s')::interval", + job_retention_secs + ) + .fetch_all(&mut *tx) + .await?; + // Use FOR UPDATE SKIP LOCKED to avoid contention between replicas // ORDER BY completed_at ensures we delete oldest jobs first let deleted_jobs: Vec = sqlx::query_scalar!( "DELETE FROM v2_job_completed WHERE id IN ( - SELECT id FROM v2_job_completed - WHERE completed_at <= now() - ($1::bigint::text || ' s')::interval - ORDER BY completed_at ASC + SELECT jc.id 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 - FOR UPDATE SKIP LOCKED + FOR UPDATE OF jc SKIP LOCKED ) RETURNING id", job_retention_secs, - batch_size + batch_size, + &active_root_job_ids ) .fetch_all(&mut *tx) .await?;