fix: add orphan adoption and atomic v2_job_completed deletion in job cleanup

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
This commit is contained in:
Ruben Fiszel
2026-02-16 07:03:45 +00:00
co-authored by Claude Opus 4.6
parent 0385996ba1
commit be43f3b3fe
9 changed files with 91 additions and 80 deletions
@@ -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"
}
@@ -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"
}
@@ -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"
}
@@ -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"
}
@@ -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"
}
@@ -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"
}
@@ -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
+27 -13
View File
@@ -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<usize> {
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<Uuid> = 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",
@@ -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",