From 9a203ca79b6ca69eff6ab79cf0b6e987c703987e Mon Sep 17 00:00:00 2001 From: Ibrahim BinAlshikh Date: Sun, 20 Sep 2026 21:55:40 +0300 Subject: [PATCH] feat(queue): add opt-in error callback and standardize README #9: expose job-processing exceptions - Add Queue::setOnError()/getOnError(): a callback invoked for every caught throwable while processing a job, with signature fn(?Job $job, Throwable $e, int $attempts, bool $willRetry). It fires for both intermediate retries and terminal failures, letting applications bridge queue failures into a centralized logger with full exception context. - Enrich the failed job's failReason with the exception class in addition to the message. - QueueFacade::setOnError() passthrough. - Invalid (non-Job) payloads are reported to the callback as terminal failures with a null job. Additive and opt-in: behavior is unchanged without a callback. #6: standardize README - Add Table of Contents and Testing/Examples/Contributing/Support/Changelog sections; document setOnError; fix an inaccurate getFailed() retry example. Closes #9, closes #6 --- README.md | 75 ++++++++++++++++++++++++++- WebFiori/Queue/Queue.php | 87 ++++++++++++++++++++++++++----- WebFiori/Queue/QueueFacade.php | 18 ++++--- tests/QueueTest.php | 93 +++++++++++++++++++++++++++++++++- 4 files changed, 249 insertions(+), 24 deletions(-) diff --git a/README.md b/README.md index 2c66661..33c8120 100644 --- a/README.md +++ b/README.md @@ -18,6 +18,20 @@ A lightweight job queue library for PHP with file-based storage, priority orderi

