Code Coverage
 
Lines
Functions and Methods
Classes and Traits
Total
90.16% covered (success)
90.16%
220 / 244
72.09% covered (success)
72.09%
31 / 43
CRAP
0.00% covered (danger)
0.00%
0 / 1
File
90.16% covered (success)
90.16%
220 / 244
72.09% covered (success)
72.09%
31 / 43
132.48
0.00% covered (danger)
0.00%
0 / 1
 __construct
72.73% covered (success)
72.73%
8 / 11
0.00% covered (danger)
0.00%
0 / 1
5.51
 create
100.00% covered (success)
100.00%
1 / 1
100.00% covered (success)
100.00%
1 / 1
1
 getFolder
100.00% covered (success)
100.00%
1 / 1
100.00% covered (success)
100.00%
1 / 1
1
 folder
100.00% covered (success)
100.00%
1 / 1
100.00% covered (success)
100.00%
1 / 1
1
 pendingPath
100.00% covered (success)
100.00%
1 / 1
100.00% covered (success)
100.00%
1 / 1
1
 reservedPath
100.00% covered (success)
100.00%
1 / 1
100.00% covered (success)
100.00%
1 / 1
1
 getEndIndex
100.00% covered (success)
100.00%
7 / 7
100.00% covered (success)
100.00%
1 / 1
3
 readVerifiedPayload
100.00% covered (success)
100.00%
2 / 2
100.00% covered (success)
100.00%
1 / 1
2
 getLeaseUntil
100.00% covered (success)
100.00%
2 / 2
100.00% covered (success)
100.00%
1 / 1
2
 isLeaseExpired
80.00% covered (success)
80.00%
4 / 5
0.00% covered (danger)
0.00%
0 / 1
4.13
 reclaimExpiredLeases
66.67% covered (warning)
66.67%
10 / 15
0.00% covered (danger)
0.00%
0 / 1
7.33
 findReservedIndexForJob
82.35% covered (success)
82.35%
14 / 17
0.00% covered (danger)
0.00%
0 / 1
9.45
 claimPendingDir
100.00% covered (success)
100.00%
1 / 1
100.00% covered (success)
100.00%
1 / 1
1
 push
94.74% covered (success)
94.74%
18 / 19
0.00% covered (danger)
0.00%
0 / 1
6.01
 reserve
92.31% covered (success)
92.31%
24 / 26
0.00% covered (danger)
0.00%
0 / 1
9.04
 release
88.89% covered (success)
88.89%
8 / 9
0.00% covered (danger)
0.00%
0 / 1
2.01
 delete
80.00% covered (success)
80.00%
4 / 5
0.00% covered (danger)
0.00%
0 / 1
2.03
 removeJobDir
83.33% covered (success)
83.33%
5 / 6
0.00% covered (danger)
0.00%
0 / 1
4.07
 bury
100.00% covered (success)
100.00%
3 / 3
100.00% covered (success)
100.00%
1 / 1
1
 hasJobs
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
 clear
100.00% covered (success)
100.00%
5 / 5
100.00% covered (success)
100.00%
1 / 1
3
 getDeadJobIds
100.00% covered (success)
100.00%
5 / 5
100.00% covered (success)
100.00%
1 / 1
3
 hasDeadJobs
100.00% covered (success)
100.00%
1 / 1
100.00% covered (success)
100.00%
1 / 1
1
 countDead
100.00% covered (success)
100.00%
1 / 1
100.00% covered (success)
100.00%
1 / 1
1
 getDeadJobs
100.00% covered (success)
100.00%
4 / 4
100.00% covered (success)
100.00%
1 / 1
2
 getDeadJob
80.00% covered (success)
80.00%
8 / 10
0.00% covered (danger)
0.00%
0 / 1
5.20
 retryDeadJob
100.00% covered (success)
100.00%
5 / 5
100.00% covered (success)
100.00%
1 / 1
2
 deleteDeadJob
100.00% covered (success)
100.00%
4 / 4
100.00% covered (success)
100.00%
1 / 1
2
 clearDead
100.00% covered (success)
100.00%
3 / 3
100.00% covered (success)
100.00%
1 / 1
2
 schedule
100.00% covered (success)
100.00%
5 / 5
100.00% covered (success)
100.00%
1 / 1
2
 getTasks
100.00% covered (success)
100.00%
6 / 6
100.00% covered (success)
100.00%
1 / 1
3
 getTask
83.33% covered (success)
83.33%
5 / 6
0.00% covered (danger)
0.00%
0 / 1
4.07
 updateTask
100.00% covered (success)
100.00%
6 / 6
100.00% covered (success)
100.00%
1 / 1
2
 removeTask
100.00% covered (success)
100.00%
5 / 5
100.00% covered (success)
100.00%
1 / 1
3
 getTaskCount
100.00% covered (success)
100.00%
1 / 1
100.00% covered (success)
100.00%
1 / 1
1
 hasTasks
100.00% covered (success)
100.00%
1 / 1
100.00% covered (success)
100.00%
1 / 1
1
 clearTasks
