Skip to content
Closed
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
49 changes: 48 additions & 1 deletion packages/queue/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -73,9 +73,56 @@ $client->enqueue([
]);
```

## Idempotent publication

Redis publishers implement `Utopia\Queue\Publisher\Idempotent`. Use
`enqueueOnce()` when a producer can retry after losing the publish response:

```php
use Utopia\Queue\Publisher\Result;
use Utopia\Queue\Queue;

$queue = new Queue('my-queue');
$result = $publisher->enqueueOnce(
$queue,
messageId: $operationId,
payload: ['operationId' => $operationId],
);

if ($result === Result::Existing) {
// The same envelope was already accepted.
}
```

`Result::Enqueued` and `Result::Existing` both acknowledge success. Message IDs
are retained per queue without an implicit expiry. Reusing an ID with a
different canonical payload or priority throws
`Utopia\Queue\Exception\Conflict`.

## Visibility leases

Redis claims atomically move a pending envelope into processing state.
Visibility reclaim is disabled by default to preserve existing long-running
workers. Set an application-specific timeout on the queue to recover messages
when a worker loses its acknowledgement:

```php
$queue = new Queue('my-queue', visibilityTimeout: $leaseSeconds);
$message = $consumer->receive($queue, timeout: 2);

if ($message !== null) {
$consumer->renew($queue, $message);
$consumer->commit($queue, $message);
}
```

Call `renew()` before the deadline while processing valid long-running work.
Each redelivery receives a new receipt, so an acknowledgement from an older
delivery cannot complete the active claim.

## System requirements

Utopia Framework requires PHP 8.0 or later. We recommend using the latest PHP version whenever possible.
Utopia Queue requires PHP 8.5 or later.

## Copyright and license

Expand Down
32 changes: 26 additions & 6 deletions packages/queue/docker-compose.yml
Original file line number Diff line number Diff line change
Expand Up @@ -10,12 +10,32 @@ services:
retries: 15

redis-cluster:
image: grokzen/redis-cluster:7.0.10
environment:
IP: "0.0.0.0"
INITIAL_PORT: 17000
MASTERS: 3
SLAVES_PER_MASTER: 0
image: redis:alpine
entrypoint: ["/bin/sh", "-c"]
command:
- |
for port in 17000 17001 17002; do
mkdir -p "/data/$${port}"
redis-server \
--port "$${port}" \
--bind 0.0.0.0 \
--protected-mode no \
--appendonly no \
--cluster-enabled yes \
--cluster-config-file "/data/$${port}/nodes.conf" \
--cluster-node-timeout 5000 \
--cluster-announce-ip 127.0.0.1 \
--cluster-announce-port "$${port}" \
--cluster-announce-bus-port "$$((port + 10000))" &
done
until redis-cli -p 17000 ping > /dev/null 2>&1; do sleep 1; done
redis-cli --cluster create \
127.0.0.1:17000 \
127.0.0.1:17001 \
127.0.0.1:17002 \
--cluster-replicas 0 \
--cluster-yes
wait
ports:
- "17000-17002:17000-17002"
healthcheck:
Expand Down
3 changes: 3 additions & 0 deletions packages/queue/phpunit.xml
Original file line number Diff line number Diff line change
Expand Up @@ -7,12 +7,15 @@
<testsuites>
<testsuite name="unit">
<file>./tests/Queue/E2E/Adapter/LockingTest.php</file>
<file>./tests/Queue/E2E/Adapter/RedisKeysTest.php</file>
<file>./tests/Queue/E2E/Adapter/RedisReconnectCallbackTest.php</file>
<file>./tests/Queue/E2E/Adapter/ServerTelemetryTest.php</file>
<file>./tests/Queue/E2E/Adapter/SwooleConcurrencyTest.php</file>
</testsuite>
<testsuite name="e2e">
<file>./tests/Queue/E2E/Adapter/PoolTest.php</file>
<file>./tests/Queue/E2E/Adapter/RedisClusterDurabilityTest.php</file>
<file>./tests/Queue/E2E/Adapter/RedisDurabilityTest.php</file>
<file>./tests/Queue/E2E/Adapter/SwooleTest.php</file>
<file>./tests/Queue/E2E/Adapter/SwooleRedisClusterTest.php</file>
<file>./tests/Queue/E2E/Adapter/WorkermanTest.php</file>
Expand Down
36 changes: 35 additions & 1 deletion packages/queue/src/Queue/Broker/Pool.php
Original file line number Diff line number Diff line change
Expand Up @@ -4,11 +4,15 @@

use Utopia\Pools\Pool as UtopiaPool;
use Utopia\Queue\Consumer;
use Utopia\Queue\Consumer\Leased as LeasedConsumer;
use Utopia\Queue\Exception\Unsupported;
use Utopia\Queue\Message;
use Utopia\Queue\Publisher;
use Utopia\Queue\Publisher\Idempotent as IdempotentPublisher;
use Utopia\Queue\Publisher\Result;
use Utopia\Queue\Queue;

readonly class Pool implements Publisher, Consumer
readonly class Pool implements IdempotentPublisher, LeasedConsumer
{
public function __construct(
private ?UtopiaPool $publisher = null,
Expand All @@ -20,6 +24,23 @@ public function enqueue(Queue $queue, array $payload, bool $priority = false): b
return $this->delegate($this->publisher, __FUNCTION__, \func_get_args());
}

public function enqueueOnce(
Queue $queue,
string $messageId,
array $payload,
bool $priority = false,
): Result {
return $this->publisher?->use(
function (Publisher $publisher) use ($queue, $messageId, $payload, $priority): Result {
if (!$publisher instanceof IdempotentPublisher) {
throw new Unsupported('idempotent publishing');
}

return $publisher->enqueueOnce($queue, $messageId, $payload, $priority);
},
) ?? throw new Unsupported('publishing');
}

public function retry(Queue $queue, ?int $limit = null): void
{
$this->delegate($this->publisher, __FUNCTION__, \func_get_args());
Expand All @@ -45,6 +66,19 @@ public function reject(Queue $queue, Message $message): void
$this->delegate($this->consumer, __FUNCTION__, \func_get_args());
}

public function renew(Queue $queue, Message $message): bool
{
return $this->consumer?->use(
function (Consumer $consumer) use ($queue, $message): bool {
if (!$consumer instanceof LeasedConsumer) {
throw new Unsupported('visibility leases');
}

return $consumer->renew($queue, $message);
},
) ?? throw new Unsupported('consuming');
}

public function close(): void
{
// TODO: Implement closing all connections in the pool
Expand Down
Loading
Loading