diff --git a/config/di.php b/config/di.php index e8ae8707..609c2f58 100644 --- a/config/di.php +++ b/config/di.php @@ -3,6 +3,7 @@ declare(strict_types=1); use Psr\Container\ContainerInterface; +use Yiisoft\Definitions\Reference; use Yiisoft\Queue\Cli\LoopInterface; use Yiisoft\Queue\Cli\SignalLoop; use Yiisoft\Queue\Cli\SimpleLoop; @@ -21,6 +22,10 @@ use Yiisoft\Queue\Middleware\Push\PushMiddlewareConfig; use Yiisoft\Queue\Middleware\Push\PushMiddlewareFactory; use Yiisoft\Queue\Middleware\Push\PushMiddlewareFactoryInterface; +use Yiisoft\Queue\Middleware\Worker\WorkerFinalHandler; +use Yiisoft\Queue\Middleware\Worker\WorkerMiddlewareDispatcher; +use Yiisoft\Queue\Middleware\Worker\WorkerMiddlewareFactory; +use Yiisoft\Queue\Middleware\Worker\WorkerMiddlewareFactoryInterface; use Yiisoft\Queue\Message\Handler\HandlerResolver; use Yiisoft\Queue\Worker\Worker as QueueWorker; use Yiisoft\Queue\Worker\WorkerInterface; @@ -40,6 +45,7 @@ PushMiddlewareFactoryInterface::class => PushMiddlewareFactory::class, ConsumeMiddlewareFactoryInterface::class => ConsumeMiddlewareFactory::class, FailureMiddlewareFactoryInterface::class => FailureMiddlewareFactory::class, + WorkerMiddlewareFactoryInterface::class => WorkerMiddlewareFactory::class, PushMiddlewareConfig::class => [ '__construct()' => ['commonMiddlewareDefinitions' => $params['yiisoft/queue']['middlewares-push']], ], @@ -49,6 +55,18 @@ FailureMiddlewareDispatcher::class => [ '__construct()' => ['middlewareDefinitions' => $params['yiisoft/queue']['middlewares-fail']], ], + WorkerMiddlewareDispatcher::class => [ + '__construct()' => [ + 'middlewareDefinitions' => $params['yiisoft/queue']['middlewares-worker'], + 'finalHandler' => Reference::to(WorkerFinalHandler::class), + ], + ], + WorkerFinalHandler::class => [ + '__construct()' => [ + 'handlerResolver' => Reference::to(HandlerResolver::class), + 'consumeMiddlewareDispatcher' => Reference::to(ConsumeMiddlewareDispatcher::class), + ], + ], MessageEncoderInterface::class => JsonMessageEncoder::class, MessageSerializerInterface::class => MessageSerializer::class, MessageClassResolverInterface::class => [ diff --git a/config/params.php b/config/params.php index e502d23a..61639cd7 100644 --- a/config/params.php +++ b/config/params.php @@ -46,6 +46,7 @@ 'middlewares-push' => [], 'middlewares-consume' => [], 'middlewares-fail' => [], + 'middlewares-worker' => [], ], 'yiisoft/yii-debug' => [ 'collectors' => [ diff --git a/src/AsyncQueueProducer.php b/src/AsyncQueueProducer.php index f2ae8145..8d7745d9 100644 --- a/src/AsyncQueueProducer.php +++ b/src/AsyncQueueProducer.php @@ -11,6 +11,7 @@ use Yiisoft\Queue\Message\MessageInterface; use Yiisoft\Queue\Middleware\Push\AdapterPushHandler; use Yiisoft\Queue\Middleware\Push\PushMiddlewareConfig; +use Yiisoft\Queue\Middleware\Push\PushRequest; use Yiisoft\Queue\Middleware\Push\PushMiddlewareDispatcher; /** @@ -50,7 +51,7 @@ public function push(MessageInterface $message): MessageInterface 'Preparing to push message with message type "{messageType}".', ['messageType' => $message->getType()], ); - $message = $this->dispatcher->dispatch($message); + $message = $this->dispatcher->dispatch(new PushRequest($message, $this->queueName))->getMessage(); $id = IdEnvelope::fromMessage($message)->getId(); $this->logger->info( $id === null diff --git a/src/Middleware/Push/AdapterPushHandler.php b/src/Middleware/Push/AdapterPushHandler.php index 20e35c7d..c6befb60 100644 --- a/src/Middleware/Push/AdapterPushHandler.php +++ b/src/Middleware/Push/AdapterPushHandler.php @@ -5,7 +5,6 @@ namespace Yiisoft\Queue\Middleware\Push; use Yiisoft\Queue\Adapter\AdapterInterface; -use Yiisoft\Queue\Message\MessageInterface; /** * @internal @@ -16,8 +15,8 @@ public function __construct( private readonly AdapterInterface $adapter, ) {} - public function handlePush(MessageInterface $message): MessageInterface + public function handlePush(PushRequest $request): PushRequest { - return $this->adapter->push($message); + return $request->withMessage($this->adapter->push($request->getMessage())); } } diff --git a/src/Middleware/Push/Implementation/IdMiddleware.php b/src/Middleware/Push/Implementation/IdMiddleware.php index d2ec2482..31071665 100644 --- a/src/Middleware/Push/Implementation/IdMiddleware.php +++ b/src/Middleware/Push/Implementation/IdMiddleware.php @@ -5,25 +5,26 @@ namespace Yiisoft\Queue\Middleware\Push\Implementation; use Yiisoft\Queue\Message\IdEnvelope; -use Yiisoft\Queue\Message\MessageInterface; use Yiisoft\Queue\Middleware\Push\PushHandlerInterface; use Yiisoft\Queue\Middleware\Push\PushMiddlewareInterface; +use Yiisoft\Queue\Middleware\Push\PushRequest; /** * A middleware for message ID setting. */ final class IdMiddleware implements PushMiddlewareInterface { - public function processPush(MessageInterface $message, PushHandlerInterface $handler): MessageInterface + public function processPush(PushRequest $request, PushHandlerInterface $handler): PushRequest { + $message = $request->getMessage(); $envelope = IdEnvelope::fromMessage($message); if ($envelope->getId() === null) { return $handler->handlePush( - new IdEnvelope($message, uniqid('yii3-message-', true)), + $request->withMessage(new IdEnvelope($message, uniqid('yii3-message-', true))), ); } - return $handler->handlePush($message); + return $handler->handlePush($request); } } diff --git a/src/Middleware/Push/PushHandlerInterface.php b/src/Middleware/Push/PushHandlerInterface.php index 3f048123..0e7043c2 100644 --- a/src/Middleware/Push/PushHandlerInterface.php +++ b/src/Middleware/Push/PushHandlerInterface.php @@ -4,9 +4,7 @@ namespace Yiisoft\Queue\Middleware\Push; -use Yiisoft\Queue\Message\MessageInterface; - interface PushHandlerInterface { - public function handlePush(MessageInterface $message): MessageInterface; + public function handlePush(PushRequest $request): PushRequest; } diff --git a/src/Middleware/Push/PushMiddlewareDispatcher.php b/src/Middleware/Push/PushMiddlewareDispatcher.php index 77c23606..ef196b6b 100644 --- a/src/Middleware/Push/PushMiddlewareDispatcher.php +++ b/src/Middleware/Push/PushMiddlewareDispatcher.php @@ -5,7 +5,6 @@ namespace Yiisoft\Queue\Middleware\Push; use Closure; -use Yiisoft\Queue\Message\MessageInterface; /** * @internal Used internally by {@see SyncQueueProducer} and {@see AsyncQueueProducer}. @@ -33,13 +32,13 @@ public function __construct( /** * Dispatch message through middleware to get response. * - * @param MessageInterface $message Message to pass to middleware. + * @param PushRequest $request Request to pass to middleware. */ - public function dispatch(MessageInterface $message): MessageInterface + public function dispatch(PushRequest $request): PushRequest { $this->stack ??= new PushMiddlewareStack($this->buildMiddlewares(), $this->finalHandler); - return $this->stack->handlePush($message); + return $this->stack->handlePush($request); } /** diff --git a/src/Middleware/Push/PushMiddlewareFactory.php b/src/Middleware/Push/PushMiddlewareFactory.php index de70f488..98adb786 100644 --- a/src/Middleware/Push/PushMiddlewareFactory.php +++ b/src/Middleware/Push/PushMiddlewareFactory.php @@ -4,7 +4,6 @@ namespace Yiisoft\Queue\Middleware\Push; -use Yiisoft\Queue\Message\MessageInterface; use Yiisoft\Queue\Middleware\InvalidMiddlewareDefinitionException; use Yiisoft\Queue\Middleware\MiddlewareFactory; @@ -25,15 +24,14 @@ final class PushMiddlewareFactory extends MiddlewareFactory implements PushMiddl * * - A middleware object. * - A name of a middleware class. The middleware instance will be obtained from container and executed. - * - A callable with `function(MessageInterface $message, PushHandlerInterface $handler): - * MessageInterface` signature. + * - A callable with `function(PushRequest $request, PushHandlerInterface $handler): PushRequest` signature. * - A controller handler action in format `[TestController::class, 'index']`. `TestController` instance will * be created and `index()` method will be executed. * - A function returning a middleware. The middleware returned will be executed. * * For handler action and callable * typed parameters are automatically injected using dependency injection container. - * Current message and handler could be obtained by type-hinting for {@see MessageInterface} + * Current request and handler could be obtained by type-hinting for {@see PushRequest} * and {@see PushHandlerInterface}. * * @throws InvalidMiddlewareDefinitionException @@ -74,15 +72,15 @@ public function __construct(callable $callback) $this->callback = $callback; } - public function processPush(MessageInterface $message, PushHandlerInterface $handler): MessageInterface + public function processPush(PushRequest $request, PushHandlerInterface $handler): PushRequest { - $response = ($this->callback)($message, $handler); - if ($response instanceof MessageInterface) { + $response = ($this->callback)($request, $handler); + if ($response instanceof PushRequest) { return $response; } if ($response instanceof PushMiddlewareInterface) { - return $response->processPush($message, $handler); + return $response->processPush($request, $handler); } throw new InvalidMiddlewareDefinitionException($this->callback); diff --git a/src/Middleware/Push/PushMiddlewareInterface.php b/src/Middleware/Push/PushMiddlewareInterface.php index 290be84b..8002ab24 100644 --- a/src/Middleware/Push/PushMiddlewareInterface.php +++ b/src/Middleware/Push/PushMiddlewareInterface.php @@ -4,9 +4,7 @@ namespace Yiisoft\Queue\Middleware\Push; -use Yiisoft\Queue\Message\MessageInterface; - interface PushMiddlewareInterface { - public function processPush(MessageInterface $message, PushHandlerInterface $handler): MessageInterface; + public function processPush(PushRequest $request, PushHandlerInterface $handler): PushRequest; } diff --git a/src/Middleware/Push/PushMiddlewareStack.php b/src/Middleware/Push/PushMiddlewareStack.php index 11c4e432..819d30d9 100644 --- a/src/Middleware/Push/PushMiddlewareStack.php +++ b/src/Middleware/Push/PushMiddlewareStack.php @@ -5,7 +5,6 @@ namespace Yiisoft\Queue\Middleware\Push; use Closure; -use Yiisoft\Queue\Message\MessageInterface; /** * @internal @@ -31,10 +30,10 @@ public function __construct( private readonly PushHandlerInterface $finalHandler, ) {} - public function handlePush(MessageInterface $message): MessageInterface + public function handlePush(PushRequest $request): PushRequest { $this->stack ??= $this->build(); - return $this->stack->handlePush($message); + return $this->stack->handlePush($request); } private function build(): PushHandlerInterface @@ -66,11 +65,11 @@ public function __construct( private readonly PushHandlerInterface $handler, ) {} - public function handlePush(MessageInterface $message): MessageInterface + public function handlePush(PushRequest $request): PushRequest { $this->middleware ??= ($this->middlewareFactory)(); - return $this->middleware->processPush($message, $this->handler); + return $this->middleware->processPush($request, $this->handler); } }; } diff --git a/src/Middleware/Push/PushRequest.php b/src/Middleware/Push/PushRequest.php new file mode 100644 index 00000000..5da95fd4 --- /dev/null +++ b/src/Middleware/Push/PushRequest.php @@ -0,0 +1,38 @@ +message; + } + + public function getQueueName(): string + { + return $this->queueName; + } + + public function withMessage(MessageInterface $message): self + { + $instance = clone $this; + $instance->message = $message; + + return $instance; + } + + public function withQueueName(string $queueName): self + { + $instance = clone $this; + $instance->queueName = $queueName; + + return $instance; + } +} diff --git a/src/Middleware/Push/SynchronousPushHandler.php b/src/Middleware/Push/SynchronousPushHandler.php index 57a59c4d..3b072d86 100644 --- a/src/Middleware/Push/SynchronousPushHandler.php +++ b/src/Middleware/Push/SynchronousPushHandler.php @@ -4,7 +4,6 @@ namespace Yiisoft\Queue\Middleware\Push; -use Yiisoft\Queue\Message\MessageInterface; use Yiisoft\Queue\Worker\WorkerInterface; /** @@ -17,10 +16,10 @@ public function __construct( private readonly string $queueName, ) {} - public function handlePush(MessageInterface $message): MessageInterface + public function handlePush(PushRequest $request): PushRequest { - $this->worker->process($message, $this->queueName); + $this->worker->process($request->getMessage(), $request->getQueueName()); - return $message; + return $request; } } diff --git a/src/Middleware/Worker/WorkerFinalHandler.php b/src/Middleware/Worker/WorkerFinalHandler.php new file mode 100644 index 00000000..14f99c68 --- /dev/null +++ b/src/Middleware/Worker/WorkerFinalHandler.php @@ -0,0 +1,27 @@ +handlerResolver->resolve($request->getMessage()->getType()); + $consumeRequest = new ConsumeRequest($request->getMessage(), $request->getQueueName()); + $this->consumeMiddlewareDispatcher->dispatch($consumeRequest, new ConsumeFinalHandler($handler->handle(...))); + + return $request; + } +} diff --git a/src/Middleware/Worker/WorkerHandlerInterface.php b/src/Middleware/Worker/WorkerHandlerInterface.php new file mode 100644 index 00000000..74ead4f6 --- /dev/null +++ b/src/Middleware/Worker/WorkerHandlerInterface.php @@ -0,0 +1,10 @@ +stack ??= new WorkerMiddlewareStack($this->buildMiddlewares(), $this->finalHandler); + + return $this->stack->handleWorker($request); + } + + /** + * @psalm-return list + */ + private function buildMiddlewares(): array + { + /** @var list $middlewares */ + $middlewares = []; + $factory = $this->middlewareFactory; + + foreach ($this->middlewareDefinitions as $middlewareDefinition) { + $middlewares[] = static fn(): WorkerMiddlewareInterface => $factory->createWorkerMiddleware( + $middlewareDefinition, + ); + } + + return $middlewares; + } +} diff --git a/src/Middleware/Worker/WorkerMiddlewareFactory.php b/src/Middleware/Worker/WorkerMiddlewareFactory.php new file mode 100644 index 00000000..0346eb95 --- /dev/null +++ b/src/Middleware/Worker/WorkerMiddlewareFactory.php @@ -0,0 +1,74 @@ + + * @template-implements WorkerMiddlewareFactoryInterface + */ +final class WorkerMiddlewareFactory extends MiddlewareFactory implements WorkerMiddlewareFactoryInterface +{ + public function createWorkerMiddleware(mixed $middlewareDefinition): WorkerMiddlewareInterface + { + if ($middlewareDefinition instanceof WorkerMiddlewareInterface) { + return $middlewareDefinition; + } + + if ( + !is_callable($middlewareDefinition) + && !is_array($middlewareDefinition) + && !is_string($middlewareDefinition) + ) { + throw new InvalidMiddlewareDefinitionException($middlewareDefinition); + } + + /** @var object $middleware */ + $middleware = $this->create($middlewareDefinition); + if (!$middleware instanceof WorkerMiddlewareInterface) { + throw new InvalidMiddlewareDefinitionException($middlewareDefinition); + } + + return $middleware; + } + + protected function getInterfaceName(): string + { + return WorkerMiddlewareInterface::class; + } + + protected function wrapMiddleware(callable $callback): WorkerMiddlewareInterface + { + return new class ($callback) implements WorkerMiddlewareInterface { + /** @var callable */ + private readonly mixed $callback; + + public function __construct(callable $callback) + { + $this->callback = $callback; + } + + public function processWorker(WorkerRequest $request, WorkerHandlerInterface $handler): WorkerRequest + { + $response = ($this->callback)($request, $handler); + if ($response instanceof WorkerRequest) { + return $response; + } + + if ($response instanceof WorkerMiddlewareInterface) { + return $response->processWorker($request, $handler); + } + + throw new InvalidMiddlewareDefinitionException($this->callback); + } + }; + } +} diff --git a/src/Middleware/Worker/WorkerMiddlewareFactoryInterface.php b/src/Middleware/Worker/WorkerMiddlewareFactoryInterface.php new file mode 100644 index 00000000..6649d464 --- /dev/null +++ b/src/Middleware/Worker/WorkerMiddlewareFactoryInterface.php @@ -0,0 +1,23 @@ + $middlewares + */ + public function __construct( + private readonly array $middlewares, + private readonly WorkerHandlerInterface $finalHandler, + ) {} + + public function handleWorker(WorkerRequest $request): WorkerRequest + { + $this->stack ??= $this->build(); + + return $this->stack->handleWorker($request); + } + + private function build(): WorkerHandlerInterface + { + $handler = $this->finalHandler; + + foreach (array_reverse($this->middlewares) as $middleware) { + $handler = $this->wrap($middleware, $handler); + } + + return $handler; + } + + /** + * @psalm-param Closure():WorkerMiddlewareInterface $middlewareFactory + */ + private function wrap(Closure $middlewareFactory, WorkerHandlerInterface $handler): WorkerHandlerInterface + { + return new class ($middlewareFactory, $handler) implements WorkerHandlerInterface { + private ?WorkerMiddlewareInterface $middleware = null; + + /** + * @psalm-param Closure():WorkerMiddlewareInterface $middlewareFactory + */ + public function __construct( + private readonly Closure $middlewareFactory, + private readonly WorkerHandlerInterface $handler, + ) {} + + public function handleWorker(WorkerRequest $request): WorkerRequest + { + $this->middleware ??= ($this->middlewareFactory)(); + + return $this->middleware->processWorker($request, $this->handler); + } + }; + } +} diff --git a/src/Middleware/Worker/WorkerRequest.php b/src/Middleware/Worker/WorkerRequest.php new file mode 100644 index 00000000..428c1ee0 --- /dev/null +++ b/src/Middleware/Worker/WorkerRequest.php @@ -0,0 +1,41 @@ +message; + } + + public function getQueueName(): string + { + return $this->queueName; + } + + public function withMessage(MessageInterface $message): self + { + $instance = clone $this; + $instance->message = $message; + + return $instance; + } + + public function withQueueName(string $queueName): self + { + $instance = clone $this; + $instance->queueName = $queueName; + + return $instance; + } +} diff --git a/src/SyncQueueProducer.php b/src/SyncQueueProducer.php index 51990f2a..2a02de9b 100644 --- a/src/SyncQueueProducer.php +++ b/src/SyncQueueProducer.php @@ -8,6 +8,7 @@ use Psr\Log\LoggerInterface; use Yiisoft\Queue\Message\MessageInterface; use Yiisoft\Queue\Middleware\Push\PushMiddlewareConfig; +use Yiisoft\Queue\Middleware\Push\PushRequest; use Yiisoft\Queue\Middleware\Push\PushMiddlewareDispatcher; use Yiisoft\Queue\Middleware\Push\SynchronousPushHandler; use Yiisoft\Queue\Worker\WorkerInterface; @@ -46,7 +47,7 @@ public function getQueueName(): string 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); + $message = $this->dispatcher->dispatch(new PushRequest($message, $this->queueName))->getMessage(); $this->logger->info('Processed message with message type "{messageType}" synchronously.', ['messageType' => $message->getType()]); return $message; } diff --git a/src/Worker/Worker.php b/src/Worker/Worker.php index 1ea4cc86..6c16e799 100644 --- a/src/Worker/Worker.php +++ b/src/Worker/Worker.php @@ -7,11 +7,9 @@ use Psr\Log\LoggerInterface; use Throwable; use Yiisoft\Queue\Exception\MessageFailureException; -use Yiisoft\Queue\Message\Handler\HandlerResolver; use Yiisoft\Queue\Message\MessageInterface; -use Yiisoft\Queue\Middleware\Consume\ConsumeFinalHandler; -use Yiisoft\Queue\Middleware\Consume\ConsumeMiddlewareDispatcher; -use Yiisoft\Queue\Middleware\Consume\ConsumeRequest; +use Yiisoft\Queue\Middleware\Worker\WorkerMiddlewareDispatcher; +use Yiisoft\Queue\Middleware\Worker\WorkerRequest; use Yiisoft\Queue\Middleware\FailureHandling\FailureFinalHandler; use Yiisoft\Queue\Middleware\FailureHandling\FailureHandlingRequest; use Yiisoft\Queue\Middleware\FailureHandling\FailureMiddlewareDispatcher; @@ -21,9 +19,8 @@ final class Worker implements WorkerInterface { public function __construct( private readonly LoggerInterface $logger, - private readonly ConsumeMiddlewareDispatcher $consumeMiddlewareDispatcher, + private readonly WorkerMiddlewareDispatcher $workerMiddlewareDispatcher, private readonly FailureMiddlewareDispatcher $failureMiddlewareDispatcher, - private readonly HandlerResolver $handlerResolver, ) {} /** @@ -38,11 +35,9 @@ public function process(MessageInterface $message, string $queueName): void $this->logger->info('Processing message #{message}.', ['message' => $messageId]); } - $request = new ConsumeRequest($message, $queueName); + $request = new WorkerRequest($message, $queueName); try { - $handler = $this->handlerResolver->resolve($message->getType()); - $finalHandler = new ConsumeFinalHandler($handler->handle(...)); - $this->consumeMiddlewareDispatcher->dispatch($request, $finalHandler); + $this->workerMiddlewareDispatcher->dispatch($request); } catch (Throwable $exception) { $request = new FailureHandlingRequest($request->getMessage(), $exception, $request->getQueueName()); diff --git a/tests/Benchmark/QueueBench.php b/tests/Benchmark/QueueBench.php index c5bc691f..42cbf187 100644 --- a/tests/Benchmark/QueueBench.php +++ b/tests/Benchmark/QueueBench.php @@ -19,6 +19,9 @@ use Yiisoft\Queue\Middleware\FailureHandling\FailureMiddlewareFactory; use Yiisoft\Queue\Middleware\Push\PushMiddlewareConfig; use Yiisoft\Queue\Middleware\Push\PushMiddlewareFactory; +use Yiisoft\Queue\Middleware\Worker\WorkerFinalHandler; +use Yiisoft\Queue\Middleware\Worker\WorkerMiddlewareDispatcher; +use Yiisoft\Queue\Middleware\Worker\WorkerMiddlewareFactory; use Yiisoft\Queue\AsyncQueueProducer; use Yiisoft\Queue\QueueConsumer; use Yiisoft\Queue\QueueConsumerInterface; @@ -42,17 +45,18 @@ public function __construct() $worker = new Worker( $logger, - new ConsumeMiddlewareDispatcher(new ConsumeMiddlewareFactory($container)), - new FailureMiddlewareDispatcher( - new FailureMiddlewareFactory($container), + new WorkerMiddlewareDispatcher( + new WorkerMiddlewareFactory($container), [], + new WorkerFinalHandler( + new HandlerResolver( + ['foo' => static function (): void {}], + $container, + ), + new ConsumeMiddlewareDispatcher(new ConsumeMiddlewareFactory($container)), + ), ), - new HandlerResolver( - [ - 'foo' => static function (): void {}, - ], - $container, - ), + new FailureMiddlewareDispatcher(new FailureMiddlewareFactory($container), []), ); $this->serializer = new MessageSerializer(new JsonMessageEncoder()); $this->adapter = new VoidAdapter($this->serializer); diff --git a/tests/Integration/MessageConsumingTest.php b/tests/Integration/MessageConsumingTest.php index e978602f..959d9f06 100644 --- a/tests/Integration/MessageConsumingTest.php +++ b/tests/Integration/MessageConsumingTest.php @@ -12,6 +12,9 @@ use Yiisoft\Queue\Middleware\Consume\ConsumeMiddlewareFactoryInterface; use Yiisoft\Queue\Middleware\FailureHandling\FailureMiddlewareDispatcher; use Yiisoft\Queue\Middleware\FailureHandling\FailureMiddlewareFactoryInterface; +use Yiisoft\Queue\Middleware\Worker\WorkerFinalHandler; +use Yiisoft\Queue\Middleware\Worker\WorkerMiddlewareDispatcher; +use Yiisoft\Queue\Middleware\Worker\WorkerMiddlewareFactory; use Yiisoft\Queue\Message\Handler\HandlerResolver; use Yiisoft\Queue\Tests\Integration\Support\TestHandler; use Yiisoft\Queue\Tests\TestCase; @@ -30,15 +33,21 @@ public function testMessagesConsumed(): void $container = $this->createMock(ContainerInterface::class); $worker = new Worker( new NullLogger(), - new ConsumeMiddlewareDispatcher($this->createMock(ConsumeMiddlewareFactoryInterface::class)), - new FailureMiddlewareDispatcher($this->createMock(FailureMiddlewareFactoryInterface::class), []), - new HandlerResolver( - [ - 'test' => fn(MessageInterface $message): mixed => $this->messagesProcessed[] = $message->getPayload(), - 'test2' => fn(MessageInterface $message): mixed => $this->messagesProcessedSecond[] = $message->getPayload(), - ], - $container, + new WorkerMiddlewareDispatcher( + new WorkerMiddlewareFactory($container), + [], + new WorkerFinalHandler( + new HandlerResolver( + [ + 'test' => fn(MessageInterface $message): mixed => $this->messagesProcessed[] = $message->getPayload(), + 'test2' => fn(MessageInterface $message): mixed => $this->messagesProcessedSecond[] = $message->getPayload(), + ], + $container, + ), + new ConsumeMiddlewareDispatcher($this->createMock(ConsumeMiddlewareFactoryInterface::class)), + ), ), + new FailureMiddlewareDispatcher($this->createMock(FailureMiddlewareFactoryInterface::class), []), ); $messages = [1, 'foo', 'bar-baz']; @@ -59,9 +68,15 @@ public function testMessagesConsumedByHandlerClass(): void $container->method('has')->with(TestHandler::class)->willReturn(true); $worker = new Worker( new NullLogger(), - new ConsumeMiddlewareDispatcher($this->createMock(ConsumeMiddlewareFactoryInterface::class)), + new WorkerMiddlewareDispatcher( + new WorkerMiddlewareFactory($container), + [], + new WorkerFinalHandler( + new HandlerResolver([], $container), + new ConsumeMiddlewareDispatcher($this->createMock(ConsumeMiddlewareFactoryInterface::class)), + ), + ), new FailureMiddlewareDispatcher($this->createMock(FailureMiddlewareFactoryInterface::class), []), - new HandlerResolver([], $container), ); $messages = [1, 'foo', 'bar-baz']; diff --git a/tests/Integration/MiddlewareTest.php b/tests/Integration/MiddlewareTest.php index 44df9211..ec776f85 100644 --- a/tests/Integration/MiddlewareTest.php +++ b/tests/Integration/MiddlewareTest.php @@ -8,26 +8,29 @@ use PHPUnit\Framework\TestCase; use Psr\Container\ContainerInterface; use Psr\Log\LoggerInterface; -use Yiisoft\Test\Support\Container\SimpleContainer; -use Yiisoft\Test\Support\Log\SimpleLogger; use Yiisoft\Queue\Message\GenericMessage; +use Yiisoft\Queue\Message\Handler\HandlerResolver; use Yiisoft\Queue\Message\MessageInterface; use Yiisoft\Queue\Middleware\Consume\ConsumeMiddlewareDispatcher; use Yiisoft\Queue\Middleware\Consume\ConsumeMiddlewareFactory; use Yiisoft\Queue\Middleware\FailureHandling\FailureFinalHandler; use Yiisoft\Queue\Middleware\FailureHandling\FailureHandlingRequest; use Yiisoft\Queue\Middleware\FailureHandling\FailureMiddlewareDispatcher; +use Yiisoft\Queue\Middleware\FailureHandling\FailureMiddlewareFactory; use Yiisoft\Queue\Middleware\FailureHandling\Implementation\ExponentialDelayMiddleware; use Yiisoft\Queue\Middleware\FailureHandling\Implementation\SendAgainMiddleware; -use Yiisoft\Queue\Middleware\FailureHandling\FailureMiddlewareFactory; use Yiisoft\Queue\Middleware\Push\PushMiddlewareConfig; use Yiisoft\Queue\Middleware\Push\PushMiddlewareFactory; -use Yiisoft\Queue\SyncQueueProducer; +use Yiisoft\Queue\Middleware\Worker\WorkerFinalHandler; +use Yiisoft\Queue\Middleware\Worker\WorkerMiddlewareDispatcher; +use Yiisoft\Queue\Middleware\Worker\WorkerMiddlewareFactory; use Yiisoft\Queue\QueueProducerInterface; -use Yiisoft\Queue\Message\Handler\HandlerResolver; +use Yiisoft\Queue\SyncQueueProducer; use Yiisoft\Queue\Tests\Integration\Support\TestMiddleware; use Yiisoft\Queue\Worker\Worker; use Yiisoft\Queue\Worker\WorkerInterface; +use Yiisoft\Test\Support\Container\SimpleContainer; +use Yiisoft\Test\Support\Log\SimpleLogger; final class MiddlewareTest extends TestCase { @@ -98,14 +101,20 @@ public function testFullStackConsume(): void $handledMessage = null; $worker = new Worker( new SimpleLogger(), - $consumeMiddlewareDispatcher, - $failureMiddlewareDispatcher, - new HandlerResolver( - ['test' => static function (MessageInterface $message) use (&$handledMessage): void { - $handledMessage = $message; - }], - $container, + new WorkerMiddlewareDispatcher( + new WorkerMiddlewareFactory($container), + [], + new WorkerFinalHandler( + new HandlerResolver( + ['test' => static function (MessageInterface $message) use (&$handledMessage): void { + $handledMessage = $message; + }], + $container, + ), + $consumeMiddlewareDispatcher, + ), ), + $failureMiddlewareDispatcher, ); $message = new GenericMessage('test', ['initial']); diff --git a/tests/Integration/Support/TestMiddleware.php b/tests/Integration/Support/TestMiddleware.php index 05eb63a7..465cc97a 100644 --- a/tests/Integration/Support/TestMiddleware.php +++ b/tests/Integration/Support/TestMiddleware.php @@ -5,23 +5,24 @@ namespace Yiisoft\Queue\Tests\Integration\Support; use Yiisoft\Queue\Message\GenericMessage; -use Yiisoft\Queue\Message\MessageInterface; use Yiisoft\Queue\Middleware\Consume\ConsumeRequest; use Yiisoft\Queue\Middleware\Consume\ConsumeHandlerInterface; use Yiisoft\Queue\Middleware\Consume\ConsumeMiddlewareInterface; use Yiisoft\Queue\Middleware\Push\PushHandlerInterface; use Yiisoft\Queue\Middleware\Push\PushMiddlewareInterface; +use Yiisoft\Queue\Middleware\Push\PushRequest; final class TestMiddleware implements PushMiddlewareInterface, ConsumeMiddlewareInterface { public function __construct(private readonly string $stage) {} - public function processPush(MessageInterface $message, PushHandlerInterface $handler): MessageInterface + public function processPush(PushRequest $request, PushHandlerInterface $handler): PushRequest { + $message = $request->getMessage(); $stack = $message->getPayload(); $stack[] = $this->stage; - return $handler->handlePush(new GenericMessage($message->getType(), $stack)); + return $handler->handlePush($request->withMessage(new GenericMessage($message->getType(), $stack))); } public function processConsume(ConsumeRequest $request, ConsumeHandlerInterface $handler): ConsumeRequest diff --git a/tests/Integration/SynchronousFailureHandlingTest.php b/tests/Integration/SynchronousFailureHandlingTest.php index 2059887b..2412e876 100644 --- a/tests/Integration/SynchronousFailureHandlingTest.php +++ b/tests/Integration/SynchronousFailureHandlingTest.php @@ -21,6 +21,9 @@ use Yiisoft\Queue\Middleware\FailureHandling\Implementation\SendAgainMiddleware; use Yiisoft\Queue\Middleware\Push\PushMiddlewareConfig; use Yiisoft\Queue\Middleware\Push\PushMiddlewareFactory; +use Yiisoft\Queue\Middleware\Worker\WorkerFinalHandler; +use Yiisoft\Queue\Middleware\Worker\WorkerMiddlewareDispatcher; +use Yiisoft\Queue\Middleware\Worker\WorkerMiddlewareFactory; use Yiisoft\Queue\Provider\InvalidQueueConfigException; use Yiisoft\Queue\Stubs\InMemoryAdapter; use Yiisoft\Queue\SyncQueueProducer; @@ -106,12 +109,18 @@ private function createSyncProducer(array $failureMiddlewareDefinitions): SyncQu $container = new SimpleContainer(); $worker = new Worker( new NullLogger(), - new ConsumeMiddlewareDispatcher(new ConsumeMiddlewareFactory($container)), + new WorkerMiddlewareDispatcher( + new WorkerMiddlewareFactory($container), + [], + new WorkerFinalHandler( + new HandlerResolver($this->getMessageHandlers(), $container), + new ConsumeMiddlewareDispatcher(new ConsumeMiddlewareFactory($container)), + ), + ), new FailureMiddlewareDispatcher( new FailureMiddlewareFactory($container), $failureMiddlewareDefinitions, ), - new HandlerResolver($this->getMessageHandlers(), $container), ); return new SyncQueueProducer( diff --git a/tests/TestCase.php b/tests/TestCase.php index 8660c167..45d4f5fe 100644 --- a/tests/TestCase.php +++ b/tests/TestCase.php @@ -19,6 +19,9 @@ use Yiisoft\Queue\Middleware\FailureHandling\FailureMiddlewareDispatcher; use Yiisoft\Queue\Middleware\FailureHandling\FailureMiddlewareFactory; use Yiisoft\Queue\Middleware\Push\PushMiddlewareConfig; +use Yiisoft\Queue\Middleware\Worker\WorkerFinalHandler; +use Yiisoft\Queue\Middleware\Worker\WorkerMiddlewareDispatcher; +use Yiisoft\Queue\Middleware\Worker\WorkerMiddlewareFactory; use Yiisoft\Queue\Middleware\Push\PushMiddlewareFactory; use Yiisoft\Queue\AsyncQueueProducer; use Yiisoft\Queue\QueueProducerInterface; @@ -110,12 +113,15 @@ protected function createWorker(): WorkerInterface { return new Worker( new NullLogger(), - $this->getConsumeMiddlewareDispatcher(), - $this->getFailureMiddlewareDispatcher(), - new HandlerResolver( - $this->getMessageHandlers(), - $this->getContainer(), + new WorkerMiddlewareDispatcher( + new WorkerMiddlewareFactory($this->getContainer()), + [], + new WorkerFinalHandler( + new HandlerResolver($this->getMessageHandlers(), $this->getContainer()), + $this->getConsumeMiddlewareDispatcher(), + ), ), + $this->getFailureMiddlewareDispatcher(), ); } diff --git a/tests/Unit/AsyncQueueProducerTest.php b/tests/Unit/AsyncQueueProducerTest.php new file mode 100644 index 00000000..851dbb1e --- /dev/null +++ b/tests/Unit/AsyncQueueProducerTest.php @@ -0,0 +1,116 @@ +createMock(LoggerInterface::class); + $logger->expects(self::once()) + ->method('debug') + ->with( + 'Preparing to push message with message type "{messageType}".', + ['messageType' => 'test'], + ); + $logger->expects(self::once()) + ->method('info') + ->with( + 'Pushed message with message type "{messageType}" to the queue. Assigned ID #{id}.', + ['messageType' => 'test', 'id' => 0], + ); + + $producer = new AsyncQueueProducer( + $logger, + new PushMiddlewareConfig($this->createMock(PushMiddlewareFactoryInterface::class), []), + new InMemoryAdapter(), + ); + + $producer->push(new GenericMessage('test', [])); + } + + public function testPushLogsCompletionWithoutAssignedId(): void + { + $logger = $this->createMock(LoggerInterface::class); + $logger->expects(self::once()) + ->method('debug') + ->with( + 'Preparing to push message with message type "{messageType}".', + ['messageType' => 'test'], + ); + $logger->expects(self::once()) + ->method('info') + ->with( + 'Pushed message with message type "{messageType}" to the queue. ID doesn\'t assigned.', + ['messageType' => 'test', 'id' => null], + ); + + $adapter = $this->createMock(AdapterInterface::class); + $adapter->expects(self::once()) + ->method('push') + ->willReturnArgument(0); + + $producer = new AsyncQueueProducer( + $logger, + new PushMiddlewareConfig($this->createMock(PushMiddlewareFactoryInterface::class), []), + $adapter, + ); + + $producer->push(new GenericMessage('test', [])); + } + + public function testCommonAndQueueSpecificMiddlewareDefinitionsAreCombined(): void + { + $definitions = []; + $factory = $this->createMock(PushMiddlewareFactoryInterface::class); + $factory->expects(self::exactly(2)) + ->method('createPushMiddleware') + ->willReturnCallback(function (mixed $definition) use (&$definitions): PushMiddlewareInterface { + $definitions[] = $definition; + + return new class ($definition) implements PushMiddlewareInterface { + public function __construct(private readonly string $stage) {} + + public function processPush(PushRequest $request, PushHandlerInterface $handler): PushRequest + { + $message = $request->getMessage(); + $payload = (array) $message->getPayload(); + $payload[] = $this->stage; + + return $handler->handlePush($request->withMessage( + new GenericMessage($message->getType(), $payload), + )); + } + }; + }); + + $adapter = new InMemoryAdapter(); + $producer = new AsyncQueueProducer( + new NullLogger(), + new PushMiddlewareConfig($factory, ['common']), + $adapter, + 'queue', + ['specific'], + ); + + $producer->push(new GenericMessage('test', [])); + + self::assertSame(['common', 'specific'], $definitions); + self::assertSame(['common', 'specific'], $adapter->getMessagesList()[0]->getPayload()); + } +} diff --git a/tests/Unit/Middleware/DispatcherCachingTest.php b/tests/Unit/Middleware/DispatcherCachingTest.php new file mode 100644 index 00000000..d420eaf7 --- /dev/null +++ b/tests/Unit/Middleware/DispatcherCachingTest.php @@ -0,0 +1,87 @@ +createMock(PushMiddlewareFactoryInterface::class); + $factoryCalls = 0; + $factory->expects(self::once())->method('createPushMiddleware')->willReturnCallback( + function () use (&$factoryCalls): PushMiddlewareInterface { + $factoryCalls++; + + return new class implements PushMiddlewareInterface { + public function processPush(PushRequest $request, PushHandlerInterface $handler): PushRequest + { + return $handler->handlePush($request); + } + }; + }, + ); + $final = new class implements PushHandlerInterface { + public function handlePush(PushRequest $request): PushRequest + { + return $request; + } + }; + $dispatcher = new PushMiddlewareDispatcher($factory, ['definition'], $final); + self::assertSame(0, $factoryCalls); + $request = new PushRequest(new GenericMessage('test', null), 'queue'); + + $dispatcher->dispatch($request); + self::assertSame(1, $factoryCalls); + + $dispatcher->dispatch($request); + self::assertSame(1, $factoryCalls); + } + + public function testWorkerMiddlewareIsCreatedOnlyOnceForRepeatedDispatches(): void + { + $factory = $this->createMock(WorkerMiddlewareFactoryInterface::class); + $factoryCalls = 0; + $factory->expects(self::once())->method('createWorkerMiddleware')->willReturnCallback( + function () use (&$factoryCalls): WorkerMiddlewareInterface { + $factoryCalls++; + + return new class implements WorkerMiddlewareInterface { + public function processWorker(WorkerRequest $request, WorkerHandlerInterface $handler): WorkerRequest + { + return $handler->handleWorker($request); + } + }; + }, + ); + $final = new class implements WorkerHandlerInterface { + public function handleWorker(WorkerRequest $request): WorkerRequest + { + return $request; + } + }; + $dispatcher = new WorkerMiddlewareDispatcher($factory, ['definition'], $final); + self::assertSame(0, $factoryCalls); + $request = new WorkerRequest(new GenericMessage('test', null), 'queue'); + + $dispatcher->dispatch($request); + self::assertSame(1, $factoryCalls); + + $dispatcher->dispatch($request); + self::assertSame(1, $factoryCalls); + } +} diff --git a/tests/Unit/Middleware/Push/AdapterPushHandlerTest.php b/tests/Unit/Middleware/Push/AdapterPushHandlerTest.php index b8e13c5b..2f792ec8 100644 --- a/tests/Unit/Middleware/Push/AdapterPushHandlerTest.php +++ b/tests/Unit/Middleware/Push/AdapterPushHandlerTest.php @@ -6,7 +6,9 @@ use PHPUnit\Framework\TestCase; use Yiisoft\Queue\Message\GenericMessage; +use Yiisoft\Queue\Message\IdEnvelope; use Yiisoft\Queue\Middleware\Push\AdapterPushHandler; +use Yiisoft\Queue\Middleware\Push\PushRequest; use Yiisoft\Queue\Stubs\InMemoryAdapter; final class AdapterPushHandlerTest extends TestCase @@ -17,8 +19,13 @@ public function testHandlePushUsesAdapter(): void $handler = new AdapterPushHandler($adapter); $message = new GenericMessage('handler', 'data'); - $handler->handlePush($message); + $result = $handler->handlePush(new PushRequest($message, 'queue')); + $enrichedMessage = $result->getMessage(); + + self::assertInstanceOf(IdEnvelope::class, $enrichedMessage); + self::assertSame(0, $enrichedMessage->getId()); + self::assertSame($message, $enrichedMessage->getMessage()); self::assertSame([$message], $adapter->getMessagesList()); } } diff --git a/tests/Unit/Middleware/Push/Implementation/IdMiddlewareTest.php b/tests/Unit/Middleware/Push/Implementation/IdMiddlewareTest.php index d3ffb71f..f34a9344 100644 --- a/tests/Unit/Middleware/Push/Implementation/IdMiddlewareTest.php +++ b/tests/Unit/Middleware/Push/Implementation/IdMiddlewareTest.php @@ -9,6 +9,7 @@ use Yiisoft\Queue\Message\GenericMessage; use Yiisoft\Queue\Middleware\Push\Implementation\IdMiddleware; use Yiisoft\Queue\Middleware\Push\PushHandlerInterface; +use Yiisoft\Queue\Middleware\Push\PushRequest; final class IdMiddlewareTest extends TestCase { @@ -22,13 +23,13 @@ public function testWithId(): void ->willReturnArgument(0); $middleware = new IdMiddleware(); - $result = $middleware->processPush($message, $handler); + $result = $middleware->processPush(new PushRequest($message, 'test-queue'), $handler); - $this->assertSame($message, $result); - $this->assertNotInstanceOf(IdEnvelope::class, $result); - $this->assertEquals('test-id', $result->getMeta()[IdEnvelope::META_ID]); - $this->assertSame($message->getPayload(), $result->getPayload()); - $this->assertSame($message->getType(), $result->getType()); + $this->assertSame($message, $result->getMessage()); + $this->assertNotInstanceOf(IdEnvelope::class, $result->getMessage()); + $this->assertEquals('test-id', $result->getMessage()->getMeta()[IdEnvelope::META_ID]); + $this->assertSame($message->getPayload(), $result->getMessage()->getPayload()); + $this->assertSame($message->getType(), $result->getMessage()->getType()); } public function testWithoutId(): void @@ -41,13 +42,13 @@ public function testWithoutId(): void ->willReturnArgument(0); $middleware = new IdMiddleware(); - $result = $middleware->processPush($message, $handler); + $result = $middleware->processPush(new PushRequest($message, 'test-queue'), $handler); - $this->assertInstanceOf(IdEnvelope::class, $result); - $this->assertNotSame($message, $result); - $this->assertNotEmpty($result->getMeta()[IdEnvelope::META_ID] ?? null); - $this->assertSame($message->getPayload(), $result->getPayload()); - $this->assertSame($message->getType(), $result->getType()); + $this->assertInstanceOf(IdEnvelope::class, $result->getMessage()); + $this->assertNotSame($message, $result->getMessage()); + $this->assertNotEmpty($result->getMessage()->getMeta()[IdEnvelope::META_ID] ?? null); + $this->assertSame($message->getPayload(), $result->getMessage()->getPayload()); + $this->assertSame($message->getType(), $result->getMessage()->getType()); } public function testWithEmptyId(): void @@ -60,13 +61,13 @@ public function testWithEmptyId(): void ->willReturnArgument(0); $middleware = new IdMiddleware(); - $result = $middleware->processPush($message, $handler); + $result = $middleware->processPush(new PushRequest($message, 'test-queue'), $handler); - $this->assertInstanceOf(IdEnvelope::class, $result); - $this->assertNotSame($message, $result); - $this->assertNotEmpty($result->getMeta()[IdEnvelope::META_ID] ?? null); - $this->assertNotSame('', $result->getMeta()[IdEnvelope::META_ID]); - $this->assertSame($message->getPayload(), $result->getPayload()); - $this->assertSame($message->getType(), $result->getType()); + $this->assertInstanceOf(IdEnvelope::class, $result->getMessage()); + $this->assertNotSame($message, $result->getMessage()); + $this->assertNotEmpty($result->getMessage()->getMeta()[IdEnvelope::META_ID] ?? null); + $this->assertNotSame('', $result->getMessage()->getMeta()[IdEnvelope::META_ID]); + $this->assertSame($message->getPayload(), $result->getMessage()->getPayload()); + $this->assertSame($message->getType(), $result->getMessage()->getType()); } } diff --git a/tests/Unit/Middleware/Push/MiddlewareDispatcherTest.php b/tests/Unit/Middleware/Push/MiddlewareDispatcherTest.php index 7ec7a0eb..b3826e7d 100644 --- a/tests/Unit/Middleware/Push/MiddlewareDispatcherTest.php +++ b/tests/Unit/Middleware/Push/MiddlewareDispatcherTest.php @@ -12,6 +12,7 @@ use Yiisoft\Queue\Message\MessageInterface; use Yiisoft\Queue\Middleware\Push\PushHandlerInterface; use Yiisoft\Queue\Middleware\Push\PushMiddlewareFactory; +use Yiisoft\Queue\Middleware\Push\PushRequest; use Yiisoft\Queue\Middleware\Push\PushMiddlewareDispatcher; use Yiisoft\Queue\Stubs\InMemoryAdapter; use Yiisoft\Queue\Tests\Unit\Middleware\Push\Support\TestCallableMiddleware; @@ -24,13 +25,14 @@ public function testCallableMiddlewareCalled(): void $message = $this->getMessage(); $dispatcher = $this->createDispatcher(middlewareDefinitions: [ - static function (MessageInterface $message, PushHandlerInterface $handler): MessageInterface { - return new GenericMessage('test', 'New closure test data'); + static function (PushRequest $request, PushHandlerInterface $handler): PushRequest { + return $request->withMessage(new GenericMessage('test', 'New closure test data')); }, ]); - $result = $dispatcher->dispatch($message); - $this->assertSame('New closure test data', $result->getPayload()); + $result = $dispatcher->dispatch(new PushRequest($message, 'queue')); + $this->assertSame('New closure test data', $result->getMessage()->getPayload()); + $this->assertSame('queue', $result->getQueueName()); } public function testArrayMiddlewareCallableDefinition(): void @@ -42,8 +44,8 @@ public function testArrayMiddlewareCallableDefinition(): void ], ); $dispatcher = $this->createDispatcher($container, [[TestCallableMiddleware::class, 'index']]); - $result = $dispatcher->dispatch($message); - $this->assertSame('New test data', $result->getPayload()); + $result = $dispatcher->dispatch(new PushRequest($message, 'queue')); + $this->assertSame('New test data', $result->getMessage()->getPayload()); } public function testFactoryArrayDefinition(): void @@ -55,38 +57,38 @@ public function testFactoryArrayDefinition(): void '__construct()' => ['message' => 'New test data from the definition'], ]; $dispatcher = $this->createDispatcher($container, [$definition]); - $result = $dispatcher->dispatch($message); - $this->assertSame('New test data from the definition', $result->getPayload()); + $result = $dispatcher->dispatch(new PushRequest($message, 'queue')); + $this->assertSame('New test data from the definition', $result->getMessage()->getPayload()); } public function testMiddlewareFullStackCalled(): void { $message = $this->getMessage(); - $middleware1 = static function (MessageInterface $message, PushHandlerInterface $handler): MessageInterface { - return $handler->handlePush(new GenericMessage($message->getType(), 'new test data')); + $middleware1 = static function (PushRequest $request, PushHandlerInterface $handler): PushRequest { + return $handler->handlePush($request->withMessage(new GenericMessage($request->getMessage()->getType(), 'new test data'))); }; - $middleware2 = static function (MessageInterface $message, PushHandlerInterface $handler): MessageInterface { - return $handler->handlePush($message); + $middleware2 = static function (PushRequest $request, PushHandlerInterface $handler): PushRequest { + return $handler->handlePush($request); }; $dispatcher = $this->createDispatcher(middlewareDefinitions: [$middleware1, $middleware2]); - $result = $dispatcher->dispatch($message); - $this->assertSame('new test data', $result->getPayload()); + $result = $dispatcher->dispatch(new PushRequest($message, 'queue')); + $this->assertSame('new test data', $result->getMessage()->getPayload()); } public function testMiddlewareStackInterrupted(): void { $message = $this->getMessage(); - $middleware1 = static fn(MessageInterface $message, PushHandlerInterface $handler): MessageInterface => new GenericMessage($message->getType(), 'first'); - $middleware2 = static fn(MessageInterface $message, PushHandlerInterface $handler): MessageInterface => new GenericMessage($message->getType(), 'second'); + $middleware1 = static fn(PushRequest $request, PushHandlerInterface $handler): PushRequest => $request->withMessage(new GenericMessage($request->getMessage()->getType(), 'first')); + $middleware2 = static fn(PushRequest $request, PushHandlerInterface $handler): PushRequest => $request->withMessage(new GenericMessage($request->getMessage()->getType(), 'second')); $dispatcher = $this->createDispatcher(middlewareDefinitions: [$middleware1, $middleware2]); - $result = $dispatcher->dispatch($message); - $this->assertSame('first', $result->getPayload()); + $result = $dispatcher->dispatch(new PushRequest($message, 'queue')); + $this->assertSame('first', $result->getMessage()->getPayload()); } private function createDispatcher( @@ -99,9 +101,9 @@ private function createDispatcher( new PushMiddlewareFactory($container), $middlewareDefinitions, new class implements PushHandlerInterface { - public function handlePush(MessageInterface $message): MessageInterface + public function handlePush(PushRequest $request): PushRequest { - return $message; + return $request; } }, ); diff --git a/tests/Unit/Middleware/Push/MiddlewareFactoryTest.php b/tests/Unit/Middleware/Push/MiddlewareFactoryTest.php index 6a30802d..4866a89b 100644 --- a/tests/Unit/Middleware/Push/MiddlewareFactoryTest.php +++ b/tests/Unit/Middleware/Push/MiddlewareFactoryTest.php @@ -14,6 +14,7 @@ use Yiisoft\Queue\Middleware\InvalidMiddlewareDefinitionException; use Yiisoft\Queue\Middleware\Push\PushHandlerInterface; use Yiisoft\Queue\Middleware\Push\PushMiddlewareFactory; +use Yiisoft\Queue\Middleware\Push\PushRequest; use Yiisoft\Queue\Middleware\Push\PushMiddlewareFactoryInterface; use Yiisoft\Queue\Middleware\Push\PushMiddlewareInterface; use Yiisoft\Queue\Stubs\InMemoryAdapter; @@ -39,9 +40,9 @@ public function testCreateCallableFromArray(): void self::assertSame( 'New test data', $middleware->processPush( - $this->getMessage(), + new PushRequest($this->getMessage(), 'queue'), $this->createMock(PushHandlerInterface::class), - )->getPayload(), + )->getMessage()->getPayload(), ); } @@ -49,16 +50,16 @@ public function testCreateFromClosureResponse(): void { $container = $this->getContainer([TestCallableMiddleware::class => new TestCallableMiddleware()]); $middleware = $this->getMiddlewareFactory($container)->createPushMiddleware( - static function (): MessageInterface { - return new GenericMessage('test', 'test data'); + static function (): PushRequest { + return new PushRequest(new GenericMessage('test', 'test data'), 'queue'); }, ); self::assertSame( 'test data', $middleware->processPush( - $this->getMessage(), + new PushRequest($this->getMessage(), 'queue'), $this->createMock(PushHandlerInterface::class), - )->getPayload(), + )->getMessage()->getPayload(), ); } @@ -73,9 +74,9 @@ static function (): PushMiddlewareInterface { self::assertSame( 'New middleware test data', $middleware->processPush( - $this->getMessage(), + new PushRequest($this->getMessage(), 'queue'), $this->createMock(PushHandlerInterface::class), - )->getPayload(), + )->getMessage()->getPayload(), ); } @@ -87,9 +88,9 @@ public function testCreateWithUseParamsMiddleware(): void self::assertSame( 'New middleware test data', $middleware->processPush( - $this->getMessage(), + new PushRequest($this->getMessage(), 'queue'), $this->getRequestHandler(), - )->getPayload(), + )->getMessage()->getPayload(), ); } @@ -101,9 +102,9 @@ public function testCreateWithTestCallableMiddleware(): void self::assertSame( 'New test data', $middleware->processPush( - $this->getMessage(), + new PushRequest($this->getMessage(), 'queue'), $this->getRequestHandler(), - )->getPayload(), + )->getMessage()->getPayload(), ); } @@ -115,9 +116,9 @@ public function testCreateFromStringCallable(): void self::assertSame( 'String callable data', $middleware->processPush( - $this->getMessage(), + new PushRequest($this->getMessage(), 'queue'), $this->createMock(PushHandlerInterface::class), - )->getPayload(), + )->getMessage()->getPayload(), ); } @@ -129,9 +130,9 @@ public function testCreateFromCallableObject(): void self::assertSame( 'Callable object data', $middleware->processPush( - $this->getMessage(), + new PushRequest($this->getMessage(), 'queue'), $this->createMock(PushHandlerInterface::class), - )->getPayload(), + )->getMessage()->getPayload(), ); } @@ -165,7 +166,7 @@ public function testInvalidMiddlewareWithWrongController(): void $this->expectException(InvalidMiddlewareDefinitionException::class); $middleware->processPush( - $this->getMessage(), + new PushRequest($this->getMessage(), 'queue'), $this->createMock(PushHandlerInterface::class), ); } @@ -185,9 +186,9 @@ private function getContainer(array $instances = []): ContainerInterface private function getRequestHandler(): PushHandlerInterface { return new class implements PushHandlerInterface { - public function handlePush(MessageInterface $message): MessageInterface + public function handlePush(PushRequest $request): PushRequest { - return $message; + return $request; } }; } diff --git a/tests/Unit/Middleware/Push/Support/CallableObjectMiddleware.php b/tests/Unit/Middleware/Push/Support/CallableObjectMiddleware.php index 996665b6..02940bfb 100644 --- a/tests/Unit/Middleware/Push/Support/CallableObjectMiddleware.php +++ b/tests/Unit/Middleware/Push/Support/CallableObjectMiddleware.php @@ -5,12 +5,13 @@ namespace Yiisoft\Queue\Tests\Unit\Middleware\Push\Support; use Yiisoft\Queue\Message\GenericMessage; -use Yiisoft\Queue\Message\MessageInterface; +use Yiisoft\Queue\Middleware\Push\PushRequest; +use Yiisoft\Queue\Middleware\Push\PushHandlerInterface; final class CallableObjectMiddleware { - public function __invoke(MessageInterface $message): MessageInterface + public function __invoke(PushRequest $request, PushHandlerInterface $handler): PushRequest { - return new GenericMessage('test', 'Callable object data'); + return $request->withMessage(new GenericMessage('test', 'Callable object data')); } } diff --git a/tests/Unit/Middleware/Push/Support/StringCallableMiddleware.php b/tests/Unit/Middleware/Push/Support/StringCallableMiddleware.php index f0c11636..1f437131 100644 --- a/tests/Unit/Middleware/Push/Support/StringCallableMiddleware.php +++ b/tests/Unit/Middleware/Push/Support/StringCallableMiddleware.php @@ -5,12 +5,13 @@ namespace Yiisoft\Queue\Tests\Unit\Middleware\Push\Support; use Yiisoft\Queue\Message\GenericMessage; -use Yiisoft\Queue\Message\MessageInterface; +use Yiisoft\Queue\Middleware\Push\PushRequest; +use Yiisoft\Queue\Middleware\Push\PushHandlerInterface; final class StringCallableMiddleware { - public static function handle(MessageInterface $message): MessageInterface + public static function handle(PushRequest $request, PushHandlerInterface $handler): PushRequest { - return new GenericMessage('test', 'String callable data'); + return $request->withMessage(new GenericMessage('test', 'String callable data')); } } diff --git a/tests/Unit/Middleware/Push/Support/TestCallableMiddleware.php b/tests/Unit/Middleware/Push/Support/TestCallableMiddleware.php index a8219ac8..d1b152a5 100644 --- a/tests/Unit/Middleware/Push/Support/TestCallableMiddleware.php +++ b/tests/Unit/Middleware/Push/Support/TestCallableMiddleware.php @@ -5,12 +5,13 @@ namespace Yiisoft\Queue\Tests\Unit\Middleware\Push\Support; use Yiisoft\Queue\Message\GenericMessage; -use Yiisoft\Queue\Message\MessageInterface; +use Yiisoft\Queue\Middleware\Push\PushRequest; +use Yiisoft\Queue\Middleware\Push\PushHandlerInterface; final class TestCallableMiddleware { - public function index(MessageInterface $message): MessageInterface + public function index(PushRequest $request, PushHandlerInterface $handler): PushRequest { - return new GenericMessage('test', 'New test data'); + return $request->withMessage(new GenericMessage('test', 'New test data')); } } diff --git a/tests/Unit/Middleware/Push/Support/TestMiddleware.php b/tests/Unit/Middleware/Push/Support/TestMiddleware.php index ba942661..907bdb6f 100644 --- a/tests/Unit/Middleware/Push/Support/TestMiddleware.php +++ b/tests/Unit/Middleware/Push/Support/TestMiddleware.php @@ -5,16 +5,16 @@ namespace Yiisoft\Queue\Tests\Unit\Middleware\Push\Support; use Yiisoft\Queue\Message\GenericMessage; -use Yiisoft\Queue\Message\MessageInterface; use Yiisoft\Queue\Middleware\Push\PushHandlerInterface; use Yiisoft\Queue\Middleware\Push\PushMiddlewareInterface; +use Yiisoft\Queue\Middleware\Push\PushRequest; final class TestMiddleware implements PushMiddlewareInterface { public function __construct(private readonly string $message = 'New middleware test data') {} - public function processPush(MessageInterface $message, PushHandlerInterface $handler): MessageInterface + public function processPush(PushRequest $request, PushHandlerInterface $handler): PushRequest { - return new GenericMessage('test', $this->message); + return $request->withMessage(new GenericMessage('test', $this->message)); } } diff --git a/tests/Unit/Middleware/RequestImmutabilityTest.php b/tests/Unit/Middleware/RequestImmutabilityTest.php new file mode 100644 index 00000000..a111c950 --- /dev/null +++ b/tests/Unit/Middleware/RequestImmutabilityTest.php @@ -0,0 +1,49 @@ +withMessage(new GenericMessage('changed', 'new')); + $withQueue = $request->withQueueName('changed-queue'); + + self::assertNotSame($request, $withMessage); + self::assertNotSame($request, $withQueue); + self::assertSame($message, $request->getMessage()); + self::assertSame('original-queue', $request->getQueueName()); + self::assertSame('changed', $withMessage->getMessage()->getType()); + self::assertSame('original-queue', $withMessage->getQueueName()); + self::assertSame('changed-queue', $withQueue->getQueueName()); + self::assertSame($message, $withQueue->getMessage()); + } + + public function testWorkerRequestWithersDoNotMutateOriginal(): void + { + $message = new GenericMessage('original', 'payload'); + $request = new WorkerRequest($message, 'original-queue'); + + $withMessage = $request->withMessage(new GenericMessage('changed', 'new')); + $withQueue = $request->withQueueName('changed-queue'); + + self::assertNotSame($request, $withMessage); + self::assertNotSame($request, $withQueue); + self::assertSame($message, $request->getMessage()); + self::assertSame('original-queue', $request->getQueueName()); + self::assertSame('changed', $withMessage->getMessage()->getType()); + self::assertSame('original-queue', $withMessage->getQueueName()); + self::assertSame('changed-queue', $withQueue->getQueueName()); + self::assertSame($message, $withQueue->getMessage()); + } +} diff --git a/tests/Unit/Middleware/Worker/WorkerMiddlewareFactoryTest.php b/tests/Unit/Middleware/Worker/WorkerMiddlewareFactoryTest.php new file mode 100644 index 00000000..014d11ee --- /dev/null +++ b/tests/Unit/Middleware/Worker/WorkerMiddlewareFactoryTest.php @@ -0,0 +1,110 @@ +factory()->createWorkerMiddleware($middleware)); + } + + public function testCallableReturningRequestIsSupported(): void + { + $request = $this->request(); + $handler = $this->createMock(WorkerHandlerInterface::class); + $handler + ->expects(self::once()) + ->method('handleWorker') + ->with($request) + ->willReturn($request); + + $middleware = $this->factory()->createWorkerMiddleware( + static function (WorkerRequest $actual, WorkerHandlerInterface $actualHandler): WorkerRequest { + $actualHandler->handleWorker($actual); + + return $actual->withQueueName('changed'); + }, + ); + + $result = $middleware->processWorker($request, $handler); + + self::assertSame('changed', $result->getQueueName()); + } + + public function testCallableReturningMiddlewareIsSupported(): void + { + $middleware = $this->factory()->createWorkerMiddleware(static fn(): WorkerMiddlewareInterface => new class implements WorkerMiddlewareInterface { + public function processWorker(WorkerRequest $request, WorkerHandlerInterface $handler): WorkerRequest + { + return $request->withQueueName('changed'); + } + }); + + self::assertSame('changed', $middleware->processWorker($this->request(), $this->handler())->getQueueName()); + } + + public static function invalidDefinitions(): array + { + return [ + 'scalar' => [42], + 'unknown string' => ['unknown'], + 'invalid array' => [['not', 'a', 'definition']], + ]; + } + + #[DataProvider('invalidDefinitions')] + public function testInvalidDefinitionThrows(mixed $definition): void + { + $this->expectException(InvalidMiddlewareDefinitionException::class); + + $this->factory()->createWorkerMiddleware($definition); + } + + public function testCallableReturningInvalidResultThrows(): void + { + $middleware = $this->factory()->createWorkerMiddleware(static fn(): string => 'invalid'); + + $this->expectException(InvalidMiddlewareDefinitionException::class); + $middleware->processWorker($this->request(), $this->handler()); + } + + private function factory(): WorkerMiddlewareFactory + { + return new WorkerMiddlewareFactory(new SimpleContainer()); + } + + private function request(): WorkerRequest + { + return new WorkerRequest(new GenericMessage('test', null), 'queue'); + } + + private function handler(): WorkerHandlerInterface + { + return new class implements WorkerHandlerInterface { + public function handleWorker(WorkerRequest $request): WorkerRequest + { + return $request; + } + }; + } +} diff --git a/tests/Unit/Middleware/Worker/WorkerMiddlewareTest.php b/tests/Unit/Middleware/Worker/WorkerMiddlewareTest.php new file mode 100644 index 00000000..dd565a5e --- /dev/null +++ b/tests/Unit/Middleware/Worker/WorkerMiddlewareTest.php @@ -0,0 +1,81 @@ +createMock(WorkerMiddlewareFactoryInterface::class); + $factory->method('createWorkerMiddleware')->willReturnCallback( + static function (mixed $definition) use (&$calls): WorkerMiddlewareInterface { + return new class ($definition, $calls) implements WorkerMiddlewareInterface { + private array $calls; + + public function __construct(private string $name, array &$calls) + { + $this->calls = & $calls; + } + + public function processWorker(WorkerRequest $request, WorkerHandlerInterface $handler): WorkerRequest + { + $this->calls[] = $this->name . ':before'; + $result = $handler->handleWorker($request); + $this->calls[] = $this->name . ':after'; + + return $result; + } + }; + }, + ); + $final = new class ($calls) implements WorkerHandlerInterface { + private array $calls; + + public function __construct(array &$calls) + { + $this->calls = & $calls; + } + + public function handleWorker(WorkerRequest $request): WorkerRequest + { + $this->calls[] = 'final'; + + return $request; + } + }; + + (new WorkerMiddlewareDispatcher($factory, ['one', 'two'], $final)) + ->dispatch(new WorkerRequest(new GenericMessage('test', null), 'queue')); + + self::assertSame(['one:before', 'two:before', 'final', 'two:after', 'one:after'], $calls); + } + + public function testMiddlewareCanShortCircuit(): void + { + $factory = $this->createMock(WorkerMiddlewareFactoryInterface::class); + $factory->method('createWorkerMiddleware')->willReturn(new class implements WorkerMiddlewareInterface { + public function processWorker(WorkerRequest $request, WorkerHandlerInterface $handler): WorkerRequest + { + return $request->withQueueName('short-circuited'); + } + }); + $final = $this->createMock(WorkerHandlerInterface::class); + $final->expects(self::never())->method('handleWorker'); + + $result = (new WorkerMiddlewareDispatcher($factory, ['stop'], $final)) + ->dispatch(new WorkerRequest(new GenericMessage('test', null), 'queue')); + + self::assertSame('short-circuited', $result->getQueueName()); + } +} diff --git a/tests/Unit/SyncQueueProducerTest.php b/tests/Unit/SyncQueueProducerTest.php new file mode 100644 index 00000000..d587cc5d --- /dev/null +++ b/tests/Unit/SyncQueueProducerTest.php @@ -0,0 +1,48 @@ +createMock(LoggerInterface::class); + $logger->expects(self::once()) + ->method('debug') + ->with( + 'Preparing to push message with message type "{messageType}".', + ['messageType' => 'test'], + ); + $logger->expects(self::once()) + ->method('info') + ->with( + 'Processed message with message type "{messageType}" synchronously.', + ['messageType' => 'test'], + ); + + $message = new GenericMessage('test', []); + $worker = $this->createMock(WorkerInterface::class); + $worker->expects(self::once()) + ->method('process') + ->with($message, 'queue'); + + $producer = new SyncQueueProducer( + $logger, + new PushMiddlewareConfig($this->createMock(PushMiddlewareFactoryInterface::class), []), + $worker, + 'queue', + ); + + self::assertSame($message, $producer->push($message)); + } +} diff --git a/tests/Unit/WorkerFailurePathTest.php b/tests/Unit/WorkerFailurePathTest.php new file mode 100644 index 00000000..d6e4023a --- /dev/null +++ b/tests/Unit/WorkerFailurePathTest.php @@ -0,0 +1,73 @@ +createMock(WorkerMiddlewareFactoryInterface::class); + $workerFactory->method('createWorkerMiddleware')->willReturn( + new class ($failure) implements WorkerMiddlewareInterface { + public function __construct(private readonly RuntimeException $failure) {} + + public function processWorker(WorkerRequest $request, WorkerHandlerInterface $handler): WorkerRequest + { + throw $this->failure; + } + }, + ); + $workerDispatcher = new WorkerMiddlewareDispatcher( + $workerFactory, + ['worker-middleware'], + $this->createMock(WorkerHandlerInterface::class), + ); + + $failureFactory = $this->createMock(FailureMiddlewareFactoryInterface::class); + $failureFactory->expects(self::once()) + ->method('createFailureMiddleware') + ->willReturn(new class ($failure) implements FailureMiddlewareInterface { + public function __construct(private readonly RuntimeException $expected) {} + + public function processFailure(FailureHandlingRequest $request, FailureHandlerInterface $handler): FailureHandlingRequest + { + TestCase::assertSame($this->expected, $request->getException()); + + return $request; + } + }); + $failureDispatcher = new FailureMiddlewareDispatcher( + $failureFactory, + ['queue' => ['failure-middleware']], + ); + $logger = new SimpleLogger(); + + (new Worker($logger, $workerDispatcher, $failureDispatcher)) + ->process(new GenericMessage('test', null), 'queue'); + + $messages = $logger->getMessages(); + self::assertCount(2, $messages); + self::assertSame(LogLevel::INFO, $messages[1]['level']); + self::assertSame('worker middleware failed', $messages[1]['message']); + } +} diff --git a/tests/Unit/WorkerTest.php b/tests/Unit/WorkerTest.php index adab4e2f..ac53f66f 100644 --- a/tests/Unit/WorkerTest.php +++ b/tests/Unit/WorkerTest.php @@ -21,6 +21,9 @@ use Yiisoft\Queue\Middleware\FailureHandling\FailureMiddlewareDispatcher; use Yiisoft\Queue\Middleware\FailureHandling\FailureMiddlewareFactoryInterface; use Yiisoft\Queue\Middleware\FailureHandling\FailureMiddlewareInterface; +use Yiisoft\Queue\Middleware\Worker\WorkerFinalHandler; +use Yiisoft\Queue\Middleware\Worker\WorkerMiddlewareDispatcher; +use Yiisoft\Queue\Middleware\Worker\WorkerMiddlewareFactory; use Yiisoft\Queue\Tests\App\FakeHandler; use Yiisoft\Queue\Tests\TestCase; use Yiisoft\Queue\Worker\Worker; @@ -163,9 +166,15 @@ private function createWorkerByParams( return new Worker( $logger ?? new NullLogger(), - $consumeMiddlewareDispatcher ?? new ConsumeMiddlewareDispatcher($consumeMiddlewareFactory), + new WorkerMiddlewareDispatcher( + new WorkerMiddlewareFactory(new SimpleContainer()), + [], + new WorkerFinalHandler( + $handlerResolver, + $consumeMiddlewareDispatcher ?? new ConsumeMiddlewareDispatcher($consumeMiddlewareFactory), + ), + ), $failureMiddlewareDispatcher ?? new FailureMiddlewareDispatcher($failureMiddlewareFactory, []), - $handlerResolver, ); } }