+## Table of Contents + +- [Supported PHP Versions](#supported-php-versions) +- [Features](#features) +- [Installation](#installation) +- [Usage](#usage) +- [API](#api) +- [Testing](#testing) +- [Examples](#examples) +- [Contributing](#contributing) +- [License](#license) +- [Support](#support) +- [Changelog](#changelog) + ## Supported PHP Versions This library requires **PHP 8.1 or higher**. @@ -104,16 +118,42 @@ QueueFacade::process(); ### Failed Jobs ```php -// View failed jobs +// View failed jobs (each is a QueuedJob) $failed = $queue->getFailed(); // Retry a specific failed job -$queue->retry($failed[0]['id']); +$queue->retry($failed[0]->getId()); // Clear all failed jobs $queue->flush(); ``` +### Observing Failures + +By default, exceptions thrown by a job's `handle()` are caught and turned into +retries or a failed-job record. Register an **opt-in error callback** to observe +every caught throwable — with full exception context — for both intermediate +retries and terminal failures. This makes it easy to bridge queue failures into +a centralized logger (e.g. a globally registered `WebFiori\Error\Handler`). + +```php +$queue->setOnError(function (?Job $job, \Throwable $e, int $attempts, bool $willRetry): void { + // $job is null if the stored payload was not a valid Job instance. + error_log(sprintf( + '[queue] %s failed on attempt %d (%s): %s', + $job !== null ? get_class($job) : 'invalid-payload', + $attempts, + $willRetry ? 'will retry' : 'terminal', + $e->getMessage() + )); +}); +``` + +The failed `QueuedJob` also records richer context — its `getFailReason()` now +includes the exception class in addition to the message. Behavior is unchanged +when no callback is registered. The same callback can be set on the facade via +`QueueFacade::setOnError(...)`. + ### Inspecting Pending Jobs ```php @@ -159,6 +199,8 @@ $count = $queue->getPendingCount(); | `getFailed(): array` | All failed jobs | | `flush(): void` | Remove all failed jobs | | `getStorage(): QueueStorage` | Get the storage backend | +| `setOnError(?callable $cb): Queue` | Register a callback invoked on every caught throwable: `fn(?Job $job, \Throwable $e, int $attempts, bool $willRetry)` | +| `getOnError(): ?callable` | Get the registered error callback | ### `QueueStorage` (interface) @@ -186,6 +228,35 @@ Backends that support listing (files, database, Redis) implement `ListableQueueS Static wrapper. Same methods as `Queue` plus `getInstance()`, `setInstance()`, `reset()`. +## Testing + +Run the test suite with: + +```bash +composer test +``` + +## Examples + +Runnable examples are available in the [`examples/`](examples) directory: + +- [`examples/01-basic-queue.php`](examples/01-basic-queue.php) — dispatching and processing jobs +- [`examples/02-custom-storage.php`](examples/02-custom-storage.php) — implementing a custom storage backend + +## Contributing + +Contributions are welcome. Please open an issue to discuss significant changes, +follow [Conventional Commits](https://www.conventionalcommits.org/) for commit +messages, and ensure `composer test` passes before opening a pull request. + ## License MIT + +## Support + +- **Issues**: [GitHub Issues](https://github.com/WebFiori/queue/issues) + +## Changelog + +See [CHANGELOG.md](CHANGELOG.md) for a full history of changes and releases. diff --git a/WebFiori/Queue/Queue.php b/WebFiori/Queue/Queue.php index 612aa46..b3f1067 100644 --- a/WebFiori/Queue/Queue.php +++ b/WebFiori/Queue/Queue.php @@ -11,6 +11,10 @@ */ namespace WebFiori\Queue; +use LogicException; +use Throwable; +use UnexpectedValueException; + /** * Core queue class that dispatches and processes jobs. * @@ -18,6 +22,14 @@ * only deals with QueuedJob value objects containing opaque payloads. */ class Queue { + /** + * Optional callback invoked for every throwable caught while processing a job. + * + * Signature: fn(?Job $job, Throwable $e, int $attempts, bool $willRetry): void + * + * @var callable|null + */ + private $onError = null; private QueueStorage $storage; /** @@ -64,12 +76,12 @@ public function getFailed(): array { return $this->storage->getFailed(); } /** - * Returns the number of pending jobs. + * Returns the configured error callback, if any. * - * @return int + * @return callable|null */ - public function getPendingCount(): int { - return $this->storage->getPendingCount(); + public function getOnError(): ?callable { + return $this->onError; } /** * Returns all pending jobs, including delayed ones not yet available. @@ -79,18 +91,26 @@ public function getPendingCount(): int { * * @return QueuedJob[] Array of all pending queued jobs. * - * @throws \LogicException If the storage backend does not implement ListableQueueStorage. + * @throws LogicException If the storage backend does not implement ListableQueueStorage. */ public function getPending(): array { if (!($this->storage instanceof ListableQueueStorage)) { - throw new \LogicException( + throw new LogicException( 'The configured storage backend does not support listing pending jobs. ' - . 'Use a ListableQueueStorage implementation (e.g. FileQueueStorage).' + .'Use a ListableQueueStorage implementation (e.g. FileQueueStorage).' ); } return $this->storage->getPending(); } + /** + * Returns the number of pending jobs. + * + * @return int + */ + public function getPendingCount(): int { + return $this->storage->getPendingCount(); + } /** * Returns the storage backend. * @@ -116,25 +136,28 @@ public function process(int $limit = 10): int { foreach ($pending as $queuedJob) { $id = $queuedJob->getId(); $attempts = $queuedJob->getAttempts() + 1; + $job = null; try { $job = unserialize($this->decrypt($queuedJob->getPayload())); if (!($job instanceof Job)) { - $queuedJob->setAttempts($attempts); - $queuedJob->setFailReason('Payload is not a valid Job instance.'); - $this->storage->markFailed($queuedJob); + $job = null; - continue; + throw new UnexpectedValueException('Payload is not a valid Job instance.'); } $job->handle(); $this->storage->markComplete($id); $processed++; - } catch (\Throwable $e) { - if ($attempts >= $job->getMaxAttempts()) { + } catch (Throwable $e) { + // A non-Job payload cannot be retried; treat it as terminal. + $maxAttempts = $job !== null ? $job->getMaxAttempts() : 1; + $willRetry = $attempts < $maxAttempts; + + if (!$willRetry) { $queuedJob->setAttempts($attempts); - $queuedJob->setFailReason($e->getMessage()); + $queuedJob->setFailReason(get_class($e).': '.$e->getMessage()); $this->storage->markFailed($queuedJob); } else { // Re-queue with updated attempt count and delay @@ -144,6 +167,8 @@ public function process(int $limit = 10): int { $queuedJob->setAvailableAt(time() + $delay); $this->storage->push($queuedJob); } + + $this->invokeOnError($job, $e, $attempts, $willRetry); } } @@ -157,6 +182,27 @@ public function process(int $limit = 10): int { public function retry(string $id): void { $this->storage->retry($id); } + /** + * Sets an optional callback invoked whenever a throwable is caught while + * processing a job. + * + * The callback lets applications observe, log, or react to job failures + * (e.g. bridge them into a globally registered error handler) with full + * exception context. It is invoked for both terminal failures (attempts + * exhausted) and intermediate failures that will be retried. + * + * Signature: fn(?Job $job, Throwable $e, int $attempts, bool $willRetry): void + * ($job is null when the stored payload is not a valid Job instance.) + * + * @param callable|null $callback The error callback, or null to disable. + * + * @return Queue This instance, for chaining. + */ + public function setOnError(?callable $callback): Queue { + $this->onError = $callback; + + return $this; + } /** * Decrypts data if it was encrypted. * @@ -211,4 +257,17 @@ private function encrypt(string $data): string { private function generateId(): string { return bin2hex(random_bytes(16)); } + /** + * Invokes the error callback if one is configured. + * + * @param Job|null $job The job that failed, or null for an invalid payload. + * @param Throwable $e The caught throwable. + * @param int $attempts The attempt count at the time of failure. + * @param bool $willRetry Whether the job will be retried. + */ + private function invokeOnError(?Job $job, Throwable $e, int $attempts, bool $willRetry): void { + if ($this->onError !== null) { + ($this->onError)($job, $e, $attempts, $willRetry); + } + } } diff --git a/WebFiori/Queue/QueueFacade.php b/WebFiori/Queue/QueueFacade.php index 6be2263..fbee5f3 100644 --- a/WebFiori/Queue/QueueFacade.php +++ b/WebFiori/Queue/QueueFacade.php @@ -50,12 +50,6 @@ public static function getInstance(): Queue { return self::$inst; } - /** - * @see Queue::getPendingCount() - */ - public static function getPendingCount(): int { - return self::getInstance()->getPendingCount(); - } /** * Returns all pending jobs, including delayed ones not yet available. * @@ -70,6 +64,12 @@ public static function getPendingCount(): int { public static function getPending(): array { return self::getInstance()->getPending(); } + /** + * @see Queue::getPendingCount() + */ + public static function getPendingCount(): int { + return self::getInstance()->getPendingCount(); + } /** * @see Queue::process() */ @@ -96,4 +96,10 @@ public static function retry(string $id): void { public static function setInstance(Queue $queue): void { self::$inst = $queue; } + /** + * @see Queue::setOnError() + */ + public static function setOnError(?callable $callback): void { + self::getInstance()->setOnError($callback); + } } diff --git a/tests/QueueTest.php b/tests/QueueTest.php index b20676e..f600266 100644 --- a/tests/QueueTest.php +++ b/tests/QueueTest.php @@ -287,7 +287,8 @@ public function testMaxAttemptsExhaustedMovesToFailed() { $this->assertEquals(0, $this->queue->getPendingCount()); $failed = $this->queue->getFailed(); $this->assertCount(1, $failed); - $this->assertEquals('Always fails', $failed[0]->getFailReason()); + $this->assertStringContainsString('RuntimeException', $failed[0]->getFailReason()); + $this->assertStringContainsString('Always fails', $failed[0]->getFailReason()); $this->assertEquals(2, $failed[0]->getAttempts()); } /** @@ -314,7 +315,7 @@ public function testGetFailedPreservesReason() { $failed = $this->queue->getFailed(); $this->assertCount(1, $failed); - $this->assertEquals('Always fails', $failed[0]->getFailReason()); + $this->assertStringContainsString('Always fails', $failed[0]->getFailReason()); $this->assertNotEmpty($failed[0]->getId()); } /** @@ -436,4 +437,92 @@ public function retry(string $id): void {} $this->expectExceptionMessage('does not support listing pending jobs'); $queue->getPending(); } + + /** + * @test + * The onError callback fires on a terminal failure with willRetry = false. + */ + public function testOnErrorInvokedOnTerminalFailure() { + $calls = []; + $this->queue->setOnError(function ($job, $e, $attempts, $willRetry) use (&$calls) { + $calls[] = [$job, $e, $attempts, $willRetry]; + }); + + $this->queue->dispatch(new AlwaysFailsJob()); // maxAttempts = 2 + $this->queue->process(); // attempt 1 -> retry + $this->queue->process(); // attempt 2 -> terminal + + $this->assertCount(2, $calls); + + // First call: intermediate retry. + $this->assertInstanceOf(AlwaysFailsJob::class, $calls[0][0]); + $this->assertInstanceOf(\RuntimeException::class, $calls[0][1]); + $this->assertSame(1, $calls[0][2]); + $this->assertTrue($calls[0][3]); + + // Second call: terminal. + $this->assertSame(2, $calls[1][2]); + $this->assertFalse($calls[1][3]); + $this->assertSame('Always fails', $calls[1][1]->getMessage()); + } + + /** + * @test + * With no callback configured, processing behaves exactly as before. + */ + public function testNoOnErrorCallbackKeepsDefaultBehavior() { + $this->queue->dispatch(new AlwaysFailsJob()); + $this->queue->process(); + $this->queue->process(); + + $this->assertEquals(0, $this->queue->getPendingCount()); + $this->assertCount(1, $this->queue->getFailed()); + } + + /** + * @test + * A non-Job payload is reported to the callback as a terminal failure + * with a null job. + */ + public function testOnErrorInvokedForInvalidPayload() { + $captured = null; + $this->queue->setOnError(function ($job, $e, $attempts, $willRetry) use (&$captured) { + $captured = [$job, $willRetry]; + }); + + // Push a payload that is not a serialized Job. + $this->queue->getStorage()->push(new QueuedJob('bad-id', serialize('not a job'), 0, 0, 0)); + $this->queue->process(); + + $this->assertNotNull($captured); + $this->assertNull($captured[0]); + $this->assertFalse($captured[1]); + $this->assertCount(1, $this->queue->getFailed()); + } + + /** + * @test + * setOnError is chainable and exposed via getOnError. + */ + public function testSetOnErrorChainableAndGettable() { + $cb = function () { + }; + $this->assertSame($this->queue, $this->queue->setOnError($cb)); + $this->assertSame($cb, $this->queue->getOnError()); + $this->queue->setOnError(null); + $this->assertNull($this->queue->getOnError()); + } + + /** + * @test + * The facade exposes setOnError, delegating to the default instance. + */ + public function testFacadeSetOnErrorPassthrough() { + QueueFacade::reset(); + $cb = function () { + }; + QueueFacade::setOnError($cb); + $this->assertSame($cb, QueueFacade::getInstance()->getOnError()); + QueueFacade::reset(); + } }