* Main queue processing loop
()
| 191 | * Main queue processing loop |
| 192 | */ |
| 193 | async function processQueue() { |
| 194 | const BATCH_SIZE = Number.parseInt(process.env.UPDATE_BATCH_SIZE || "100"); |
| 195 | |
| 196 | console.log(`[Queue] Starting with batch size ${BATCH_SIZE}`); |
| 197 | |
| 198 | while (!isShuttingDown) { |
| 199 | try { |
| 200 | const jobs = await dequeueJobs(BATCH_SIZE); |
| 201 | |
| 202 | if (jobs.length > 0) { |
| 203 | await processBatch(jobs); |
| 204 | } else { |
| 205 | await handleIdleQueue(BATCH_SIZE); |
| 206 | } |
| 207 | |
| 208 | // Wait before next poll |
| 209 | await new Promise((resolve) => |
| 210 | setTimeout(resolve, QUEUE_POLL_INTERVAL_MS) |
| 211 | ); |
| 212 | } catch (error) { |
| 213 | console.error("[Queue] Unexpected error in main loop:", error); |
| 214 | // Continue processing despite errors in individual iterations |
| 215 | } |
| 216 | } |
| 217 | |
| 218 | console.log("[Queue] Shutdown complete"); |
| 219 | } |
| 220 | |
| 221 | if (import.meta.url === `file://${process.argv[1]}`) { |
| 222 | processQueue().catch((error) => { |
no test coverage detected