Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
75 changes: 73 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,20 @@ A lightweight job queue library for PHP with file-based storage, priority orderi
</a>
</p>

## 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**.
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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)

Expand Down Expand Up @@ -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.
87 changes: 73 additions & 14 deletions WebFiori/Queue/Queue.php
Original file line number Diff line number Diff line change
Expand Up @@ -11,13 +11,25 @@
*/
namespace WebFiori\Queue;

use LogicException;
use Throwable;
use UnexpectedValueException;

/**
* Core queue class that dispatches and processes jobs.
*
* Handles serialization, encryption, and retry logic. The storage layer
* 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;

/**
Expand Down Expand Up @@ -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.
Expand All @@ -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.
*
Expand All @@ -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
Expand All @@ -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);
}
}

Expand All @@ -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.
*
Expand Down Expand Up @@ -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);
}
}
}
18 changes: 12 additions & 6 deletions WebFiori/Queue/QueueFacade.php
Original file line number Diff line number Diff line change
Expand Up @@ -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.
*
Expand All @@ -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()
*/
Expand All @@ -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);
}
}
Loading
Loading