100.00% covered (success)
100.00%
7 / 7
100.00% covered (success)
100.00%
1 / 1
4
 taskClaimPath
100.00% covered (success)
100.00%
1 / 1
100.00% covered (success)
100.00%
1 / 1
1
 claimTaskRun
86.36% covered (success)
86.36%
19 / 22
0.00% covered (danger)
0.00%
0 / 1
7.12
 getFolders
100.00% covered (success)
100.00%
1 / 1
100.00% covered (success)
100.00%
1 / 1
1
 getFiles
100.00% covered (success)
100.00%
1 / 1
100.00% covered (success)
100.00%
1 / 1
1
 readDirectory
100.00% covered (success)
100.00%
9 / 9
100.00% covered (success)
100.00%
1 / 1
5
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\Adapter;
16
17use Pop\Queue\Process\AbstractJob;
18use Pop\Queue\Process\PayloadSigner;
19use Pop\Queue\Process\Task;
20
21/**
22 * File adapter class
23 *
24 * @category   Pop
25 * @package    Pop\Queue
26 * @author     Nick Sagona, III <nick@popphp.org>
27 * @copyright  Copyright (c) 2009-2026 Nick Sagona, III
28 * @license    https://www.popphp.org/license     New BSD License
29 * @version    3.0.0
30 */
31class File extends AbstractTaskAdapter
32{
33
34    /**
35     * Folder
36     * @var ?string
37     */
38    protected ?string $folder = null;
39
40    /**
41     * Reservation lease length, in seconds
42     * @var int
43     */
44    protected int $leaseSeconds = 60;
45
46    /**
47     * Highest job index this adapter instance knows to be taken, or null when
48     * it hasn't looked yet. Purely a starting hint for push() - never a source
49     * of truth. See getEndIndex() for why being wrong in either direction is
50     * safe.
51     * @var ?int
52     */
53    protected ?int $endIndex = null;
54
55    /**
56     * Constructor
57     *
58     * Instantiate the file object
59     *
60     * @param  string  $folder
61     * @param  ?string $priority
62     * @param  int     $leaseSeconds
63     * @throws Exception
64     */
65    public function __construct(string $folder, ?string $priority = null, int $leaseSeconds = 60)
66    {
67        if (!file_exists($folder)) {
68            throw new Exception("Error: The folder '" . $folder . "' does not exist.");
69        }
70        if (!is_writable($folder)) {
71            throw new Exception("Error: The folder '" . $folder . "' is not writable.");
72        }
73
74        $this->folder       = $folder;
75        $this->leaseSeconds = $leaseSeconds;
76
77        if (!file_exists($this->pendingPath())) {
78            mkdir($this->pendingPath());
79        }
80        if (!file_exists($this->reservedPath())) {
81            mkdir($this->reservedPath());
82        }
83
84        parent::__construct($priority);
85    }
86
87    /**
88     * Create file adapter
89     *
90     * @param  string  $folder
91     * @param  ?string $priority
92     * @param  int     $leaseSeconds
93     * @throws Exception
94     * @return File
95     */
96    public static function create(string $folder, ?string $priority = null, int $leaseSeconds = 60): File
97    {
98        return new self($folder, $priority, $leaseSeconds);
99    }
100
101    /**
102     * Get folder
103     *
104     * @return ?string
105     */
106    public function getFolder(): ?string
107    {
108        return $this->folder;
109    }
110
111    /**
112     * Get folder (alias)
113     *
114     * @return ?string
115     */
116    public function folder(): ?string
117    {
118        return $this->folder;
119    }
120
121    /**
122     * Get the pending-jobs subdirectory path
123     *
124     * @return string
125     */
126    protected function pendingPath(): string
127    {
128        return $this->folder . DIRECTORY_SEPARATOR . 'pending';
129    }
130
131    /**
132     * Get the reserved-jobs subdirectory path
133     *
134     * @return string
135     */
136    protected function reservedPath(): string
137    {
138        return $this->folder . DIRECTORY_SEPARATOR . 'reserved';
139    }
140
141    /**
142     * Get queue end index across both pending and reserved jobs (indices must
143     * stay unique across both so a reclaimed reserved job can never collide
144     * with a newly-pushed one)
145     *
146     * Scans the directories once per adapter instance and then tracks the
147     * index forward in memory, because the scan is what made push() quadratic:
148     * it walked both directories on every single push, so the cost of pushing
149     * the Nth job grew with the number of jobs already queued.
150     *
151     * Caching it is safe precisely because push() never trusted this value in
152     * the first place - mkdir() is the real allocator, and push()'s retry loop
153     * already handles the index being taken. A cached value that is too low
154     * (another process pushed since the scan) costs one wasted mkdir() attempt
155     * per collision and the loop walks up; a cached value that is too high
156     * (jobs were cleared elsewhere) just leaves a gap in the numbering, which
157     * nothing depends on - indices only ever need to be unique and ordered,
158     * never contiguous.
159     *
160     * @return int
161     */
162    protected function getEndIndex(): int
163    {
164        if ($this->endIndex === null) {
165            $indices = array_merge(
166                array_map('intval', $this->getFolders($this->pendingPath())),
167                array_map('intval', $this->getFolders($this->reservedPath()))
168            );
169
170            $this->endIndex = !empty($indices) ? max($indices) : 0;
171        }
172
173        return $this->endIndex;
174    }
175
176    /**
177     * Read a payload file and verify its signature, or false if it can't be read
178     *
179     * Every caller already treats a false return as "unusable payload, skip it",
180     * so an unreadable file joins the corrupt and tampered ones on that path.
181     * The explicit check matters because file_get_contents() signals failure
182     * with false rather than '', and under declare(strict_types=1) passing that
183     * to PayloadSigner::verify(string) is a TypeError - which would turn an
184     * everyday race (another worker claiming and unlinking the same payload
185     * between the file_exists() check and the read) into a crashed worker.
186     *
187     * @param  string $path
188     * @return string|false
189     */
190    protected function readVerifiedPayload(string $path): string|false
191    {
192        // Suppressed for the same reason as the @rename() in reclaimExpiredLeases():
193        // losing this read to a concurrent worker is expected operation, not a
194        // fault worth writing to the worker's stderr on every occurrence. The
195        // false return is what callers act on.
196        $payload = @file_get_contents($path);
197
198        return ($payload !== false) ? PayloadSigner::verify($payload) : false;
199    }
200
201    /**
202     * Read a reserved job's lease expiry, if any
203     *
204     * @param  string $reservedDir
205     * @return ?int
206     */
207    protected function getLeaseUntil(string $reservedDir): ?int
208    {
209        $leaseFile = $reservedDir . DIRECTORY_SEPARATOR . 'lease';
210        return file_exists($leaseFile) ? (int)file_get_contents($leaseFile) : null;
211    }
212
213    /**
214     * Determine whether a reserved job's lease has expired. The directory's
215     * own mtime is the primary signal - reserve() freshens it at claim time
216     * (rename() itself doesn't update it), so a fresh claim's mtime alone is
217     * sufficient to prove it isn't expired, regardless of what stale lease
218     * content happens to still be sitting in the directory (a job that was
219     * previously release()d or reclaimed carries its old, already-expired
220     * lease file back into pending/ with it - nothing unlinks it, so the
221     * next claim's directory starts out holding a stale expired lease next
222     * to a brand new mtime). The lease file is only consulted to confirm
223     * expiry once the mtime already looks stale, never to override a fresh
224     * mtime.
225     *
226     * @param  string $dir
227     * @param  int    $now
228     * @return bool
229     */
230    protected function isLeaseExpired(string $dir, int $now): bool
231    {
232        $mtime = @filemtime($dir);
233        if ($mtime === false) {
234            return false;
235        }
236
237        $leaseUntil = $this->getLeaseUntil($dir);
238
239        return ($mtime + $this->leaseSeconds <= $now) && (($leaseUntil === null) || ($leaseUntil <= $now));
240    }
241
242    /**
243     * Move any reserved job whose lease has expired back to pending, so a
244     * crashed worker's claim self-heals instead of being stuck forever.
245     * Reclaimed jobs are eligible on this or a later reserve() call, not
246     * necessarily returned by this one.
247     *
248     * Deliberately not throttled to one sweep per second, which would otherwise
249     * look free: every input to the expiry decision has one-second resolution,
250     * so two sweeps within the same second agree whenever the only thing moving
251     * is the clock. That is not the only thing that can move. A lease can be
252     * expired by writing to it, and a caller that does so and then calls
253     * reserve() is entitled to see the reclaim happen on that call rather than
254     * on whichever one lands in the next second. The sweep is cheap in any case,
255     * because reserved/ only ever holds jobs currently in flight - it is bounded
256     * by the number of live workers, not by queue depth.
257     *
258     * @return void
259     */
260    protected function reclaimExpiredLeases(): void
261    {
262        $now = time();
263
264        foreach ($this->getFolders($this->reservedPath()) as $index) {
265            $reservedDir = $this->reservedPath() . DIRECTORY_SEPARATOR . $index;
266
267            if (!$this->isLeaseExpired($reservedDir, $now)) {
268                continue;
269            }
270
271            // Stage through a uniquely-named path so at most one worker can win
272            // the rename off this exact directory, and so we can re-verify what
273            // actually got moved (not what we read a moment ago) before deciding
274            // to reclaim it.
275            $stagingDir = $this->reservedPath() . DIRECTORY_SEPARATOR . $index . '.reclaim-' . getmypid() . '-' . bin2hex(random_bytes(4));
276            if (!@rename($reservedDir, $stagingDir)) {
277                // Lost the race - someone else is reclaiming or already re-claimed this index.
278                continue;
279            }
280
281            // Re-check against what actually got staged (not the stale pre-read)
282            // using the exact same lease-or-mtime rule, so this can't drift out
283            // of sync with the check above and reclaim someone's fresh claim.
284            if (!$this->isLeaseExpired($stagingDir, $now)) {
285                // Whatever we staged turned out to be a fresh claim (another
286                // worker reclaimed-and-re-reserved this same index between our
287                // stale read and our rename winning) - put it back rather than
288                // reclaim someone's active work.
289                @rename($stagingDir, $reservedDir);
290                continue;
291            }
292
293            $pendingDir = $this->pendingPath() . DIRECTORY_SEPARATOR . $index;
294            if (!@rename($stagingDir, $pendingDir)) {
295                // Target collision or other failure - put it back to reserved/
296                // rather than permanently orphaning it in a staging directory
297                // nothing else will ever look at again.
298                @rename($stagingDir, $reservedDir);
299                continue;
300            }
301        }
302    }
303
304    /**
305     * Find the reserved-job directory index holding a given job, if any
306     *
307     * Called by release(), delete() and bury() - so once per job a worker
308     * finishes, whatever the outcome. It used to answer by reading and
309     * unserializing every reserved payload in turn, which meant reconstructing
310     * whole job objects (closures included) purely to read one string off each.
311     *
312     * reserve() now writes the claimed job's ID into a 'job-id' file beside the
313     * payload, so the common path compares a short string read against a string,
314     * and the payload is only unserialized for directories written before this
315     * file existed - or by a reserve() that died between its rename() and its
316     * sidecar write. Keeping that fallback is what makes the sidecar a pure
317     * optimization: its absence costs speed, never correctness.
318     *
319     * @param  AbstractJob $job
320     * @return ?int
321     */
322    protected function findReservedIndexForJob(AbstractJob $job): ?int
323    {
324        $jobId = $job->getJobId();
325
326        foreach ($this->getFolders($this->reservedPath()) as $index) {
327            // Skip staging directories (e.g. "5.reclaim-1234-abcd") left behind
328            // mid-reclaim - (int) casting one of those resolves to a path that
329            // doesn't actually exist, so treat only purely-numeric folder names
330            // as real job slots.
331            if (!ctype_digit((string)$index)) {
332                continue;
333            }
334
335            $dir = $this->reservedPath() . DIRECTORY_SEPARATOR . $index;
336
337            // Suppressed for the same reason as every other read on this path:
338            // losing the file to a concurrent reclaim is ordinary operation.
339            $storedId = @file_get_contents($dir . DIRECTORY_SEPARATOR . 'job-id');
340            if ($storedId !== false) {
341                if ($storedId === $jobId) {
342                    return (int)$index;
343                }
344                continue;
345            }
346
347            $payloadFile = $dir . DIRECTORY_SEPARATOR . 'payload';
348            if (file_exists($payloadFile)) {
349                $raw = $this->readVerifiedPayload($payloadFile);
350                // Suppressed: a corrupt/tampered payload makes unserialize() emit
351                // a warning and return false, which the instanceof check below
352                // handles.
353                $stored = ($raw !== false) ? @unserialize($raw) : false;
354                if (($stored instanceof AbstractJob) && ($stored->getJobId() === $jobId)) {
355                    return (int)$index;
356                }
357            }
358        }
359
360        return null;
361    }
362
363    /**
364     * Atomically claim a pending job directory by moving it into reserved/.
365     * rename() is the whole claim: exactly one worker's rename off a given
366     * source directory can succeed. Isolated into its own method (rather than
367     * inlined in reserve()) purely so a test can deterministically simulate a
368     * concurrent reclaim landing in the narrow window between this rename
369     * winning and reserve()'s follow-up touch()/lease write - see the guard in
370     * reserve() for what that window can otherwise damage.
371     *
372     * @param  string $pendingDir
373     * @param  string $reservedDir
374     * @return bool
375     */
376    protected function claimPendingDir(string $pendingDir, string $reservedDir): bool
377    {
378        return @rename($pendingDir, $reservedDir);
379    }
380
381    /**
382     * Push job on to queue
383     *
384     * @param  AbstractJob $job
385     * @throws Exception
386     * @return File
387     */
388    public function push(AbstractJob $job): File
389    {
390        // Force job ID generation before persisting so identity survives the
391        // serialize/unserialize round-trip on subsequent reserve/release/delete calls.
392        $job->getJobId();
393
394        // mkdir() itself is the atomic allocator: if two pushers compute the
395        // same next index, only one mkdir() wins and the loser retries the
396        // next index instead of silently clobbering the winner's payload. A
397        // failed mkdir() only means "keep trying" when the collision is with
398        // an existing directory (someone else's job) - any other failure
399        // (permissions, disk full, pending/ missing) must not spin forever.
400        //
401        // getEndIndex() is only a hint now that it's cached per instance, so
402        // this loop carries the two guards that hint can't provide on its own:
403        //
404        //  - A reserved/<index> check after the mkdir wins. An index is taken
405        //    if a directory bearing it exists in *either* pending/ or reserved/,
406        //    and mkdir() only knows about pending/. Checking after rather than
407        //    before also closes the race where another worker reserves that
408        //    index (moving it out of pending/ and into reserved/) in the gap
409        //    between the two calls. Backing out is safe: the directory we
410        //    created is still empty, and reserve() skips payload-less
411        //    directories, so nothing can have claimed it in between.
412        //
413        //  - A single re-scan on the first collision. A stale hint means the
414        //    real high-water mark has moved on, and walking up to it one
415        //    mkdir() at a time would reintroduce exactly the per-push cost the
416        //    cache exists to remove. Re-deriving it once jumps straight past
417        //    every taken index; only if that *also* collides does this fall
418        //    back to incrementing, which is the genuinely-contended case.
419        $index     = $this->getEndIndex() + 1;
420        $rescanned = false;
421
422        while (true) {
423            $dir = $this->pendingPath() . DIRECTORY_SEPARATOR . $index;
424
425            if (@mkdir($dir)) {
426                if (!is_dir($this->reservedPath() . DIRECTORY_SEPARATOR . $index)) {
427                    break;
428                }
429                @rmdir($dir);
430            } else if (!is_dir($dir)) {
431                throw new Exception('Error: Unable to create a new job folder in ' . $this->pendingPath() . '.');
432            }
433
434            if (!$rescanned) {
435                $rescanned      = true;
436                $this->endIndex = null;
437                $index          = $this->getEndIndex() + 1;
438            } else {
439                $index++;
440            }
441        }
442
443        $this->endIndex = $index;
444
445        file_put_contents($dir . DIRECTORY_SEPARATOR . 'payload', PayloadSigner::sign(serialize(clone $job)));
446
447        return $this;
448    }
449
450    /**
451     * Atomically claim the next eligible job. Reclaims any reserved job whose
452     * lease has expired first, then scans pending jobs in FIFO/FILO order,
453     * skipping any that aren't yet available, and atomically claims the first
454     * eligible one via rename() - if the rename fails, another worker won the
455     * race and this moves on to the next candidate.
456     *
457     * @return ?AbstractJob
458     */
459    public function reserve(): ?AbstractJob
460    {
461        $this->reclaimExpiredLeases();
462
463        // sort()/rsort() rather than usort() with a closure: the ordering is a
464        // property of the queue, not of any pair of indices, so re-asking
465        // isFifo() inside the comparator was doing O(n log n) method calls to
466        // re-derive one value that cannot change mid-sort. These also sort ints
467        // natively instead of calling back into PHP userland per comparison.
468        $indices = array_map('intval', $this->getFolders($this->pendingPath()));
469
470        if ($this->isFifo()) {
471            sort($indices, SORT_NUMERIC);
472        } else {
473            rsort($indices, SORT_NUMERIC);
474        }
475
476        foreach ($indices as $index) {
477            $pendingDir  = $this->pendingPath() . DIRECTORY_SEPARATOR . $index;
478            $payloadFile = $pendingDir . DIRECTORY_SEPARATOR . 'payload';
479
480            if (!file_exists($payloadFile)) {
481                continue;
482            }
483
484            // Suppressed: a corrupt/truncated payload makes unserialize() emit a
485            // warning and return false, which the instanceof check below handles.
486            // A payload that fails PayloadSigner::verify() (tampered, or written
487            // by something other than this application when a signing key is
488            // configured) is treated identically - $raw is false, so
489            // unserialize() is never called on it at all.
490            $raw = $this->readVerifiedPayload($payloadFile);
491            $job = ($raw !== false) ? @unserialize($raw) : false;
492
493            if (!($job instanceof AbstractJob)) {
494                // Corrupt/tampered payload - skip rather than crash or claim garbage.
495                continue;
496            }
497
498            if (!$job->isAvailable()) {
499                continue;
500            }
501
502            $reservedDir = $this->reservedPath() . DIRECTORY_SEPARATOR . $index;
503            if (!$this->claimPendingDir($pendingDir, $reservedDir)) {
504                // Lost the race to another worker - move on.
505                continue;
506            }
507
508            // A concurrent reclaim can move reserved/<index> away in the narrow
509            // window between the rename() above winning and the touch() below.
510            // Without this check, touch() on the now-vacant path would create a
511            // stray *regular file* where a job directory is expected, making
512            // that index permanently unclaimable (every future reserve() scan
513            // reaches it, the claiming rename() fails against a file, so it is
514            // skipped forever) and invisible to clear() (which only walks real
515            // job directories). Bailing out here degrades that into an ordinary
516            // lost race - the same safe failure mode every other race path in
517            // this class has. It does not close the underlying window; encoding
518            // the claim time into the rename itself is the real fix, deferred
519            // to a future design pass.
520            if (!is_dir($reservedDir)) {
521                continue;
522            }
523
524            // Freshen mtime immediately: rename() doesn't update it, so without
525            // this the directory's mtime would still reflect when the job was
526            // originally pushed, breaking isLeaseExpired()'s mtime fallback for
527            // any job that sat pending longer than the lease window before
528            // being claimed.
529            touch($reservedDir);
530            file_put_contents($reservedDir . DIRECTORY_SEPARATOR . 'lease', (string)(time() + $this->leaseSeconds));
531
532            // Sidecar for findReservedIndexForJob(), which release()/delete()/
533            // bury() all go through when this job finishes. Written after the
534            // lease so it can never be the thing that makes a claim look
535            // complete before it is; a failure to write it costs nothing but
536            // the fast path, since that lookup falls back to the payload.
537            file_put_contents($reservedDir . DIRECTORY_SEPARATOR . 'job-id', $job->getJobId());
538
539            return $job;
540        }
541
542        return null;
543    }
544
545    /**
546     * Put a job back to pending, honoring its backoff schedule unless an
547     * explicit delay is given
548     *
549     * @param  AbstractJob $job
550     * @param  ?int        $delay
551     * @return File
552     */
553    public function release(AbstractJob $job, ?int $delay = null): File
554    {
555        $index = $this->findReservedIndexForJob($job);
556        if ($index === null) {
557            return $this;
558        }
559
560        $job->delay($delay ?? $job->getBackoffDelay());
561
562        $reservedDir = $this->reservedPath() . DIRECTORY_SEPARATOR . $index;
563        $pendingDir  = $this->pendingPath() . DIRECTORY_SEPARATOR . $index;
564
565        file_put_contents($reservedDir . DIRECTORY_SEPARATOR . 'payload', PayloadSigner::sign(serialize(clone $job)));
566        rename($reservedDir, $pendingDir);
567
568        return $this;
569    }
570
571    /**
572     * Permanently remove a job
573     *
574     * @param  AbstractJob $job
575     * @return File
576     */
577    public function delete(AbstractJob $job): File
578    {
579        $index = $this->findReservedIndexForJob($job);
580        if ($index === null) {
581            return $this;
582        }
583
584        $this->removeJobDir($this->reservedPath() . DIRECTORY_SEPARATOR . $index);
585
586        return $this;
587    }
588
589    /**
590     * Empty and remove a job directory.
591     *
592     * Deliberately generic rather than unlinking 'payload' and 'lease' by name:
593     * rmdir() fails on a non-empty directory, so every file a job directory can
594     * hold has to be accounted for here, and naming them individually means any
595     * future addition silently turns delete() into a no-op that leaves the
596     * directory behind. Clearing whatever is actually in there cannot drift out
597     * of sync that way.
598     *
599     * @param  string $dir
600     * @return void
601     */
602    protected function removeJobDir(string $dir): void
603    {
604        if (!is_dir($dir)) {
605            return;
606        }
607
608        // Iterated raw rather than through getFiles(), which filters out
609        // '.empty' - anything left behind, whatever its name, keeps rmdir()
610        // from succeeding.
611        foreach (new \FilesystemIterator($dir, \FilesystemIterator::SKIP_DOTS) as $entry) {
612            if ($entry->isFile()) {
613                @unlink($entry->getPathname());
614            }
615        }
616
617        @rmdir($dir);
618    }
619
620    /**
621     * Move a job to the dead-letter store
622     *
623     * @param  AbstractJob $job
624     * @param  ?string     $reason
625     * @return File
626     */
627    public function bury(AbstractJob $job, ?string $reason = null): File
628    {
629        $this->delete($job);
630        file_put_contents($this->folder . DIRECTORY_SEPARATOR . 'dead-' . $job->getJobId(), PayloadSigner::sign(serialize(clone $job)));
631
632        return $this;
633    }
634
635    /**
636     * Check if adapter has jobs
637     *
638     * @return bool
639     */
640    public function hasJobs(): bool
641    {
642        return ($this->count() > 0);
643    }
644
645    /**
646     * Count of pending + reserved jobs
647     *
648     * @return int
649     */
650    public function count(): int
651    {
652        return count($this->getFolders($this->pendingPath())) + count($this->getFolders($this->reservedPath()));
653    }
654
655    /**
656     * Clear pending and reserved jobs (not tasks or dead-letter jobs)
657     *
658     * @return File
659     */
660    public function clear(): File
661    {
662        foreach ([$this->pendingPath(), $this->reservedPath()] as $path) {
663            foreach ($this->getFolders($path) as $index) {
664                $this->removeJobDir($path . DIRECTORY_SEPARATOR . $index);
665            }
666        }
667
668        // Both directories are empty now, so the cached high-water index is
669        // stale in the one direction worth correcting: leaving it set would
670        // keep numbering new jobs from wherever the cleared queue left off,
671        // where a fresh scan restarts from 1 the way it always has.
672        $this->endIndex = null;
673
674        return $this;
675    }
676
677    /**
678     * Get the dead-letter job IDs
679     *
680     * @return array
681     */
682    protected function getDeadJobIds(): array
683    {
684        $ids = [];
685        foreach ($this->getFiles($this->folder) as $file) {
686            if (str_starts_with($file, 'dead-')) {
687                $ids[] = substr($file, 5);
688            }
689        }
690
691        return $ids;
692    }
693
694    /**
695     * Check if adapter has dead-letter jobs
696     *
697     * @return bool
698     */
699    public function hasDeadJobs(): bool
700    {
701        return !empty($this->getDeadJobIds());
702    }
703
704    /**
705     * Count of dead-letter jobs
706     *
707     * @return int
708     */
709    public function countDead(): int
710    {
711        return count($this->getDeadJobIds());
712    }
713
714    /**
715     * Get dead-letter jobs
716     *
717     * @param  bool $unserialize
718     * @return array
719     */
720    public function getDeadJobs(bool $unserialize = true): array
721    {
722        $jobs = [];
723        foreach ($this->getDeadJobIds() as $jobId) {
724            $jobs[$jobId] = $this->getDeadJob($jobId, $unserialize);
725        }
726
727        return $jobs;
728    }
729
730    /**
731     * Get a dead-letter job
732     *
733     * @param  string $jobId
734     * @param  bool   $unserialize
735     * @return mixed
736     */
737    public function getDeadJob(string $jobId, bool $unserialize = true): mixed
738    {
739        $path = $this->folder . DIRECTORY_SEPARATOR . 'dead-' . $jobId;
740        if (!file_exists($path)) {
741            return null;
742        }
743
744        // Read once and share it with both branches rather than going through
745        // readVerifiedPayload(), which would re-read the file for the
746        // unserialize case. A failed read is reported the same way a missing
747        // file is, above - the job is unreadable either way.
748        $payload = @file_get_contents($path);
749        if ($payload === false) {
750            return null;
751        }
752
753        if (!$unserialize) {
754            return $payload;
755        }
756
757        $raw = PayloadSigner::verify($payload);
758        return ($raw !== false) ? unserialize($raw) : false;
759    }
760
761    /**
762     * Move a dead-letter job back to pending
763     *
764     * @param  string $jobId
765     * @return File
766     */
767    public function retryDeadJob(string $jobId): File
768    {
769        $job = $this->getDeadJob($jobId);
770        if ($job instanceof AbstractJob) {
771            $this->push($job);
772            $this->deleteDeadJob($jobId);
773        }
774
775        return $this;
776    }
777
778    /**
779     * Permanently remove a dead-letter job
780     *
781     * @param  string $jobId
782     * @return File
783     */
784    public function deleteDeadJob(string $jobId): File
785    {
786        $path = $this->folder . DIRECTORY_SEPARATOR . 'dead-' . $jobId;
787        if (file_exists($path)) {
788            unlink($path);
789        }
790
791        return $this;
792    }
793
794    /**
795     * Clear all dead-letter jobs
796     *
797     * @return File
798     */
799    public function clearDead(): File
800    {
801        foreach ($this->getDeadJobIds() as $jobId) {
802            $this->deleteDeadJob($jobId);
803        }
804
805        return $this;
806    }
807
808    /**
809     * Schedule job with queue
810     *
811     * @param  Task $task
812     * @return File
813     */
814    public function schedule(Task $task): File
815    {
816        if ($task->isValid()) {
817            file_put_contents(
818                $this->folder . DIRECTORY_SEPARATOR . 'task-' . $task->getJobId(), PayloadSigner::sign(serialize(clone $task))
819            );
820        }
821        return $this;
822    }
823
824    /**
825     * Get scheduled tasks
826     *
827     * @return array
828     */
829    public function getTasks(): array
830    {
831        $files = $this->getFiles($this->folder);
832        $tasks = [];
833
834        foreach ($files as $file) {
835            if (str_starts_with($file, 'task-')) {
836                $tasks[] = substr($file, 5);
837            }
838        }
839
840        return $tasks;
841    }
842
843    /**
844     * Get scheduled task
845     *
846     * @param  string $taskId
847     * @return ?Task
848     */
849    public function getTask(string $taskId): ?Task
850    {
851        $path = $this->folder . DIRECTORY_SEPARATOR . 'task-' . $taskId;
852        if (!file_exists($path)) {
853            return null;
854        }
855
856        // Guarded with instanceof rather than returning unserialize()'s result
857        // directly, the same way reserve() does: a corrupt, truncated or
858        // tampered payload makes unserialize() return false, and false out of
859        // a ": ?Task" method is a TypeError, not a null. A bad task file should
860        // read as "no such task", never as a crash.
861        $raw  = $this->readVerifiedPayload($path);
862        $task = ($raw !== false) ? @unserialize($raw) : false;
863
864        return ($task instanceof Task) ? $task : null;
865    }
866
867    /**
868     * Update scheduled task
869     *
870     * @param  Task $task
871     * @return File
872     */
873    public function updateTask(Task $task): File
874    {
875        if ($task->isValid()) {
876            file_put_contents(
877                $this->folder . DIRECTORY_SEPARATOR . 'task-' . $task->getJobId(), PayloadSigner::sign(serialize(clone $task))
878            );
879        } else {
880            $this->removeTask($task->getJobId());
881        }
882
883        return $this;
884    }
885
886    /**
887     * Remove scheduled task
888     *
889     * @param  string $taskId
890     * @return File
891     */
892    public function removeTask(string $taskId): File
893    {
894        if (file_exists($this->folder . DIRECTORY_SEPARATOR . 'task-' . $taskId)) {
895            unlink($this->folder . DIRECTORY_SEPARATOR . 'task-' . $taskId);
896        }
897        if (file_exists($this->taskClaimPath($taskId))) {
898            unlink($this->taskClaimPath($taskId));
899        }
900        return $this;
901    }
902
903    /**
904     * Get scheduled tasks count
905     *
906     * @return int
907     */
908    public function getTaskCount(): int
909    {
910        return count($this->getTasks());
911    }
912
913    /**
914     * Has scheduled tasks
915     *
916     * @return bool
917     */
918    public function hasTasks(): bool
919    {
920        return ($this->getTaskCount() > 0);
921    }
922
923    /**
924     * Clear all scheduled task
925     *
926     * @return File
927     */
928    public function clearTasks(): File
929    {
930        $tasks = $this->getTasks();
931
932        foreach ($tasks as $taskId) {
933            $this->removeTask($taskId);
934        }
935
936        foreach ($this->getFiles($this->folder) as $file) {
937            if (str_starts_with($file, 'claim-task-')) {
938                unlink($this->folder . DIRECTORY_SEPARATOR . $file);
939            }
940        }
941
942        return $this;
943    }
944
945    /**
946     * Get the claim-marker file path for a task
947     *
948     * @param  string $taskId
949     * @return string
950     */
951    protected function taskClaimPath(string $taskId): string
952    {
953        return $this->folder . DIRECTORY_SEPARATOR . 'claim-task-' . $taskId;
954    }
955
956    /**
957     * Atomically claim a task's current due-window via a small sidecar
958     * file, deliberately separate from the task's own serialized
959     * definition file (task-<taskId>) so a claim attempt never touches or
960     * re-serializes the closure-bearing Task object. Content format is
961     * "<window>:<expiresAtUnixTimestamp>". flock() provides real
962     * cross-process mutual exclusion for the read-decide-write.
963     *
964     * @param  string $taskId
965     * @param  string $window
966     * @return bool
967     */
968    public function claimTaskRun(string $taskId, string $window): bool
969    {
970        $path = $this->taskClaimPath($taskId);
971
972        $fh = fopen($path, 'c+');
973        if ($fh === false) {
974            return false;
975        }
976
977        if (!flock($fh, LOCK_EX)) {
978            fclose($fh);
979            return false;
980        }
981
982        $contents = stream_get_contents($fh);
983        $now      = time();
984        $claimed  = true;
985
986        if (!empty($contents)) {
987            [$storedWindow, $storedExpiry] = array_pad(explode(':', $contents, 2), 2, '0');
988            if (($storedWindow === $window) && ((int)$storedExpiry > $now)) {
989                $claimed = false;
990            }
991        }
992
993        if ($claimed) {
994            ftruncate($fh, 0);
995            rewind($fh);
996            fwrite($fh, $window . ':' . ($now + self::TASK_CLAIM_TTL));
997            fflush($fh);
998        }
999
1000        flock($fh, LOCK_UN);
1001        fclose($fh);
1002
1003        return $claimed;
1004    }
1005
1006    /**
1007     * Get folders
1008     *
1009     * @param  string $folder
1010     * @return array
1011     */
1012    public function getFolders(string $folder): array
1013    {
1014        return $this->readDirectory($folder, true);
1015    }
1016
1017    /**
1018     * Get files from folder
1019     *
1020     * @param  string $folder
1021     * @return array
1022     */
1023    public function getFiles(string $folder): array
1024    {
1025        return $this->readDirectory($folder, false);
1026    }
1027
1028    /**
1029     * List the entries of a directory, keeping either the subdirectories or the
1030     * plain files.
1031     *
1032     * Both public listers route through here rather than each running their own
1033     * scandir(). Two things made that pairing expensive on the hot paths -
1034     * reserve() and count() call it on every invocation, once per pending or
1035     * reserved job:
1036     *
1037     *  - scandir() returns names only, so deciding what each entry *is* meant an
1038     *    is_dir() stat syscall per entry. FilesystemIterator carries the type
1039     *    along with the directory read, so isDir() answers from what the OS
1040     *    already handed back.
1041     *  - scandir() sorts alphabetically by default, and every caller here either
1042     *    wants numeric order (reserve() re-sorts these index names itself) or no
1043     *    order at all, so that sort was pure waste.
1044     *
1045     * A missing directory reads as empty rather than raising, matching the
1046     * is_dir() guard both listers carried before.
1047     *
1048     * @param  string $folder
1049     * @param  bool   $directories
1050     * @return array
1051     */
1052    protected function readDirectory(string $folder, bool $directories): array
1053    {
1054        if (!is_dir($folder)) {
1055            return [];
1056        }
1057
1058        $entries = [];
1059
1060        // SKIP_DOTS drops '.' and '..'; '.empty' is a repo placeholder that is
1061        // neither a job nor a task and has always been filtered out here.
1062        $iterator = new \FilesystemIterator($folder, \FilesystemIterator::SKIP_DOTS | \FilesystemIterator::CURRENT_AS_FILEINFO);
1063
1064        foreach ($iterator as $entry) {
1065            $name = $entry->getFilename();
1066            if (($name !== '.empty') && ($entry->isDir() === $directories)) {
1067                $entries[] = $name;
1068            }
1069        }
1070
1071        return $entries;
1072    }
1073
1074}