diff --git a/README.md b/README.md index 5b593890..5c11189f 100644 --- a/README.md +++ b/README.md @@ -38,7 +38,7 @@ See the [adapter list](docs/guide/en/adapter-list.md) and follow the adapter-spe > If you don't have an external broker — whether for development, testing, or because you want to > design around `QueueProducerInterface` from day one and add a real broker later — you can run the queue -> in [synchronous mode](docs/guide/en/synchronous-mode.md) (the adapter argument is optional). +> in [synchronous mode](docs/guide/en/synchronous-mode.md) using `SyncQueueProducer` instead of `AsyncQueueProducer`. > In this mode messages are processed immediately in the same process, so it won't provide true > async execution, but the code stays the same when you switch to a real adapter. diff --git a/docs/guide/en/configuration-manual.md b/docs/guide/en/configuration-manual.md index 58947487..5ad61cc7 100644 --- a/docs/guide/en/configuration-manual.md +++ b/docs/guide/en/configuration-manual.md @@ -8,7 +8,7 @@ To use the queue, you need to create instances of the following classes: 1. **Adapter** - handles the actual queue backend like AMQP, Redis, etc. 2. **Worker** - processes messages from the queue -3. **QueueProducer** - pushes messages; **QueueConsumer** consumes them when needed +3. **SyncQueueProducer** / **AsyncQueueProducer** - pushes messages; **QueueConsumer** consumes them when needed ### Example @@ -25,7 +25,7 @@ use Yiisoft\Queue\Middleware\FailureHandling\FailureMiddlewareFactory; use Yiisoft\Queue\Middleware\Push\PushMiddlewareConfig; use Yiisoft\Queue\Middleware\Push\PushMiddlewareFactory; use Yiisoft\Queue\QueueConsumer; -use Yiisoft\Queue\QueueProducer; +use Yiisoft\Queue\SyncQueueProducer; use Yiisoft\Queue\Worker\Worker; // A PSR-11 container is required for resolving dependencies of middleware and handlers. @@ -69,16 +69,17 @@ $worker = new Worker( // Create loop (SignalLoop requires ext-pcntl; SimpleLoop works without it) $loop = new SimpleLoop(); -// Create queue. Without an adapter the queue runs in synchronous mode (messages are processed -// immediately on push). Pass an adapter (e.g., AMQP, Redis) for asynchronous processing. -$producer = new QueueProducer( +// Create queue. SyncQueueProducer runs in synchronous mode (messages are processed +// immediately on push). Use AsyncQueueProducer with an adapter (e.g., AMQP, Redis) instead +// for asynchronous processing. +$producer = new SyncQueueProducer( $logger, $pushMiddlewareConfig, - worker: $worker, + $worker, ); $consumer = new QueueConsumer($worker, $loop, $logger); -// Now you can push messages. With no adapter, the producer dispatches directly to the worker. +// Now you can push messages. SyncQueueProducer dispatches directly to the worker. $message = new DownloadFileMessage(url: 'https://example.com/file.pdf', destinationPath: '/tmp/file.pdf'); $producer->push($message); ``` @@ -100,7 +101,7 @@ $provider = new PredefinedQueueProvider([ ## Running the queue Message consumption methods are available on `Yiisoft\Queue\QueueConsumerInterface`. -`QueueProducer` and `QueueConsumer` are separate capabilities. Obtain or construct the consumer role before calling these methods. +The producer and `QueueConsumer` are separate capabilities. Obtain or construct the consumer role before calling these methods. ### Processing existing messages diff --git a/docs/guide/en/middleware-pipelines.md b/docs/guide/en/middleware-pipelines.md index cf4a55fb..9828c19d 100644 --- a/docs/guide/en/middleware-pipelines.md +++ b/docs/guide/en/middleware-pipelines.md @@ -149,6 +149,6 @@ See [Configuration with yiisoft/config](configuration-with-config.md) for exampl ### Manual configuration (without yiisoft/config) -When configuring the component manually, you instantiate the middleware dispatchers and pass them to `QueueProducer`, `QueueConsumer`, and `Worker` as appropriate. +When configuring the component manually, you instantiate the middleware dispatchers and pass them to `SyncQueueProducer` / `AsyncQueueProducer`, `QueueConsumer`, and `Worker` as appropriate. See [Manual configuration](configuration-manual.md) for a full runnable example. diff --git a/docs/guide/en/performance-tuning.md b/docs/guide/en/performance-tuning.md index 8d6723be..63602e9a 100644 --- a/docs/guide/en/performance-tuning.md +++ b/docs/guide/en/performance-tuning.md @@ -127,15 +127,15 @@ return [ 'yiisoft/queue' => [ 'queues' => [ 'critical' => [ - 'producer' => ['class' => \Yiisoft\Queue\QueueProducer::class, '__construct()' => ['adapter' => AmqpAdapter::class]], + 'producer' => ['class' => \Yiisoft\Queue\AsyncQueueProducer::class, '__construct()' => ['adapter' => AmqpAdapter::class]], 'consumer' => ['class' => \Yiisoft\Queue\QueueConsumer::class, '__construct()' => ['adapter' => AmqpAdapter::class]], ], 'normal' => [ - 'producer' => ['class' => \Yiisoft\Queue\QueueProducer::class, '__construct()' => ['adapter' => AmqpAdapter::class]], + 'producer' => ['class' => \Yiisoft\Queue\AsyncQueueProducer::class, '__construct()' => ['adapter' => AmqpAdapter::class]], 'consumer' => ['class' => \Yiisoft\Queue\QueueConsumer::class, '__construct()' => ['adapter' => AmqpAdapter::class]], ], 'low' => [ - 'producer' => ['class' => \Yiisoft\Queue\QueueProducer::class, '__construct()' => ['adapter' => AmqpAdapter::class]], + 'producer' => ['class' => \Yiisoft\Queue\AsyncQueueProducer::class, '__construct()' => ['adapter' => AmqpAdapter::class]], 'consumer' => ['class' => \Yiisoft\Queue\QueueConsumer::class, '__construct()' => ['adapter' => AmqpAdapter::class]], ], ], @@ -165,19 +165,19 @@ return [ 'yiisoft/queue' => [ 'queues' => [ 'fast' => [ // Quick tasks (< 1s) - 'producer' => ['class' => \Yiisoft\Queue\QueueProducer::class, '__construct()' => ['adapter' => AmqpAdapter::class]], + 'producer' => ['class' => \Yiisoft\Queue\AsyncQueueProducer::class, '__construct()' => ['adapter' => AmqpAdapter::class]], 'consumer' => ['class' => \Yiisoft\Queue\QueueConsumer::class, '__construct()' => ['adapter' => AmqpAdapter::class]], ], 'slow' => [ // Long tasks (> 10s) - 'producer' => ['class' => \Yiisoft\Queue\QueueProducer::class, '__construct()' => ['adapter' => AmqpAdapter::class]], + 'producer' => ['class' => \Yiisoft\Queue\AsyncQueueProducer::class, '__construct()' => ['adapter' => AmqpAdapter::class]], 'consumer' => ['class' => \Yiisoft\Queue\QueueConsumer::class, '__construct()' => ['adapter' => AmqpAdapter::class]], ], 'cpu-bound' => [ // CPU-intensive - 'producer' => ['class' => \Yiisoft\Queue\QueueProducer::class, '__construct()' => ['adapter' => AmqpAdapter::class]], + 'producer' => ['class' => \Yiisoft\Queue\AsyncQueueProducer::class, '__construct()' => ['adapter' => AmqpAdapter::class]], 'consumer' => ['class' => \Yiisoft\Queue\QueueConsumer::class, '__construct()' => ['adapter' => AmqpAdapter::class]], ], 'io-bound' => [ // I/O-intensive - 'producer' => ['class' => \Yiisoft\Queue\QueueProducer::class, '__construct()' => ['adapter' => AmqpAdapter::class]], + 'producer' => ['class' => \Yiisoft\Queue\AsyncQueueProducer::class, '__construct()' => ['adapter' => AmqpAdapter::class]], 'consumer' => ['class' => \Yiisoft\Queue\QueueConsumer::class, '__construct()' => ['adapter' => AmqpAdapter::class]], ], ], diff --git a/docs/guide/en/queue-capabilities.md b/docs/guide/en/queue-capabilities.md index 189a0dae..5ee4ac8d 100644 --- a/docs/guide/en/queue-capabilities.md +++ b/docs/guide/en/queue-capabilities.md @@ -6,18 +6,18 @@ Named providers use a strict nested role map. `getProducerNames()` and `getConsu ```php use Yiisoft\Queue\QueueConsumer; -use Yiisoft\Queue\QueueProducer; +use Yiisoft\Queue\AsyncQueueProducer; $definitions = [ 'orders' => [ - 'producer' => ['class' => QueueProducer::class], + 'producer' => ['class' => AsyncQueueProducer::class], 'consumer' => ['class' => QueueConsumer::class], ], - 'outbound-events' => ['producer' => ['class' => QueueProducer::class]], + 'outbound-events' => ['producer' => ['class' => AsyncQueueProducer::class]], 'inbound-events' => ['consumer' => ['class' => QueueConsumer::class]], ]; ``` -`QueueFactoryProvider` accepts factory definitions in each role. `PredefinedQueueProvider` uses the same outer shape but each role value must already be its respective interface instance. A raw definition such as `'orders' => ['class' => QueueProducer::class]`, an empty role map, and unknown role keys are invalid. +`QueueFactoryProvider` accepts factory definitions in each role. `PredefinedQueueProvider` uses the same outer shape but each role value must already be its respective interface instance. A raw definition such as `'orders' => ['class' => AsyncQueueProducer::class]`, an empty role map, and unknown role keys are invalid. -`QueueInterface`, `Queue`, and `QueueProviderInterface` were removed before release. Replace them with `QueueProducerInterface`, `QueueProducer` / `QueueConsumer`, and the relevant typed provider. Synchronous consumers retain no-op `run()` and `listen()` behavior when no adapter is configured. Default retry of an asynchronously consumed message resolves a producer for the execution queue name through a configured producer provider; if none is available it fails with an actionable configuration error rather than dropping the message. +`QueueInterface`, `Queue`, and `QueueProviderInterface` were removed before release. Replace them with `QueueProducerInterface`, `SyncQueueProducer` / `AsyncQueueProducer` / `QueueConsumer`, and the relevant typed provider. Synchronous consumers retain no-op `run()` and `listen()` behavior when no adapter is configured. Default retry of an asynchronously consumed message resolves a producer for the execution queue name through a configured producer provider; if none is available it fails with an actionable configuration error rather than dropping the message. diff --git a/docs/guide/en/queue-names-advanced.md b/docs/guide/en/queue-names-advanced.md index 4e71b0e0..f67ce37b 100644 --- a/docs/guide/en/queue-names-advanced.md +++ b/docs/guide/en/queue-names-advanced.md @@ -30,15 +30,15 @@ Choose the provider by how the roles are created: ```php use Yiisoft\Queue\Provider\QueueFactoryProvider; use Yiisoft\Queue\QueueConsumer; -use Yiisoft\Queue\QueueProducer; +use Yiisoft\Queue\AsyncQueueProducer; $provider = new QueueFactoryProvider([ 'emails' => [ - 'producer' => ['class' => QueueProducer::class], + 'producer' => ['class' => AsyncQueueProducer::class], 'consumer' => ['class' => QueueConsumer::class], ], 'audit' => [ - 'producer' => ['class' => QueueProducer::class], + 'producer' => ['class' => AsyncQueueProducer::class], ], ], $container); diff --git a/docs/guide/en/queue-names.md b/docs/guide/en/queue-names.md index e5b0e1a6..33ee34d2 100644 --- a/docs/guide/en/queue-names.md +++ b/docs/guide/en/queue-names.md @@ -20,19 +20,19 @@ Named queues use a strict role map under `yiisoft/queue.queues`. Each name must use Yiisoft\Queue\Adapter\AdapterInterface; use Yiisoft\Queue\Provider\QueueProducerProviderInterface; use Yiisoft\Queue\QueueConsumer; -use Yiisoft\Queue\QueueProducer; +use Yiisoft\Queue\AsyncQueueProducer; return [ 'yiisoft/queue' => [ 'queues' => [ // A queue with both capabilities. QueueProducerProviderInterface::DEFAULT_QUEUE => [ - 'producer' => ['class' => QueueProducer::class, '__construct()' => ['adapter' => AdapterInterface::class]], + 'producer' => ['class' => AsyncQueueProducer::class, '__construct()' => ['adapter' => AdapterInterface::class]], 'consumer' => ['class' => QueueConsumer::class, '__construct()' => ['adapter' => AdapterInterface::class]], ], // Produce-only and consume-only names are valid. 'outbound-events' => [ - 'producer' => ['class' => QueueProducer::class, '__construct()' => ['adapter' => AdapterInterface::class]], + 'producer' => ['class' => AsyncQueueProducer::class, '__construct()' => ['adapter' => AdapterInterface::class]], ], 'inbound-events' => [ 'consumer' => ['class' => QueueConsumer::class, '__construct()' => ['adapter' => AdapterInterface::class]], diff --git a/docs/guide/en/synchronous-mode.md b/docs/guide/en/synchronous-mode.md index 4d4a7af0..efe73cdf 100644 --- a/docs/guide/en/synchronous-mode.md +++ b/docs/guide/en/synchronous-mode.md @@ -8,7 +8,7 @@ Run tasks synchronously in the same process. Useful for: doesn't have an external broker yet — you can switch to a real adapter later without touching the call sites. -To enable it, create the queue instance without an adapter (the `adapter` argument defaults to `null`): +To enable it, use `SyncQueueProducer` instead of `AsyncQueueProducer` — it takes a worker instead of an adapter: ```php $logger = $DIContainer->get(\Psr\Log\LoggerInterface::class); @@ -18,10 +18,10 @@ $pushMiddlewareConfig = $DIContainer->get( \Yiisoft\Queue\Middleware\Push\PushMiddlewareConfig::class ); -$producer = new \Yiisoft\Queue\QueueProducer( +$producer = new \Yiisoft\Queue\SyncQueueProducer( $logger, $pushMiddlewareConfig, - worker: $worker, + $worker, ); ``` diff --git a/src/QueueProducer.php b/src/AsyncQueueProducer.php similarity index 52% rename from src/QueueProducer.php rename to src/AsyncQueueProducer.php index 8e6338a1..22e785dc 100644 --- a/src/QueueProducer.php +++ b/src/AsyncQueueProducer.php @@ -12,38 +12,31 @@ use Yiisoft\Queue\Middleware\Push\AdapterPushHandler; use Yiisoft\Queue\Middleware\Push\PushMiddlewareConfig; use Yiisoft\Queue\Middleware\Push\PushMiddlewareDispatcher; -use Yiisoft\Queue\Middleware\Push\SynchronousPushHandler; use Yiisoft\Queue\Provider\QueueProducerProviderInterface; -use Yiisoft\Queue\Worker\WorkerInterface; -use InvalidArgumentException; -/** Produces messages for one logical queue. */ -final class QueueProducer implements QueueProducerInterface +/** + * Produces messages for one logical queue, pushing them to an adapter-backed broker. + */ +final class AsyncQueueProducer implements QueueProducerInterface { private string $name; private PushMiddlewareDispatcher $dispatcher; /** - * @param mixed ...$middlewareDefinitions Queue-specific push middleware definitions. + * @param mixed[] $middlewareDefinitions Queue-specific push middleware definitions. */ public function __construct( private readonly LoggerInterface $logger, PushMiddlewareConfig $middlewareConfig, - private readonly ?AdapterInterface $adapter = null, + private readonly AdapterInterface $adapter, string|BackedEnum $name = QueueProducerProviderInterface::DEFAULT_QUEUE, - ?WorkerInterface $worker = null, - mixed ...$middlewareDefinitions, + array $middlewareDefinitions = [], ) { $this->name = StringNormalizer::normalize($name); - if ($adapter === null && $worker === null) { - throw new InvalidArgumentException('A synchronous queue producer requires a worker.'); - } $this->dispatcher = new PushMiddlewareDispatcher( middlewareFactory: $middlewareConfig->middlewareFactory, middlewareDefinitions: [...$middlewareConfig->commonMiddlewareDefinitions, ...$middlewareDefinitions], - finishHandler: $adapter === null - ? new SynchronousPushHandler($worker, $this) - : new AdapterPushHandler($adapter), + finishHandler: new AdapterPushHandler($adapter), ); } @@ -54,15 +47,16 @@ public function getName(): string public function push(MessageInterface $message): MessageInterface { - $this->logger->debug('Preparing to push message with message type "{messageType}".', ['messageType' => $message->getType()]); + $this->logger->debug( + 'Preparing to push message with message type "{messageType}".', + ['messageType' => $message->getType()], + ); $message = $this->dispatcher->dispatch($message); - if ($this->adapter === null) { - $this->logger->info('Processed message with message type "{messageType}" synchronously.', ['messageType' => $message->getType()]); - return $message; - } $id = IdEnvelope::fromMessage($message)->getId(); $this->logger->info( - $id === null ? 'Pushed message with message type "{messageType}" to the queue. ID doesn\'t assigned.' : 'Pushed message with message type "{messageType}" to the queue. Assigned ID #{id}.', + $id === null + ? 'Pushed message with message type "{messageType}" to the queue. ID doesn\'t assigned.' + : 'Pushed message with message type "{messageType}" to the queue. Assigned ID #{id}.', ['messageType' => $message->getType(), 'id' => $id], ); return $message; @@ -70,6 +64,6 @@ public function push(MessageInterface $message): MessageInterface public function status(string|int $id): MessageStatus { - return $this->adapter?->status($id) ?? MessageStatus::NOT_FOUND; + return $this->adapter->status($id); } } diff --git a/src/Middleware/Push/PushMiddlewareDispatcher.php b/src/Middleware/Push/PushMiddlewareDispatcher.php index ce2a89dc..4347d8a5 100644 --- a/src/Middleware/Push/PushMiddlewareDispatcher.php +++ b/src/Middleware/Push/PushMiddlewareDispatcher.php @@ -8,7 +8,7 @@ use Yiisoft\Queue\Message\MessageInterface; /** - * @internal Used internally by {@see QueueProducer}. + * @internal Used internally by {@see SyncQueueProducer} and {@see AsyncQueueProducer}. */ final class PushMiddlewareDispatcher { diff --git a/src/SyncQueueProducer.php b/src/SyncQueueProducer.php new file mode 100644 index 00000000..dcfdde5c --- /dev/null +++ b/src/SyncQueueProducer.php @@ -0,0 +1,59 @@ +name = StringNormalizer::normalize($name); + $this->dispatcher = new PushMiddlewareDispatcher( + middlewareFactory: $middlewareConfig->middlewareFactory, + middlewareDefinitions: [...$middlewareConfig->commonMiddlewareDefinitions, ...$middlewareDefinitions], + finishHandler: new SynchronousPushHandler($worker, $this), + ); + } + + public function getName(): string + { + return $this->name; + } + + public function push(MessageInterface $message): MessageInterface + { + $this->logger->debug('Preparing to push message with message type "{messageType}".', ['messageType' => $message->getType()]); + $message = $this->dispatcher->dispatch($message); + $this->logger->info('Processed message with message type "{messageType}" synchronously.', ['messageType' => $message->getType()]); + return $message; + } + + public function status(string|int $id): MessageStatus + { + return MessageStatus::NOT_FOUND; + } +} diff --git a/tests/Benchmark/QueueBench.php b/tests/Benchmark/QueueBench.php index 4803c019..7f56b83b 100644 --- a/tests/Benchmark/QueueBench.php +++ b/tests/Benchmark/QueueBench.php @@ -21,7 +21,7 @@ use Yiisoft\Queue\Middleware\FailureHandling\FailureMiddlewareFactory; use Yiisoft\Queue\Middleware\Push\PushMiddlewareConfig; use Yiisoft\Queue\Middleware\Push\PushMiddlewareFactory; -use Yiisoft\Queue\QueueProducer; +use Yiisoft\Queue\AsyncQueueProducer; use Yiisoft\Queue\QueueConsumer; use Yiisoft\Queue\QueueConsumerInterface; use Yiisoft\Queue\QueueProducerInterface; @@ -59,7 +59,7 @@ public function __construct() $this->serializer = new MessageSerializer(new JsonMessageEncoder()); $this->adapter = new VoidAdapter($this->serializer); - $this->producer = new QueueProducer( + $this->producer = new AsyncQueueProducer( $logger, new PushMiddlewareConfig(new PushMiddlewareFactory($container, $callableFactory)), $this->adapter, diff --git a/tests/Integration/MiddlewareTest.php b/tests/Integration/MiddlewareTest.php index 4ee51af5..3ef03f6f 100644 --- a/tests/Integration/MiddlewareTest.php +++ b/tests/Integration/MiddlewareTest.php @@ -24,7 +24,7 @@ use Yiisoft\Queue\Middleware\FailureHandling\FailureMiddlewareFactory; use Yiisoft\Queue\Middleware\Push\PushMiddlewareConfig; use Yiisoft\Queue\Middleware\Push\PushMiddlewareFactory; -use Yiisoft\Queue\QueueProducer; +use Yiisoft\Queue\SyncQueueProducer; use Yiisoft\Queue\QueueProducerInterface; use Yiisoft\Queue\Tests\Integration\Support\TestMiddleware; use Yiisoft\Queue\Worker\Worker; @@ -58,16 +58,17 @@ public function testFullStackPush(): void ); $worker = $this->createMock(WorkerInterface::class); $worker->method('process')->willReturnArgument(0); - $queue = new QueueProducer( + $queue = new SyncQueueProducer( $this->createMock(LoggerInterface::class), $pushMiddlewareConfig, - null, - 'test', $worker, - new TestMiddleware('channel 1'), - new TestMiddleware('channel 2'), - new TestMiddleware('channel 3'), - new TestMiddleware('channel 4'), + 'test', + [ + new TestMiddleware('channel 1'), + new TestMiddleware('channel 2'), + new TestMiddleware('channel 3'), + new TestMiddleware('channel 4'), + ], ); $message = new GenericMessage('test', ['initial']); diff --git a/tests/TestCase.php b/tests/TestCase.php index 592e821d..765f5dc9 100644 --- a/tests/TestCase.php +++ b/tests/TestCase.php @@ -22,7 +22,9 @@ use Yiisoft\Queue\Middleware\FailureHandling\FailureMiddlewareFactory; use Yiisoft\Queue\Middleware\Push\PushMiddlewareConfig; use Yiisoft\Queue\Middleware\Push\PushMiddlewareFactory; -use Yiisoft\Queue\QueueProducer; +use Yiisoft\Queue\AsyncQueueProducer; +use Yiisoft\Queue\QueueProducerInterface; +use Yiisoft\Queue\SyncQueueProducer; use Yiisoft\Queue\Worker\Worker; use Yiisoft\Queue\Worker\WorkerInterface; @@ -32,7 +34,7 @@ abstract class TestCase extends BaseTestCase { protected ?ContainerInterface $container = null; - protected ?QueueProducer $queue = null; + protected ?QueueProducerInterface $queue = null; protected ?LoopInterface $loop = null; protected ?WorkerInterface $worker = null; protected array $eventHandlers = []; @@ -51,9 +53,9 @@ protected function setUp(): void } /** - * @return QueueProducer The same object every time + * @return QueueProducerInterface The same object every time */ - protected function getQueue(): QueueProducer + protected function getQueue(): QueueProducerInterface { if ($this->queue === null) { $this->queue = $this->createQueue(); @@ -92,14 +94,20 @@ protected function getContainer(): ContainerInterface protected function createQueue( ?AdapterInterface $adapter = null, string|BackedEnum $name = QueueProducerProviderInterface::DEFAULT_QUEUE, - ): QueueProducer { - return new QueueProducer( - new NullLogger(), - $this->getPushMiddlewareConfig(), - $adapter, - $name, - $adapter === null ? $this->getWorker() : null, - ); + ): QueueProducerInterface { + return $adapter === null + ? new SyncQueueProducer( + new NullLogger(), + $this->getPushMiddlewareConfig(), + $this->getWorker(), + $name, + ) + : new AsyncQueueProducer( + new NullLogger(), + $this->getPushMiddlewareConfig(), + $adapter, + $name, + ); } protected function createLoop(): LoopInterface