diff --git a/apps/nuq-postgres/nuq.sql b/apps/nuq-postgres/nuq.sql index 79b7626b4..78c0a2d81 100644 --- a/apps/nuq-postgres/nuq.sql +++ b/apps/nuq-postgres/nuq.sql @@ -53,8 +53,46 @@ SELECT cron.schedule('nuq_queue_scrape_clean_failed', '*/5 * * * *', $$ $$); SELECT cron.schedule('nuq_queue_scrape_lock_reaper', '15 seconds', $$ - UPDATE nuq.queue_scrape SET status = 'queued'::nuq.job_status, lock = null, locked_at = null, stalls = COALESCE(stalls, 0) + 1 WHERE nuq.queue_scrape.locked_at <= now() - interval '1 minute' AND nuq.queue_scrape.status = 'active'::nuq.job_status AND COALESCE(nuq.queue_scrape.stalls, 0) < 9; - WITH stallfail AS (UPDATE nuq.queue_scrape SET status = 'failed'::nuq.job_status, lock = null, locked_at = null, stalls = COALESCE(stalls, 0) + 1 WHERE nuq.queue_scrape.locked_at <= now() - interval '1 minute' AND nuq.queue_scrape.status = 'active'::nuq.job_status AND COALESCE(nuq.queue_scrape.stalls, 0) >= 9 RETURNING id) + WITH requeued AS ( + UPDATE nuq.queue_scrape + SET status = 'queued'::nuq.job_status, lock = null, locked_at = null, stalls = COALESCE(stalls, 0) + 1 + WHERE nuq.queue_scrape.locked_at <= now() - interval '1 minute' + AND nuq.queue_scrape.status = 'active'::nuq.job_status + AND COALESCE(nuq.queue_scrape.stalls, 0) < 9 + RETURNING id, owner_id + ), + requeued_counts AS ( + SELECT owner_id, COUNT(*) as job_count + FROM requeued + WHERE owner_id IS NOT NULL + GROUP BY owner_id + ), + requeue_concurrency_update AS ( + UPDATE nuq.queue_scrape_owner_concurrency + SET current_concurrency = GREATEST(0, current_concurrency - requeued_counts.job_count) + FROM requeued_counts + WHERE nuq.queue_scrape_owner_concurrency.id = requeued_counts.owner_id + ), + stallfail AS ( + UPDATE nuq.queue_scrape + SET status = 'failed'::nuq.job_status, lock = null, locked_at = null, stalls = COALESCE(stalls, 0) + 1 + WHERE nuq.queue_scrape.locked_at <= now() - interval '1 minute' + AND nuq.queue_scrape.status = 'active'::nuq.job_status + AND COALESCE(nuq.queue_scrape.stalls, 0) >= 9 + RETURNING id, owner_id + ), + stallfail_counts AS ( + SELECT owner_id, COUNT(*) as job_count + FROM stallfail + WHERE owner_id IS NOT NULL + GROUP BY owner_id + ), + stallfail_concurrency_update AS ( + UPDATE nuq.queue_scrape_owner_concurrency + SET current_concurrency = GREATEST(0, current_concurrency - stallfail_counts.job_count) + FROM stallfail_counts + WHERE nuq.queue_scrape_owner_concurrency.id = stallfail_counts.owner_id + ) SELECT pg_notify('nuq.queue_scrape', (id::text || '|' || 'failed'::text)) FROM stallfail; $$);