Code Coverage
 
Lines
Functions and Methods
Classes and Traits
Total
97.33% covered (success)
97.33%
182 / 187
94.23% covered (success)
94.23%
49 / 52
CRAP
0.00% covered (danger)
0.00%
0 / 1
Worker
97.33% covered (success)
97.33%
182 / 187
94.23% covered (success)
94.23%
49 / 52
102
0.00% covered (danger)
0.00%
0 / 1
 __construct
100.00% covered (success)
100.00%
6 / 6
100.00% covered (success)
100.00%
1 / 1
4
 create
100.00% covered (success)
100.00%
1 / 1
100.00% covered (success)
100.00%
1 / 1
1
 getApplication
100.00% covered (success)
100.00%
1 / 1
100.00% covered (success)
100.00%
1 / 1
1
 application
100.00% covered (success)
100.00%
1 / 1
100.00% covered (success)
100.00%
1 / 1
1
 hasApplication
100.00% covered (success)
100.00%
1 / 1
100.00% covered (success)
100.00%
1 / 1
1
 setEvents
100.00% covered (success)
100.00%
2 / 2
100.00% covered (success)
100.00%
1 / 1
1
 getEvents
100.00% covered (success)
100.00%
1 / 1
100.00% covered (success)
100.00%
1 / 1
1
 events
100.00% covered (success)
100.00%
1 / 1
100.00% covered (success)
100.00%
1 / 1
1
 hasEvents
100.00% covered (success)
100.00%
1 / 1
100.00% covered (success)
100.00%
1 / 1
1
 setRegistry
100.00% covered (success)
100.00%
3 / 3
100.00% covered (success)
100.00%
1 / 1
1
 getRegistry
100.00% covered (success)
100.00%
1 / 1
100.00% covered (success)
100.00%
1 / 1
1
 registry
100.00% covered (success)
100.00%
1 / 1
100.00% covered (success)
100.00%
1 / 1
1
 hasRegistry
100.00% covered (success)
100.00%
1 / 1
100.00% covered (success)
100.00%
1 / 1
1
 setName
100.00% covered (success)
100.00%
2 / 2
100.00% covered (success)
100.00%
1 / 1
1
 getName
100.00% covered (success)
100.00%
1 / 1
100.00% covered (success)
100.00%
1 / 1
1
 hasName
100.00% covered (success)
100.00%
1 / 1
100.00% covered (success)
100.00%
1 / 1
1
 triggerEvent
100.00% covered (success)
100.00%
4 / 4
100.00% covered (success)
100.00%
1 / 1
4
 stop
100.00% covered (success)
100.00%
2 / 2
100.00% covered (success)
100.00%
1 / 1
1
 isStopped
100.00% covered (success)
100.00%
1 / 1
100.00% covered (success)
100.00%
1 / 1
1
 addQueue
100.00% covered (success)
100.00%
5 / 5
100.00% covered (success)
100.00%
1 / 1
2
 addQueues
100.00% covered (success)
100.00%
3 / 3
100.00% covered (success)
100.00%
1 / 1
2
 getQueues
100.00% covered (success)
100.00%
1 / 1
100.00% covered (success)
100.00%
1 / 1
1
 getQueuesByWeight
100.00% covered (success)
100.00%
5 / 5
100.00% covered (success)
100.00%
1 / 1
1
 getQueue
100.00% covered (success)
100.00%
1 / 1
100.00% covered (success)
100.00%
1 / 1
1
 hasQueue
100.00% covered (success)
100.00%
1 / 1
100.00% covered (success)
100.00%
1 / 1
1
 getWeight
100.00% covered (success)
100.00%
1 / 1
100.00% covered (success)
100.00%
1 / 1
1
 ensureRegistered
100.00% covered (success)
100.00%
6 / 6
100.00% covered (success)
100.00%
1 / 1
4
 heartbeat
75.00% covered (success)
75.00%
3 / 4
0.00% covered (danger)
0.00%
0 / 1
3.14
 deregisterIfOwned
75.00% covered (success)
75.00%
3 / 4
0.00% covered (danger)
0.00%
0 / 1
4.25
 work
