From d5911fc1ee0a93f6da97b3dcfbe37fc7638dee07 Mon Sep 17 00:00:00 2001 From: Sajeeb Ahamed Date: Mon, 28 Sep 2026 19:42:32 +0600 Subject: [PATCH] Add database-backed job queue with WP-Cron and loopback workers Adds a Laravel-style queue for plugins built on the framework, designed for WordPress hosting where no `queue:work` daemon can run. Job API - `ShouldQueue` contract and `Queueable` trait; a job is a class with `handle()`, whose class-typed parameters are resolved from the container. - `Job::dispatch(...)` returns a pending dispatch written on destruction, with `delay()`, `on_queue()` and `with_priority()`; plus `dispatch_if`, `dispatch_unless` and `dispatch_sync`. - Inside `handle()`: `attempts()`, `release($delay)` and `fail($e)`. - `$tries`, `$backoff` (seconds or per-attempt array) and `failed(Throwable)`. Storage and processing - Jobs live in app-prefixed `{prefix}jobs` / `{prefix}failed_jobs` tables, so several framework plugins on one site never share a queue. - Workers claim batches atomically with a single tokened `UPDATE ... ORDER BY priority DESC ... LIMIT n`, work within a time budget, hand back unstarted jobs without using an attempt, and reclaim reservations older than `retry_after` so a crashed worker cannot strand a job. - An every-minute WP-Cron sweep and a shutdown spawn after undelayed dispatches start an HMAC-signed, non-blocking admin-ajax loopback worker that daisy-chains until the queue is empty. One chain runs at a time. Opt-in and tooling - Enabled only by registering `QueueServiceProvider`; without it nothing is hooked and dispatch throws a helpful `QueueException`. - `Queue` facade (`push`, `later`, `size`, `clear`) and `Queue::fake()` with `assert_pushed`, `assert_pushed_times`, `assert_not_pushed`, `assert_nothing_pushed`. - `JobQueued`, `JobProcessing`, `JobProcessed`, `JobFailed` events; failures are also logged. - WP-CLI: `queue:table` (generates migrations, runs no DDL), `make:job`, `queue:work`, `queue:failed`, `queue:retry`, `queue:forget`, `queue:flush`, `queue:clear`. - `docs/queues.md`, including the blocked-loopback / low-traffic fallback and where this differs from Laravel; optional `config/queue.php` in the example. Also records the grilling preference in CLAUDE.md. New public API is tagged @since 3.2.0. Co-Authored-By: Claude Opus 5.5 --- CLAUDE.md | 6 + docs/queues.md | 361 ++++++++++++ example/config/queue.php | 59 ++ src/Application.php | 6 + src/Console/Commands/MakeJobCommand.php | 89 +++ src/Console/Commands/QueueClearCommand.php | 57 ++ src/Console/Commands/QueueFailedCommand.php | 71 +++ src/Console/Commands/QueueFlushCommand.php | 49 ++ src/Console/Commands/QueueForgetCommand.php | 74 +++ src/Console/Commands/QueueRetryCommand.php | 80 +++ src/Console/Commands/QueueTableCommand.php | 115 ++++ src/Console/Commands/QueueWorkCommand.php | 80 +++ src/Console/stubs/job.stub | 30 + .../stubs/queue-failed-jobs-table.stub | 29 + src/Console/stubs/queue-jobs-table.stub | 33 ++ src/Contracts/ShouldQueue.php | 18 + src/Exceptions/QueueException.php | 82 +++ src/Queue/Concerns/Queueable.php | 309 ++++++++++ src/Queue/DatabaseQueue.php | 541 ++++++++++++++++++ src/Queue/Events/JobFailed.php | 64 +++ src/Queue/Events/JobProcessed.php | 52 ++ src/Queue/Events/JobProcessing.php | 52 ++ src/Queue/Events/JobQueued.php | 51 ++ src/Queue/JobRecord.php | 255 +++++++++ src/Queue/Payload.php | 127 ++++ src/Queue/PendingDispatch.php | 125 ++++ src/Queue/QueueFake.php | 196 +++++++ src/Queue/QueueManager.php | 322 +++++++++++ src/Queue/QueueServiceProvider.php | 141 +++++ src/Queue/Spawner.php | 242 ++++++++ src/Queue/Sweeper.php | 66 +++ src/Queue/Worker.php | 411 +++++++++++++ src/Supports/Facades/Queue.php | 47 ++ tests/Support/Queue/ArrayDatabaseQueue.php | 222 +++++++ tests/Support/Queue/Jobs/AttemptsJob.php | 23 + tests/Support/Queue/Jobs/DefaultsJob.php | 21 + tests/Support/Queue/Jobs/FailedThrowsJob.php | 22 + tests/Support/Queue/Jobs/InjectedJob.php | 16 + tests/Support/Queue/Jobs/Journal.php | 28 + tests/Support/Queue/Jobs/Mailer.php | 8 + tests/Support/Queue/Jobs/NotAJob.php | 11 + tests/Support/Queue/Jobs/ReleasingJob.php | 21 + tests/Support/Queue/Jobs/SelfFailingJob.php | 25 + tests/Support/Queue/Jobs/SendEmail.php | 31 + tests/Support/Queue/Jobs/ThrowingJob.php | 29 + tests/Support/Queue/RecordingLogger.php | 16 + tests/Support/Queue/RecordingSpawner.php | 36 ++ tests/Support/Queue/TestLock.php | 44 ++ tests/Support/Queue/TestWorker.php | 39 ++ tests/Support/StubsWordPressFunctions.php | 63 ++ tests/Unit/Queue/DatabaseQueueSqlTest.php | 134 +++++ tests/Unit/Queue/DispatchTest.php | 192 +++++++ tests/Unit/Queue/OptInTest.php | 53 ++ tests/Unit/Queue/PayloadTest.php | 98 ++++ tests/Unit/Queue/QueueCommandsTest.php | 221 +++++++ tests/Unit/Queue/QueueFakeTest.php | 88 +++ tests/Unit/Queue/QueueTestCase.php | 93 +++ tests/Unit/Queue/SpawnerTest.php | 168 ++++++ tests/Unit/Queue/StaleRecoveryTest.php | 73 +++ tests/Unit/Queue/WorkerTest.php | 224 ++++++++ 60 files changed, 6239 insertions(+) create mode 100644 docs/queues.md create mode 100644 example/config/queue.php create mode 100644 src/Console/Commands/MakeJobCommand.php create mode 100644 src/Console/Commands/QueueClearCommand.php create mode 100644 src/Console/Commands/QueueFailedCommand.php create mode 100644 src/Console/Commands/QueueFlushCommand.php create mode 100644 src/Console/Commands/QueueForgetCommand.php create mode 100644 src/Console/Commands/QueueRetryCommand.php create mode 100644 src/Console/Commands/QueueTableCommand.php create mode 100644 src/Console/Commands/QueueWorkCommand.php create mode 100644 src/Console/stubs/job.stub create mode 100644 src/Console/stubs/queue-failed-jobs-table.stub create mode 100644 src/Console/stubs/queue-jobs-table.stub create mode 100644 src/Contracts/ShouldQueue.php create mode 100644 src/Exceptions/QueueException.php create mode 100644 src/Queue/Concerns/Queueable.php create mode 100644 src/Queue/DatabaseQueue.php create mode 100644 src/Queue/Events/JobFailed.php create mode 100644 src/Queue/Events/JobProcessed.php create mode 100644 src/Queue/Events/JobProcessing.php create mode 100644 src/Queue/Events/JobQueued.php create mode 100644 src/Queue/JobRecord.php create mode 100644 src/Queue/Payload.php create mode 100644 src/Queue/PendingDispatch.php create mode 100644 src/Queue/QueueFake.php create mode 100644 src/Queue/QueueManager.php create mode 100644 src/Queue/QueueServiceProvider.php create mode 100644 src/Queue/Spawner.php create mode 100644 src/Queue/Sweeper.php create mode 100644 src/Queue/Worker.php create mode 100644 src/Supports/Facades/Queue.php create mode 100644 tests/Support/Queue/ArrayDatabaseQueue.php create mode 100644 tests/Support/Queue/Jobs/AttemptsJob.php create mode 100644 tests/Support/Queue/Jobs/DefaultsJob.php create mode 100644 tests/Support/Queue/Jobs/FailedThrowsJob.php create mode 100644 tests/Support/Queue/Jobs/InjectedJob.php create mode 100644 tests/Support/Queue/Jobs/Journal.php create mode 100644 tests/Support/Queue/Jobs/Mailer.php create mode 100644 tests/Support/Queue/Jobs/NotAJob.php create mode 100644 tests/Support/Queue/Jobs/ReleasingJob.php create mode 100644 tests/Support/Queue/Jobs/SelfFailingJob.php create mode 100644 tests/Support/Queue/Jobs/SendEmail.php create mode 100644 tests/Support/Queue/Jobs/ThrowingJob.php create mode 100644 tests/Support/Queue/RecordingLogger.php create mode 100644 tests/Support/Queue/RecordingSpawner.php create mode 100644 tests/Support/Queue/TestLock.php create mode 100644 tests/Support/Queue/TestWorker.php create mode 100644 tests/Unit/Queue/DatabaseQueueSqlTest.php create mode 100644 tests/Unit/Queue/DispatchTest.php create mode 100644 tests/Unit/Queue/OptInTest.php create mode 100644 tests/Unit/Queue/PayloadTest.php create mode 100644 tests/Unit/Queue/QueueCommandsTest.php create mode 100644 tests/Unit/Queue/QueueFakeTest.php create mode 100644 tests/Unit/Queue/QueueTestCase.php create mode 100644 tests/Unit/Queue/SpawnerTest.php create mode 100644 tests/Unit/Queue/StaleRecoveryTest.php create mode 100644 tests/Unit/Queue/WorkerTest.php diff --git a/CLAUDE.md b/CLAUDE.md index b0c2109..9c148d4 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -415,3 +415,9 @@ behaviour is worse than no doc. Don't commit or push unless I ask. When I do ask commit the changes with a inferred commit message that is good enough for PR title and description and also push on behalf of me. + + +## 8. Grilling behavior + +When using the /grill-me skill ask me questions one by one and use the graphical interface +so that I can select my answer graphically. Always mention your recommendation while questioning. \ No newline at end of file diff --git a/docs/queues.md b/docs/queues.md new file mode 100644 index 0000000..47f9b4c --- /dev/null +++ b/docs/queues.md @@ -0,0 +1,361 @@ +# Queues + +This guide covers deferring work to the background: publishing a product at a future date, or sending 500 emails without making a visitor wait. The API follows Laravel's queued jobs (`ShouldQueue`, `Queueable`, `Job::dispatch()`), adapted to the framework's snake_case naming and to the fact that WordPress has no `queue:work` daemon. Jobs are stored in a database table, and WP-Cron plus self-chaining background requests run them. + +## Table of contents + +1. [Quick start](#1-quick-start) +2. [Writing jobs](#2-writing-jobs) +3. [Dispatching jobs](#3-dispatching-jobs) +4. [How jobs run on WordPress](#4-how-jobs-run-on-wordpress) +5. [Failures and retries](#5-failures-and-retries) +6. [Configuration](#6-configuration) +7. [Loopback blocked or low traffic](#7-loopback-blocked-or-low-traffic) +8. [Events](#8-events) +9. [Testing](#9-testing) +10. [Where this differs from Laravel](#10-where-this-differs-from-laravel) +11. [CLI reference](#11-cli-reference) + +--- + +## 1. Quick start + +The queue is **off** until your plugin opts in. A plugin that never queues anything pays nothing for it: no cron event, no AJAX endpoint, no query. + +**1. Enable it** by adding the provider to `bootstrap/providers.php`: + +```php +return [ + AppServiceProvider::class, + Framework\Queue\QueueServiceProvider::class, +]; +``` + +**2. Create the tables.** `queue:table` writes two migration classes. It does not run any SQL. + +```bash +wp kirki queue:table +``` + +It prints the two lines to add to `config/migrations.php`. Add them, then run: + +```bash +wp kirki migrate +``` + +**3. Write a job:** + +```bash +wp kirki make:job PublishScheduledProduct +``` + +```php +namespace Acme\Shop\Jobs; + +use Framework\Contracts\ShouldQueue; +use Framework\Queue\Concerns\Queueable; + +class PublishScheduledProduct implements ShouldQueue +{ + use Queueable; + + public $product_id; + + public function __construct(int $product_id) + { + $this->product_id = $product_id; + } + + public function handle() + { + wp_update_post(['ID' => $this->product_id, 'post_status' => 'publish']); + } +} +``` + +**4. Dispatch it:** + +```php +PublishScheduledProduct::dispatch(123)->delay($publish_at); // a DateTimeInterface +``` + +--- + +## 2. Writing jobs + +A job is a class that implements `Framework\Contracts\ShouldQueue`, uses the `Framework\Queue\Concerns\Queueable` trait, and has a `handle()` method. + +**Pass IDs and scalars, not objects.** When a job is dispatched, the job object is serialized into the queue table. It is restored when it runs, which may be minutes or days later. A `WP_Post` or a model on a property is frozen as it was at dispatch time. Store the ID and load it fresh in `handle()`. + +**`handle()` can ask for services.** Class-typed parameters are resolved from the container: + +```php +public function handle(Mailer $mailer) +{ + $mailer->send($this->cart_id); +} +``` + +**Job settings** are optional properties on the class: + +| Property | Meaning | Default | +|---|---|---| +| `$tries` | How many times the job may be attempted | `config('queue.tries')`, 1 | +| `$backoff` | Seconds to wait before a retry, or an array of seconds per attempt | `config('queue.backoff')`, 0 | + +```php +class SendAbandonedCartEmail implements ShouldQueue +{ + use Queueable; + + protected $tries = 3; + protected $backoff = [60, 300]; // 1 minute, then 5 minutes (and 5 minutes after that) +} +``` + +**A job can set its own dispatch defaults** in its constructor, and a dispatch can still override them: + +```php +public function __construct(int $cart_id) +{ + $this->cart_id = $cart_id; + $this->on_queue('emails')->with_priority(5); +} +``` + +**Inside `handle()`**, a job can inspect and control itself: + +| Method | Effect | +|---|---| +| `$this->attempts()` | The attempt this run is, counting from 1 | +| `$this->release($delay)` | Put the job back on the queue after `$delay` seconds. `failed()` is not called. The attempt still counts toward `$tries` | +| `$this->fail($exception)` | Send the job straight to the failed jobs table, regardless of the tries it has left | + +```php +public function handle() +{ + if (!$this->api->is_available()) { + $this->release(120); + + return; + } + + // ... +} +``` + +**Jobs must be idempotent.** A job can run more than once: after a retry, or if it outlives `retry_after` (see [section 6](#6-configuration)). Running it a second time must be harmless. + +--- + +## 3. Dispatching jobs + +```php +SendAbandonedCartEmail::dispatch($cart_id); // as soon as possible +SendAbandonedCartEmail::dispatch($cart_id)->delay(3600); // in an hour +SendAbandonedCartEmail::dispatch($cart_id)->delay(new DateTime('+1 day')); +SendAbandonedCartEmail::dispatch($cart_id)->delay(new DateInterval('PT30M')); +SendAbandonedCartEmail::dispatch($cart_id)->on_queue('emails')->with_priority(10); + +SendAbandonedCartEmail::dispatch_if($cart->is_abandoned(), $cart_id); +SendAbandonedCartEmail::dispatch_unless($cart->is_recovered(), $cart_id); + +SendAbandonedCartEmail::dispatch_sync($cart_id); // run now, in this request +``` + +`dispatch()` returns a pending dispatch. The row is written when that object is destroyed, which is what allows the chained modifiers. If the queue is not enabled or its table is missing, `dispatch()` throws a `Framework\Exceptions\QueueException` that says what to do. + +**Priority** is a number, and higher runs first. Among equal priorities, the job that became available first runs first, then the one dispatched first. + +**Queue names** group jobs. The background worker drains *every* queue in priority order. Names let you filter jobs with `queue:work --queue=` and `queue:clear --queue=`, and you can see them in `queue:failed`. + +The `Queue` facade covers job instances you have already built, and inspecting the queue: + +```php +use Framework\Supports\Facades\Queue; + +Queue::push(new SendAbandonedCartEmail($id)); +Queue::later(600, new SendAbandonedCartEmail($id)); +Queue::size(); // pending, all queues +Queue::size('emails'); +Queue::clear('emails'); // delete pending jobs +``` + +`dispatch_sync()` runs the job immediately without storing it. If the job throws, its `failed()` method is called and the exception is rethrown. `dispatch_sync()` works without the provider. + +--- + +## 4. How jobs run on WordPress + +There is no worker process waiting for jobs. Three things start one: + +1. **An every-minute WP-Cron event** checks whether any job is due. If none is, it stops there, so a quiet queue costs one indexed query per cron tick. If a job is due, it sends a **non-blocking loopback request** to `admin-ajax.php` and returns immediately. The visitor whose request triggered WP-Cron does not wait. +2. **Dispatching an undelayed job** registers one `shutdown` callback for that request, which sends the same loopback request. "Dispatch now" does not wait for the next cron tick. This happens at most once per request, however many jobs were dispatched. +3. **A worker that runs out of time with jobs still due** sends the loopback request itself. This is the daisy-chain, and it keeps going without any traffic until the queue is empty. + +The worker that receives the request works like this: + +- It checks the request's signature. The signature is an HMAC keyed with your site's salts and valid for 60 seconds, so strangers cannot start workers. +- It takes the **chain lock**. Only one worker chain runs at a time. A second spawn finds the lock held and does nothing. +- It **claims** a batch of due jobs with a single atomic `UPDATE`, so no job is ever reserved by two workers at once. +- It runs jobs, and claims more batches, until its **time budget** runs out. Jobs it claimed but did not start are handed back without using up an attempt. +- It releases the lock. If jobs are still due, it spawns its successor. + +So 500 abandoned-cart emails are processed in a relay of short requests, each bounded by the time budget, rather than in one request that times out. + +--- + +## 5. Failures and retries + +When `handle()` throws: + +- If the job has tries left, it goes back on the queue after its backoff delay. +- If not, it is moved to the **failed jobs table**, its `failed()` method is called, `JobFailed` is dispatched, and an error is written to the framework log. + +```php +public function failed(Throwable $exception) +{ + // Notify someone, undo partial work, ... +} +``` + +An exception thrown by `failed()` itself is ignored, and one failing job never stops the others in the batch. + +**Crashed workers.** If PHP dies in the middle of a job (a fatal error, or a host that kills the process), the job stays reserved. After `retry_after` seconds the reservation counts as abandoned and the job can be claimed again. Every claim uses up an attempt, so a job that keeps crashing its worker ends up in the failed jobs table instead of looping forever. + +Inspect and recover failures with the CLI (see [section 11](#11-cli-reference)): + +```bash +wp kirki queue:failed +wp kirki queue:retry 5 +wp kirki queue:retry all +``` + +--- + +## 6. Configuration + +Every key is optional. `config/queue.php`: + +```php +return [ + 'table' => 'kirki_jobs', // default: {app prefix}jobs + 'failed_table' => 'kirki_failed_jobs', // default: {app prefix}failed_jobs + 'batch_size' => 10, // jobs per claim + 'time_limit' => 20, // seconds a web worker spends starting jobs + 'retry_after' => 300, // seconds before a reservation counts as abandoned + 'tries' => 1, // default $tries + 'backoff' => 0, // default $backoff +]; +``` + +Table names are given without the WordPress table prefix. They default to your app prefix, so two plugins built on the framework on the same site never share a queue. + +The **time budget** of a web worker is the smaller of `time_limit` and 80% of `max_execution_time`. A job that is already running is never cut short, so keep each job well under the budget. + +**`retry_after` must be longer than your slowest job.** If a job takes longer than that, another worker can claim it while it is still running. + +--- + +## 7. Loopback blocked or low traffic + +The background worker relies on the site being able to send HTTP requests to itself. This fails on: + +- staging sites behind HTTP Basic Auth, +- some firewalls and security plugins, +- hosts that block loopback requests. + +When that happens, jobs are stored but never run. WordPress's own Site Health screen reports "loopback request failed" in this situation. + +WP-Cron also only runs when someone visits. On a quiet site, a delayed job runs at the first visit after its time, not at its time. + +The fix for both is a real cron job that runs the queue from WP-CLI: + +``` +# wp-config.php +define('DISABLE_WP_CRON', true); + +# crontab: every minute +* * * * * cd /path/to/site && wp cron event run --due-now >/dev/null 2>&1 +* * * * * cd /path/to/site && wp kirki queue:work >/dev/null 2>&1 +``` + +`queue:work` runs jobs directly in the CLI process. It needs no loopback, has no time budget, and stops when the queue is empty. + +--- + +## 8. Events + +Dispatched through the framework event system, and only when something listens for them: + +| Event | When | Carries | +|---|---|---| +| `Framework\Queue\Events\JobQueued` | After a job is stored | `$id`, `$job` | +| `Framework\Queue\Events\JobProcessing` | Before `handle()` | `$record`, `$job` | +| `Framework\Queue\Events\JobProcessed` | After a successful `handle()` | `$record`, `$job` | +| `Framework\Queue\Events\JobFailed` | When a job is moved to failed jobs | `$record`, `$job` (null if it could not be restored), `$exception` | + +Failures are also logged whether or not anything listens. + +--- + +## 9. Testing + +`Queue::fake()` replaces the queue with an in-memory recorder. Nothing is stored and no worker is spawned. It does not need the provider or the tables. + +```php +use Framework\Supports\Facades\Queue; + +Queue::fake(); + +$this->checkout->abandon($cart); + +Queue::assert_pushed(SendAbandonedCartEmail::class); +Queue::assert_pushed(SendAbandonedCartEmail::class, function ($job) use ($cart) { + return $job->cart_id === $cart->id; +}); +Queue::assert_pushed_times(SendAbandonedCartEmail::class, 1); +Queue::assert_not_pushed(PublishScheduledProduct::class); +Queue::assert_nothing_pushed(); +``` + +A failed assertion throws PHP's `AssertionError`, which PHPUnit reports as a failure. Successful assertions are not counted by PHPUnit, so a test that only makes queue assertions should call `$this->addToAssertionCount()` to avoid being marked risky. + +To test a job's own logic, call `handle()` directly or use `dispatch_sync()`. + +--- + +## 10. Where this differs from Laravel + +**There is no daemon.** Laravel's `queue:work` is a long-running process. Here, work is driven by WP-Cron and loopback requests, and `queue:work` is a one-shot drain for system cron. Latency depends on traffic or on your cron (see [section 7](#7-loopback-blocked-or-low-traffic)). + +**One worker chain at a time.** Laravel scales by running more workers. Here, a single chain keeps shared hosting from being overwhelmed, and throughput is bounded accordingly. + +**The database driver only.** There is no Redis, SQS, or `sync` connection. `dispatch_sync()` covers the synchronous case. + +**No per-job `$timeout`.** PHP-FPM cannot interrupt a running `handle()` the way Laravel does with `pcntl`. Stale reservations are recovered with the global `retry_after` instead. + +**Priority is a number on the job.** Laravel orders work by the list of queue names a worker is given. Here, the background worker drains every queue by numeric priority, and queue names are for grouping and filtering. + +**No `SerializesModels`.** Models are serialized as they are, not re-fetched when the job runs. Pass IDs. + +**Not implemented:** `ShouldBeUnique`, job chains (`Bus::chain`), job batches (`Bus::batch`), job middleware, rate-limited jobs, `dispatch_after_response`, and failed-job pruning (`queue:prune-failed`). + +**Failed jobs are identified by numeric ID** in `queue:retry` and `queue:forget`, not by UUID. + +--- + +## 11. CLI reference + +| Command | What it does | +|---|---| +| `queue:table` | Generate the `CreateJobsTable` and `CreateFailedJobsTable` migrations. It never overwrites files and runs no SQL. Available without the provider | +| `make:job ` | Generate a job class in `app/Jobs/`. Available without the provider | +| `queue:work [--once] [--max-jobs=] [--queue=]` | Process due jobs in the CLI process | +| `queue:failed` | List failed jobs | +| `queue:retry ` | Put failed jobs back on the queue with attempts reset | +| `queue:forget ` | Delete one failed job | +| `queue:flush` | Delete all failed jobs | +| `queue:clear [--queue=]` | Delete pending jobs | + +Every command other than `queue:table` and `make:job` is available only when `QueueServiceProvider` is registered. diff --git a/example/config/queue.php b/example/config/queue.php new file mode 100644 index 0000000..3f9c91e --- /dev/null +++ b/example/config/queue.php @@ -0,0 +1,59 @@ + FRAMEWORK_EXAMPLE_PREFIX . 'jobs', + + /* + * The table jobs are moved to once they have used up their tries. + */ + 'failed_table' => FRAMEWORK_EXAMPLE_PREFIX . 'failed_jobs', + + /* + * How many jobs one claim reserves. A worker keeps claiming batches while + * its time budget lasts, so this bounds a single claim, not a request. + */ + 'batch_size' => 10, + + /* + * The most seconds a background worker spends starting jobs before it + * hands over to a fresh request. The effective budget is the smaller of + * this and 80% of max_execution_time. A job already running is never cut + * short; keep individual jobs well under this. + */ + 'time_limit' => 20, + + /* + * Seconds after which a reserved job whose worker vanished (a fatal error, + * a killed process) is considered abandoned and may be claimed again. + * Must be longer than your slowest job, or that job can run twice. + */ + 'retry_after' => 300, + + /* + * Default number of attempts for a job that does not declare $tries. + */ + 'tries' => 1, + + /* + * Default seconds to wait before retrying a failed attempt, for a job that + * does not declare $backoff. An array gives one delay per attempt, with the + * last value reused: [10, 60, 300]. + */ + 'backoff' => 0, +]; diff --git a/src/Application.php b/src/Application.php index cbc3181..0f260c5 100644 --- a/src/Application.php +++ b/src/Application.php @@ -24,6 +24,7 @@ use Framework\Cache\CacheServiceProvider; use Framework\RateLimiting\RateLimiter; use Framework\RateLimiting\RateLimiterServiceProvider; +use Framework\Queue\QueueManager; use Framework\Console\Commands\ClearCacheCommand; use Framework\Console\Commands\ForgetCacheCommand; use Framework\Console\Commands\GcCacheCommand; @@ -32,6 +33,8 @@ use Framework\Console\Commands\MakeProviderCommand; use Framework\Console\Commands\MakeRequestCommand; use Framework\Console\Commands\MakeSeederCommand; +use Framework\Console\Commands\MakeJobCommand; +use Framework\Console\Commands\QueueTableCommand; use Framework\Console\Commands\MigrateCommand; use Framework\Console\Commands\RollbackCommand; use Framework\Console\Commands\SeedCommand; @@ -324,6 +327,8 @@ protected function register_base_cli_commands() 'cache:clear' => ClearCacheCommand::class, 'cache:forget' => ForgetCacheCommand::class, 'cache:gc' => GcCacheCommand::class, + 'queue:table' => QueueTableCommand::class, + 'make:job' => MakeJobCommand::class, ]; foreach ($commands as $command => $class) { @@ -343,6 +348,7 @@ protected function register_base_aliases() foreach ( [ 'cache' => CacheManager::class, + 'queue' => QueueManager::class, 'limiter' => RateLimiter::class, 'db' => DatabaseManager::class, 'schema' => SchemaManager::class, diff --git a/src/Console/Commands/MakeJobCommand.php b/src/Console/Commands/MakeJobCommand.php new file mode 100644 index 0000000..c591198 --- /dev/null +++ b/src/Console/Commands/MakeJobCommand.php @@ -0,0 +1,89 @@ +summary('Create a new queued job class') + ->description("## EXAMPLES \n\n wp kirki make:job SendAbandonedCartEmail") + ->synopsis( + Synopsis::type('positional') + ->name('name') + ->description('The job class name') + ); + } + + /** + * Check if the command passed the validation. + * + * @param mixed $args The positional arguments. + * @param mixed $assoc The associative arguments. + * + * @return bool + * + * @since 3.2.0 + */ + protected function passed($args, $assoc) + { + return !empty($args[0]); + } + + /** + * Run the command. + * + * @param mixed $args The positional arguments. + * @param mixed $assoc The associative arguments. + * + * @return void + * + * @since 3.2.0 + */ + public function run($args, $assoc) + { + $class = Str::pascal((string) $args[0]); + $namespace = app()->qualify_app_namespace('Jobs'); + $output_file = app_path('Jobs/' . $class . '.php'); + + if (File::exists($output_file)) { + $this->cli_error(sprintf('Job file already exists: %s', $output_file)); + + return; + } + + File::make_dir($output_file); + File::put($output_file, Str::replace( + ['{{namespace}}', '{{class_name}}'], + [$namespace, $class], + File::get($this->stub_path() . '/job.stub') + )); + + $this->cli_success(sprintf('Job [%s] created.', $namespace . '\\' . $class)); + } +} diff --git a/src/Console/Commands/QueueClearCommand.php b/src/Console/Commands/QueueClearCommand.php new file mode 100644 index 0000000..0f1a4d8 --- /dev/null +++ b/src/Console/Commands/QueueClearCommand.php @@ -0,0 +1,57 @@ +summary('Delete the pending jobs on the queue') + ->description("## EXAMPLES \n\n wp kirki queue:clear\n\n wp kirki queue:clear --queue=emails") + ->synopsis( + Synopsis::type('assoc') + ->name('queue') + ->description('Only clear this queue') + ->optional() + ); + } + + /** + * Run the command. + * + * @param mixed $args The positional arguments. + * @param mixed $assoc The associative arguments. + * + * @return void + * + * @since 3.2.0 + */ + public function run($args, $assoc) + { + $deleted = app(DatabaseQueue::class)->clear($assoc['queue'] ?? null); + + $this->cli_success(sprintf('Deleted %d pending %s.', $deleted, $deleted === 1 ? 'job' : 'jobs')); + } +} diff --git a/src/Console/Commands/QueueFailedCommand.php b/src/Console/Commands/QueueFailedCommand.php new file mode 100644 index 0000000..2fee242 --- /dev/null +++ b/src/Console/Commands/QueueFailedCommand.php @@ -0,0 +1,71 @@ +summary('List the failed queue jobs') + ->description("## EXAMPLES \n\n wp kirki queue:failed"); + } + + /** + * Run the command. + * + * @param mixed $args The positional arguments. + * @param mixed $assoc The associative arguments. + * + * @return void + * + * @since 3.2.0 + */ + public function run($args, $assoc) + { + $rows = app(DatabaseQueue::class)->failed_all(); + + if (empty($rows)) { + $this->cli_line('No failed jobs.'); + + return; + } + + call_user_func( + '\\WP_CLI\\Utils\\format_items', + 'table', + array_map(function (array $row) { + $exception = strtok((string) $row['exception'], "\n"); + + return [ + 'ID' => $row['id'], + 'Queue' => $row['queue'], + 'Job' => Payload::display_name((string) $row['payload']), + 'Failed At' => gmdate('Y-m-d H:i:s', (int) $row['failed_at']), + 'Exception' => $exception === false ? '' : $exception, + ]; + }, $rows), + ['ID', 'Queue', 'Job', 'Failed At', 'Exception'] + ); + } +} diff --git a/src/Console/Commands/QueueFlushCommand.php b/src/Console/Commands/QueueFlushCommand.php new file mode 100644 index 0000000..06cd146 --- /dev/null +++ b/src/Console/Commands/QueueFlushCommand.php @@ -0,0 +1,49 @@ +summary('Delete all failed queue jobs') + ->description("## EXAMPLES \n\n wp kirki queue:flush"); + } + + /** + * Run the command. + * + * @param mixed $args The positional arguments. + * @param mixed $assoc The associative arguments. + * + * @return void + * + * @since 3.2.0 + */ + public function run($args, $assoc) + { + $deleted = app(DatabaseQueue::class)->flush_failed(); + + $this->cli_success(sprintf('Deleted %d failed %s.', $deleted, $deleted === 1 ? 'job' : 'jobs')); + } +} diff --git a/src/Console/Commands/QueueForgetCommand.php b/src/Console/Commands/QueueForgetCommand.php new file mode 100644 index 0000000..e6f5466 --- /dev/null +++ b/src/Console/Commands/QueueForgetCommand.php @@ -0,0 +1,74 @@ +summary('Delete a failed queue job') + ->description("## EXAMPLES \n\n wp kirki queue:forget 5") + ->synopsis( + Synopsis::type('positional') + ->name('id') + ->description('The failed job ID') + ); + } + + /** + * Check if the command passed the validation. + * + * @param mixed $args The positional arguments. + * @param mixed $assoc The associative arguments. + * + * @return bool + * + * @since 3.2.0 + */ + protected function passed($args, $assoc) + { + return !empty($args[0]) && ctype_digit((string) $args[0]); + } + + /** + * Run the command. + * + * @param mixed $args The positional arguments. + * @param mixed $assoc The associative arguments. + * + * @return void + * + * @since 3.2.0 + */ + public function run($args, $assoc) + { + if (!app(DatabaseQueue::class)->forget((int) $args[0])) { + $this->cli_error(sprintf('No failed job with ID [%s].', $args[0])); + + return; + } + + $this->cli_success(sprintf('Failed job [%s] deleted.', $args[0])); + } +} diff --git a/src/Console/Commands/QueueRetryCommand.php b/src/Console/Commands/QueueRetryCommand.php new file mode 100644 index 0000000..6ef5a53 --- /dev/null +++ b/src/Console/Commands/QueueRetryCommand.php @@ -0,0 +1,80 @@ +summary('Retry failed queue jobs') + ->description("## EXAMPLES \n\n wp kirki queue:retry 5\n\n wp kirki queue:retry all") + ->synopsis( + Synopsis::type('positional') + ->name('id') + ->description('The failed job ID, or "all"') + ); + } + + /** + * Check if the command passed the validation. + * + * @param mixed $args The positional arguments. + * @param mixed $assoc The associative arguments. + * + * @return bool + * + * @since 3.2.0 + */ + protected function passed($args, $assoc) + { + return !empty($args[0]) && ($args[0] === 'all' || ctype_digit((string) $args[0])); + } + + /** + * Run the command. + * + * @param mixed $args The positional arguments. + * @param mixed $assoc The associative arguments. + * + * @return void + * + * @since 3.2.0 + */ + public function run($args, $assoc) + { + $retried = app(DatabaseQueue::class)->retry($args[0] === 'all' ? 'all' : (int) $args[0]); + + if ($retried === 0) { + $this->cli_warning('No matching failed jobs.'); + + return; + } + + $this->cli_success(sprintf( + 'Pushed %d failed %s back onto the queue.', + $retried, + $retried === 1 ? 'job' : 'jobs' + )); + } +} diff --git a/src/Console/Commands/QueueTableCommand.php b/src/Console/Commands/QueueTableCommand.php new file mode 100644 index 0000000..8304fdd --- /dev/null +++ b/src/Console/Commands/QueueTableCommand.php @@ -0,0 +1,115 @@ + ['queue-jobs-table.stub', 'get_table_name'], + 'CreateFailedJobsTable' => ['queue-failed-jobs-table.stub', 'get_failed_table_name'], + ]; + + /** + * Prepare the command's synopsis and other metadata. + * + * @return void + * + * @since 3.2.0 + */ + protected function prepare() + { + $this->summary('Create the migrations for the queue tables') + ->description("## EXAMPLES \n\n wp kirki queue:table"); + } + + /** + * Run the command. + * + * @param mixed $args The positional arguments. + * @param mixed $assoc The associative arguments. + * + * @return void + * + * @since 3.2.0 + */ + public function run($args, $assoc) + { + $directory = database_path('migrations'); + $existing = []; + + foreach (array_keys($this->migrations) as $class) { + if (File::exists($this->output_file($directory, $class))) { + $existing[] = $this->output_file($directory, $class); + } + } + + if (!empty($existing)) { + $this->cli_error(sprintf('Migration files already exist: %s', implode(', ', $existing))); + + return; + } + + $queue = app(DatabaseQueue::class); + $namespace = app()->get_migrations_namespace(); + + foreach ($this->migrations as $class => [$stub, $table_getter]) { + $output_file = $this->output_file($directory, $class); + + File::make_dir($output_file); + File::put($output_file, Str::replace( + ['{{migrations_namespace}}', '{{class_name}}', '{{table}}'], + [$namespace, $class, $queue->$table_getter()], + File::get($this->stub_path() . '/' . $stub) + )); + + $this->cli_success(sprintf('Migration [%s] created.', $class)); + } + + $this->cli_line('Add these to config/migrations.php, then run `migrate`:'); + + foreach (array_keys($this->migrations) as $class) { + $this->cli_line(sprintf(' \\%s\\%s::class,', $namespace, $class)); + } + } + + /** + * Get the path a migration class is written to. + * + * @param string $directory The migrations directory. + * @param string $class The migration class name. + * + * @return string + * + * @since 3.2.0 + */ + protected function output_file(string $directory, string $class) + { + return sprintf('%s/%s.php', rtrim($directory, '/'), $class); + } +} diff --git a/src/Console/Commands/QueueWorkCommand.php b/src/Console/Commands/QueueWorkCommand.php new file mode 100644 index 0000000..2cd3cc8 --- /dev/null +++ b/src/Console/Commands/QueueWorkCommand.php @@ -0,0 +1,80 @@ +summary('Process the jobs that are due on the queue') + ->description("## EXAMPLES \n\n wp kirki queue:work\n\n wp kirki queue:work --once --queue=emails") + ->synopsis( + Synopsis::type('flag') + ->name('once') + ->description('Process a single job and stop') + ->optional() + ) + ->synopsis( + Synopsis::type('assoc') + ->name('max-jobs') + ->description('Stop after processing this many jobs') + ->optional() + ) + ->synopsis( + Synopsis::type('assoc') + ->name('queue') + ->description('Only process these queues, comma separated') + ->optional() + ); + } + + /** + * Run the command. + * + * @param mixed $args The positional arguments. + * @param mixed $assoc The associative arguments. + * + * @return void + * + * @since 3.2.0 + */ + public function run($args, $assoc) + { + $max_jobs = !empty($assoc['once']) ? 1 : ($assoc['max-jobs'] ?? null); + $queues = isset($assoc['queue']) + ? array_filter(array_map('trim', explode(',', (string) $assoc['queue']))) + : null; + + $remaining = app(Worker::class)->run([ + 'budget' => null, + 'max_jobs' => is_null($max_jobs) ? null : (int) $max_jobs, + 'queues' => empty($queues) ? null : array_values($queues), + ]); + + $this->cli_success($remaining ? 'Stopped with jobs still due.' : 'No jobs are due.'); + } +} diff --git a/src/Console/stubs/job.stub b/src/Console/stubs/job.stub new file mode 100644 index 0000000..819d82f --- /dev/null +++ b/src/Console/stubs/job.stub @@ -0,0 +1,30 @@ +id(); + $table->string('uuid', 36); + $table->string('queue', 191); + $table->long_text('payload'); + $table->long_text('exception'); + $table->unsigned_integer('failed_at'); + + $table->unique('uuid'); + }); + } + + public function down() + { + Schema::drop_if_exists('{{table}}'); + } +} diff --git a/src/Console/stubs/queue-jobs-table.stub b/src/Console/stubs/queue-jobs-table.stub new file mode 100644 index 0000000..cb54e73 --- /dev/null +++ b/src/Console/stubs/queue-jobs-table.stub @@ -0,0 +1,33 @@ +id(); + $table->string('queue', 191)->default('default'); + $table->integer('priority')->default(0); + $table->long_text('payload'); + $table->unsigned_tiny_integer('attempts')->default(0); + $table->unsigned_integer('reserved_at')->nullable(); + $table->string('reserved_by', 32)->nullable(); + $table->unsigned_integer('available_at'); + $table->unsigned_integer('created_at'); + + $table->index(['reserved_at', 'available_at', 'priority']); + $table->index('reserved_by'); + }); + } + + public function down() + { + Schema::drop_if_exists('{{table}}'); + } +} diff --git a/src/Contracts/ShouldQueue.php b/src/Contracts/ShouldQueue.php new file mode 100644 index 0000000..d1863c9 --- /dev/null +++ b/src/Contracts/ShouldQueue.php @@ -0,0 +1,18 @@ +ensure_ready(); + + return new PendingDispatch(new static(...$arguments), $manager); + } + + /** + * Dispatch the job when the condition is true. + * + * @param bool|\Closure $boolean The condition. + * @param mixed $arguments The job's constructor arguments. + * + * @return \Framework\Queue\PendingDispatch|null + * + * @since 3.2.0 + */ + public static function dispatch_if($boolean, ...$arguments) + { + $boolean = $boolean instanceof Closure ? $boolean() : $boolean; + + return $boolean ? static::dispatch(...$arguments) : null; + } + + /** + * Dispatch the job unless the condition is true. + * + * @param bool|\Closure $boolean The condition. + * @param mixed $arguments The job's constructor arguments. + * + * @return \Framework\Queue\PendingDispatch|null + * + * @since 3.2.0 + */ + public static function dispatch_unless($boolean, ...$arguments) + { + $boolean = $boolean instanceof Closure ? $boolean() : $boolean; + + return $boolean ? null : static::dispatch(...$arguments); + } + + /** + * Run the job immediately in the current request without storing it. + * + * @param mixed $arguments The job's constructor arguments. + * + * @return mixed The value handle() returned. + * + * @throws \Throwable Whatever handle() threw, after failed() has been called. + * + * @since 3.2.0 + */ + public static function dispatch_sync(...$arguments) + { + return app(Worker::class)->run_sync(new static(...$arguments)); + } + + /** + * Set the queue name the job is dispatched to. + * + * @param string|null $queue The queue name. + * + * @return $this + * + * @since 3.2.0 + */ + public function on_queue($queue) + { + $this->queue = $queue; + + return $this; + } + + /** + * Set how long to wait before the job becomes available. + * + * @param int|\DateTimeInterface|\DateInterval|null $delay Seconds, or the moment it becomes available. + * + * @return $this + * + * @since 3.2.0 + */ + public function delay($delay) + { + $this->delay = $delay; + + return $this; + } + + /** + * Set the job's priority; higher runs first. + * + * @param int $priority The priority. + * + * @return $this + * + * @since 3.2.0 + */ + public function with_priority(int $priority) + { + $this->priority = $priority; + + return $this; + } + + /** + * Get the queue name the job is dispatched to. + * + * @return string + * + * @since 3.2.0 + */ + public function get_queue() + { + return $this->queue ?? 'default'; + } + + /** + * Get the delay before the job becomes available. + * + * @return int|\DateTimeInterface|\DateInterval|null + * + * @since 3.2.0 + */ + public function get_delay() + { + return $this->delay; + } + + /** + * Get the job's priority. + * + * @return int + * + * @since 3.2.0 + */ + public function get_priority() + { + return (int) ($this->priority ?? 0); + } + + /** + * Get the number of times the job may be attempted, when the job class defines it. + * + * @return int|null + * + * @since 3.2.0 + */ + public function get_tries() + { + return property_exists($this, 'tries') ? $this->tries : null; + } + + /** + * Get the seconds to wait before retrying, when the job class defines it. + * + * @return int|array|null + * + * @since 3.2.0 + */ + public function get_backoff() + { + return property_exists($this, 'backoff') ? $this->backoff : null; + } + + /** + * Get the attempt this run is, counting from one. + * + * @return int + * + * @since 3.2.0 + */ + public function attempts() + { + return $this->job_record ? $this->job_record->attempts() : 1; + } + + /** + * Put the job back on the queue once handle() returns, without treating it as a failure. + * + * @param int|\DateTimeInterface|\DateInterval $delay How long before it is available again. + * + * @return void + * + * @since 3.2.0 + */ + public function release($delay = 0) + { + if ($this->job_record) { + $this->job_record->release($delay); + } + } + + /** + * Mark the job as failed once handle() returns, regardless of its remaining tries. + * + * @param \Throwable|string|null $exception Why the job failed. + * + * @return void + * + * @since 3.2.0 + */ + public function fail($exception = null) + { + if (!$exception instanceof Throwable) { + $exception = new RuntimeException($exception ?: sprintf('The job [%s] failed itself.', static::class)); + } + + if ($this->job_record) { + $this->job_record->fail($exception); + } + } + + /** + * Attach the record of the run in progress. + * + * @param \Framework\Queue\JobRecord|null $record The record, or null to detach it. + * + * @return $this + * + * @since 3.2.0 + */ + public function set_job_record(?JobRecord $record) + { + $this->job_record = $record; + + return $this; + } +} diff --git a/src/Queue/DatabaseQueue.php b/src/Queue/DatabaseQueue.php new file mode 100644 index 0000000..53acfdb --- /dev/null +++ b/src/Queue/DatabaseQueue.php @@ -0,0 +1,541 @@ +options)) { + return $this->options[$key]; + } + + return config('queue.' . $key, $default); + } + + /** + * Get the jobs table name without the WordPress table prefix. + * + * @return string + * + * @since 3.2.0 + */ + public function get_table_name() + { + return (string) $this->option('table', app()->prefix() . 'jobs'); + } + + /** + * Get the failed jobs table name without the WordPress table prefix. + * + * @return string + * + * @since 3.2.0 + */ + public function get_failed_table_name() + { + return (string) $this->option('failed_table', app()->prefix() . 'failed_jobs'); + } + + /** + * Get the full name of the jobs table. + * + * @return string + * + * @since 3.2.0 + */ + public function get_table() + { + return DB::get_table_prefix() . $this->get_table_name(); + } + + /** + * Get the full name of the failed jobs table. + * + * @return string + * + * @since 3.2.0 + */ + public function get_failed_table() + { + return DB::get_table_prefix() . $this->get_failed_table_name(); + } + + /** + * Get the current time as a UNIX timestamp. + * + * @return int + * + * @since 3.2.0 + */ + public function now() + { + return (int) $this->current_timestamp(); + } + + /** + * Determine whether a job with the given delay would be available immediately. + * + * @param int|\DateTimeInterface|\DateInterval|null $delay The delay. + * + * @return bool + * + * @since 3.2.0 + */ + public function is_immediate($delay) + { + return $this->available_at($delay) <= $this->now(); + } + + /** + * Determine whether the jobs table exists. + * + * Only a positive answer is remembered, so a table created mid-request is still found. + * + * @return bool + * + * @since 3.2.0 + */ + public function table_exists() + { + if ($this->table_confirmed) { + return true; + } + + $rows = DB::select( + 'SELECT 1 FROM information_schema.tables WHERE table_schema = DATABASE() AND table_name = %s LIMIT 1', + [$this->get_table()] + ); + + $this->table_confirmed = !empty($rows); + + return $this->table_confirmed; + } + + /** + * Store a job. + * + * @param string $payload The encoded payload. + * @param string $queue The queue name. + * @param int $priority The priority; higher runs first. + * @param int|\DateTimeInterface|\DateInterval|null $delay The delay before it is available. + * + * @return int The new row id. + * + * @since 3.2.0 + */ + public function push(string $payload, string $queue, int $priority = 0, $delay = null) + { + DB::insert( + "INSERT INTO {$this->get_table()} (queue, priority, payload, attempts, available_at, created_at) + VALUES (%s, %d, %s, 0, %d, %d)", + [$queue, $priority, $payload, $this->available_at($delay), $this->now()] + ); + + return (int) DB::get_db()->insert_id; + } + + /** + * Reserve up to the given number of due jobs for one worker. + * + * The attempt count is incremented as part of the claim, so a job that keeps killing its + * worker still uses up its tries. + * + * @param int $limit The most jobs to reserve. + * @param string $token The token that marks this claim's rows. + * @param array|null $queues Restrict the claim to these queue names. + * + * @return \Framework\Queue\JobRecord[] + * + * @since 3.2.0 + */ + public function claim(int $limit, string $token, ?array $queues = null) + { + $now = $this->now(); + $bindings = [$now, $token, $now, $this->stale_before()]; + + DB::update( + "UPDATE {$this->get_table()} + SET reserved_at = %d, reserved_by = %s, attempts = attempts + 1 + WHERE available_at <= %d + AND (reserved_at IS NULL OR reserved_at <= %d)" + . $this->queue_constraint($queues, $bindings) . " + ORDER BY priority DESC, available_at ASC, id ASC + LIMIT %d", + array_merge($bindings, [$limit]) + ); + + return $this->reserved($token); + } + + /** + * Get the rows reserved under a claim token. + * + * @param string $token The claim token. + * + * @return \Framework\Queue\JobRecord[] + * + * @since 3.2.0 + */ + public function reserved(string $token) + { + $rows = DB::select( + "SELECT id, queue, payload, attempts FROM {$this->get_table()} + WHERE reserved_by = %s + ORDER BY priority DESC, available_at ASC, id ASC", + [$token] + ); + + return array_map([JobRecord::class, 'from_row'], is_array($rows) ? $rows : []); + } + + /** + * Delete a finished job. + * + * @param int $id The row id. + * + * @return void + * + * @since 3.2.0 + */ + public function delete(int $id) + { + DB::delete("DELETE FROM {$this->get_table()} WHERE id = %d", [$id]); + } + + /** + * Put a job back on the queue after a delay; its attempt stays counted. + * + * @param int $id The row id. + * @param int|\DateTimeInterface|\DateInterval|null $delay The delay before it is available. + * + * @return void + * + * @since 3.2.0 + */ + public function release(int $id, $delay = 0) + { + DB::update( + "UPDATE {$this->get_table()} + SET reserved_at = NULL, reserved_by = NULL, available_at = %d + WHERE id = %d", + [$this->available_at($delay), $id] + ); + } + + /** + * Hand back claimed jobs the worker never started, refunding the attempt the claim took. + * + * @param int[] $ids The row ids. + * + * @return void + * + * @since 3.2.0 + */ + public function release_unstarted(array $ids) + { + $ids = array_values(array_map('intval', $ids)); + + if (empty($ids)) { + return; + } + + DB::update( + "UPDATE {$this->get_table()} + SET reserved_at = NULL, reserved_by = NULL, attempts = IF(attempts > 0, attempts - 1, 0) + WHERE id IN (" . implode(', ', array_fill(0, count($ids), '%d')) . ')', + $ids + ); + } + + /** + * Move a job to the failed jobs table. + * + * @param \Framework\Queue\JobRecord $record The claimed record. + * @param \Throwable $exception Why it failed. + * + * @return void + * + * @since 3.2.0 + */ + public function fail(JobRecord $record, Throwable $exception) + { + $envelope = json_decode($record->payload(), true); + + DB::insert( + "INSERT INTO {$this->get_failed_table()} (uuid, queue, payload, exception, failed_at) + VALUES (%s, %s, %s, %s, %d)", + [ + is_array($envelope) && !empty($envelope['uuid']) ? (string) $envelope['uuid'] : (string) uuid(), + $record->queue(), + $record->payload(), + (string) $exception, + $this->now(), + ] + ); + + $this->delete($record->id()); + } + + /** + * Determine whether any job is due, optionally within the given queues. + * + * @param array|null $queues Restrict the check to these queue names. + * + * @return bool + * + * @since 3.2.0 + */ + public function has_due(?array $queues = null) + { + $bindings = [$this->now(), $this->stale_before()]; + + $rows = DB::select( + "SELECT 1 FROM {$this->get_table()} + WHERE available_at <= %d + AND (reserved_at IS NULL OR reserved_at <= %d)" + . $this->queue_constraint($queues, $bindings) . ' + LIMIT 1', + $bindings + ); + + return !empty($rows); + } + + /** + * Count the jobs not currently being worked on, optionally on one queue. + * + * @param string|null $queue The queue name. + * + * @return int + * + * @since 3.2.0 + */ + public function size(?string $queue = null) + { + $bindings = [$this->stale_before()]; + + $rows = DB::select( + "SELECT COUNT(*) AS aggregate FROM {$this->get_table()} + WHERE (reserved_at IS NULL OR reserved_at <= %d)" + . $this->queue_constraint(is_null($queue) ? null : [$queue], $bindings), + $bindings + ); + + return (int) ($rows[0]['aggregate'] ?? 0); + } + + /** + * Delete the jobs not currently being worked on, optionally on one queue. + * + * @param string|null $queue The queue name. + * + * @return int The number of jobs deleted. + * + * @since 3.2.0 + */ + public function clear(?string $queue = null) + { + $bindings = [$this->stale_before()]; + + return (int) DB::delete( + "DELETE FROM {$this->get_table()} + WHERE (reserved_at IS NULL OR reserved_at <= %d)" + . $this->queue_constraint(is_null($queue) ? null : [$queue], $bindings), + $bindings + ); + } + + /** + * Get every failed job, newest first. + * + * @return array + * + * @since 3.2.0 + */ + public function failed_all() + { + $rows = DB::select( + "SELECT id, uuid, queue, payload, exception, failed_at FROM {$this->get_failed_table()} ORDER BY id DESC" + ); + + return is_array($rows) ? $rows : []; + } + + /** + * Find one failed job. + * + * @param int $id The failed job's id. + * + * @return array|null + * + * @since 3.2.0 + */ + public function failed_find(int $id) + { + $rows = DB::select( + "SELECT id, uuid, queue, payload, exception, failed_at FROM {$this->get_failed_table()} WHERE id = %d", + [$id] + ); + + return empty($rows) ? null : $rows[0]; + } + + /** + * Move failed jobs back onto the queue with their attempts reset. + * + * @param int|string $id A failed job's id, or "all". + * + * @return int The number of jobs moved back. + * + * @since 3.2.0 + */ + public function retry($id) + { + if ($id === 'all') { + $failed = $this->failed_all(); + } else { + $failed = array_filter([$this->failed_find((int) $id)]); + } + + foreach ($failed as $row) { + $envelope = json_decode((string) $row['payload'], true); + + $this->push( + (string) $row['payload'], + (string) $row['queue'], + is_array($envelope) ? (int) ($envelope['priority'] ?? 0) : 0 + ); + + $this->forget((int) $row['id']); + } + + return count($failed); + } + + /** + * Delete one failed job. + * + * @param int $id The failed job's id. + * + * @return bool Whether a row was deleted. + * + * @since 3.2.0 + */ + public function forget(int $id) + { + return (int) DB::delete("DELETE FROM {$this->get_failed_table()} WHERE id = %d", [$id]) > 0; + } + + /** + * Delete every failed job. + * + * @return int The number of rows deleted. + * + * @since 3.2.0 + */ + public function flush_failed() + { + return (int) DB::delete("DELETE FROM {$this->get_failed_table()}"); + } + + /** + * Get the reservation time before which a reservation is considered abandoned. + * + * @return int + * + * @since 3.2.0 + */ + protected function stale_before() + { + return $this->now() - (int) $this->option('retry_after', 300); + } + + /** + * Get the moment a job with the given delay becomes available. + * + * @param int|\DateTimeInterface|\DateInterval|null $delay The delay. + * + * @return int + * + * @since 3.2.0 + */ + protected function available_at($delay) + { + $seconds = $this->seconds_until($delay); + + return $this->now() + max(0, (int) $seconds); + } + + /** + * Build the queue-name constraint and append its bindings. + * + * @param array|null $queues The queue names, or null for every queue. + * @param array $bindings The bindings to append to. + * + * @return string + * + * @since 3.2.0 + */ + protected function queue_constraint(?array $queues, array &$bindings) + { + if (empty($queues)) { + return ''; + } + + $queues = array_values(array_map('strval', $queues)); + array_push($bindings, ...$queues); + + return ' AND queue IN (' . implode(', ', array_fill(0, count($queues), '%s')) . ')'; + } +} diff --git a/src/Queue/Events/JobFailed.php b/src/Queue/Events/JobFailed.php new file mode 100644 index 0000000..c65c169 --- /dev/null +++ b/src/Queue/Events/JobFailed.php @@ -0,0 +1,64 @@ +record = $record; + $this->job = $job; + $this->exception = $exception; + } +} diff --git a/src/Queue/Events/JobProcessed.php b/src/Queue/Events/JobProcessed.php new file mode 100644 index 0000000..6960da3 --- /dev/null +++ b/src/Queue/Events/JobProcessed.php @@ -0,0 +1,52 @@ +record = $record; + $this->job = $job; + } +} diff --git a/src/Queue/Events/JobProcessing.php b/src/Queue/Events/JobProcessing.php new file mode 100644 index 0000000..18f5c4e --- /dev/null +++ b/src/Queue/Events/JobProcessing.php @@ -0,0 +1,52 @@ +record = $record; + $this->job = $job; + } +} diff --git a/src/Queue/Events/JobQueued.php b/src/Queue/Events/JobQueued.php new file mode 100644 index 0000000..a75cf0d --- /dev/null +++ b/src/Queue/Events/JobQueued.php @@ -0,0 +1,51 @@ +id = $id; + $this->job = $job; + } +} diff --git a/src/Queue/JobRecord.php b/src/Queue/JobRecord.php new file mode 100644 index 0000000..b343541 --- /dev/null +++ b/src/Queue/JobRecord.php @@ -0,0 +1,255 @@ +id = $id; + $this->queue = $queue; + $this->payload = $payload; + $this->attempts = $attempts; + } + + /** + * Create a record from a row of the jobs table. + * + * @param array $row The row as an associative array. + * + * @return static + * + * @since 3.2.0 + */ + public static function from_row(array $row) + { + return new static( + (int) $row['id'], + (string) $row['queue'], + (string) $row['payload'], + (int) $row['attempts'] + ); + } + + /** + * Get the row id. + * + * @return int + * + * @since 3.2.0 + */ + public function id() + { + return $this->id; + } + + /** + * Get the queue name. + * + * @return string + * + * @since 3.2.0 + */ + public function queue() + { + return $this->queue; + } + + /** + * Get the raw JSON payload. + * + * @return string + * + * @since 3.2.0 + */ + public function payload() + { + return $this->payload; + } + + /** + * Get the attempt this run is. + * + * @return int + * + * @since 3.2.0 + */ + public function attempts() + { + return $this->attempts; + } + + /** + * Record that the job asked to go back on the queue. + * + * @param int|\DateTimeInterface|\DateInterval $delay How long before it is available again. + * + * @return void + * + * @since 3.2.0 + */ + public function release($delay = 0) + { + $this->released = true; + $this->release_delay = $delay; + } + + /** + * Determine whether the job released itself. + * + * @return bool + * + * @since 3.2.0 + */ + public function is_released() + { + return $this->released; + } + + /** + * Get the delay the job asked to be released with. + * + * @return int|\DateTimeInterface|\DateInterval|null + * + * @since 3.2.0 + */ + public function release_delay() + { + return $this->release_delay; + } + + /** + * Record that the job failed itself. + * + * @param \Throwable $exception Why the job failed. + * + * @return void + * + * @since 3.2.0 + */ + public function fail(Throwable $exception) + { + $this->failed = true; + $this->failure = $exception; + } + + /** + * Determine whether the job failed itself. + * + * @return bool + * + * @since 3.2.0 + */ + public function has_failed() + { + return $this->failed; + } + + /** + * Get the exception the job failed itself with. + * + * @return \Throwable|null + * + * @since 3.2.0 + */ + public function failure() + { + return $this->failure; + } +} diff --git a/src/Queue/Payload.php b/src/Queue/Payload.php new file mode 100644 index 0000000..3c92073 --- /dev/null +++ b/src/Queue/Payload.php @@ -0,0 +1,127 @@ +set_job_record(null); + + $tries = $job->get_tries(); + $backoff = $job->get_backoff(); + + return (string) wp_json_encode([ + 'uuid' => (string) uuid(), + 'display_name' => get_class($job), + 'job' => get_class($job), + 'max_tries' => max(1, (int) ($tries ?? $default_tries)), + 'backoff' => $backoff ?? $default_backoff, + 'priority' => $job->get_priority(), + 'data' => serialize($clone), + ]); + } + + /** + * Decode a stored envelope without restoring the job. + * + * @param string $payload The raw JSON payload. + * + * @return array + * + * @throws \Framework\Exceptions\QueueException When the envelope is malformed. + * + * @since 3.2.0 + */ + public static function decode(string $payload) + { + $envelope = json_decode($payload, true); + + if (!is_array($envelope) || !isset($envelope['job'], $envelope['data']) || !is_string($envelope['data'])) { + throw QueueException::invalid_payload('the envelope is not valid JSON or is missing its job data.'); + } + + return $envelope; + } + + /** + * Restore the job instance from a decoded envelope. + * + * @param array $envelope The decoded envelope. + * + * @return \Framework\Contracts\ShouldQueue + * + * @throws \Framework\Exceptions\QueueException When the class is unknown or not a queued job. + * + * @since 3.2.0 + */ + public static function restore(array $envelope) + { + $class = (string) $envelope['job']; + + if (!class_exists($class)) { + throw QueueException::invalid_payload(sprintf('the class [%s] does not exist.', $class)); + } + + if (!is_subclass_of($class, ShouldQueue::class)) { + throw QueueException::invalid_payload(sprintf( + 'the class [%s] does not implement %s.', + $class, + ShouldQueue::class + )); + } + + $job = unserialize($envelope['data'], ['allowed_classes' => true]); + + if (!$job instanceof $class) { + throw QueueException::invalid_payload(sprintf('the data does not restore to [%s].', $class)); + } + + return $job; + } + + /** + * Get the display name recorded in an envelope, falling back to the raw payload's class. + * + * @param string $payload The raw JSON payload. + * + * @return string + * + * @since 3.2.0 + */ + public static function display_name(string $payload) + { + $envelope = json_decode($payload, true); + + return is_array($envelope) ? (string) ($envelope['display_name'] ?? $envelope['job'] ?? 'unknown') : 'unknown'; + } +} diff --git a/src/Queue/PendingDispatch.php b/src/Queue/PendingDispatch.php new file mode 100644 index 0000000..707693c --- /dev/null +++ b/src/Queue/PendingDispatch.php @@ -0,0 +1,125 @@ +delay(60)->on_queue('emails')` read + * as one statement. The checks that can fail (queue enabled, table present) already ran in + * dispatch(), so the destructor only has to insert. + * + * @package Framework + * @subpackage Queue + * @since 3.2.0 + */ +namespace Framework\Queue; + +defined('ABSPATH') || exit; + +use Framework\Contracts\ShouldQueue; + +class PendingDispatch +{ + /** + * The job being dispatched. + * + * @var \Framework\Contracts\ShouldQueue + * + * @since 3.2.0 + */ + protected $job; + + /** + * The queue manager the job is pushed to. + * + * @var \Framework\Queue\QueueManager + * + * @since 3.2.0 + */ + protected $manager; + + /** + * Create a pending dispatch. + * + * @param \Framework\Contracts\ShouldQueue $job The job being dispatched. + * @param \Framework\Queue\QueueManager $manager The queue manager. + * + * @return void + * + * @since 3.2.0 + */ + public function __construct(ShouldQueue $job, QueueManager $manager) + { + $this->job = $job; + $this->manager = $manager; + } + + /** + * Set how long to wait before the job becomes available. + * + * @param int|\DateTimeInterface|\DateInterval|null $delay Seconds, or the moment it becomes available. + * + * @return $this + * + * @since 3.2.0 + */ + public function delay($delay) + { + $this->job->delay($delay); + + return $this; + } + + /** + * Set the queue name the job is dispatched to. + * + * @param string|null $queue The queue name. + * + * @return $this + * + * @since 3.2.0 + */ + public function on_queue($queue) + { + $this->job->on_queue($queue); + + return $this; + } + + /** + * Set the job's priority; higher runs first. + * + * @param int $priority The priority. + * + * @return $this + * + * @since 3.2.0 + */ + public function with_priority(int $priority) + { + $this->job->with_priority($priority); + + return $this; + } + + /** + * Get the job being dispatched. + * + * @return \Framework\Contracts\ShouldQueue + * + * @since 3.2.0 + */ + public function get_job() + { + return $this->job; + } + + /** + * Push the job to the queue. + * + * @return void + * + * @since 3.2.0 + */ + public function __destruct() + { + $this->manager->push($this->job); + } +} diff --git a/src/Queue/QueueFake.php b/src/Queue/QueueFake.php new file mode 100644 index 0000000..19c511a --- /dev/null +++ b/src/Queue/QueueFake.php @@ -0,0 +1,196 @@ +jobs[] = $job; + + return count($this->jobs); + } + + /** + * Get the pushed jobs of a class that pass the callback. + * + * @param string $class The job class. + * @param callable|null $callback Receives the job; return true to keep it. + * + * @return \Framework\Contracts\ShouldQueue[] + * + * @since 3.2.0 + */ + public function pushed(string $class, ?callable $callback = null) + { + return array_values(array_filter($this->jobs, function ($job) use ($class, $callback) { + return $job instanceof $class && (is_null($callback) || $callback($job)); + })); + } + + /** + * Count the recorded jobs, optionally on one queue. + * + * @param string|null $queue The queue name. + * + * @return int + * + * @since 3.2.0 + */ + public function size(?string $queue = null) + { + return count(array_filter($this->jobs, function ($job) use ($queue) { + return is_null($queue) || $job->get_queue() === $queue; + })); + } + + /** + * Forget the recorded jobs, optionally on one queue. + * + * @param string|null $queue The queue name. + * + * @return int The number of jobs forgotten. + * + * @since 3.2.0 + */ + public function clear(?string $queue = null) + { + $before = count($this->jobs); + + $this->jobs = array_values(array_filter($this->jobs, function ($job) use ($queue) { + return !is_null($queue) && $job->get_queue() !== $queue; + })); + + return $before - count($this->jobs); + } + + /** + * Assert that a job of the class was pushed and passes the callback. + * + * @param string $class The job class. + * @param callable|null $callback Receives the job; return true when it matches. + * + * @return void + * + * @throws \AssertionError When no matching job was pushed. + * + * @since 3.2.0 + */ + public function assert_pushed(string $class, ?callable $callback = null) + { + $this->assert( + count($this->pushed($class, $callback)) > 0, + sprintf('The expected [%s] job was not pushed.', $class) + ); + } + + /** + * Assert that a job of the class was pushed exactly the given number of times. + * + * @param string $class The job class. + * @param int $times The expected count. + * + * @return void + * + * @throws \AssertionError When the count differs. + * + * @since 3.2.0 + */ + public function assert_pushed_times(string $class, int $times) + { + $count = count($this->pushed($class)); + + $this->assert( + $count === $times, + sprintf('The expected [%s] job was pushed %d times instead of %d times.', $class, $count, $times) + ); + } + + /** + * Assert that no job of the class passing the callback was pushed. + * + * @param string $class The job class. + * @param callable|null $callback Receives the job; return true when it matches. + * + * @return void + * + * @throws \AssertionError When a matching job was pushed. + * + * @since 3.2.0 + */ + public function assert_not_pushed(string $class, ?callable $callback = null) + { + $this->assert( + count($this->pushed($class, $callback)) === 0, + sprintf('The unexpected [%s] job was pushed.', $class) + ); + } + + /** + * Assert that no job was pushed at all. + * + * @return void + * + * @throws \AssertionError When any job was pushed. + * + * @since 3.2.0 + */ + public function assert_nothing_pushed() + { + $this->assert( + empty($this->jobs), + sprintf('%d unexpected jobs were pushed.', count($this->jobs)) + ); + } + + /** + * Fail with the message unless the condition holds. + * + * @param bool $condition The condition. + * @param string $message The failure message. + * + * @return void + * + * @throws \AssertionError When the condition is false. + * + * @since 3.2.0 + */ + protected function assert(bool $condition, string $message) + { + if (!$condition) { + throw new AssertionError($message); + } + } +} diff --git a/src/Queue/QueueManager.php b/src/Queue/QueueManager.php new file mode 100644 index 0000000..a494189 --- /dev/null +++ b/src/Queue/QueueManager.php @@ -0,0 +1,322 @@ +database = $database; + $this->spawner = $spawner; + } + + /** + * Get the bound queue manager. + * + * @return static + * + * @throws \Framework\Exceptions\QueueException When the queue service provider is not registered. + * + * @since 3.2.0 + */ + public static function resolve() + { + if (!app()->bound(QueueManager::class)) { + throw QueueException::not_enabled(); + } + + return app(QueueManager::class); + } + + /** + * Make sure a job can be stored. + * + * @return void + * + * @throws \Framework\Exceptions\QueueException When the jobs table does not exist. + * + * @since 3.2.0 + */ + public function ensure_ready() + { + if ($this->fake) { + return; + } + + if (!$this->database->table_exists()) { + throw QueueException::table_missing($this->database->get_table()); + } + } + + /** + * Push a job onto the queue. + * + * @param \Framework\Contracts\ShouldQueue $job The job. + * + * @return int The stored row id. + * + * @since 3.2.0 + */ + public function push(ShouldQueue $job) + { + if ($this->fake) { + return $this->fake->push($job); + } + + $id = $this->database->push( + Payload::encode( + $job, + (int) $this->database->option('tries', 1), + $this->database->option('backoff', 0) + ), + $job->get_queue(), + $job->get_priority(), + $job->get_delay() + ); + + $events = app('event'); + + if ($events->has_listeners(JobQueued::class)) { + $events->dispatch(new JobQueued($id, $job)); + } + + if ($this->database->is_immediate($job->get_delay())) { + $this->schedule_spawn(); + } + + return $id; + } + + /** + * Push a job onto the queue after a delay. + * + * @param int|\DateTimeInterface|\DateInterval $delay Seconds, or the moment it becomes available. + * @param \Framework\Contracts\ShouldQueue $job The job. + * + * @return int The stored row id. + * + * @since 3.2.0 + */ + public function later($delay, ShouldQueue $job) + { + $job->delay($delay); + + return $this->push($job); + } + + /** + * Count the jobs not currently being worked on, optionally on one queue. + * + * @param string|null $queue The queue name. + * + * @return int + * + * @since 3.2.0 + */ + public function size(?string $queue = null) + { + return $this->fake ? $this->fake->size($queue) : $this->database->size($queue); + } + + /** + * Delete the jobs not currently being worked on, optionally on one queue. + * + * @param string|null $queue The queue name. + * + * @return int The number of jobs deleted. + * + * @since 3.2.0 + */ + public function clear(?string $queue = null) + { + return $this->fake ? $this->fake->clear($queue) : $this->database->clear($queue); + } + + /** + * Replace the queue with a recorder for the rest of the request. + * + * Also binds this manager, so a test can fake the queue without registering the provider. + * + * @return \Framework\Queue\QueueFake + * + * @since 3.2.0 + */ + public function fake() + { + $this->fake = new QueueFake(); + + app()->instance(QueueManager::class, $this); + + return $this->fake; + } + + /** + * Determine whether the queue is faked. + * + * @return bool + * + * @since 3.2.0 + */ + public function is_faked() + { + return !is_null($this->fake); + } + + /** + * Assert that a job of the class was pushed and passes the callback. + * + * @param string $class The job class. + * @param callable|null $callback Receives the job; return true when it matches. + * + * @return void + * + * @since 3.2.0 + */ + public function assert_pushed(string $class, ?callable $callback = null) + { + $this->faked()->assert_pushed($class, $callback); + } + + /** + * Assert that a job of the class was pushed exactly the given number of times. + * + * @param string $class The job class. + * @param int $times The expected count. + * + * @return void + * + * @since 3.2.0 + */ + public function assert_pushed_times(string $class, int $times) + { + $this->faked()->assert_pushed_times($class, $times); + } + + /** + * Assert that no job of the class passing the callback was pushed. + * + * @param string $class The job class. + * @param callable|null $callback Receives the job; return true when it matches. + * + * @return void + * + * @since 3.2.0 + */ + public function assert_not_pushed(string $class, ?callable $callback = null) + { + $this->faked()->assert_not_pushed($class, $callback); + } + + /** + * Assert that no job was pushed at all. + * + * @return void + * + * @since 3.2.0 + */ + public function assert_nothing_pushed() + { + $this->faked()->assert_nothing_pushed(); + } + + /** + * Get the fake, failing loudly when the queue was not faked first. + * + * @return \Framework\Queue\QueueFake + * + * @throws \Framework\Exceptions\QueueException When Queue::fake() was not called. + * + * @since 3.2.0 + */ + protected function faked() + { + if (!$this->fake) { + throw new QueueException('Call Queue::fake() before making queue assertions.'); + } + + return $this->fake; + } + + /** + * Start a worker when the current request finishes, at most once per request. + * + * @return void + * + * @since 3.2.0 + */ + protected function schedule_spawn() + { + if ($this->spawn_scheduled || !function_exists('add_action')) { + return; + } + + $this->spawn_scheduled = true; + + add_action('shutdown', function () { + $this->spawner->spawn_if_idle(); + }, PHP_INT_MAX); + } +} diff --git a/src/Queue/QueueServiceProvider.php b/src/Queue/QueueServiceProvider.php new file mode 100644 index 0000000..b87c4be --- /dev/null +++ b/src/Queue/QueueServiceProvider.php @@ -0,0 +1,141 @@ +app->singleton(DatabaseQueue::class); + $this->app->singleton(Worker::class); + $this->app->singleton(Spawner::class); + $this->app->singleton(QueueManager::class); + } + + /** + * Wire the sweep schedule, the worker endpoint, and the queue commands. + * + * @return void + * + * @since 3.2.0 + */ + public function boot() + { + if (function_exists('add_action')) { + $this->register_sweep(); + $this->register_worker_endpoint(); + } + + if ($this->app->is_cli_available()) { + $this->register_commands(); + } + } + + /** + * Add the every-minute schedule and the recurring sweep event. + * + * @return void + * + * @since 3.2.0 + */ + protected function register_sweep() + { + $schedule = $this->app->prefix() . 'queue_every_minute'; + $hook = $this->app->prefix() . 'queue_sweep'; + + add_filter('cron_schedules', function ($schedules) use ($schedule) { + $schedules[$schedule] = [ + 'interval' => MINUTE_IN_SECONDS, + 'display' => __('Every minute (queue sweep)'), + ]; + + return $schedules; + }); + + add_action($hook, function () { + try { + $this->app->make(Sweeper::class)->sweep(); + } catch (Throwable $exception) { + // A sweep that cannot run must never fail the request that triggered it. + } + }); + + if (function_exists('wp_next_scheduled') && !wp_next_scheduled($hook)) { + wp_schedule_event(time() + MINUTE_IN_SECONDS, $schedule, $hook); + } + } + + /** + * Register the AJAX action the loopback request is sent to. + * + * Both the logged-in and logged-out variants are hooked because which one fires depends on + * the cookies the request carries; the signature, not the session, is the credential. + * + * @return void + * + * @since 3.2.0 + */ + protected function register_worker_endpoint() + { + $callback = function () { + $accepted = $this->app->make(Spawner::class)->handle_request(Superglobals::post()); + + wp_die('', '', ['response' => $accepted ? 200 : 403]); + }; + + add_action('wp_ajax_nopriv_' . Spawner::action(), $callback); + add_action('wp_ajax_' . Spawner::action(), $callback); + } + + /** + * Register the queue commands that need the queue to be enabled. + * + * @return void + * + * @since 3.2.0 + */ + protected function register_commands() + { + $commands = [ + 'queue:work' => QueueWorkCommand::class, + 'queue:failed' => QueueFailedCommand::class, + 'queue:retry' => QueueRetryCommand::class, + 'queue:forget' => QueueForgetCommand::class, + 'queue:flush' => QueueFlushCommand::class, + 'queue:clear' => QueueClearCommand::class, + ]; + + foreach ($commands as $command => $class) { + Command::register($command, $class); + } + } +} diff --git a/src/Queue/Spawner.php b/src/Queue/Spawner.php new file mode 100644 index 0000000..43609ad --- /dev/null +++ b/src/Queue/Spawner.php @@ -0,0 +1,242 @@ +queue = $queue; + $this->worker = $worker; + } + + /** + * Get the AJAX action the worker endpoint listens on. + * + * @return string + * + * @since 3.2.0 + */ + public static function action() + { + return app()->prefix() . 'queue_work'; + } + + /** + * Start a worker unless a chain is already running. + * + * @return bool Whether a worker was started. + * + * @since 3.2.0 + */ + public function spawn_if_idle() + { + $lock = $this->lock(); + + if (!$lock->acquire()) { + return false; + } + + $lock->release(); + $this->spawn(); + + return true; + } + + /** + * Send the signed, non-blocking loopback request that starts a worker. + * + * @return void + * + * @since 3.2.0 + */ + public function spawn() + { + $timestamp = (int) $this->current_timestamp(); + + try { + Http::as_form() + ->timeout(0.01) + ->with_options([ + 'blocking' => false, + 'sslverify' => apply_filters('https_local_ssl_verify', false), + ]) + ->post(admin_url('admin-ajax.php'), [ + 'action' => static::action(), + 'ts' => $timestamp, + 'sig' => $this->sign($timestamp), + ]); + } catch (Throwable $exception) { + // A spawn that cannot be sent is retried by the next sweep; it must not break the request. + } + } + + /** + * Handle an incoming worker request. + * + * @param array $input The request body. + * + * @return bool Whether the request was accepted and a worker ran. + * + * @since 3.2.0 + */ + public function handle_request(array $input) + { + if (!$this->verify($input)) { + return false; + } + + $lock = $this->lock(); + + if (!$lock->acquire()) { + return false; + } + + $this->prepare_runtime(); + + try { + $remaining = $this->worker->run(); + } catch (Throwable $exception) { + $remaining = true; + } + + $lock->release(); + + if ($remaining) { + $this->spawn(); + } + + return true; + } + + /** + * Sign a worker request. + * + * @param int $timestamp The request timestamp. + * + * @return string + * + * @since 3.2.0 + */ + public function sign(int $timestamp) + { + return hash_hmac('sha256', static::action() . '|' . $timestamp, wp_salt('auth')); + } + + /** + * Verify a worker request's signature and freshness. + * + * @param array $input The request body. + * + * @return bool + * + * @since 3.2.0 + */ + public function verify(array $input) + { + if (!isset($input['ts'], $input['sig']) || !is_scalar($input['ts']) || !is_string($input['sig'])) { + return false; + } + + $timestamp = (int) $input['ts']; + + if (abs((int) $this->current_timestamp() - $timestamp) > static::SIGNATURE_WINDOW) { + return false; + } + + return hash_equals($this->sign($timestamp), $input['sig']); + } + + /** + * Get the lock that allows one worker chain at a time. + * + * It outlives a worker's budget by a margin, so a worker that dies leaves it to expire. + * + * @return \Framework\Contracts\Lock + * + * @since 3.2.0 + */ + protected function lock() + { + return app(CacheManager::class)->lock( + app()->prefix() . 'queue_chain', + (int) ceil($this->worker->budget()) + 60 + ); + } + + /** + * Keep the worker running after the loopback caller disconnects, with room for its budget. + * + * @return void + * + * @since 3.2.0 + */ + protected function prepare_runtime() + { + ignore_user_abort(true); + + if (function_exists('set_time_limit')) { + @set_time_limit((int) ceil($this->worker->budget()) + 30); + } + } +} diff --git a/src/Queue/Sweeper.php b/src/Queue/Sweeper.php new file mode 100644 index 0000000..191a7f9 --- /dev/null +++ b/src/Queue/Sweeper.php @@ -0,0 +1,66 @@ +queue = $queue; + $this->spawner = $spawner; + } + + /** + * Start a worker when jobs are due and no chain is running. + * + * @return bool Whether a worker was started. + * + * @since 3.2.0 + */ + public function sweep() + { + if (!$this->queue->table_exists() || !$this->queue->has_due()) { + return false; + } + + return $this->spawner->spawn_if_idle(); + } +} diff --git a/src/Queue/Worker.php b/src/Queue/Worker.php new file mode 100644 index 0000000..b790cc0 --- /dev/null +++ b/src/Queue/Worker.php @@ -0,0 +1,411 @@ +queue = $queue; + } + + /** + * Process due jobs until the budget, the job limit, or the queue runs out. + * + * Options: `budget` (seconds, or null for no limit), `max_jobs` (int or null), and `queues` + * (queue names to restrict to, or null for every queue). + * + * @param array $options The run options. + * + * @return bool Whether due jobs remain. + * + * @since 3.2.0 + */ + public function run(array $options = []) + { + $budget = array_key_exists('budget', $options) ? $options['budget'] : $this->budget(); + $max_jobs = isset($options['max_jobs']) ? max(1, (int) $options['max_jobs']) : null; + $queues = $options['queues'] ?? null; + $batch_size = max(1, (int) $this->queue->option('batch_size', 10)); + + $this->started_at = microtime(true); + $processed = 0; + + while (!$this->out_of_time($budget) && (is_null($max_jobs) || $processed < $max_jobs)) { + $limit = is_null($max_jobs) ? $batch_size : min($batch_size, $max_jobs - $processed); + $records = array_values($this->queue->claim($limit, $this->token(), $queues)); + + if (empty($records)) { + break; + } + + foreach ($records as $index => $record) { + if ($this->out_of_time($budget)) { + $this->queue->release_unstarted(array_map(function (JobRecord $unstarted) { + return $unstarted->id(); + }, array_slice($records, $index))); + + break 2; + } + + $this->process($record); + $processed++; + } + } + + return $this->queue->has_due($queues); + } + + /** + * Get the time budget for a web worker, in seconds. + * + * @return float + * + * @since 3.2.0 + */ + public function budget() + { + $limit = (float) $this->queue->option('time_limit', 20); + $max_execution_time = (int) ini_get('max_execution_time'); + + if ($max_execution_time > 0) { + $limit = min($limit, floor($max_execution_time * 0.8)); + } + + return max(1.0, $limit); + } + + /** + * Run one claimed job and settle it. + * + * @param \Framework\Queue\JobRecord $record The claimed record. + * + * @return void + * + * @since 3.2.0 + */ + public function process(JobRecord $record) + { + try { + $envelope = Payload::decode($record->payload()); + } catch (QueueException $exception) { + $this->fail_job($record, null, $exception); + + return; + } + + if ($record->attempts() > max(1, (int) ($envelope['max_tries'] ?? 1))) { + $this->fail_job( + $record, + null, + QueueException::max_attempts_exceeded(Payload::display_name($record->payload())) + ); + + return; + } + + try { + $job = Payload::restore($envelope); + } catch (QueueException $exception) { + $this->fail_job($record, null, $exception); + + return; + } + + $job->set_job_record($record); + $this->fire(JobProcessing::class, [$record, $job]); + + try { + $this->call_handle($job); + } catch (Throwable $exception) { + $this->handle_exception($record, $job, $envelope, $exception); + + return; + } + + if ($record->has_failed()) { + $this->fail_job($record, $job, $record->failure()); + + return; + } + + if ($record->is_released()) { + $this->queue->release($record->id(), $record->release_delay()); + + return; + } + + $this->queue->delete($record->id()); + $this->fire(JobProcessed::class, [$record, $job]); + } + + /** + * Run a job in the current request without storing it. + * + * @param \Framework\Contracts\ShouldQueue $job The job. + * + * @return mixed The value handle() returned. + * + * @throws \Throwable Whatever handle() threw, after failed() has been called. + * + * @since 3.2.0 + */ + public function run_sync(ShouldQueue $job) + { + $record = new JobRecord(0, $job->get_queue(), '', 1); + $job->set_job_record($record); + + try { + $result = $this->call_handle($job); + } catch (Throwable $exception) { + $this->call_failed($job, $exception); + + throw $exception; + } + + if ($record->has_failed()) { + $this->call_failed($job, $record->failure()); + } + + return $result; + } + + /** + * Decide between retrying and failing a job whose handle() threw. + * + * @param \Framework\Queue\JobRecord $record The claimed record. + * @param \Framework\Contracts\ShouldQueue $job The job. + * @param array $envelope The decoded envelope. + * @param \Throwable $exception What handle() threw. + * + * @return void + * + * @since 3.2.0 + */ + protected function handle_exception(JobRecord $record, ShouldQueue $job, array $envelope, Throwable $exception) + { + if (!$record->has_failed() && $record->attempts() < max(1, (int) ($envelope['max_tries'] ?? 1))) { + $this->queue->release($record->id(), $this->backoff_for($envelope['backoff'] ?? 0, $record->attempts())); + + return; + } + + $this->fail_job($record, $job, $exception); + } + + /** + * Get the retry delay for the attempt that just failed. + * + * @param int|array $backoff Seconds, or seconds per attempt with the last value reused. + * @param int $attempts The attempt that just failed. + * + * @return int + * + * @since 3.2.0 + */ + protected function backoff_for($backoff, int $attempts) + { + if (!is_array($backoff)) { + return max(0, (int) $backoff); + } + + $backoff = array_values($backoff); + + if (empty($backoff)) { + return 0; + } + + return max(0, (int) $backoff[min(max(1, $attempts), count($backoff)) - 1]); + } + + /** + * Record a job as failed and tell everyone who needs to know. + * + * @param \Framework\Queue\JobRecord $record The claimed record. + * @param \Framework\Contracts\ShouldQueue|null $job The job, if it was restored. + * @param \Throwable $exception Why it failed. + * + * @return void + * + * @since 3.2.0 + */ + protected function fail_job(JobRecord $record, ?ShouldQueue $job, Throwable $exception) + { + $this->queue->fail($record, $exception); + + if ($job) { + $this->call_failed($job, $exception); + } + + $this->fire(JobFailed::class, [$record, $job, $exception]); + + try { + Log::error(sprintf( + 'Queued job [%s] failed on attempt %d: %s', + Payload::display_name($record->payload()), + $record->attempts(), + $exception->getMessage() + )); + } catch (Throwable $ignored) { + // A log that cannot be written must not stop the worker. + } + } + + /** + * Call the job's failed() method, if it has one, without letting it stop the worker. + * + * @param \Framework\Contracts\ShouldQueue $job The job. + * @param \Throwable $exception Why it failed. + * + * @return void + * + * @since 3.2.0 + */ + protected function call_failed(ShouldQueue $job, Throwable $exception) + { + if (!method_exists($job, 'failed')) { + return; + } + + try { + $job->failed($exception); + } catch (Throwable $ignored) { + // The job has already failed; its failure hook failing too changes nothing. + } + } + + /** + * Call the job's handle() method with its class-typed parameters resolved from the container. + * + * @param \Framework\Contracts\ShouldQueue $job The job. + * + * @return mixed + * + * @since 3.2.0 + */ + protected function call_handle(ShouldQueue $job) + { + $dependencies = []; + + foreach ((new ReflectionMethod($job, 'handle'))->getParameters() as $parameter) { + $type = $parameter->getType(); + + if ($type instanceof ReflectionNamedType && !$type->isBuiltin()) { + $dependencies[] = app()->make($type->getName()); + continue; + } + + $dependencies[] = $parameter->isDefaultValueAvailable() ? $parameter->getDefaultValue() : null; + } + + return $job->handle(...$dependencies); + } + + /** + * Dispatch a queue event when anything listens for it. + * + * @param string $event_class The event class. + * @param array $arguments The event's constructor arguments. + * + * @return void + * + * @since 3.2.0 + */ + protected function fire(string $event_class, array $arguments) + { + $events = app('event'); + + if ($events->has_listeners($event_class)) { + $events->dispatch(new $event_class(...$arguments)); + } + } + + /** + * Determine whether the run has used up its budget. + * + * @param float|null $budget The budget in seconds, or null for no limit. + * + * @return bool + * + * @since 3.2.0 + */ + protected function out_of_time($budget) + { + return !is_null($budget) && $this->elapsed() >= $budget; + } + + /** + * Get the seconds elapsed since the run started. + * + * @return float + * + * @since 3.2.0 + */ + protected function elapsed() + { + return microtime(true) - $this->started_at; + } + + /** + * Generate a fresh claim token. + * + * @return string + * + * @since 3.2.0 + */ + protected function token() + { + return bin2hex(random_bytes(16)); + } +} diff --git a/src/Supports/Facades/Queue.php b/src/Supports/Facades/Queue.php new file mode 100644 index 0000000..fe300f2 --- /dev/null +++ b/src/Supports/Facades/Queue.php @@ -0,0 +1,47 @@ +options = array_merge($this->options, $options); + + return $this; + } + + public function get_table() + { + return 'wp_' . $this->get_table_name(); + } + + public function get_failed_table() + { + return 'wp_' . $this->get_failed_table_name(); + } + + public function table_exists() + { + return $this->exists; + } + + public function push(string $payload, string $queue, int $priority = 0, $delay = null) + { + $id = $this->next_id++; + + $this->rows[$id] = [ + 'id' => $id, + 'queue' => $queue, + 'priority' => $priority, + 'payload' => $payload, + 'attempts' => 0, + 'reserved_at' => null, + 'reserved_by' => null, + 'available_at' => $this->available_at($delay), + 'created_at' => $this->now(), + ]; + + return $id; + } + + public function claim(int $limit, string $token, ?array $queues = null) + { + $due = $this->due_rows($queues); + + usort($due, function ($a, $b) { + return [$b['priority'], $a['available_at'], $a['id']] <=> [$a['priority'], $b['available_at'], $b['id']]; + }); + + $claimed = array_slice($due, 0, $limit); + $this->claims[] = count($claimed); + + foreach ($claimed as $row) { + $this->rows[$row['id']]['reserved_at'] = $this->now(); + $this->rows[$row['id']]['reserved_by'] = $token; + $this->rows[$row['id']]['attempts']++; + } + + return $this->reserved($token); + } + + public function reserved(string $token) + { + $rows = array_values(array_filter($this->rows, function ($row) use ($token) { + return $row['reserved_by'] === $token; + })); + + usort($rows, function ($a, $b) { + return [$b['priority'], $a['available_at'], $a['id']] <=> [$a['priority'], $b['available_at'], $b['id']]; + }); + + return array_map([JobRecord::class, 'from_row'], $rows); + } + + public function delete(int $id) + { + unset($this->rows[$id]); + } + + public function release(int $id, $delay = 0) + { + $this->rows[$id]['reserved_at'] = null; + $this->rows[$id]['reserved_by'] = null; + $this->rows[$id]['available_at'] = $this->available_at($delay); + } + + public function release_unstarted(array $ids) + { + foreach ($ids as $id) { + $this->rows[$id]['reserved_at'] = null; + $this->rows[$id]['reserved_by'] = null; + $this->rows[$id]['attempts'] = max(0, $this->rows[$id]['attempts'] - 1); + } + } + + public function fail(JobRecord $record, Throwable $exception) + { + $id = $this->next_failed_id++; + $envelope = json_decode($record->payload(), true); + + $this->failed[$id] = [ + 'id' => $id, + 'uuid' => $envelope['uuid'] ?? 'uuid-' . $id, + 'queue' => $record->queue(), + 'payload' => $record->payload(), + 'exception' => (string) $exception, + 'failed_at' => $this->now(), + ]; + + $this->delete($record->id()); + } + + public function has_due(?array $queues = null) + { + return !empty($this->due_rows($queues)); + } + + public function size(?string $queue = null) + { + return count($this->waiting_rows($queue)); + } + + public function clear(?string $queue = null) + { + $waiting = $this->waiting_rows($queue); + + foreach ($waiting as $row) { + unset($this->rows[$row['id']]); + } + + return count($waiting); + } + + public function failed_all() + { + return array_reverse(array_values($this->failed)); + } + + public function failed_find(int $id) + { + return $this->failed[$id] ?? null; + } + + public function forget(int $id) + { + $existed = isset($this->failed[$id]); + unset($this->failed[$id]); + + return $existed; + } + + public function flush_failed() + { + $count = count($this->failed); + $this->failed = []; + + return $count; + } + + public function row(int $id): ?array + { + return $this->rows[$id] ?? null; + } + + public function payload_of(int $id): array + { + return json_decode($this->rows[$id]['payload'], true); + } + + protected function is_claimable(array $row): bool + { + return is_null($row['reserved_at']) || $row['reserved_at'] <= $this->stale_before(); + } + + protected function due_rows(?array $queues): array + { + return array_values(array_filter($this->rows, function ($row) use ($queues) { + return $row['available_at'] <= $this->now() + && $this->is_claimable($row) + && (empty($queues) || in_array($row['queue'], $queues, true)); + })); + } + + protected function waiting_rows(?string $queue): array + { + return array_values(array_filter($this->rows, function ($row) use ($queue) { + return $this->is_claimable($row) && (is_null($queue) || $row['queue'] === $queue); + })); + } +} diff --git a/tests/Support/Queue/Jobs/AttemptsJob.php b/tests/Support/Queue/Jobs/AttemptsJob.php new file mode 100644 index 0000000..4b9206b --- /dev/null +++ b/tests/Support/Queue/Jobs/AttemptsJob.php @@ -0,0 +1,23 @@ +attempts()); + + if ($this->attempts() === 1) { + throw new RuntimeException('first try fails'); + } + } +} diff --git a/tests/Support/Queue/Jobs/DefaultsJob.php b/tests/Support/Queue/Jobs/DefaultsJob.php new file mode 100644 index 0000000..fb21add --- /dev/null +++ b/tests/Support/Queue/Jobs/DefaultsJob.php @@ -0,0 +1,21 @@ +on_queue('emails')->with_priority(5)->delay(60); + } + + public function handle() + { + Journal::write('defaults'); + } +} diff --git a/tests/Support/Queue/Jobs/FailedThrowsJob.php b/tests/Support/Queue/Jobs/FailedThrowsJob.php new file mode 100644 index 0000000..9be7064 --- /dev/null +++ b/tests/Support/Queue/Jobs/FailedThrowsJob.php @@ -0,0 +1,22 @@ +name); + } +} diff --git a/tests/Support/Queue/Jobs/Journal.php b/tests/Support/Queue/Jobs/Journal.php new file mode 100644 index 0000000..9e5d32f --- /dev/null +++ b/tests/Support/Queue/Jobs/Journal.php @@ -0,0 +1,28 @@ +release(30); + } + + public function failed() + { + Journal::write('failed'); + } +} diff --git a/tests/Support/Queue/Jobs/SelfFailingJob.php b/tests/Support/Queue/Jobs/SelfFailingJob.php new file mode 100644 index 0000000..7624b24 --- /dev/null +++ b/tests/Support/Queue/Jobs/SelfFailingJob.php @@ -0,0 +1,25 @@ +fail(new RuntimeException('gave up')); + } + + public function failed(Throwable $exception) + { + Journal::write('failed', $exception->getMessage()); + } +} diff --git a/tests/Support/Queue/Jobs/SendEmail.php b/tests/Support/Queue/Jobs/SendEmail.php new file mode 100644 index 0000000..328fe6c --- /dev/null +++ b/tests/Support/Queue/Jobs/SendEmail.php @@ -0,0 +1,31 @@ +user_id = $user_id; + $this->subject = $subject; + $this->tags = $tags; + } + + public function handle() + { + Journal::write('handled', $this->user_id); + + return 'sent:' . $this->user_id; + } +} diff --git a/tests/Support/Queue/Jobs/ThrowingJob.php b/tests/Support/Queue/Jobs/ThrowingJob.php new file mode 100644 index 0000000..c647296 --- /dev/null +++ b/tests/Support/Queue/Jobs/ThrowingJob.php @@ -0,0 +1,29 @@ +attempts()); + + throw new RuntimeException('boom'); + } + + public function failed(Throwable $exception) + { + Journal::write('failed', $exception->getMessage()); + } +} diff --git a/tests/Support/Queue/RecordingLogger.php b/tests/Support/Queue/RecordingLogger.php new file mode 100644 index 0000000..74d1543 --- /dev/null +++ b/tests/Support/Queue/RecordingLogger.php @@ -0,0 +1,16 @@ +errors[] = $message; + } +} diff --git a/tests/Support/Queue/RecordingSpawner.php b/tests/Support/Queue/RecordingSpawner.php new file mode 100644 index 0000000..84cedf2 --- /dev/null +++ b/tests/Support/Queue/RecordingSpawner.php @@ -0,0 +1,36 @@ +spawned++; + } + + public function chain_lock(): TestLock + { + return $this->lock(); + } + + protected function lock() + { + return new TestLock('queue_chain'); + } + + protected function prepare_runtime() + { + // Raising the time limit would leak into the rest of the test run. + } +} diff --git a/tests/Support/Queue/TestLock.php b/tests/Support/Queue/TestLock.php new file mode 100644 index 0000000..074db35 --- /dev/null +++ b/tests/Support/Queue/TestLock.php @@ -0,0 +1,44 @@ +name = $name; + } + + public function acquire() + { + if (!empty(static::$held[$this->name])) { + return false; + } + + static::$held[$this->name] = true; + $this->owned = true; + + return true; + } + + public function release() + { + if (!$this->owned) { + return false; + } + + unset(static::$held[$this->name]); + $this->owned = false; + + return true; + } +} diff --git a/tests/Support/Queue/TestWorker.php b/tests/Support/Queue/TestWorker.php new file mode 100644 index 0000000..3a7de21 --- /dev/null +++ b/tests/Support/Queue/TestWorker.php @@ -0,0 +1,39 @@ +fake_elapsed = 0.0; + + return parent::run($options); + } + + public function process(JobRecord $record) + { + $this->processed[] = $record->id(); + + parent::process($record); + + $this->fake_elapsed += $this->seconds_per_job; + } + + protected function elapsed() + { + return $this->fake_elapsed; + } +} diff --git a/tests/Support/StubsWordPressFunctions.php b/tests/Support/StubsWordPressFunctions.php index 52dd8f4..4cdfedd 100644 --- a/tests/Support/StubsWordPressFunctions.php +++ b/tests/Support/StubsWordPressFunctions.php @@ -953,3 +953,66 @@ public function __construct($args = []) } } } + +if (!function_exists('admin_url')) { + function admin_url($path = '', $scheme = 'admin') + { + return 'https://example.test/wp-admin/' . ltrim((string) $path, '/'); + } +} + +if (!function_exists('wp_next_scheduled')) { + function wp_next_scheduled($hook, $args = []) + { + return $GLOBALS['framework_test_cron'][$hook]['timestamp'] ?? false; + } +} + +if (!function_exists('wp_schedule_event')) { + function wp_schedule_event($timestamp, $recurrence, $hook, $args = [], $wp_error = false) + { + $GLOBALS['framework_test_cron'][$hook] = [ + 'timestamp' => $timestamp, + 'recurrence' => $recurrence, + ]; + + return true; + } +} + +if (!class_exists('WP_CLI')) { + /** + * Records what a command printed. error() throws, standing in for the real one's exit. + */ + class WP_CLI + { + public static $messages = []; + + public static function add_command($name, $callable, $args = []) + { + return true; + } + + public static function success($message) + { + static::$messages[] = ['success', $message]; + } + + public static function warning($message) + { + static::$messages[] = ['warning', $message]; + } + + public static function line($message = '') + { + static::$messages[] = ['line', $message]; + } + + public static function error($message, $exit = true) + { + static::$messages[] = ['error', $message]; + + throw new RuntimeException((string) $message); + } + } +} diff --git a/tests/Unit/Queue/DatabaseQueueSqlTest.php b/tests/Unit/Queue/DatabaseQueueSqlTest.php new file mode 100644 index 0000000..3e8ebe5 --- /dev/null +++ b/tests/Unit/Queue/DatabaseQueueSqlTest.php @@ -0,0 +1,134 @@ +bootstrap_application(); + $app->use_prefix('kirki_'); + + $wpdb = new TestWpdb(); + $wpdb->insert_id = 7; + $wpdb->rows_affected = 0; + $this->wpdb = $wpdb; + $app->instance(Connection::class, new Connection()); + + $this->queue = new class extends DatabaseQueue { + use FreezesTime; + }; + $this->queue->freeze(self::NOW); + } + + public function test_table_names_carry_the_wordpress_and_app_prefixes(): void + { + $this->assertSame('wp_kirki_jobs', $this->queue->get_table()); + $this->assertSame('wp_kirki_failed_jobs', $this->queue->get_failed_table()); + $this->assertSame('kirki_jobs', $this->queue->get_table_name()); + } + + public function test_push_binds_the_computed_availability_and_returns_the_insert_id(): void + { + $id = $this->queue->push('{"job":"X"}', 'emails', 5, 60); + + $this->assertSame(7, $id); + $this->assertStringContainsString('INSERT INTO wp_kirki_jobs', $this->last_query()); + $this->assertStringContainsString("VALUES ('emails', 5, '{\\\"job\\\":\\\"X\\\"}', 0, 1700000060, 1700000000)", $this->last_query()); + } + + public function test_claim_is_one_ordered_limited_update_then_a_select_by_token(): void + { + $this->queue->claim(10, 'tok'); + + [$update, $select] = $this->wpdb->queries; + + $this->assertStringContainsString('UPDATE wp_kirki_jobs', $update); + $this->assertStringContainsString( + "SET reserved_at = 1700000000, reserved_by = 'tok', attempts = attempts + 1", + $update + ); + $this->assertStringContainsString('WHERE available_at <= 1700000000', $update); + $this->assertStringContainsString('AND (reserved_at IS NULL OR reserved_at <= 1699999700)', $update); + $this->assertStringContainsString('ORDER BY priority DESC, available_at ASC, id ASC', $update); + $this->assertStringContainsString('LIMIT 10', $update); + $this->assertStringNotContainsString('queue IN', $update); + + $this->assertStringContainsString("WHERE reserved_by = 'tok'", $select); + } + + public function test_claim_can_be_restricted_to_queues(): void + { + $this->queue->claim(5, 'tok', ['emails', 'default']); + + $this->assertStringContainsString("AND queue IN ('emails', 'default')", $this->wpdb->queries[0]); + $this->assertStringContainsString('LIMIT 5', $this->wpdb->queries[0]); + } + + public function test_has_due_is_a_single_row_probe(): void + { + $this->queue->has_due(); + + $this->assertStringContainsString('SELECT 1 FROM wp_kirki_jobs', $this->last_query()); + $this->assertStringContainsString('LIMIT 1', $this->last_query()); + } + + public function test_release_unstarted_refunds_the_attempt(): void + { + $this->queue->release_unstarted([3, 4]); + + $this->assertStringContainsString('attempts = IF(attempts > 0, attempts - 1, 0)', $this->last_query()); + $this->assertStringContainsString('WHERE id IN (3, 4)', $this->last_query()); + } + + public function test_release_unstarted_with_nothing_runs_no_query(): void + { + $this->queue->release_unstarted([]); + + $this->assertSame([], $this->wpdb->queries); + } + + public function test_fail_records_the_payload_uuid_and_deletes_the_job(): void + { + $record = new JobRecord(9, 'emails', '{"uuid":"abc-123","job":"X"}', 2); + + $this->queue->fail($record, new RuntimeException('nope')); + + [$insert, $delete] = $this->wpdb->queries; + $this->assertStringContainsString("INSERT INTO wp_kirki_failed_jobs", $insert); + $this->assertStringContainsString("'abc-123', 'emails'", $insert); + $this->assertStringContainsString('DELETE FROM wp_kirki_jobs WHERE id = 9', $delete); + } + + public function test_a_configured_table_name_is_used(): void + { + $queue = new class extends DatabaseQueue { + protected $options = ['table' => 'custom_jobs']; + }; + + $this->assertSame('wp_custom_jobs', $queue->get_table()); + } + + protected function last_query(): string + { + return (string) end($this->wpdb->queries); + } +} diff --git a/tests/Unit/Queue/DispatchTest.php b/tests/Unit/Queue/DispatchTest.php new file mode 100644 index 0000000..d4bbdb4 --- /dev/null +++ b/tests/Unit/Queue/DispatchTest.php @@ -0,0 +1,192 @@ +enable_queue(); + + SendEmail::dispatch(42); + + $this->assertCount(1, $this->queue->rows); + + $row = $this->queue->row(1); + $this->assertSame('default', $row['queue']); + $this->assertSame(0, $row['priority']); + $this->assertSame(0, $row['attempts']); + $this->assertNull($row['reserved_at']); + $this->assertSame(self::NOW, $row['available_at']); + $this->assertSame(SendEmail::class, $this->queue->payload_of(1)['job']); + } + + public function test_the_row_is_written_only_when_the_pending_dispatch_is_destroyed(): void + { + $this->enable_queue(); + + $pending = SendEmail::dispatch(42)->delay(10); + $this->assertCount(0, $this->queue->rows); + + unset($pending); + $this->assertCount(1, $this->queue->rows); + } + + public function test_delay_in_seconds(): void + { + $this->enable_queue(); + + SendEmail::dispatch(1)->delay(3600); + + $this->assertSame(self::NOW + 3600, $this->queue->row(1)['available_at']); + } + + public function test_delay_until_a_date(): void + { + $this->enable_queue(); + + SendEmail::dispatch(1)->delay(new DateTimeImmutable('@' . (self::NOW + 120))); + + $this->assertSame(self::NOW + 120, $this->queue->row(1)['available_at']); + } + + public function test_delay_by_an_interval(): void + { + $this->enable_queue(); + + SendEmail::dispatch(1)->delay(new DateInterval('PT5M')); + + $this->assertSame(self::NOW + 300, $this->queue->row(1)['available_at']); + } + + public function test_a_past_date_is_available_now(): void + { + $this->enable_queue(); + + SendEmail::dispatch(1)->delay(new DateTimeImmutable('@' . (self::NOW - 500))); + + $this->assertSame(self::NOW, $this->queue->row(1)['available_at']); + } + + public function test_queue_name_and_priority(): void + { + $this->enable_queue(); + + SendEmail::dispatch(1)->on_queue('emails')->with_priority(10); + + $this->assertSame('emails', $this->queue->row(1)['queue']); + $this->assertSame(10, $this->queue->row(1)['priority']); + } + + public function test_conditional_dispatch(): void + { + $this->enable_queue(); + + SendEmail::dispatch_if(false, 1); + SendEmail::dispatch_unless(true, 2); + $this->assertCount(0, $this->queue->rows); + + SendEmail::dispatch_if(function () { + return true; + }, 3); + SendEmail::dispatch_unless(false, 4); + $this->assertCount(2, $this->queue->rows); + } + + public function test_dispatch_sync_runs_immediately_without_storing_or_needing_the_provider(): void + { + $result = SendEmail::dispatch_sync(7); + + $this->assertSame('sent:7', $result); + $this->assertSame([7], Journal::values('handled')); + $this->assertCount(0, $this->queue->rows); + } + + public function test_dispatch_sync_calls_failed_then_rethrows(): void + { + try { + ThrowingJob::dispatch_sync(); + $this->fail('The exception was swallowed.'); + } catch (RuntimeException $exception) { + $this->assertSame('boom', $exception->getMessage()); + } + + $this->assertSame(['boom'], Journal::values('failed')); + } + + public function test_a_job_can_set_its_own_defaults_and_the_dispatch_can_override_them(): void + { + $this->enable_queue(); + + DefaultsJob::dispatch(); + DefaultsJob::dispatch()->on_queue('other')->with_priority(1)->delay(0); + + $this->assertSame('emails', $this->queue->row(1)['queue']); + $this->assertSame(5, $this->queue->row(1)['priority']); + $this->assertSame(self::NOW + 60, $this->queue->row(1)['available_at']); + + $this->assertSame('other', $this->queue->row(2)['queue']); + $this->assertSame(1, $this->queue->row(2)['priority']); + $this->assertSame(self::NOW, $this->queue->row(2)['available_at']); + } + + public function test_dispatch_without_the_provider_names_the_provider(): void + { + $this->expectException(QueueException::class); + $this->expectExceptionMessage('QueueServiceProvider'); + + SendEmail::dispatch(1); + } + + public function test_dispatch_with_a_missing_table_says_how_to_create_it(): void + { + $this->enable_queue(); + $this->queue->exists = false; + + try { + SendEmail::dispatch(1); + $this->fail('No exception was thrown.'); + } catch (QueueException $exception) { + $this->assertStringContainsString('wp_jobs', $exception->getMessage()); + $this->assertStringContainsString('queue:table', $exception->getMessage()); + } + + $this->assertCount(0, $this->queue->rows); + } + + public function test_the_payload_records_the_jobs_tries_or_the_configured_default(): void + { + $this->enable_queue(); + $this->queue->with_options(['tries' => 4, 'backoff' => 20]); + + ThrowingJob::dispatch(); + SendEmail::dispatch(1); + + $this->assertSame(3, $this->queue->payload_of(1)['max_tries']); + $this->assertSame([10, 60], $this->queue->payload_of(1)['backoff']); + $this->assertSame(4, $this->queue->payload_of(2)['max_tries']); + $this->assertSame(20, $this->queue->payload_of(2)['backoff']); + } + + public function test_job_queued_is_dispatched_when_listened_for(): void + { + $this->enable_queue(); + $this->events->listen_for(JobQueued::class); + + SendEmail::dispatch(1); + + $this->assertSame([JobQueued::class], $this->events->dispatched_classes()); + $this->assertSame(1, $this->events->dispatched[0]->id); + $this->assertInstanceOf(SendEmail::class, $this->events->dispatched[0]->job); + } +} diff --git a/tests/Unit/Queue/OptInTest.php b/tests/Unit/Queue/OptInTest.php new file mode 100644 index 0000000..376e1b3 --- /dev/null +++ b/tests/Unit/Queue/OptInTest.php @@ -0,0 +1,53 @@ +assertSame([], $this->queue_hooks()); + $this->assertSame([], $GLOBALS['framework_test_cron']); + + $this->expectException(QueueException::class); + + SendEmail::dispatch(1); + } + + public function test_the_provider_hooks_the_sweep_and_the_worker_endpoint(): void + { + $provider = $this->app->register(new QueueServiceProvider($this->app)); + $provider->boot(); + + $this->assertSame( + ['queue_sweep', 'wp_ajax_nopriv_queue_work', 'wp_ajax_queue_work'], + $this->queue_hooks() + ); + $this->assertSame('queue_every_minute', $GLOBALS['framework_test_cron']['queue_sweep']['recurrence']); + + SendEmail::dispatch(1); + $this->assertCount(1, $this->queue->rows); + } + + public function test_the_sweep_is_scheduled_only_once(): void + { + $provider = $this->app->register(new QueueServiceProvider($this->app)); + $provider->boot(); + $scheduled = $GLOBALS['framework_test_cron']['queue_sweep']['timestamp']; + + $provider->boot(); + + $this->assertSame($scheduled, $GLOBALS['framework_test_cron']['queue_sweep']['timestamp']); + } + + protected function queue_hooks(): array + { + return array_values(array_filter(array_keys($GLOBALS['framework_test_actions']), function ($hook) { + return strpos($hook, 'queue') !== false; + })); + } +} diff --git a/tests/Unit/Queue/PayloadTest.php b/tests/Unit/Queue/PayloadTest.php new file mode 100644 index 0000000..7b7d9c1 --- /dev/null +++ b/tests/Unit/Queue/PayloadTest.php @@ -0,0 +1,98 @@ +on_queue('emails')->with_priority(3); + + $restored = Payload::restore(Payload::decode(Payload::encode($job))); + + $this->assertInstanceOf(SendEmail::class, $restored); + $this->assertSame(42, $restored->user_id); + $this->assertSame('Welcome', $restored->subject); + $this->assertSame(['vip', 'new'], $restored->tags); + $this->assertSame('emails', $restored->get_queue()); + $this->assertSame(3, $restored->get_priority()); + } + + public function test_the_runtime_record_is_not_persisted(): void + { + $job = new SendEmail(1); + $job->set_job_record(new JobRecord(9, 'default', '{"secret":true}', 4)); + + $envelope = Payload::decode(Payload::encode($job)); + + $this->assertStringNotContainsString('JobRecord', $envelope['data']); + $this->assertSame(1, Payload::restore($envelope)->attempts()); + $this->assertSame(4, $job->attempts()); + } + + public function test_a_malformed_envelope_is_rejected(): void + { + $this->expectException(QueueException::class); + + Payload::decode('not json'); + } + + public function test_a_class_that_does_not_exist_fails_without_running(): void + { + $this->assert_row_fails_without_running(json_encode([ + 'job' => 'Framework\\Tests\\Support\\Queue\\Jobs\\Missing', + 'max_tries' => 1, + 'data' => 'O:0:"":0:{}', + ]), 'does not exist'); + } + + public function test_a_class_that_is_not_a_queued_job_fails_without_running(): void + { + $this->assert_row_fails_without_running(json_encode([ + 'job' => NotAJob::class, + 'max_tries' => 1, + 'data' => serialize(new NotAJob()), + ]), 'does not implement'); + } + + public function test_data_that_restores_to_a_different_class_fails_without_running(): void + { + $this->assert_row_fails_without_running(json_encode([ + 'job' => SendEmail::class, + 'max_tries' => 1, + 'data' => serialize(new NotAJob()), + ]), 'does not restore'); + } + + public function test_processing_continues_after_an_invalid_payload(): void + { + $this->enable_queue(); + $this->queue->push('garbage', 'default'); + SendEmail::dispatch(5); + + $this->worker->run(); + + $this->assertCount(1, $this->queue->failed); + $this->assertSame([5], Journal::values('handled')); + } + + protected function assert_row_fails_without_running(string $payload, string $reason): void + { + $this->queue->push($payload, 'default'); + + $this->worker->run(); + + $this->assertSame([], Journal::names()); + $this->assertCount(0, $this->queue->rows); + $this->assertCount(1, $this->queue->failed); + $this->assertStringContainsString($reason, $this->queue->failed[1]['exception']); + } +} diff --git a/tests/Unit/Queue/QueueCommandsTest.php b/tests/Unit/Queue/QueueCommandsTest.php new file mode 100644 index 0000000..928e825 --- /dev/null +++ b/tests/Unit/Queue/QueueCommandsTest.php @@ -0,0 +1,221 @@ +base = sys_get_temp_dir() . '/framework-queue-commands-' . uniqid(); + + mkdir($this->base . '/app', 0777, true); + mkdir($this->base . '/database/migrations', 0777, true); + file_put_contents($this->base . '/composer.json', json_encode([ + 'autoload' => [ + 'psr-4' => [ + 'Acme\\Shop\\' => 'app/', + 'Acme\\Shop\\Database\\Migrations\\' => 'database/migrations/', + ], + ], + ])); + + $this->reset_container_instance(); + + $app = Application::get_instance($this->base); + $app->use_prefix('shop_'); + $app->instance(Filesystem::class, new TestFilesystem()); + + return $app; + } + + protected function tearDown(): void + { + (new TestFilesystem())->delete($this->base, true); + + parent::tearDown(); + } + + public function test_queue_table_writes_two_migrations_without_running_any_sql(): void + { + (new QueueTableCommand())->run([], []); + + $jobs = $this->read('database/migrations/CreateJobsTable.php'); + $failed = $this->read('database/migrations/CreateFailedJobsTable.php'); + + $this->assertStringContainsString('namespace Acme\\Shop\\Database\\Migrations;', $jobs); + $this->assertStringContainsString('class CreateJobsTable implements Migration', $jobs); + $this->assertStringContainsString("Schema::create('shop_jobs'", $jobs); + $this->assertStringContainsString("\$table->string('reserved_by', 32)->nullable();", $jobs); + $this->assertStringContainsString("Schema::drop_if_exists('shop_jobs')", $jobs); + + $this->assertStringContainsString('class CreateFailedJobsTable implements Migration', $failed); + $this->assertStringContainsString("Schema::create('shop_failed_jobs'", $failed); + + $this->assert_valid_php($jobs); + $this->assert_valid_php($failed); + $this->assertContains( + ['line', ' \\Acme\\Shop\\Database\\Migrations\\CreateJobsTable::class,'], + \WP_CLI::$messages + ); + } + + public function test_queue_table_refuses_to_overwrite(): void + { + (new QueueTableCommand())->run([], []); + file_put_contents($this->base . '/database/migrations/CreateJobsTable.php', 'edited'); + + try { + (new QueueTableCommand())->run([], []); + $this->fail('The command did not stop.'); + } catch (RuntimeException $exception) { + $this->assertStringContainsString('CreateJobsTable.php', $exception->getMessage()); + } + + $this->assertSame('edited', $this->read('database/migrations/CreateJobsTable.php')); + } + + public function test_make_job_writes_a_dispatchable_job_class(): void + { + (new MakeJobCommand())->run(['send_abandoned_cart_email'], []); + + $job = $this->read('app/Jobs/SendAbandonedCartEmail.php'); + + $this->assertStringContainsString('namespace Acme\\Shop\\Jobs;', $job); + $this->assertStringContainsString('class SendAbandonedCartEmail implements ShouldQueue', $job); + $this->assertStringContainsString('use Queueable;', $job); + $this->assertStringContainsString('public function handle()', $job); + $this->assert_valid_php($job); + + $this->expectException(RuntimeException::class); + + (new MakeJobCommand())->run(['SendAbandonedCartEmail'], []); + } + + public function test_queue_work_once_processes_a_single_job(): void + { + $this->enable_queue(); + SendEmail::dispatch(1); + SendEmail::dispatch(2); + SendEmail::dispatch(3); + + (new QueueWorkCommand())->run([], ['once' => true]); + + $this->assertSame([1], Journal::values('handled')); + $this->assertCount(2, $this->queue->rows); + } + + public function test_queue_work_can_be_restricted_to_queues(): void + { + $this->enable_queue(); + SendEmail::dispatch('a'); + SendEmail::dispatch('b')->on_queue('emails'); + SendEmail::dispatch('c')->on_queue('reports'); + + (new QueueWorkCommand())->run([], ['queue' => 'emails, reports']); + + $this->assertSame(['b', 'c'], Journal::values('handled')); + $this->assertCount(1, $this->queue->rows); + } + + public function test_queue_retry_all_moves_failed_jobs_back_with_attempts_reset(): void + { + $this->enable_queue(); + SendEmail::dispatch(1)->with_priority(7); + SendEmail::dispatch(2)->on_queue('emails'); + $this->fail_row(1); + $this->fail_row(2); + + (new QueueRetryCommand())->run(['all'], []); + + $this->assertSame([], $this->queue->failed); + $this->assertCount(2, $this->queue->rows); + + $rows = array_values($this->queue->rows); + $this->assertSame([0, 0], array_column($rows, 'attempts')); + $this->assertEqualsCanonicalizing(['default', 'emails'], array_column($rows, 'queue')); + $this->assertContains(7, array_column($rows, 'priority')); + } + + public function test_queue_retry_one_leaves_the_others(): void + { + $this->enable_queue(); + SendEmail::dispatch(1); + SendEmail::dispatch(2); + $this->fail_row(1); + $this->fail_row(2); + + (new QueueRetryCommand())->run(['2'], []); + + $this->assertSame([1], array_keys($this->queue->failed)); + $this->assertCount(1, $this->queue->rows); + } + + public function test_queue_forget_and_flush(): void + { + $this->enable_queue(); + SendEmail::dispatch(1); + SendEmail::dispatch(2); + SendEmail::dispatch(3); + $this->fail_row(1); + $this->fail_row(2); + $this->fail_row(3); + + (new QueueForgetCommand())->run(['1'], []); + $this->assertSame([2, 3], array_keys($this->queue->failed)); + + (new QueueFlushCommand())->run([], []); + $this->assertSame([], $this->queue->failed); + + $this->expectException(RuntimeException::class); + + (new QueueForgetCommand())->run(['99'], []); + } + + public function test_queue_clear_deletes_pending_jobs_optionally_by_queue(): void + { + $this->enable_queue(); + SendEmail::dispatch(1); + SendEmail::dispatch(2)->on_queue('emails'); + SendEmail::dispatch(3)->on_queue('emails'); + + (new QueueClearCommand())->run([], ['queue' => 'emails']); + $this->assertCount(1, $this->queue->rows); + + (new QueueClearCommand())->run([], []); + $this->assertCount(0, $this->queue->rows); + } + + protected function fail_row(int $id): void + { + $this->queue->fail(JobRecord::from_row($this->queue->row($id)), new RuntimeException('failed')); + } + + protected function read(string $relative): string + { + return (string) file_get_contents($this->base . '/' . $relative); + } + + protected function assert_valid_php(string $code): void + { + token_get_all($code, TOKEN_PARSE); + + $this->addToAssertionCount(1); + } +} diff --git a/tests/Unit/Queue/QueueFakeTest.php b/tests/Unit/Queue/QueueFakeTest.php new file mode 100644 index 0000000..17bf597 --- /dev/null +++ b/tests/Unit/Queue/QueueFakeTest.php @@ -0,0 +1,88 @@ +on_queue('emails'); + + $this->assertCount(0, $this->queue->rows); + $this->assertArrayNotHasKey('shutdown', $GLOBALS['framework_test_actions']); + $this->assertSame(1, Queue::size('emails')); + $this->assertTrue(Queue::is_faked()); + } + + public function test_assertions_pass_for_what_was_pushed(): void + { + Queue::fake(); + + SendEmail::dispatch(123); + + Queue::assert_pushed(SendEmail::class); + Queue::assert_pushed(SendEmail::class, function (SendEmail $job) { + return $job->user_id === 123; + }); + Queue::assert_pushed_times(SendEmail::class, 1); + Queue::assert_not_pushed(ThrowingJob::class); + Queue::assert_not_pushed(SendEmail::class, function (SendEmail $job) { + return $job->user_id === 999; + }); + + $this->addToAssertionCount(5); + } + + public function test_assert_pushed_fails_when_nothing_matches(): void + { + Queue::fake(); + SendEmail::dispatch(1); + + $this->expectException(AssertionError::class); + + Queue::assert_pushed(SendEmail::class, function (SendEmail $job) { + return $job->user_id === 2; + }); + } + + public function test_assert_pushed_times_fails_on_a_different_count(): void + { + Queue::fake(); + SendEmail::dispatch(1); + SendEmail::dispatch(2); + + $this->expectException(AssertionError::class); + $this->expectExceptionMessage('pushed 2 times instead of 1 times'); + + Queue::assert_pushed_times(SendEmail::class, 1); + } + + public function test_assert_not_pushed_fails_when_it_was(): void + { + Queue::fake(); + SendEmail::dispatch(1); + + $this->expectException(AssertionError::class); + + Queue::assert_not_pushed(SendEmail::class); + } + + public function test_assert_nothing_pushed(): void + { + Queue::fake(); + Queue::assert_nothing_pushed(); + + SendEmail::dispatch(1); + + $this->expectException(AssertionError::class); + + Queue::assert_nothing_pushed(); + } +} diff --git a/tests/Unit/Queue/QueueTestCase.php b/tests/Unit/Queue/QueueTestCase.php new file mode 100644 index 0000000..63c008a --- /dev/null +++ b/tests/Unit/Queue/QueueTestCase.php @@ -0,0 +1,93 @@ +reset_queue_globals(); + $this->app = $this->make_application(); + + $this->queue = (new ArrayDatabaseQueue())->freeze(self::NOW); + $this->worker = new TestWorker($this->queue); + $this->spawner = (new RecordingSpawner($this->queue, $this->worker))->freeze(self::NOW); + $this->events = new RecordingEventManager(); + $this->log = new RecordingLogger(); + + $this->app->instance(DatabaseQueue::class, $this->queue); + $this->app->instance(Worker::class, $this->worker); + $this->app->instance(Spawner::class, $this->spawner); + $this->app->instance(EventManager::class, $this->events); + $this->app->instance(LogManager::class, $this->log); + } + + protected function tearDown(): void + { + $this->reset_queue_globals(); + + parent::tearDown(); + } + + protected function make_application(): Application + { + return $this->bootstrap_application(); + } + + protected function enable_queue(): QueueManager + { + $manager = new QueueManager($this->queue, $this->spawner); + $this->app->instance(QueueManager::class, $manager); + + return $manager; + } + + protected function travel(int $seconds): void + { + $this->queue->travel($seconds); + $this->spawner->travel($seconds); + } + + protected function reset_queue_globals(): void + { + Journal::$entries = []; + TestLock::$held = []; + \WP_CLI::$messages = []; + + $GLOBALS['framework_test_actions'] = []; + $GLOBALS['framework_test_cron'] = []; + $GLOBALS['framework_test_salt'] = 'framework-test-salt'; + } +} diff --git a/tests/Unit/Queue/SpawnerTest.php b/tests/Unit/Queue/SpawnerTest.php new file mode 100644 index 0000000..8ac3c13 --- /dev/null +++ b/tests/Unit/Queue/SpawnerTest.php @@ -0,0 +1,168 @@ +enable_queue(); + } + + public function test_a_correctly_signed_request_runs_a_worker(): void + { + SendEmail::dispatch(1); + + $this->assertTrue($this->spawner->handle_request($this->signed(self::NOW))); + $this->assertSame([1], Journal::values('handled')); + } + + public function test_an_unsigned_request_is_rejected_before_claiming(): void + { + SendEmail::dispatch(1); + + $this->assertFalse($this->spawner->handle_request(['action' => 'queue_work'])); + $this->assert_nothing_claimed(); + } + + public function test_a_tampered_signature_is_rejected(): void + { + SendEmail::dispatch(1); + + $request = $this->signed(self::NOW); + $request['sig'] = str_repeat('0', 64); + + $this->assertFalse($this->spawner->handle_request($request)); + $this->assert_nothing_claimed(); + } + + public function test_a_signature_for_another_timestamp_is_rejected(): void + { + SendEmail::dispatch(1); + + $request = $this->signed(self::NOW); + $request['ts'] = self::NOW + 1; + + $this->assertFalse($this->spawner->handle_request($request)); + $this->assert_nothing_claimed(); + } + + public function test_an_expired_signature_is_rejected(): void + { + SendEmail::dispatch(1); + + $this->assertFalse($this->spawner->handle_request($this->signed(self::NOW - 61))); + $this->assertTrue($this->spawner->handle_request($this->signed(self::NOW - 60))); + } + + public function test_the_sweeper_does_nothing_when_nothing_is_due(): void + { + SendEmail::dispatch(1)->delay(3600); + + $this->assertFalse($this->sweeper()->sweep()); + $this->assertSame(0, $this->spawner->spawned); + } + + public function test_the_sweeper_spawns_once_a_delayed_job_becomes_due(): void + { + SendEmail::dispatch(1)->delay(3600); + $this->travel(3600); + + $this->assertTrue($this->sweeper()->sweep()); + $this->assertSame(1, $this->spawner->spawned); + $this->assertSame([], TestLock::$held); + } + + public function test_the_sweeper_skips_a_missing_table(): void + { + $this->queue->exists = false; + + $this->assertFalse($this->sweeper()->sweep()); + $this->assertSame(0, $this->spawner->spawned); + } + + public function test_no_worker_is_spawned_or_run_while_a_chain_holds_the_lock(): void + { + SendEmail::dispatch(1); + $this->spawner->chain_lock()->acquire(); + + $this->assertFalse($this->sweeper()->sweep()); + $this->assertFalse($this->spawner->handle_request($this->signed(self::NOW))); + $this->assertSame(0, $this->spawner->spawned); + $this->assert_nothing_claimed(); + } + + public function test_a_worker_chains_a_successor_when_work_remains(): void + { + $this->queue->with_options(['time_limit' => 2]); + $this->worker->seconds_per_job = 1.0; + + for ($i = 1; $i <= 5; $i++) { + SendEmail::dispatch($i); + } + + $this->spawner->handle_request($this->signed(self::NOW)); + + $this->assertSame([1, 2], Journal::values('handled')); + $this->assertSame(1, $this->spawner->spawned); + $this->assertSame([], TestLock::$held); + } + + public function test_the_last_worker_releases_the_lock_without_spawning(): void + { + SendEmail::dispatch(1); + + $this->spawner->handle_request($this->signed(self::NOW)); + + $this->assertSame(0, $this->spawner->spawned); + $this->assertSame([], TestLock::$held); + } + + public function test_undelayed_dispatches_spawn_once_on_shutdown(): void + { + SendEmail::dispatch(1); + SendEmail::dispatch(2); + SendEmail::dispatch(3); + + $this->assertCount(1, $GLOBALS['framework_test_actions']['shutdown']); + $this->assertSame(0, $this->spawner->spawned); + + framework_test_do_action('shutdown'); + + $this->assertSame(1, $this->spawner->spawned); + } + + public function test_delayed_dispatches_do_not_spawn_on_shutdown(): void + { + SendEmail::dispatch(1)->delay(60); + + $this->assertArrayNotHasKey('shutdown', $GLOBALS['framework_test_actions']); + } + + protected function signed(int $timestamp): array + { + return [ + 'action' => 'queue_work', + 'ts' => (string) $timestamp, + 'sig' => $this->spawner->sign($timestamp), + ]; + } + + protected function sweeper(): Sweeper + { + return new Sweeper($this->queue, $this->spawner); + } + + protected function assert_nothing_claimed(): void + { + $this->assertSame([], $this->queue->claims); + $this->assertSame([], Journal::names()); + } +} diff --git a/tests/Unit/Queue/StaleRecoveryTest.php b/tests/Unit/Queue/StaleRecoveryTest.php new file mode 100644 index 0000000..fdb53d2 --- /dev/null +++ b/tests/Unit/Queue/StaleRecoveryTest.php @@ -0,0 +1,73 @@ +enable_queue(); + $this->queue->with_options(['retry_after' => 300]); + } + + public function test_a_live_reservation_is_not_claimed_again(): void + { + SendEmail::dispatch(1); + $this->queue->claim(10, 'crashed-worker'); + + $this->travel(299); + $this->worker->run(); + + $this->assertSame([], Journal::values('handled')); + $this->assertSame('crashed-worker', $this->queue->row(1)['reserved_by']); + } + + public function test_an_abandoned_reservation_is_reclaimed_and_counts_an_attempt(): void + { + AttemptsJob::dispatch(); + $this->queue->claim(10, 'crashed-worker'); + + $this->travel(300); + $this->assertTrue($this->queue->has_due()); + + $this->worker->run(); + + $this->assertSame([2], Journal::values('attempts')); + $this->assertCount(0, $this->queue->rows); + } + + public function test_a_single_try_job_whose_worker_crashed_is_failed_not_rerun(): void + { + SendEmail::dispatch(1); + $this->queue->claim(10, 'crashed-worker'); + + $this->travel(300); + $this->worker->run(); + + $this->assertSame([], Journal::values('handled')); + $this->assertCount(1, $this->queue->failed); + } + + public function test_a_crash_loop_ends_in_the_failed_jobs_table_without_running_again(): void + { + AttemptsJob::dispatch(); + + $this->queue->claim(10, 'crash-1'); + $this->travel(300); + $this->queue->claim(10, 'crash-2'); + $this->travel(300); + + $this->worker->run(); + + $this->assertSame([], Journal::values('attempts')); + $this->assertCount(0, $this->queue->rows); + $this->assertCount(1, $this->queue->failed); + $this->assertStringContainsString('attempted too many times', $this->queue->failed[1]['exception']); + } +} diff --git a/tests/Unit/Queue/WorkerTest.php b/tests/Unit/Queue/WorkerTest.php new file mode 100644 index 0000000..21727de --- /dev/null +++ b/tests/Unit/Queue/WorkerTest.php @@ -0,0 +1,224 @@ +enable_queue(); + } + + public function test_jobs_run_highest_priority_first(): void + { + SendEmail::dispatch('low')->with_priority(0); + SendEmail::dispatch('high')->with_priority(10); + SendEmail::dispatch('middle')->with_priority(5); + + $this->worker->run(); + + $this->assertSame(['high', 'middle', 'low'], Journal::values('handled')); + } + + public function test_equal_priority_runs_in_availability_then_insertion_order(): void + { + SendEmail::dispatch('b'); + SendEmail::dispatch('a'); + SendEmail::dispatch('c'); + $this->queue->rows[2]['available_at'] = self::NOW - 10; + + $this->worker->run(); + + $this->assertSame(['a', 'b', 'c'], Journal::values('handled')); + } + + public function test_a_claim_reserves_at_most_the_batch_size_and_one_run_takes_several_batches(): void + { + $this->queue->with_options(['batch_size' => 10]); + + for ($i = 1; $i <= 25; $i++) { + SendEmail::dispatch($i); + } + + $remaining = $this->worker->run(); + + $this->assertSame([10, 10, 5, 0], $this->queue->claims); + $this->assertCount(25, Journal::values('handled')); + $this->assertFalse($remaining); + $this->assertCount(0, $this->queue->rows); + } + + public function test_an_exhausted_budget_hands_back_unstarted_jobs_without_using_an_attempt(): void + { + $this->queue->with_options(['batch_size' => 10, 'time_limit' => 3]); + $this->worker->seconds_per_job = 1.0; + + for ($i = 1; $i <= 10; $i++) { + SendEmail::dispatch($i); + } + + $remaining = $this->worker->run(); + + $this->assertTrue($remaining); + $this->assertSame([1, 2, 3], Journal::values('handled')); + $this->assertCount(7, $this->queue->rows); + + foreach ($this->queue->rows as $row) { + $this->assertNull($row['reserved_at']); + $this->assertNull($row['reserved_by']); + $this->assertSame(0, $row['attempts']); + } + } + + public function test_a_job_limit_stops_the_run(): void + { + for ($i = 1; $i <= 3; $i++) { + SendEmail::dispatch($i); + } + + $remaining = $this->worker->run(['max_jobs' => 1]); + + $this->assertTrue($remaining); + $this->assertSame([1], Journal::values('handled')); + $this->assertSame([1], $this->queue->claims); + } + + public function test_a_successful_job_is_deleted_and_announced(): void + { + $this->events->listen_for(JobProcessing::class)->listen_for(JobProcessed::class); + + SendEmail::dispatch(1); + $this->worker->run(); + + $this->assertCount(0, $this->queue->rows); + $this->assertSame([JobProcessing::class, JobProcessed::class], $this->events->dispatched_classes()); + } + + public function test_a_throwing_job_is_retried_with_its_backoff_then_failed(): void + { + $this->events->listen_for(JobFailed::class); + + ThrowingJob::dispatch(); + + $this->worker->run(); + $this->assertSame(1, $this->queue->row(1)['attempts']); + $this->assertSame(self::NOW + 10, $this->queue->row(1)['available_at']); + $this->assertNull($this->queue->row(1)['reserved_at']); + + $this->travel(10); + $this->worker->run(); + $this->assertSame(2, $this->queue->row(1)['attempts']); + $this->assertSame(self::NOW + 10 + 60, $this->queue->row(1)['available_at']); + + $this->travel(60); + $this->worker->run(); + + $this->assertSame([1, 2, 3], Journal::values('attempt')); + $this->assertSame(['boom'], Journal::values('failed')); + $this->assertCount(0, $this->queue->rows); + $this->assertCount(1, $this->queue->failed); + + $failed = $this->queue->failed[1]; + $this->assertSame('default', $failed['queue']); + $this->assertSame(self::NOW + 70, $failed['failed_at']); + $this->assertStringContainsString('boom', $failed['exception']); + $this->assertSame($this->queue_uuid($failed['payload']), $failed['uuid']); + + $this->assertSame([JobFailed::class], $this->events->dispatched_classes()); + $this->assertCount(1, $this->log->errors); + $this->assertStringContainsString(ThrowingJob::class, $this->log->errors[0]); + $this->assertStringContainsString('boom', $this->log->errors[0]); + } + + public function test_a_single_backoff_value_and_default_tries_come_from_config(): void + { + $this->queue->with_options(['tries' => 2, 'backoff' => 15]); + + FailedThrowsJob::dispatch(); + $this->worker->run(); + + $this->assertSame(self::NOW + 15, $this->queue->row(1)['available_at']); + + $this->travel(15); + $this->worker->run(); + + $this->assertCount(1, $this->queue->failed); + } + + public function test_one_failing_job_does_not_stop_the_batch(): void + { + SendEmail::dispatch(1); + FailedThrowsJob::dispatch(); + SendEmail::dispatch(3); + + $this->worker->run(); + + $this->assertSame([1, 3], Journal::values('handled')); + $this->assertCount(1, $this->queue->failed); + } + + public function test_a_job_can_release_itself(): void + { + ReleasingJob::dispatch(); + + $this->worker->run(); + + $row = $this->queue->row(1); + $this->assertNotNull($row); + $this->assertNull($row['reserved_at']); + $this->assertSame(self::NOW + 30, $row['available_at']); + $this->assertSame([], Journal::names()); + $this->assertCount(0, $this->queue->failed); + } + + public function test_a_job_can_fail_itself_with_tries_remaining(): void + { + SelfFailingJob::dispatch(); + + $this->worker->run(); + + $this->assertCount(0, $this->queue->rows); + $this->assertCount(1, $this->queue->failed); + $this->assertSame(['gave up'], Journal::values('failed')); + } + + public function test_attempts_is_visible_inside_handle(): void + { + AttemptsJob::dispatch(); + + $this->worker->run(); + $this->worker->run(); + + $this->assertSame([1, 2], Journal::values('attempts')); + $this->assertCount(0, $this->queue->rows); + $this->assertCount(0, $this->queue->failed); + } + + public function test_handle_receives_container_resolved_dependencies(): void + { + InjectedJob::dispatch(); + + $this->worker->run(); + + $this->assertSame(['default mailer'], Journal::values('injected')); + } + + protected function queue_uuid(string $payload): string + { + return json_decode($payload, true)['uuid']; + } +}