'Keep running until stopped', '--sleep' => 'Seconds between polls when idle (default 2)', '--max-jobs' => 'Exit after N successful jobs', '--stale' => 'Also reap jobs stuck in processing beyond max runtime', ]; public function run(array $params): void { $loop = array_key_exists('loop', $params) || CLI::getOption('loop'); $sleep = (int) ($params['sleep'] ?? 2); $maxJobs = isset($params['max-jobs']) ? (int) $params['max-jobs'] : 0; $jobs = model(JobModel::class); $done = 0; do { $job = null; // claim with retry: SKIP LOCKED can race under contention try { $job = $jobs->claimNext(); } catch (\Throwable $e) { log_message('error', 'worker claim failed: {m}', ['m' => $e->getMessage()]); usleep(500_000); } if ($job === null) { if (! $loop) { break; } sleep(max(1, $sleep)); continue; } CLI::write("Processing job {$job['id']} [{$job['operation']}]", 'green'); self::heartbeat(); service('pipeline')->process($job); ++$done; if ($done % 10 === 0) { self::reapStale($jobs); // periodic housekeeping inside long-running workers } if ($maxJobs > 0 && $done >= $maxJobs) { break; } } while (true); CLI::write("Worker finished after {$done} job(s).", 'yellow'); } /** Touch a cache key so the admin dashboard can see workers are alive. */ public static function heartbeat(): void { try { cache()->save('tv_worker_heartbeat', gmdate('H:i:s'), 120); } catch (\Throwable) { // never let telemetry break processing } } /** Fail jobs whose worker died mid-processing. */ public static function reapStale(JobModel $jobs): void { foreach ($jobs->findStale((int) config('Site')->maxProcessingSeconds) as $stale) { Pipeline::failJob($stale['id'], 'Processing timed out.'); } } }