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

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
18 changes: 14 additions & 4 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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'
Expand Down Expand Up @@ -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
Expand Down
37 changes: 37 additions & 0 deletions Docs/Documentation/Caveats.md
Original file line number Diff line number Diff line change
@@ -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.<queue>`, 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)
119 changes: 119 additions & 0 deletions Docs/Documentation/Configuration.md
Original file line number Diff line number Diff line change
@@ -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://<connection>?table_name=<table>&polling_interval=<ms>&lazy=<bool>
```

```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=<name>` 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`.
108 changes: 108 additions & 0 deletions Docs/Documentation/Database-Schema.md
Original file line number Diff line number Diff line change
@@ -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.
80 changes: 80 additions & 0 deletions Docs/Documentation/Enqueue-Client.md
Original file line number Diff line number Diff line change
@@ -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.
Loading
Loading