100.00% covered (success)
100.00%
9 / 9
100.00% covered (success)
100.00%
1 / 1
5
 workAll
100.00% covered (success)
100.00%
6 / 6
100.00% covered (success)
100.00%
1 / 1
2
 run
100.00% covered (success)
100.00%
6 / 6
100.00% covered (success)
100.00%
1 / 1
2
 runAll
100.00% covered (success)
100.00%
25 / 25
100.00% covered (success)
100.00%
1 / 1
11
 installSignalHandlers
66.67% covered (warning)
66.67%
6 / 9
0.00% covered (danger)
0.00%
0 / 1
2.15
 workLoop
100.00% covered (success)
100.00%
20 / 20
100.00% covered (success)
100.00%
1 / 1
6
 runLoop
100.00% covered (success)
100.00%
20 / 20
100.00% covered (success)
100.00%
1 / 1
6
 clear
100.00% covered (success)
100.00%
3 / 3
100.00% covered (success)
100.00%
1 / 1
2
 clearFailed
100.00% covered (success)
100.00%
3 / 3
100.00% covered (success)
100.00%
1 / 1
2
 clearTasks
100.00% covered (success)
100.00%
3 / 3
100.00% covered (success)
100.00%
1 / 1
2
 clearAll
100.00% covered (success)
100.00%
3 / 3
100.00% covered (success)
100.00%
1 / 1
2
 clearAllFailed
100.00% covered (success)
100.00%
3 / 3
100.00% covered (success)
100.00%
1 / 1
2
 clearAllTasks
100.00% covered (success)
100.00%
3 / 3
100.00% covered (success)
100.00%
1 / 1
2
 __set
100.00% covered (success)
100.00%
1 / 1
100.00% covered (success)
100.00%
1 / 1
1
 __get
100.00% covered (success)
100.00%
1 / 1
100.00% covered (success)
100.00%
1 / 1
1
 __isset
100.00% covered (success)
100.00%
1 / 1
100.00% covered (success)
100.00%
1 / 1
1
 __unset
100.00% covered (success)
100.00%
2 / 2
100.00% covered (success)
100.00%
1 / 1
2
 offsetSet
100.00% covered (success)
100.00%
1 / 1
100.00% covered (success)
100.00%
1 / 1
1
 offsetGet
100.00% covered (success)
100.00%
1 / 1
100.00% covered (success)
100.00%
1 / 1
1
 offsetExists
100.00% covered (success)
100.00%
1 / 1
100.00% covered (success)
100.00%
1 / 1
1
 offsetUnset
100.00% covered (success)
100.00%
1 / 1
100.00% covered (success)
100.00%
1 / 1
1
 count
100.00% covered (success)
100.00%
1 / 1
100.00% covered (success)
100.00%
1 / 1
1
 getIterator
