diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 9e93f71..adf8c21 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -18,7 +18,17 @@ jobs: steps: - name: Setup MySQL latest if: matrix.db-type == 'mysql' - run: docker run --rm --name=mysqld -e MYSQL_ROOT_PASSWORD=root -e MYSQL_DATABASE=cakephp -p 3306:3306 -d mysql --default-authentication-plugin=mysql_native_password --disable-log-bin + run: docker run --rm --name=mysqld -e MYSQL_ROOT_PASSWORD=root -e MYSQL_DATABASE=cakephp -p 3306:3306 -d mysql --disable-log-bin + + - name: Wait for MySQL + if: matrix.db-type == 'mysql' + run: | + for i in $(seq 1 60); do + docker exec mysqld mysqladmin ping -h 127.0.0.1 -uroot -proot --silent && exit 0 + sleep 2 + done + docker logs mysqld + exit 1 - name: Setup PostgreSQL latest if: matrix.db-type == 'pgsql' @@ -62,9 +72,9 @@ jobs: - name: Run PHPUnit run: | - if [[ ${{ matrix.db-type }} == 'sqlite' ]]; then export DB_URL='sqlite:///:memory:'; fi - if [[ ${{ matrix.db-type }} == 'mysql' ]]; then export DB_URL='mysql://root:root@127.0.0.1/cakephp'; fi - if [[ ${{ matrix.db-type }} == 'pgsql' ]]; then export DB_URL='postgres://postgres:postgres@127.0.0.1/postgres'; fi + if [[ ${{ matrix.db-type }} == 'sqlite' ]]; then export db_dsn='sqlite:///:memory:'; fi + if [[ ${{ matrix.db-type }} == 'mysql' ]]; then export db_dsn='mysql://root:root@127.0.0.1/cakephp'; fi + if [[ ${{ matrix.db-type }} == 'pgsql' ]]; then export db_dsn='postgres://postgres:postgres@127.0.0.1/postgres'; fi if [[ ${{ matrix.php-version }} == '8.2' ]]; then export CODECOVERAGE=1 && vendor/bin/phpunit --coverage-clover=coverage.xml --configuration phpunit.xml.dist else diff --git a/Docs/Documentation/Caveats.md b/Docs/Documentation/Caveats.md new file mode 100644 index 0000000..9951e42 --- /dev/null +++ b/Docs/Documentation/Caveats.md @@ -0,0 +1,37 @@ +Known Caveats +============= + +Things that are easy to get wrong with this transport. Each item links to the page with the details. + +Configuration +------------- + +* `polling_interval` does not affect `bin/cake queue worker`, which uses `subscription_polling_interval` (200 ms). [Configuration](Configuration.md#options) +* With an array `url`, set the `queue` key (or a `client` key), otherwise the Enqueue client cannot be built. [Configuration](Configuration.md#array-configuration) +* With an array `transport`, the DSN host and query are ignored; set `connection` and `table_name` as keys. [Configuration](Configuration.md#array-configuration) +* Only `cakephp://` DSNs are accepted; `cakephpenqueue://` is registered but rejected by the factory. [Overview](Overview.md) +* `Cake/Queue` must find the `Queue` configuration when it is loaded. [Installation](Installation.md#configure-the-queue) + +Queue names and routing +----------------------- + +* Rows are stored under the prefixed queue name `enqueue.app.`, not under the `queue` option passed to `push()`. [Queue names](Usage.md#queue-names) +* A job pushed with a `queue` option different from the worker's queue is routed to nobody and dropped. [Queue names](Usage.md#queue-names) +* `purgeQueue($context->createQueue('default'))` removes nothing; build the destination through the driver. [Emptying a queue](Usage.md#emptying-a-queue) +* `delay`, `expires` and `priority` only take effect once a worker has routed the job. [Queue names](Usage.md#queue-names) + +Processing +---------- + +* A job running longer than `redelivery_delay` (20 minutes by default) is delivered again to another worker. [Message Lifecycle](Message-Lifecycle.md#failed-and-rejected-messages) +* Without `$maxAttempts` or `--max-attempts` a failing job is requeued forever. [Message Lifecycle](Message-Lifecycle.md#failed-and-rejected-messages) +* `--max-jobs` is enforced per consume cycle, so a worker may process more jobs than the limit. [Process jobs](Usage.md#process-jobs) +* An expired message claimed before the housekeeping run is kept and delivered after `redelivery_delay`. [Time to live](Message-Lifecycle.md#time-to-live) + +Database +-------- + +* Automatic table creation only supports MySQL/MariaDB, PostgreSQL and SQLite; other drivers get a warning. [Database Schema](Database-Schema.md) +* The table is created when the configuration is first used, so the database user needs `CREATE TABLE` then; the low level transport never creates it. [Installation](Installation.md#the-queue-table), [Enqueue Client](Enqueue-Client.md#low-level-transport) +* On SQLite a second queue table in the same database cannot be created automatically (index name collision). [Database Schema](Database-Schema.md#indexes) +* Creating the table also creates the migrations bookkeeping table (`cake_migrations`). [Database Schema](Database-Schema.md#indexes) diff --git a/Docs/Documentation/Configuration.md b/Docs/Documentation/Configuration.md new file mode 100644 index 0000000..3a09692 --- /dev/null +++ b/Docs/Documentation/Configuration.md @@ -0,0 +1,119 @@ +Configuration +============= + +The transport is configured through the `url` of a `cakephp/queue` configuration, either as a DSN string or as an array. + +DSN +--- + +``` +cakephp://?table_name=&polling_interval=&lazy= +``` + +```php +// config/app.php +'Queue' => [ + 'default' => [ + 'url' => 'cakephp://default?table_name=enqueue', + 'queue' => 'default', + ], + 'reports' => [ + 'url' => 'cakephp://default?table_name=enqueue&lazy=false', + 'queue' => 'reports', + ], +], +``` + +The host part is the name of the CakePHP datasource connection (`ConnectionManager::get()`); when it is empty `default` is used. Only the `cakephp` scheme is accepted, any other scheme throws `Wrong dsn schema passed`. + +Options +------- + +| Option | Type | Default | Description | +| :----- | :--- | :------ | :---------- | +| connection | string | `default` | Name of the CakePHP datasource connection. In a DSN it is the host part. | +| table_name | string | `enqueue` | Table where messages are stored. | +| lazy | bool | `true` | Connect to the database only when it is first needed. | +| polling_interval | int (ms) | `1000` | Sleep between polls of `CakeConsumer::receive()`. Not used by the worker, see below. | +| subscription_polling_interval | int (ms) | `200` | Sleep between polls of the subscription consumer, which is what `bin/cake queue worker` uses when the queue is empty. | +| redelivery_delay | int (ms) | `1200000` (20 minutes) | Time a received message can stay unacknowledged before it is redelivered to another consumer. Values that are not a multiple of 1000 are truncated to whole seconds. | + +NOTE: from a DSN, `polling_interval` and `lazy` are cast when parsed; `redelivery_delay` and `subscription_polling_interval` stay strings and are cast to integers when the consumers are created, so all options can be set in the DSN. + +Array configuration +------------------- + +`cakephp/queue` also accepts an array as `url`. The `transport` key can be an array with the options above: + +```php +'Queue' => [ + 'default' => [ + 'url' => [ + 'transport' => [ + 'dsn' => 'cakephp:', + 'connection' => 'default', + 'table_name' => 'enqueue', + 'redelivery_delay' => 600000, + 'subscription_polling_interval' => 500, + ], + ], + 'queue' => 'default', + ], +], +``` + +NOTE: when `transport` is an array with keys other than `dsn`, the whole array is passed to `CakeConnectionFactory` and the DSN is only used to select the transport: its host and query string are ignored. Set `connection` and `table_name` as keys. + +NOTE: with the array form, `cakephp/queue` only fills in the Enqueue `client` section when the `queue` key is set. Without `queue` you must add `'client' => true` (or an array of client options) next to `transport`, otherwise `SimpleClient` fails with a `TypeError` on `Config::__construct()`. + +If you use the transport classes directly, `CakeConnectionFactory` accepts the same options array or a DSN string: + +```php +use Cake\Enqueue\CakeConnectionFactory; + +$factory = new CakeConnectionFactory([ + 'connection' => 'default', + 'table_name' => 'enqueue', + 'redelivery_delay' => 600000, +]); +``` + +Dedicated connection +-------------------- + +Because the queue table is read and written constantly, you may want it in its own database. Define an additional datasource and point the DSN to it: + +```php +// config/app_local.php +'Datasources' => [ + 'default' => [ /* ... */ ], + 'queue' => [ + 'className' => Connection::class, + 'driver' => Postgres::class, + 'host' => 'localhost', + 'database' => 'my_app_queue', + ], +], +``` + +```php +// config/app.php +'Queue' => [ + 'default' => [ + 'url' => 'cakephp://queue?table_name=enqueue', + 'queue' => 'default', + ], +], +``` + +With a dedicated connection the message insert is not part of transactions opened on your application connection. Sharing the connection means a rollback also removes the queued message, which may or may not be what you want. + +Several queue configurations +---------------------------- + +Different configurations can share the same table as long as they use different `queue` names: each one stores its messages under its own queue (prefixed, see [Queue names](Usage.md#queue-names)). They can also use different tables or connections. One worker process serves one configuration, start one `bin/cake queue worker --config=` per configuration. + +Plugin name +----------- + +The plugin name is `Cake/Enqueue`, derived from its namespace, not the composer package name. Use it in `addPlugin()` and in table registry aliases, for example `Cake/Enqueue.Enqueue`. diff --git a/Docs/Documentation/Database-Schema.md b/Docs/Documentation/Database-Schema.md new file mode 100644 index 0000000..4ed7f49 --- /dev/null +++ b/Docs/Documentation/Database-Schema.md @@ -0,0 +1,108 @@ +Database Schema +=============== + +The table is created by `CakeContext::createDataBaseTable()`, which the client driver calls from `setupBroker()` the first time a queue configuration is used (first `QueueManager::push()` or worker start). The name comes from the `table_name` option (default `enqueue`). The same structure is created on MySQL/MariaDB, PostgreSQL and SQLite using the `cakephp/migrations` adapters; any other driver triggers a warning and the table must be created by hand. + +Columns +------- + +| Column | Type | Null | Description | +| :----- | :--- | :--- | :---------- | +| id | uuid | no | Primary key. | +| published_at | biginteger | no | Publish time as `microtime(true) * 10000` (tenths of a millisecond). Used for ordering. | +| body | text | yes | Message body. | +| headers | text | yes | JSON encoded headers. | +| properties | text | yes | JSON encoded properties. `cakephp/queue` keeps the topic, delay, expire and priority here. | +| redelivered | boolean | yes | `true` after the claim of an unacknowledged message was released. | +| queue | string | no | Queue name, prefixed by the Enqueue client, for example `enqueue.app.default`. | +| priority | integer(5) | yes | Stored negated, so the default ascending order delivers the highest priority first. | +| delayed_until | biginteger | yes | Unix time before which the message is not delivered. | +| time_to_live | biginteger | yes | Unix time after which an unclaimed message is removed. | +| delivery_id | uuid | yes | Set while a consumer holds the message. | +| redeliver_after | biginteger | yes | Unix time after which a claimed message is made available again. | + +Indexes +------- + +| Name | Columns | +| :--- | :------ | +| priority_idx | priority, published_at, queue, delivery_id, delayed_until, id | +| redeliver_idx | redeliver_after, delivery_id | +| ttl_idx | time_to_live, delivery_id | +| delivery_id_idx | delivery_id | + +NOTE: SQLite index names are global per database. Automatic creation of a second queue table in the same SQLite database fails on `priority_idx already exists` with a warning, leaving the second table without indexes. Use one table per SQLite database, or create the second table yourself with other index names. + +NOTE: the migrations adapter also creates its own bookkeeping table (`cake_migrations` with `cakephp/migrations` 5) in the queue connection if it does not exist yet. + +Creating the table yourself +--------------------------- + +To create the table outside of the worker, for example while deploying, either trigger the normal setup: + +```php +use Cake\Queue\QueueManager; + +QueueManager::engine('default'); // calls setupBroker(), creates the table if missing +``` + +or use the transport directly: + +```php +use Cake\Enqueue\CakeConnectionFactory; + +$context = (new CakeConnectionFactory('cakephp://default?table_name=enqueue'))->createContext(); +$context->createDataBaseTable(); // does nothing if the table exists +``` + +If you prefer a migration in your application, this one produces the same structure: + +```php +use Migrations\BaseMigration; + +class CreateEnqueue extends BaseMigration +{ + public function change(): void + { + $table = $this->table('enqueue', ['id' => false, 'primary_key' => ['id']]); + $table + ->addColumn('id', 'uuid') + ->addColumn('published_at', 'biginteger') + ->addColumn('body', 'text', ['null' => true]) + ->addColumn('headers', 'text', ['null' => true]) + ->addColumn('properties', 'text', ['null' => true]) + ->addColumn('redelivered', 'boolean', ['null' => true]) + ->addColumn('queue', 'string') + ->addColumn('priority', 'integer', ['limit' => 5, 'null' => true]) + ->addColumn('delayed_until', 'biginteger', ['null' => true]) + ->addColumn('time_to_live', 'biginteger', ['null' => true]) + ->addColumn('delivery_id', 'uuid', ['null' => true]) + ->addColumn('redeliver_after', 'biginteger', ['null' => true]) + ->addIndex(['priority', 'published_at', 'queue', 'delivery_id', 'delayed_until', 'id'], ['name' => 'priority_idx']) + ->addIndex(['redeliver_after', 'delivery_id'], ['name' => 'redeliver_idx']) + ->addIndex(['time_to_live', 'delivery_id'], ['name' => 'ttl_idx']) + ->addIndex(['delivery_id'], ['name' => 'delivery_id_idx']) + ->create(); + } +} +``` + +Reading messages +---------------- + +The `Cake/Enqueue.Enqueue` table class can be used to inspect the queue, for example to build a monitoring page. Remember that queue names are prefixed: + +```php +$table = $this->fetchTable('Cake/Enqueue.Enqueue'); + +$pending = $table->find() + ->where(['queue' => 'enqueue.app.default', 'delivery_id IS' => null]) + ->count(); + +$table->getSize('enqueue.app.default'); // rows in that queue, claimed ones included +$table->getAllSizes(); // [['queue' => 'enqueue.app.default', 'job_count' => 3], ...] +``` + +If `table_name` is not `enqueue`, call `$table->setTable('my_queue_table')` first, or get the configured instance from the context: `QueueManager::engine('default')->getDriver()->getContext()->getTable()`. + +NOTE: Do not modify rows directly while workers are running. Use the producer and consumer API instead. diff --git a/Docs/Documentation/Enqueue-Client.md b/Docs/Documentation/Enqueue-Client.md new file mode 100644 index 0000000..4eb5c37 --- /dev/null +++ b/Docs/Documentation/Enqueue-Client.md @@ -0,0 +1,80 @@ +Using the Enqueue Client +======================== + +The plugin also registers a driver for the Enqueue client, so you can use `enqueue/simple-client` with the database transport without `cakephp/queue`. The plugin must be loaded (its `bootstrap()` registers the transport), the client is a separate package: + +``` +composer require enqueue/simple-client +``` + +```php +use Enqueue\SimpleClient\SimpleClient; +use Interop\Queue\Message; +use Interop\Queue\Processor; + +$client = new SimpleClient('cakephp://default?table_name=enqueue'); + +$client->bindTopic('orders', function (Message $message) { + // process $message->getBody() + + return Processor::ACK; +}); + +$client->setupBroker(); // creates the table if missing +$client->sendEvent('orders', 'order #1'); +$client->consume(); +``` + +`cakephp/queue` wraps this same client: `QueueManager::engine('default')` returns a `SimpleClient` built from your `Queue` configuration, with the client's `router_topic`, `router_queue` and `default_queue` set to the configuration's `queue`. + +NOTE: the client routes each event through its router queue: `sendEvent()` stores one row in the router queue, the consumer's router processor then inserts the routed copy in the processor queue and deletes the router row. Both queues are `enqueue.app.default` by default (prefix `enqueue`, app name `app`, queue `default`). A consumer therefore handles two messages per event; keep that in mind when limiting consumption, for example with `--max-jobs` or `LimitConsumedMessagesExtension`. + +Low level transport +------------------- + +You can also work with the `queue-interop` API directly. Nothing creates the table for you at this level, call `createDataBaseTable()` once: + +```php +use Cake\Enqueue\CakeConnectionFactory; + +$factory = new CakeConnectionFactory('cakephp://default?table_name=enqueue'); +$context = $factory->createContext(); +$context->createDataBaseTable(); + +$queue = $context->createQueue('emails'); + +$producer = $context->createProducer(); +$producer->setDeliveryDelay(5000); // milliseconds, optional +$producer->setTimeToLive(60000); // milliseconds, optional +$producer->setPriority(3); // any integer, optional +$producer->send($queue, $context->createMessage('hello', ['key' => 'value'])); + +$consumer = $context->createConsumer($queue); +if ($message = $consumer->receive(2000)) { // wait up to 2000 ms + // ... + $consumer->acknowledge($message); + // or $consumer->reject($message); // drop + // or $consumer->reject($message, true); // requeue as a new row +} +``` + +At this level queue names are used as given (no `enqueue.app.` prefix), delays and TTLs are milliseconds (truncated to whole seconds when stored) and priorities are plain integers, higher first. + +To consume several queues in one loop use the subscription consumer: + +```php +$subscription = $context->createSubscriptionConsumer(); +$subscription->subscribe($context->createConsumer($context->createQueue('emails')), function ($message, $consumer) { + $consumer->acknowledge($message); + + return true; // return false to stop consuming +}); +$subscription->subscribe($context->createConsumer($context->createQueue('reports')), function ($message, $consumer) { + $consumer->acknowledge($message); + + return true; +}); +$subscription->consume(10000); // run for 10000 ms, 0 = forever +``` + +The subscription consumer polls every `subscription_polling_interval` ms (200 by default) when all queues are empty, see [Configuration](Configuration.md#options). Only one consumer can be subscribed per queue name. diff --git a/Docs/Documentation/Installation.md b/Docs/Documentation/Installation.md new file mode 100644 index 0000000..c3b75da --- /dev/null +++ b/Docs/Documentation/Installation.md @@ -0,0 +1,83 @@ +Installation +============ + +Composer +-------- + +``` +composer require cakedc/cakephp-enqueue +``` + +This also installs `cakephp/queue` and `cakephp/migrations`. The plugin uses the `cakephp/migrations` database adapters to create the queue table. + +Configure the queue +------------------- + +Add a queue configuration that uses the `cakephp://` DSN to your `config/app.php`: + +```php + 'Queue' => [ + 'default' => [ + 'url' => 'cakephp://default?table_name=enqueue', + 'queue' => 'default', + ], + ], +``` + +`default` in the DSN is the name of the CakePHP datasource connection to use. See [Configuration](Configuration.md) for all the available options. + +NOTE: do this before loading the plugins. `Cake/Queue` throws `Missing 'Queue' configuration key` from its bootstrap if the `Queue` key is not set. + +Load the Plugin +--------------- + +Ensure the plugin is loaded in your `src/Application.php` file, together with `Cake/Queue`: + +```php + /** + * {@inheritdoc} + */ + public function bootstrap(): void + { + parent::bootstrap(); + + $this->addPlugin('Cake/Enqueue'); + $this->addPlugin('Cake/Queue'); + } +``` + +Or using the CLI: + +``` +bin/cake plugin load Cake/Enqueue +bin/cake plugin load Cake/Queue +``` + +The plugin name is `Cake/Enqueue` (from the `Cake\Enqueue` namespace), not the composer package name. + +The queue table +--------------- + +You don't need to write a migration. The table is created the first time a queue configuration is used, that is on the first `QueueManager::push()` or when the worker starts, through the driver's `setupBroker()`. If the table already exists nothing happens. + +NOTE: automatic creation works on MySQL/MariaDB, PostgreSQL and SQLite. The database user needs the `CREATE TABLE` privilege the first time. If creation fails (other drivers, missing privilege) a `E_USER_WARNING` with the reason is triggered and the first query against the table fails. See [Database Schema](Database-Schema.md) to create it yourself. + +Run a worker +------------ + +``` +bin/cake queue worker +``` + +Continue with [Usage](Usage.md) to send your first job. + +Failed jobs storage (optional) +------------------------------ + +Failed jobs are not kept in the queue table. If you want `cakephp/queue` to store them, set `'storeFailedJobs' => true` in the queue configuration and create its `queue_failed_jobs` table: + +``` +bin/cake migrations migrate --plugin Cake/Queue +``` + +The `queue requeue` and `queue purge_failed` commands work with that table, see the [cakephp/queue book](https://book.cakephp.org/queue/2/). diff --git a/Docs/Documentation/Message-Lifecycle.md b/Docs/Documentation/Message-Lifecycle.md new file mode 100644 index 0000000..d5edc27 --- /dev/null +++ b/Docs/Documentation/Message-Lifecycle.md @@ -0,0 +1,58 @@ +Message Lifecycle +================= + +Sending +------- + +`CakeProducer::send()` inserts one row in the queue table with a generated UUID, the body, JSON encoded headers and properties, the queue name, the negated priority and `published_at` (`microtime(true) * 10000`). No `delivery_id` is set, which marks the message as available. + +With `cakephp/queue` a job is written twice: first to the router queue by `QueueManager::push()`, then again to the processor queue by the worker's router processor, which acknowledges (deletes) the router row. Both rows use the same queue name by default, see [Queue names](Usage.md#queue-names). + +Receiving +--------- + +A consumer selects the first message that: + +* belongs to the requested queue(s) +* has no `delivery_id`, meaning nobody is processing it +* is not delayed, or its delay has already passed + +Messages are ordered by priority first and then by `published_at`. The consumer claims the message by writing a new `delivery_id` and `redeliver_after = now + redelivery_delay` with an `UPDATE ... WHERE delivery_id IS NULL`, so two workers never receive the same message. If another worker won the race the consumer retries for up to 200 ms before giving up for this poll. + +When the queue is empty: + +* `bin/cake queue worker` (through `CakeSubscriptionConsumer`) sleeps `subscription_polling_interval` ms (200 by default) and polls again until `receiveTimeout` ms (10000 by default) have passed. +* `CakeConsumer::receive($timeout)` sleeps `polling_interval` ms (1000 by default) between polls. + +Acknowledging +------------- + +When a job returns `Processor::ACK` (or `null`) the row is deleted. + +Failed and rejected messages +---------------------------- + +* `Processor::REJECT` deletes the row. +* `Processor::REQUEUE`, or an exception thrown by the job, inserts a copy of the message as a new available row (new id, new `published_at`, `redelivered = false`) and deletes the original. With `cakephp/queue` the copy carries an incremented `attempts` property; once `attempts` reaches `$maxAttempts` (or `--max-attempts`) the job is rejected instead and, if `storeFailedJobs` is enabled, stored in `queue_failed_jobs`. Without a limit a failing job is requeued forever. +* If the worker dies, or never acknowledges the message, the row stays claimed until `redeliver_after` is reached (`redelivery_delay`, 20 minutes by default). Then the claim is cleared and the row is marked `redelivered = true` so another worker can take it. The job sees `$message->getOriginalMessage()->isRedelivered() === true`. + +NOTE: redelivery is time based. A job that runs longer than `redelivery_delay` is handed to a second worker while the first one is still running it. Raise `redelivery_delay` above your longest job, see [Configuration](Configuration.md#options). + +Delayed messages +---------------- + +`delay` is given in seconds to `QueueManager::push()` and converted to milliseconds by the Enqueue client. The transport stores `delayed_until = time() + delay` (whole seconds, sub-second delays are truncated) and consumers skip the row until then. The delay must be a positive integer, otherwise `CakeProducer` throws a `LogicException`. + +With `cakephp/queue` the delay is applied when the router processor re-sends the message, so it starts counting when a worker routes it, not when it is pushed. + +Time to live +------------ + +`expires` is given in seconds and stored as `time_to_live = time() + expires` on the processor row. Before each poll, consumers delete the rows whose `time_to_live` has passed and that are neither claimed nor redelivered. + +NOTE: the fetch query does not filter by `time_to_live`. If a row expires between the housekeeping run and the claim, the consumer claims it, skips it and leaves it claimed; after `redelivery_delay` it is redelivered with `redelivered = true` and, because redelivered rows are never expired, it is then delivered normally. + +Housekeeping +------------ + +Redelivery and expiration checks run inside every consumer, before each poll and at most once per second, over the whole table (all queues), so no cron job is needed. Errors in these queries are ignored and retried on the next poll. diff --git a/Docs/Documentation/Overview.md b/Docs/Documentation/Overview.md new file mode 100644 index 0000000..e28a7af --- /dev/null +++ b/Docs/Documentation/Overview.md @@ -0,0 +1,48 @@ +Overview +======== + +The plugin registers a `cakephp` transport (connection factory) and a client driver with Enqueue. Anything built on Enqueue, including `cakephp/queue`, can then use a CakePHP database connection as its message broker. + +The plugin is capable of: + +* Storing messages in a single database table, using the CakePHP ORM +* Creating the queue table automatically on MySQL/MariaDB, PostgreSQL and SQLite +* Multiple named queues in the same table +* Message priorities +* Delayed delivery (`delay`) +* Message expiration (`expires` / time to live) +* Automatic redelivery of messages that were received but never acknowledged +* Rejecting messages, with or without requeueing +* Lazy connections, so the database is only touched when a message is sent or received +* Running side by side with other Enqueue transports (Redis, filesystem, ...) in different queue configurations + +What it is not: + +* A full broker. Consumers poll the table; there is no push notification, no fan-out and no temporary queues (`createTemporaryQueue()` throws `TemporaryQueueNotSupportedException`). +* A migration. The table is created on the fly by the transport, see [Database Schema](Database-Schema.md). + +Requirements +------------ + +| Plugin | CakePHP | PHP | cakephp/queue | cakephp/migrations | +| :----: | :-----: | :----: | :-----------: | :----------------: | +| 2.x | ^5.1 | >= 8.2 | ^2.0 | ^4.0 or ^5.0 | +| 1.x | ^4.3 | >= 7.2 | any | not used | + +Classes +------- + +| Class | Responsibility | +| :---- | :------------- | +| `Cake\Enqueue\EnqueuePlugin` | Registers the transport and the client driver with Enqueue | +| `Cake\Enqueue\CakeConnectionFactory` | Parses the DSN or config array and creates the context | +| `Cake\Enqueue\CakeContext` | Creates queues, producers, consumers and the queue table | +| `Cake\Enqueue\CakeProducer` | Stores messages in the table | +| `Cake\Enqueue\CakeConsumer` | Fetches, acknowledges and rejects messages of one queue | +| `Cake\Enqueue\CakeSubscriptionConsumer` | Consumes from several queues at once, used by the worker | +| `Cake\Enqueue\CakeMessage` | Transport message (body, headers, properties and delivery data) | +| `Cake\Enqueue\CakeDestination` | Queue / topic name | +| `Cake\Enqueue\Client\Driver\CakephpDriver` | Enqueue client driver, creates the table in `setupBroker()` | +| `Cake\Enqueue\Model\Table\EnqueueTable` | ORM table used to read and write messages | + +NOTE: the transport is registered under the schemes `cakephp` and `cakephpenqueue`, but the connection factory only accepts `cakephp://` DSNs. A `cakephpenqueue://` DSN throws `Wrong dsn schema passed`. diff --git a/Docs/Documentation/Testing.md b/Docs/Documentation/Testing.md new file mode 100644 index 0000000..6967d43 --- /dev/null +++ b/Docs/Documentation/Testing.md @@ -0,0 +1,72 @@ +Testing +======= + +Running the plugin tests +------------------------ + +``` +composer install +composer test +``` + +By default the tests run against an in-memory SQLite database. To use another database set the `db_dsn` environment variable before running PHPUnit: + +``` +export db_dsn='mysql://root:root@127.0.0.1/cakephp' +vendor/bin/phpunit +``` + +``` +export db_dsn='postgres://postgres:postgres@127.0.0.1/cakephp' +vendor/bin/phpunit +``` + +The GitHub workflow runs the suite on PHP 8.2, 8.3 and 8.4 with a SQLite, MySQL and PostgreSQL matrix. + +Coding standards +---------------- + +``` +composer cs-check +composer cs-fix +``` + +Testing your application jobs +----------------------------- + +To assert that jobs were queued without touching the database, use the `QueueTrait` shipped with `cakephp/queue`, which swaps every queue client for an in-memory one: + +```php +use App\Job\InvoiceJob; +use Cake\Queue\QueueManager; +use Cake\Queue\TestSuite\QueueTrait; +use Cake\TestSuite\TestCase; + +class OrdersControllerTest extends TestCase +{ + use QueueTrait; + + public function testOrderQueuesJob(): void + { + QueueManager::push(InvoiceJob::class, ['id' => 1]); + + $this->assertJobQueued(InvoiceJob::class); + $this->assertJobQueuedWith(InvoiceJob::class, ['id' => 1]); + } +} +``` + +To run the real transport in tests point a configuration at the `test` connection; the table is created on first use: + +```php +use Cake\Queue\QueueManager; + +QueueManager::setConfig('default', [ + 'url' => 'cakephp://test', + 'queue' => 'default', +]); +``` + +To clean the table between tests, purge the prefixed queue (see [Emptying a queue](Usage.md#emptying-a-queue)) or truncate the `enqueue` table. `QueueManager::drop('default')` discards the client so that the next test builds a fresh one. + +NOTE: `QueueManager::setConfig()` throws `Cannot reconfigure existing key` if the key is already configured; drop it first in `tearDown()`. diff --git a/Docs/Documentation/Usage.md b/Docs/Documentation/Usage.md new file mode 100644 index 0000000..dda43cc --- /dev/null +++ b/Docs/Documentation/Usage.md @@ -0,0 +1,114 @@ +Usage +===== + +The plugin does not add its own API. You use the `cakephp/queue` package as usual and the plugin acts as the broker. See the [cakephp/queue book](https://book.cakephp.org/queue/2/) for everything about jobs. + +Create a job +------------ + +```php +// src/Job/ExampleJob.php +namespace App\Job; + +use Cake\Queue\Job\JobInterface; +use Cake\Queue\Job\Message; +use Interop\Queue\Processor; + +class ExampleJob implements JobInterface +{ + public static ?int $maxAttempts = 3; + + public function execute(Message $message): ?string + { + $id = $message->getArgument('id'); + + // do the work... + + return Processor::ACK; + } +} +``` + +Returning `null` is the same as `Processor::ACK`. Returning `Processor::REQUEUE`, or throwing, sends the job back to the queue; `Processor::REJECT` drops it. See [Message Lifecycle](Message-Lifecycle.md). + +Push a job +---------- + +```php +use App\Job\ExampleJob; +use Cake\Queue\QueueManager; + +$data = ['id' => 7, 'is_premium' => true]; +$options = ['config' => 'default']; + +QueueManager::push(ExampleJob::class, $data, $options); +``` + +Options +------- + +`QueueManager::push()` accepts these options: + +| Option | Description | +| :----- | :---------- | +| config | Name of the queue configuration to use. Defaults to `default`. | +| queue | Topic the job is published on. Defaults to the `queue` key of the configuration, or `default`. It must match the worker's queue, see [Queue names](#queue-names). | +| delay | Integer seconds to wait before the job becomes available. See [Delayed messages](Message-Lifecycle.md#delayed-messages). | +| expires | Integer seconds after which an unconsumed job is removed. See [Time to live](Message-Lifecycle.md#time-to-live). | +| priority | One of the `Enqueue\Client\MessagePriority` constants: `VERY_LOW`, `LOW`, `NORMAL`, `HIGH`, `VERY_HIGH` (stored as 0 to 4). Higher priority is delivered first. | + +```php +use Enqueue\Client\MessagePriority; + +QueueManager::push(ExampleJob::class, $data, [ + 'config' => 'default', + 'delay' => 60, + 'expires' => 3600, + 'priority' => MessagePriority::HIGH, +]); +``` + +Queue names +----------- + +`cakephp/queue` publishes jobs as Enqueue *events*: the job is stored in the configuration's router queue with the `queue` option as its topic, and the worker's router re-sends it to the processor queue, which is the same queue by default. Enqueue prefixes queue names with `enqueue.app.`, so a configuration with `'queue' => 'default'` stores all its rows with `queue = 'enqueue.app.default'` in the table, whatever the `queue` option of `push()` was. + +Consequences: + +* The `queue` push option is a topic name. A worker only processes jobs whose topic matches its `--queue` (defaulting to the configuration's `queue`). Jobs with another topic are routed to zero subscribers and silently acknowledged, that is, dropped. Keep the push option and the worker queue equal, or use one configuration per queue. +* `delay`, `expires` and `priority` are stored as message properties on the router row and only applied (as `delayed_until`, `time_to_live` and `priority` columns) when the worker routes the job to the processor queue. The delay starts counting at that moment, which requires a running worker. +* Anything that addresses the table by queue name, such as `purgeQueue()`, must use the prefixed name. Use the driver to build it, see [Emptying a queue](#emptying-a-queue). + +Process jobs +------------ + +``` +bin/cake queue worker +``` + +Options: `--config` (`-c`), `--queue` (`-Q`), `--processor` (`-p`), `--logger` (`-l`, only used together with `--verbose`), `--max-jobs` (`-i`), `--max-runtime` (`-r`, seconds) and `--max-attempts` (`-a`). A `$maxAttempts` property on a job overrides `--max-attempts`. + +``` +bin/cake queue worker --config=default --max-jobs=100 --max-runtime=3600 --max-attempts=3 --verbose +``` + +Run as many workers as you need; each message is claimed by a single worker. See [Message Lifecycle](Message-Lifecycle.md). + +NOTE: `--max-jobs` is checked between consume cycles, and one cycle lasts `receiveTimeout` milliseconds (10000 by default, configurable per queue configuration). With this transport every job that becomes available during a cycle is processed, so a worker can process more than `--max-jobs` jobs before it stops. + +NOTE: In production keep the worker alive with a process manager such as supervisor or systemd, and restart it on deployments so it picks up new code. + +Emptying a queue +---------------- + +To remove every row of a configuration's queue, build the destination through the driver so the prefixed name is used: + +```php +use Cake\Queue\QueueManager; + +$driver = QueueManager::engine('default')->getDriver(); +$queue = $driver->createQueue($driver->getConfig()->getRouterQueue()); +$driver->getContext()->purgeQueue($queue); +``` + +`$context->purgeQueue($context->createQueue('default'))` does nothing, because no row has `queue = 'default'`; the rows are stored as `enqueue.app.default`. `purgeQueue()` deletes claimed and delayed rows too. diff --git a/Docs/Home.md b/Docs/Home.md new file mode 100644 index 0000000..f67bdd4 --- /dev/null +++ b/Docs/Home.md @@ -0,0 +1,104 @@ +Home +==== + +The **CakePHP Enqueue** plugin is a database transport for [cakephp/queue](https://github.com/cakephp/queue). It implements the [Enqueue](https://github.com/php-enqueue/enqueue) `queue-interop` transport on top of a CakePHP database connection, so you can run background jobs without installing a dedicated message broker. + +That it works out of the box doesn't mean it is meant to replace a full broker in every scenario. It is a good fit when you already have a database and want simple, reliable job processing with zero extra infrastructure. + +Documentation +------------- + +* [Overview](Documentation/Overview.md) +* [Installation](Documentation/Installation.md) +* [Configuration](Documentation/Configuration.md) +* [Usage](Documentation/Usage.md) +* [Message Lifecycle](Documentation/Message-Lifecycle.md) +* [Database Schema](Documentation/Database-Schema.md) +* [Using the Enqueue Client](Documentation/Enqueue-Client.md) +* [Testing](Documentation/Testing.md) +* [Known Caveats](Documentation/Caveats.md) + +I want to +--------- + +* get started quickly + *
+ install and send my first job + + ``` + composer require cakedc/cakephp-enqueue + ``` + + ```php + // config/app.php + 'Queue' => [ + 'default' => [ + 'url' => 'cakephp://default?table_name=enqueue', + 'queue' => 'default', + ], + ], + ``` + + ```php + // src/Application.php + $this->addPlugin('Cake/Enqueue'); + $this->addPlugin('Cake/Queue'); + ``` + + ```php + QueueManager::push(ExampleJob::class, ['id' => 7]); + ``` + + ``` + bin/cake queue worker + ``` + + See [Installation](Documentation/Installation.md) and [Usage](Documentation/Usage.md) for the details. +
+* configure + *
+ a different database connection or table + + ```php + 'url' => 'cakephp://my_connection?table_name=my_queue_table', + ``` +
+ *
+ how often the worker polls the database + + ```php + 'url' => [ + 'transport' => [ + 'dsn' => 'cakephp:', + 'connection' => 'default', + 'subscription_polling_interval' => 500, + ], + ], + 'queue' => 'default', + ``` + + NOTE: the worker uses `subscription_polling_interval`, not `polling_interval`, and it must be an integer, so it can only be set with the array form. See [Configuration](Documentation/Configuration.md#options). +
+ *
+ how long before an unacknowledged message is redelivered + + ```php + 'url' => [ + 'transport' => [ + 'dsn' => 'cakephp:', + 'connection' => 'default', + 'redelivery_delay' => 600000, + ], + ], + 'queue' => 'default', + ``` + + NOTE: milliseconds, truncated to whole seconds. See [Configuration](Documentation/Configuration.md#options). +
+* understand + * [where my job is stored and why the `queue` column says `enqueue.app.default`](Documentation/Usage.md#queue-names) + * [what happens to a message that fails](Documentation/Message-Lifecycle.md#failed-and-rejected-messages) + * [how delayed messages work](Documentation/Message-Lifecycle.md#delayed-messages) + * [how expired messages are removed](Documentation/Message-Lifecycle.md#time-to-live) + * [which table columns are used](Documentation/Database-Schema.md) + * [what can go wrong](Documentation/Caveats.md) diff --git a/README.md b/README.md index f42b51f..190dbe0 100644 --- a/README.md +++ b/README.md @@ -11,20 +11,19 @@ Versions and branches | CakePHP | CakePHP Enqueue Plugin | Tag | Notes | | :-------------: | :------------------------: | :--: | :---- | | ^5.1 | [2.x](https://github.com/CakeDC/cakephp-enqueue/tree/2.x) | 2.0.1 | stable | -| ^4.5 | [1.x](https://github.com/CakeDC/cakephp-enqueue/tree/1.x) | 1.0.0 | stable | +| ^4.3 | [1.x](https://github.com/CakeDC/cakephp-enqueue/tree/1.x) | 1.0.0 | stable | The **CakePHP Enqueue** plugin provides message queue integration for CakePHP applications using the Enqueue library and database as a message broker. Requirements ------------ -* CakePHP 4.5+ or 5.1+ -* PHP 8.0+ +* CakePHP 5.1+ and PHP 8.2+ (2.x), or CakePHP 4.3+ and PHP 7.2+ (1.x) Documentation ------------- -For documentation, see the [Docs](docs/index.md) directory of this repository. +For documentation, see the [Docs](Docs/Home.md) directory of this repository. Support ------- diff --git a/docs/index.md b/docs/index.md deleted file mode 100644 index d7afcc2..0000000 --- a/docs/index.md +++ /dev/null @@ -1,53 +0,0 @@ -Home -==== - -The **CakePHP Enqueue** plugin provides message queue integration for CakePHP applications and uses database as a message broker. - -Quick Start ------------ - -1. Install the plugin: - -```bash - composer require cakedc/cakephp-enqueue -``` - -2. Load the plugin in your `Application.php`: - -```php - $this->addPlugin('CakephpEnqueue'); -``` - -3. Configure your queue in `config/app.php`: - -```php - 'Queue' => [ - 'default' => [ - 'url' => 'cakephp://default?table_name=queue' - ] - ] -``` - -4. Create a job: - -```php - use App\Job\ExampleJob; - use Cake\Queue\QueueManager; - - $data = ['id' => 7, 'is_premium' => true]; - $options = ['config' => 'default']; - - QueueManager::push(ExampleJob::class, $data, $options); -``` - -5. Process jobs: - -```bash - bin/cake queue:worker -``` - -DSN Configuration ------------------ - -The plugin supports standard DSN format: -- `cakephp://connection_name?table_name=queue&polling_interval=1000` diff --git a/src/CakeConsumer.php b/src/CakeConsumer.php index d0f3cc9..3b7ac5c 100644 --- a/src/CakeConsumer.php +++ b/src/CakeConsumer.php @@ -99,7 +99,7 @@ public function getQueue(): Queue */ public function receiveNoWait(): ?Message { - $redeliveryDelay = $this->getRedeliveryDelay() / 1000; + $redeliveryDelay = (int)($this->getRedeliveryDelay() / 1000); $this->removeExpiredMessages(); $this->redeliverMessages(); diff --git a/src/CakeContext.php b/src/CakeContext.php index 223777c..930e2f7 100644 --- a/src/CakeContext.php +++ b/src/CakeContext.php @@ -169,11 +169,11 @@ public function createSubscriptionConsumer(): SubscriptionConsumer $consumer = new CakeSubscriptionConsumer($this); if (isset($this->config['redelivery_delay'])) { - $consumer->setRedeliveryDelay($this->config['redelivery_delay']); + $consumer->setRedeliveryDelay((int)$this->config['redelivery_delay']); } if (isset($this->config['subscription_polling_interval'])) { - $consumer->setPollingInterval($this->config['subscription_polling_interval']); + $consumer->setPollingInterval((int)$this->config['subscription_polling_interval']); } return $consumer; diff --git a/src/CakeSubscriptionConsumer.php b/src/CakeSubscriptionConsumer.php index 3ead9a8..680f562 100644 --- a/src/CakeSubscriptionConsumer.php +++ b/src/CakeSubscriptionConsumer.php @@ -127,7 +127,7 @@ public function consume(int $timeout = 0): void $timeout /= 1000; $now = time(); - $redeliveryDelay = $this->getRedeliveryDelay() / 1000; // milliseconds to seconds + $redeliveryDelay = (int)($this->getRedeliveryDelay() / 1000); // milliseconds to seconds $currentQueueNames = []; while (true) { diff --git a/tests/TestCase/CakeConnectionFactoryTest.php b/tests/TestCase/CakeConnectionFactoryTest.php index 10fa998..ba13a09 100644 --- a/tests/TestCase/CakeConnectionFactoryTest.php +++ b/tests/TestCase/CakeConnectionFactoryTest.php @@ -102,4 +102,25 @@ public function testConstructorWithArrayConfig(): void $this->assertEquals(3000, $actualConfig['polling_interval']); $this->assertFalse($actualConfig['lazy']); } + + /** + * Options given in the DSN must not break consumer creation. + * + * @return void + */ + public function testDsnConsumerOptionsAreCast(): void + { + $factory = new CakeConnectionFactory( + 'cakephp://test?redelivery_delay=1500&subscription_polling_interval=300', + ); + $context = $factory->createContext(); + + $subscriptionConsumer = $context->createSubscriptionConsumer(); + $this->assertSame(1500, $subscriptionConsumer->getRedeliveryDelay()); + + $context->createDataBaseTable(); + $consumer = $context->createConsumer($context->createQueue('test')); + $this->assertSame(1500, $consumer->getRedeliveryDelay()); + $this->assertNull($consumer->receiveNoWait()); + } }