Coroutine-aware priority queue with pluggable backoff strategies for Momo Framework.
🇬🇧 English · 🇷🇺 Русская версия
momo-framework/queue provides a coroutine-aware job queue engine backed by a priority min-heap. Jobs declare their retry budget and dispatch priority through a single #[AsJob] attribute. The Worker drains the queue in a tight loop, retrying failures with configurable backoff and moving exhausted jobs to a failed-job store.
The package integrates with momo-framework/events: JobProcessing, JobProcessed, and JobFailed domain events are published on every lifecycle transition, enabling telemetry and notifications without coupling them to the worker itself.
InMemoryQueue is the built-in driver. It uses two binary heaps — one for ready jobs ordered by priority, one for delayed jobs ordered by availability time — and is safe under Swoole's cooperative scheduler because every mutating operation completes without an I/O suspension point.
- PHP >= 8.5
momo-framework/kernelmomo-framework/eventsext-swoole(optional — required only forCoroutineSleeper)
composer require momo-framework/queueuse Momo\Queue\Attributes\AsJob;
#[AsJob(maxAttempts: 5, priority: 10)]| Parameter | Type | Default | Description |
|---|---|---|---|
maxAttempts |
positive-int |
3 |
Total attempts including the initial run |
priority |
int |
0 |
Higher value dispatches first |
Extend Momo\Queue\Job\AbstractJob and implement handle(): void. Optionally override failed(Throwable $exception): void for compensating actions after the retry budget is exhausted.
Backed by two binary heaps. push() enqueues immediately; later() enqueues with a delay. pop() promotes due delayed jobs before returning the highest-priority ready envelope.
Three built-in policies implement BackoffPolicyInterface:
| Class | Behaviour |
|---|---|
ExponentialBackoffPolicy |
base * multiplier^(attempt-2), capped at maxSeconds |
ConstantBackoffPolicy |
Fixed delay for every retry |
FullJitterBackoffPolicy |
Random delay in [0, exponential_delay] |
Jobs may override the global policy per-class by implementing HasBackoffPolicyInterface.
Worker::run(WorkerOptions) drains the queue until stopped, the job limit is reached, or memory is exceeded. WorkerOptions::once() is a convenience factory that processes a single job and exits.
When ext-swoole is loaded, idle waits call Swoole\Coroutine::sleep() to yield the scheduler. Without Swoole it falls back to usleep().
use Momo\Queue\Attributes\AsJob;
use Momo\Queue\Job\AbstractJob;
#[AsJob(maxAttempts: 5, priority: 10)]
final class SendOrderReceipt extends AbstractJob
{
public function __construct(
private readonly string $orderId,
) {}
public function handle(): void
{
// deliver the receipt
}
public function failed(\Throwable $e): void
{
// compensating action after retry budget is exhausted
Log::critical("Receipt delivery failed for order {$this->orderId}", ['error' => $e->getMessage()]);
}
}use Momo\Queue\Contracts\QueueInterface;
// Enqueue immediately with explicit priority
$queue->push(new SendOrderReceipt($orderId), priority: 10);
// Enqueue with a delay (seconds)
$queue->later(delaySeconds: 30, job: new SendOrderReceipt($orderId));
// Priority from the #[AsJob] attribute is used when not overriding here
$queue->push(new SendOrderReceipt($orderId));php momo queue:work # Long-running worker
php momo queue:work --once # Process a single job and exit
php momo queue:failed # Inspect the dead-letter storeOr programmatically:
use Momo\Queue\Worker\Worker;
use Momo\Queue\Worker\WorkerOptions;
$worker->run(new WorkerOptions(
maxJobs: 0, // 0 = unlimited
sleepSeconds: 1, // idle polling interval
stopWhenEmpty: false, // keep running after queue drains
memoryLimitMb: 256, // restart when RSS exceeds this
));
// Process exactly one job and return:
$worker->run(WorkerOptions::once());use Momo\Queue\Backoff\ExponentialBackoffPolicy;
use Momo\Queue\Contracts\BackoffPolicyInterface;
use Momo\Queue\Contracts\HasBackoffPolicyInterface;
use Momo\Queue\Job\AbstractJob;
#[AsJob(maxAttempts: 10)]
final class ImportLargeFile extends AbstractJob implements HasBackoffPolicyInterface
{
public function handle(): void { /* ... */ }
public function backoffPolicy(): BackoffPolicyInterface
{
return new ExponentialBackoffPolicy(baseSeconds: 30, maxSeconds: 3600, multiplier: 2.0);
}
}use Momo\Events\Contracts\EventBusInterface;
use Momo\Queue\Events\JobFailed;
use Momo\Queue\Events\JobProcessed;
use Momo\Queue\Events\JobProcessing;
// Subscribe in a ServiceProvider::boot()
$bus->subscribe(JobProcessing::class, new LogJobStartedListener());
$bus->subscribe(JobProcessed::class, new RecordMetricsListener());
$bus->subscribe(JobFailed::class, new AlertOnFailureListener());composer lint
composer stan
composer rector:check
composer test
composer ci