100.00% covered (success)
100.00%
1 / 1
100.00% covered (success)
100.00%
1 / 1
1
1<?php
2declare(strict_types=1);
3/**
4 * Pop PHP Framework (https://www.popphp.org/)
5 *
6 * @link       https://github.com/popphp/popphp-framework
7 * @author     Nick Sagona, III <nick@popphp.org>
8 * @copyright  Copyright (c) 2009-2026 Nick Sagona, III
9 * @license    https://www.popphp.org/license     New BSD License
10 */
11
12/**
13 * @namespace
14 */
15namespace Pop\Queue;
16
17use ArrayIterator;
18use Pop\Application;
19use Pop\Event\Manager as EventManager;
20use Pop\Queue\Process\AbstractJob;
21use Pop\Queue\Registry\WorkerRecord;
22use Pop\Queue\Registry\WorkerRegistry;
23
24/**
25 * Queue worker class
26 *
27 * @category   Pop
28 * @package    Pop\Queue
29 * @author     Nick Sagona, III <nick@popphp.org>
30 * @copyright  Copyright (c) 2009-2026 Nick Sagona, III
31 * @license    https://www.popphp.org/license     New BSD License
32 * @version    3.0.0
33 */
34class Worker implements \ArrayAccess, \Countable, \IteratorAggregate
35{
36
37    /**
38     * Queues
39     * @var array
40     */
41    protected array $queues = [];
42
43    /**
44     * Queue weights, keyed by queue name. Higher services first. Not named
45     * "priority" - that word is already used elsewhere in this codebase for
46     * the unrelated FIFO/FILO adapter job-ordering setting.
47     * @var array
48     */
49    protected array $weights = [];
50
51    /**
52     * Application object
53     * @var ?Application
54     */
55    protected ?Application $application = null;
56
57    /**
58     * Event manager, for worker-level lifecycle observability hooks
59     * (worker.work_loop.*, worker.run_loop.*). If not set, and this
60     * Worker was constructed with an Application that has its own event
61     * manager, that Application's manager is used instead - see
62     * triggerEvent(). If neither is available, event firing is a silent
63     * no-op.
64     * @var ?EventManager
65     */
66    protected ?EventManager $events = null;
67
68    /**
69     * Whether a graceful shutdown has been requested, via stop() directly
70     * or via a caught SIGTERM/SIGINT (see installSignalHandlers()).
71     * workLoop()/runLoop() each reset this to false at their own start.
72     * @var bool
73     */
74    protected bool $stopped = false;
75
76    /**
77     * Optional operator-facing label for this worker, surfaced in its
78     * registry record
79     * @var ?string
80     */
81    protected ?string $name = null;
82
83    /**
84     * Worker registry, for observability. Null (the default) means no
85     * registry work happens at all and behavior is identical to a build
86     * without this feature.
87     * @var ?WorkerRegistry
88     */
89    protected ?WorkerRegistry $registry = null;
90
91    /**
92     * Constructor
93     *
94     * Instantiate the queue worker object.
95     *
96     * @param mixed $queues
97     * @param ?Application $application
98     */
99    public function __construct(mixed $queues = null, ?Application $application = null)
100    {
101        if (!empty($queues)) {
102            if (is_array($queues)) {
103                $this->addQueues($queues);
104            } else if ($queues instanceof Queue) {
105                $this->addQueue($queues);
106            }
107        }
108
109        $this->application = $application;
110    }
111
112    /**
113     * Create queue worker worker
114     *
115     * @param  mixed $queues
116     * @param  ?Application $application
117     * @return Worker
118     */
119    public static function create(mixed $queues = null, ?Application $application = null): Worker
120    {
121        return new self($queues, $application);
122    }
123
124    /**
125     * Get the application
126     *
127     * @return ?Application
128     */
129    public function getApplication(): ?Application
130    {
131        return $this->application;
132    }
133
134    /**
135     * Get the application (alias)
136     *
137     * @return ?Application
138     */
139    public function application(): ?Application
140    {
141        return $this->application;
142    }
143
144    /**
145     * Has application
146     *
147     * @return bool
148     */
149    public function hasApplication(): bool
150    {
151        return ($this->application !== null);
152    }
153
154    /**
155     * Set event manager
156     *
157     * @param  EventManager $events
158     * @return Worker
159     */
160    public function setEvents(EventManager $events): Worker
161    {
162        $this->events = $events;
163        return $this;
164    }
165
166    /**
167     * Get event manager
168     *
169     * @return ?EventManager
170     */
171    public function getEvents(): ?EventManager
172    {
173        return $this->events;
174    }
175
176    /**
177     * Get event manager (alias)
178     *
179     * @return ?EventManager
180     */
181    public function events(): ?EventManager
182    {
183        return $this->events;
184    }
185
186    /**
187     * Has event manager
188     *
189     * @return bool
190     */
191    public function hasEvents(): bool
192    {
193        return ($this->events !== null);
194    }
195
196    /**
197     * Set the worker registry, and wire its tracking onto this worker's
198     * queues
199     *
200     * @param  WorkerRegistry $registry
201     * @return Worker
202     */
203    public function setRegistry(WorkerRegistry $registry): Worker
204    {
205        $this->registry = $registry;
206        $registry->attachTo($this);
207
208        return $this;
209    }
210
211    /**
212     * Get the worker registry
213     *
214     * @return ?WorkerRegistry
215     */
216    public function getRegistry(): ?WorkerRegistry
217    {
218        return $this->registry;
219    }
220
221    /**
222     * Get the worker registry (alias)
223     *
224     * @return ?WorkerRegistry
225     */
226    public function registry(): ?WorkerRegistry
227    {
228        return $this->registry;
229    }
230
231    /**
232     * Has a worker registry
233     *
234     * @return bool
235     */
236    public function hasRegistry(): bool
237    {
238        return ($this->registry !== null);
239    }
240
241    /**
242     * Set this worker's operator-facing label
243     *
244     * @param  ?string $name
245     * @return Worker
246     */
247    public function setName(?string $name): Worker
248    {
249        $this->name = $name;
250        return $this;
251    }
252
253    /**
254     * Get this worker's operator-facing label
255     *
256     * @return ?string
257     */
258    public function getName(): ?string
259    {
260        return $this->name;
261    }
262
263    /**
264     * Has an operator-facing label
265     *
266     * @return bool
267     */
268    public function hasName(): bool
269    {
270        return ($this->name !== null);
271    }
272
273    /**
274     * Trigger a worker-level lifecycle event. Uses this Worker's own event
275     * manager if one is set via setEvents(); otherwise falls back to the
276     * event manager of the Application this Worker was constructed with,
277     * if it has one; otherwise does nothing. Never throws on its own
278     * account - if the resolved manager's trigger() call throws (e.g. a
279     * listener's own code throws), that exception propagates to the
280     * caller exactly as any other uncaught exception would.
281     *
282     * @param  string $name
283     * @param  array  $params
284     * @return void
285     */
286    protected function triggerEvent(string $name, array $params): void
287    {
288        if ($this->hasEvents()) {
289            $this->events->trigger($name, $params);
290        } else if (($this->application !== null) && ($this->application->events() !== null)) {
291            $this->application->events()->trigger($name, $params);
292        }
293    }
294
295    /**
296     * Request a graceful shutdown of a running workLoop()/runLoop() call.
297     * Takes effect at that loop's next iteration boundary - never mid-job
298     * or mid-task-evaluation.
299     *
300     * @return Worker
301     */
302    public function stop(): Worker
303    {
304        $this->stopped = true;
305        return $this;
306    }
307
308    /**
309     * Whether a graceful shutdown has been requested
310     *
311     * @return bool
312     */
313    public function isStopped(): bool
314    {
315        return $this->stopped;
316    }
317
318    /**
319     * Add queue
320     *
321     * @param  Queue $queue
322     * @param  int   $weight
323     * @return Worker
324     */
325    public function addQueue(Queue $queue, int $weight = 0): Worker
326    {
327        $this->queues[$queue->getName()]  = $queue;
328        $this->weights[$queue->getName()] = $weight;
329
330        // A queue added after setRegistry() would otherwise never be wired.
331        // attachTo() is idempotent per event manager, so re-attaching costs
332        // nothing for queues already covered.
333        if ($this->hasRegistry()) {
334            $this->registry->attachTo($this);
335        }
336
337        return $this;
338    }
339
340    /**
341     * Add queues
342     *
343     * @param  array $queues
344     * @return Worker
345     */
346    public function addQueues(array $queues): Worker
347    {
348        foreach ($queues as $queue) {
349            $this->addQueue($queue);
350        }
351        return $this;
352    }
353
354    /**
355     * Get queues
356     *
357     * @return array
358     */
359    public function getQueues(): array
360    {
361        return $this->getQueuesByWeight();
362    }
363
364    /**
365     * Get queues ordered by weight, highest first. PHP's sort functions
366     * are stable since 8.0, so queues with equal weight (including the
367     * default-zero case when no weight was ever set) keep their original
368     * insertion order automatically.
369     *
370     * @return array
371     */
372    protected function getQueuesByWeight(): array
373    {
374        $queues = $this->queues;
375        uksort($queues, function($a, $b) {
376            return ($this->weights[$b] ?? 0) <=> ($this->weights[$a] ?? 0);
377        });
378        return $queues;
379    }
380
381    /**
382     * Get queue
383     *
384     * @param  string $queue
385     * @return ?Queue
386     */
387    public function getQueue(string $queue): ?Queue
388    {
389        return $this->queues[$queue] ?? null;
390    }
391
392    /**
393     * Has queue
394     *
395     * @param  string $queue
396     * @return bool
397     */
398    public function hasQueue(string $queue): bool
399    {
400        return (isset($this->queues[$queue]));
401    }
402
403    /**
404     * Get a queue's weight (0 if never set)
405     *
406     * @param  string $queueName
407     * @return int
408     */
409    public function getWeight(string $queueName): int
410    {
411        return $this->weights[$queueName] ?? 0;
412    }
413
414    /**
415     * Register this process with the registry if nothing has already,
416     * returning whether THIS call performed the registration.
417     *
418     * That return value is the ownership rule that lets one mechanism serve
419     * both deployment models: a single-pass work()/run() that registered
420     * deregisters itself on the way out, while the same call made from
421     * inside a daemon loop finds the loop's registration already in place,
422     * takes no ownership, and leaves the daemon's record alone.
423     *
424     * @param  string $mode
425     * @return bool
426     */
427    protected function ensureRegistered(string $mode): bool
428    {
429        if (!$this->hasRegistry() || $this->registry->isRegistered()) {
430            return false;
431        }
432
433        try {
434            $this->registry->register($this->name, array_keys($this->queues), $mode);
435            return true;
436        } catch (\Throwable $e) {
437            // Observability must never stop the queue. Returning false is
438            // load-bearing: nothing was written, so the finally must not
439            // later try to delete a record that does not exist.
440            return false;
441        }
442    }
443
444    /**
445     * Refresh this worker's heartbeat, if it has a registry. A no-op
446     * otherwise, and a no-op inside the registry when not registered.
447     *
448     * @return void
449     */
450    protected function heartbeat(): void
451    {
452        if (!$this->hasRegistry()) {
453            return;
454        }
455
456        try {
457            $this->registry->heartbeat();
458        } catch (\Throwable $e) {
459            // Best-effort: a missed heartbeat is a stale-looking worker, not
460            // a stopped queue.
461        }
462    }
463
464    /**
465     * Deregister this process, but only if the given flag says this call
466     * owned the registration
467     *
468     * @param  bool $owned
469     * @return void
470     */
471    protected function deregisterIfOwned(bool $owned): void
472    {
473        if (!$owned || !$this->hasRegistry()) {
474            return;
475        }
476
477        try {
478            $this->registry->deregister();
479        } catch (\Throwable $e) {
480            // Swallowed deliberately: this runs in a finally, so throwing
481            // here would discard the return value of work that actually
482            // succeeded. A stranded record is reaped by prune().
483        }
484    }
485
486    /**
487     * Work next job. Pass a queue name to work that specific queue (exactly
488     * today's behavior). Pass nothing to try every registered queue in
489     * weight order (highest first), returning the first job successfully
490     * claimed - the highest-weight-first worker model, since workAll()
491     * fans out to every queue regardless of weight and doesn't need this.
492     *
493     * @param  ?string $queueName
494     * @return ?AbstractJob
495     */
496    public function work(?string $queueName = null): ?AbstractJob
497    {
498        $owned = $this->ensureRegistered(WorkerRecord::MODE_SINGLE_PASS);
499
500        try {
501            if ($queueName !== null) {
502                return isset($this->queues[$queueName]) ? $this->queues[$queueName]->work($this->application) : null;
503            }
504
505            foreach ($this->getQueuesByWeight() as $queue) {
506                $job = $queue->work($this->application);
507                if ($job !== null) {
508                    return $job;
509                }
510            }
511
512            return null;
513        } finally {
514            $this->deregisterIfOwned($owned);
515        }
516    }
517
518    /**
519     * Work next job across in all queues
520     *
521     * @return array
522     */
523    public function workAll(): array
524    {
525        $owned = $this->ensureRegistered(WorkerRecord::MODE_SINGLE_PASS);
526
527        try {
528            $jobs = [];
529            foreach ($this->getQueuesByWeight() as $queueName => $queue) {
530                $jobs[$queueName] = $queue->work($this->application);
531            }
532            return $jobs;
533        } finally {
534            $this->deregisterIfOwned($owned);
535        }
536    }
537
538    /**
539     * Run next scheduled task in queue
540     *
541     * @param  string $queueName
542     * @return array
543     */
544    public function run(string $queueName): array
545    {
546        $owned = $this->ensureRegistered(WorkerRecord::MODE_SINGLE_PASS);
547
548        try {
549            $tasks = [];
550            if (isset($this->queues[$queueName])) {
551                $tasks[$queueName] = $this->queues[$queueName]->run($this->application);
552            }
553            return $tasks;
554        } finally {
555            $this->deregisterIfOwned($owned);
556        }
557    }
558
559    /**
560     * Run next scheduled task across all queues, fairly: every queue gets
561     * one shared evaluation pass immediately, then - only if at least one
562     * queue has a sub-minute task - up to 59 more shared passes (one per
563     * second, sleeping once per tick, not once per queue per tick),
564     * mirroring Queue::run()'s own pass-1-then-tick-loop shape one level
565     * up. No single queue's tick-loop work can block another queue's
566     * evaluation on the same tick.
567     *
568     * @return array
569     */
570    public function runAll(): array
571    {
572        $owned = $this->ensureRegistered(WorkerRecord::MODE_SINGLE_PASS);
573
574        try {
575            $tasks         = [];
576            $queueTaskSets = [];
577            $queues        = $this->getQueuesByWeight();
578
579            foreach ($queues as $queueName => $queue) {
580                $tasks[$queueName]         = [];
581                $queueTaskSets[$queueName] = $queue->getScheduledTasks();
582            }
583
584            foreach ($queueTaskSets as $queueName => $scheduledTasks) {
585                foreach ($queues[$queueName]->evaluateTasksOnce($scheduledTasks, $this->application) as $jobId => $task) {
586                    $tasks[$queueName][$jobId] = $task;
587                }
588            }
589
590            $hasSubMinute = false;
591            foreach ($queueTaskSets as $scheduledTasks) {
592                foreach ($scheduledTasks as $task) {
593                    if ($task->cron()->hasSeconds()) {
594                        $hasSubMinute = true;
595                        break 2;
596                    }
597                }
598            }
599
600            if ($hasSubMinute) {
601                for ($tick = 1; $tick < 60; $tick++) {
602                    sleep(1);
603                    $this->heartbeat();
604                    foreach ($queueTaskSets as $queueName => $scheduledTasks) {
605                        foreach ($queues[$queueName]->evaluateTasksOnce($scheduledTasks, $this->application, true) as $jobId => $task) {
606                            $tasks[$queueName][$jobId] = $task;
607                        }
608                    }
609                }
610            }
611
612            return $tasks;
613        } finally {
614            $this->deregisterIfOwned($owned);
615        }
616    }
617
618    /**
619     * Install SIGTERM/SIGINT handlers that request a graceful stop(),
620     * rather than letting the OS terminate the process immediately. A
621     * no-op when ext-pcntl isn't loaded - workLoop()/runLoop() still run
622     * correctly without it, just without OS-signal-based shutdown
623     * available (only stop() can end them in that case).
624     *
625     * Unlike Queue::runWithTimeout()'s SIGALRM handler next door, these
626     * handlers are never restored to SIG_DFL once installed - not even
627     * after workLoop()/runLoop() returns. This is a deliberate, accepted
628     * tradeoff for this pass rather than an oversight: a second
629     * SIGTERM/SIGINT sent after a loop has already ended is simply inert
630     * (caught by this handler and ignored) instead of terminating the
631     * process via the OS default.
632     *
633     * @return void
634     */
635    protected function installSignalHandlers(): void
636    {
637        if (!extension_loaded('pcntl')) {
638            return;
639        }
640
641        pcntl_async_signals(true);
642        pcntl_signal(SIGTERM, function() {
643            $this->stop();
644        });
645        pcntl_signal(SIGINT, function() {
646            $this->stop();
647        });
648    }
649
650    /**
651     * Work jobs across all registered queues, forever, until stopped.
652     * Calls workAll() every iteration; sleeps $sleepSeconds only when a
653     * full pass finds nothing anywhere (every queue returned null),
654     * looping again immediately otherwise. Stoppable via stop() directly
655     * or, when ext-pcntl is loaded, via SIGTERM/SIGINT - either way, the
656     * current iteration's job is never torn down mid-execution and its
657     * remaining code always runs to completion - though a blocking call
658     * inside that code (e.g. sleep()) can itself be interrupted early if
659     * a signal lands during it; see README's Daemon mode section.
660     * A negative $sleepSeconds is silently clamped to 0.
661     *
662     * @param  int $sleepSeconds
663     * @return void
664     */
665    public function workLoop(int $sleepSeconds = 1): void
666    {
667        $sleepSeconds  = max(0, $sleepSeconds);
668        $this->stopped = false;
669        $this->installSignalHandlers();
670
671        $owned = $this->ensureRegistered(WorkerRecord::MODE_DAEMON);
672
673        try {
674            while (!$this->stopped) {
675                $jobs = $this->workAll();
676
677                $anyWorked = false;
678                foreach ($jobs as $job) {
679                    if ($job !== null) {
680                        $anyWorked = true;
681                        break;
682                    }
683                }
684
685                $this->heartbeat();
686                $this->triggerEvent('worker.work_loop.tick', ['jobs' => $jobs, 'worker' => $this]);
687
688                if ($this->stopped) {
689                    break;
690                }
691
692                if (!$anyWorked) {
693                    $this->triggerEvent('worker.work_loop.idle', ['worker' => $this]);
694                    sleep($sleepSeconds);
695                }
696            }
697
698            $this->triggerEvent('worker.work_loop.shutdown', ['worker' => $this]);
699        } finally {
700            $this->deregisterIfOwned($owned);
701        }
702    }
703
704    /**
705     * Run scheduled tasks across all registered queues, forever, until
706     * stopped. Calls runAll() every iteration; sleeps $sleepSeconds only
707     * when a full pass finds nothing due anywhere, looping again
708     * immediately otherwise. runAll() already blocks appropriately on its
709     * own whenever a sub-minute task exists, so this backoff sleep only
710     * ever triggers on the coarse-only-or-nothing-scheduled case, which
711     * returns near-instantly and would otherwise busy-loop. A stop signal
712     * arriving mid-runAll() (during its own internal sub-minute tick loop)
713     * isn't noticed until that call returns - see the design spec's
714     * Non-goals for why this latency is accepted rather than fixed here.
715     * A negative $sleepSeconds is silently clamped to 0.
716     *
717     * @param  int $sleepSeconds
718     * @return void
719     */
720    public function runLoop(int $sleepSeconds = 1): void
721    {
722        $sleepSeconds  = max(0, $sleepSeconds);
723        $this->stopped = false;
724        $this->installSignalHandlers();
725
726        $owned = $this->ensureRegistered(WorkerRecord::MODE_DAEMON);
727
728        try {
729            while (!$this->stopped) {
730                $tasks = $this->runAll();
731
732                $anyRan = false;
733                foreach ($tasks as $queueTasks) {
734                    if (!empty($queueTasks)) {
735                        $anyRan = true;
736                        break;
737                    }
738                }
739
740                $this->heartbeat();
741                $this->triggerEvent('worker.run_loop.tick', ['tasks' => $tasks, 'worker' => $this]);
742
743                if ($this->stopped) {
744                    break;
745                }
746
747                if (!$anyRan) {
748                    $this->triggerEvent('worker.run_loop.idle', ['worker' => $this]);
749                    sleep($sleepSeconds);
750                }
751            }
752
753            $this->triggerEvent('worker.run_loop.shutdown', ['worker' => $this]);
754        } finally {
755            $this->deregisterIfOwned($owned);
756        }
757    }
758
759    /**
760     * Clear jobs from queue
761     *
762     * @param  string $queueName
763     * @return Worker
764     */
765    public function clear(string $queueName): Worker
766    {
767        if (isset($this->queues[$queueName])) {
768            $this->queues[$queueName]->clear();
769        }
770        return $this;
771    }
772
773    /**
774     * Clear failed jobs from queue
775     *
776     * @param  string $queueName
777     * @return Worker
778     */
779    public function clearFailed(string $queueName): Worker
780    {
781        if (isset($this->queues[$queueName])) {
782            $this->queues[$queueName]->clearFailed();
783        }
784        return $this;
785    }
786
787    /**
788     * Clear tasks from queue
789     *
790     * @param  string $queueName
791     * @return Worker
792     */
793    public function clearTasks(string $queueName): Worker
794    {
795        if (isset($this->queues[$queueName])) {
796            $this->queues[$queueName]->clearTasks();
797        }
798        return $this;
799    }
800
801    /**
802     * Clear all jobs from queues
803     *
804     * @return Worker
805     */
806    public function clearAll(): Worker
807    {
808        foreach ($this->getQueuesByWeight() as $queue) {
809            $queue->clear();
810        }
811        return $this;
812    }
813
814    /**
815     * Clear all failed jobs from queues
816     *
817     * @return Worker
818     */
819    public function clearAllFailed(): Worker
820    {
821        foreach ($this->getQueuesByWeight() as $queue) {
822            $queue->clearFailed();
823        }
824        return $this;
825    }
826
827    /**
828     * Clear all tasks from queues
829     *
830     * @return Worker
831     */
832    public function clearAllTasks(): Worker
833    {
834        foreach ($this->getQueuesByWeight() as $queue) {
835            $queue->clearTasks();
836        }
837        return $this;
838    }
839
840    /**
841     * Register a queue with the worker
842     *
843     * @param  string $name
844     * @param  mixed $value
845     * @return void
846     */
847    public function __set(string $name, mixed $value): void
848    {
849        $this->addQueue($value);
850    }
851
852    /**
853     * Get a queue
854     *
855     * @param  string $name
856     * @return ?Queue
857     */
858    public function __get(string $name): ?Queue
859    {
860        return $this->getQueue($name);
861    }
862
863    /**
864     * Determine if a queue is registered with the worker object
865     *
866     * @param  string $name
867     * @return bool
868     */
869    public function __isset(string $name): bool
870    {
871        return isset($this->queues[$name]);
872    }
873
874    /**
875     * Unset a queue with the worker
876     *
877     * @param  string $name
878     * @return void
879     */
880    public function __unset(string $name): void
881    {
882        if (isset($this->queues[$name])) {
883            unset($this->queues[$name], $this->weights[$name]);
884        }
885    }
886
887    /**
888     * Set a queue with the worker
889     *
890     * @param  mixed $offset
891     * @param  mixed $value
892     * @return void
893     */
894    public function offsetSet(mixed $offset, mixed $value): void
895    {
896        $this->__set($offset, $value);
897    }
898
899    /**
900     * Get a queue
901     *
902     * @param  mixed $offset
903     * @return ?Queue
904     */
905    public function offsetGet(mixed $offset): ?Queue
906    {
907        return $this->__get($offset);
908    }
909
910    /**
911     * Determine if a queue is registered with the worker object
912     *
913     * @param  mixed $offset
914     * @return bool
915     */
916    public function offsetExists(mixed $offset): bool
917    {
918        return $this->__isset($offset);
919    }
920
921    /**
922     * Unset a queue from the worker
923     *
924     * @param  string $offset
925     * @return void
926     */
927    public function offsetUnset(mixed $offset): void
928    {
929        $this->__unset($offset);
930    }
931
932    /**
933     * Return count
934     *
935     * @return int
936     */
937    public function count(): int
938    {
939        return count($this->queues);
940    }
941
942    /**
943     * Get iterator
944     *
945     * @return ArrayIterator
946     */
947    public function getIterator(): ArrayIterator
948    {
949        return new ArrayIterator($this->getQueuesByWeight());
950    }
951
952}