diff --git a/config/di.php b/config/di.php index 86224d7c..1c8f3b75 100644 --- a/config/di.php +++ b/config/di.php @@ -21,16 +21,19 @@ use Yiisoft\Queue\Middleware\Push\PushMiddlewareConfig; use Yiisoft\Queue\Middleware\Push\PushMiddlewareFactory; use Yiisoft\Queue\Middleware\Push\PushMiddlewareFactoryInterface; +use Yiisoft\Queue\Message\Handler\Resolver\HandlerResolver; +use Yiisoft\Queue\Message\Handler\Resolver\HandlerResolverInterface; use Yiisoft\Queue\Worker\Worker as QueueWorker; use Yiisoft\Queue\Worker\WorkerInterface; /* @var array $params */ return [ - QueueWorker::class => [ - 'class' => QueueWorker::class, + HandlerResolver::class => [ + 'class' => HandlerResolver::class, '__construct()' => [$params['yiisoft/queue']['handlers']], ], + HandlerResolverInterface::class => HandlerResolver::class, WorkerInterface::class => QueueWorker::class, LoopInterface::class => static function (ContainerInterface $container): LoopInterface { return \extension_loaded('pcntl') diff --git a/config/params.php b/config/params.php index 03956a11..e502d23a 100644 --- a/config/params.php +++ b/config/params.php @@ -9,7 +9,7 @@ use Yiisoft\Queue\Debug\QueueConsumerProviderProxy; use Yiisoft\Queue\Debug\QueueProducerProviderProxy; use Yiisoft\Queue\Debug\QueueWorkerInterfaceProxy; -use Yiisoft\Queue\Message\MessageHandlerInterface; +use Yiisoft\Queue\Message\Handler\HandlerInterface; use Yiisoft\Queue\Message\Serializer\MessageSerializer; use Yiisoft\Queue\Provider\QueueConsumerProviderInterface; use Yiisoft\Queue\Provider\QueueProducerProviderInterface; @@ -35,7 +35,7 @@ 'messages' => [], /** * Map of message type to handler. The worker uses this to find the handler for a received message. - * A handler may be a class name implementing {@see MessageHandlerInterface}, a callable, or any definition + * A handler may be a class name implementing {@see HandlerInterface}, a callable, or any definition * supported by yiisoft/injector. Example: * [ * 'send-email' => SendEmailHandler::class, diff --git a/src/Message/Handler/CallableHandler.php b/src/Message/Handler/CallableHandler.php new file mode 100644 index 00000000..25a22e91 --- /dev/null +++ b/src/Message/Handler/CallableHandler.php @@ -0,0 +1,30 @@ +handler)($message); + } +} diff --git a/src/Message/Handler/HandlerInterface.php b/src/Message/Handler/HandlerInterface.php new file mode 100644 index 00000000..79278411 --- /dev/null +++ b/src/Message/Handler/HandlerInterface.php @@ -0,0 +1,18 @@ + Cache of resolved handlers */ + private array $cache = []; + + public function __construct( + /** @var array */ + private readonly array $handlers, + private readonly ContainerInterface $container, + private readonly CallableFactory $callableFactory, + ) {} + + public function resolve(string $messageType): HandlerInterface + { + if ($messageType === '') { + throw new HandlerNotFoundException($messageType); + } + + if (array_key_exists($messageType, $this->cache)) { + return $this->cache[$messageType]; + } + + $definition = $this->handlers[$messageType] ?? $messageType; + + if (is_string($definition) && $this->container->has($definition)) { + $resolved = $this->container->get($definition); + + if ($resolved instanceof HandlerInterface) { + return $this->cache[$messageType] = $resolved; + } + } + + try { + /** @psalm-var callable(MessageInterface): void $callable */ + $callable = $this->callableFactory->create($definition); + + return $this->cache[$messageType] = new CallableHandler($callable); + } catch (InvalidCallableConfigurationException $exception) { + throw new HandlerNotFoundException($messageType, 0, $exception); + } + } +} diff --git a/src/Message/Handler/Resolver/HandlerResolverInterface.php b/src/Message/Handler/Resolver/HandlerResolverInterface.php new file mode 100644 index 00000000..ed9bdf85 --- /dev/null +++ b/src/Message/Handler/Resolver/HandlerResolverInterface.php @@ -0,0 +1,22 @@ + Cache of resolved handlers */ - private array $handlersCached = []; - public function __construct( - /** @var array */ - private readonly array $handlers, private readonly LoggerInterface $logger, - private readonly Injector $injector, - private readonly ContainerInterface $container, private readonly ConsumeMiddlewareDispatcher $consumeMiddlewareDispatcher, private readonly FailureMiddlewareDispatcher $failureMiddlewareDispatcher, - private readonly CallableFactory $callableFactory, + private readonly HandlerResolverInterface $handlerResolver, ) {} /** @@ -61,26 +42,17 @@ public function process( $this->logger->info('Processing message #{message}.', ['message' => $messageId]); } - $messageType = $message->getType(); - try { - $handler = $this->getHandler($messageType); - } catch (InvalidCallableConfigurationException $exception) { - throw new RuntimeException(sprintf('Queue handler for message type "%s" does not exist.', $messageType), 0, $exception); - } - - if ($handler === null) { - throw new RuntimeException(sprintf('Queue handler for message type "%s" does not exist.', $messageType)); - } + $handler = $this->handlerResolver->resolve($message->getType()); $request = new ConsumeRequest($message, $queueName); - $closure = fn(MessageInterface $message): mixed => $this->injector->invoke($handler, [$message]); + $finishHandler = new ConsumeFinalHandler($handler->handle(...)); try { - return $this->consumeMiddlewareDispatcher->dispatch($request, $this->createConsumeHandler($closure))->getMessage(); + return $this->consumeMiddlewareDispatcher->dispatch($request, $finishHandler)->getMessage(); } catch (Throwable $exception) { $request = new FailureHandlingRequest($request->getMessage(), $exception, $request->getQueueName(), $retryProducer); try { - $result = $this->failureMiddlewareDispatcher->dispatch($request, $this->createFailureHandler()); + $result = $this->failureMiddlewareDispatcher->dispatch($request, new FailureFinalHandler()); $this->logger->info($exception->getMessage()); return $result->getMessage(); @@ -91,39 +63,4 @@ public function process( } } } - - private function getHandler(string $messageType): ?callable - { - if ($messageType === '') { - return null; - } - - if (!array_key_exists($messageType, $this->handlersCached)) { - $definition = $this->handlers[$messageType] ?? $messageType; - - if (is_string($definition) && $this->container->has($definition)) { - $resolved = $this->container->get($definition); - - if ($resolved instanceof MessageHandlerInterface) { - $this->handlersCached[$messageType] = $resolved->handle(...); - - return $this->handlersCached[$messageType]; - } - } - - $this->handlersCached[$messageType] = $this->callableFactory->create($definition); - } - - return $this->handlersCached[$messageType]; - } - - private function createConsumeHandler(Closure $handler): ConsumeHandlerInterface - { - return new ConsumeFinalHandler($handler); - } - - private function createFailureHandler(): FailureHandlerInterface - { - return new FailureFinalHandler(); - } } diff --git a/tests/Benchmark/QueueBench.php b/tests/Benchmark/QueueBench.php index 7f56b83b..51485bc5 100644 --- a/tests/Benchmark/QueueBench.php +++ b/tests/Benchmark/QueueBench.php @@ -7,7 +7,6 @@ use Generator; use PhpBench\Attributes\ParamProviders; use Psr\Log\NullLogger; -use Yiisoft\Injector\Injector; use Yiisoft\Queue\Cli\SimpleLoop; use Yiisoft\Queue\Message\IdEnvelope; use Yiisoft\Queue\Message\GenericMessage; @@ -25,6 +24,7 @@ use Yiisoft\Queue\QueueConsumer; use Yiisoft\Queue\QueueConsumerInterface; use Yiisoft\Queue\QueueProducerInterface; +use Yiisoft\Queue\Message\Handler\Resolver\HandlerResolver; use Yiisoft\Queue\Tests\Benchmark\Support\VoidAdapter; use Yiisoft\Queue\Worker\Worker; use Yiisoft\Test\Support\Container\SimpleContainer; @@ -43,18 +43,19 @@ public function __construct() $logger = new NullLogger(); $worker = new Worker( - [ - 'foo' => static function (): void {}, - ], $logger, - new Injector($container), - $container, new ConsumeMiddlewareDispatcher(new ConsumeMiddlewareFactory($container, $callableFactory)), new FailureMiddlewareDispatcher( new FailureMiddlewareFactory($container, $callableFactory), [], ), - $callableFactory, + new HandlerResolver( + [ + 'foo' => static function (): void {}, + ], + $container, + $callableFactory, + ), ); $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 f928f6f4..66146726 100644 --- a/tests/Integration/MessageConsumingTest.php +++ b/tests/Integration/MessageConsumingTest.php @@ -6,7 +6,6 @@ use Psr\Container\ContainerInterface; use Psr\Log\NullLogger; -use Yiisoft\Injector\Injector; use Yiisoft\Queue\Message\GenericMessage; use Yiisoft\Queue\Message\MessageInterface; use Yiisoft\Queue\Middleware\CallableFactory; @@ -14,6 +13,7 @@ use Yiisoft\Queue\Middleware\Consume\ConsumeMiddlewareFactoryInterface; use Yiisoft\Queue\Middleware\FailureHandling\FailureMiddlewareDispatcher; use Yiisoft\Queue\Middleware\FailureHandling\FailureMiddlewareFactoryInterface; +use Yiisoft\Queue\Message\Handler\Resolver\HandlerResolver; use Yiisoft\Queue\Tests\Integration\Support\TestHandler; use Yiisoft\Queue\Tests\TestCase; use Yiisoft\Queue\Worker\Worker; @@ -31,16 +31,17 @@ public function testMessagesConsumed(): void $container = $this->createMock(ContainerInterface::class); $callableFactory = new CallableFactory($container); $worker = new Worker( - [ - 'test' => fn(MessageInterface $message): mixed => $this->messagesProcessed[] = $message->getPayload(), - 'test2' => fn(MessageInterface $message): mixed => $this->messagesProcessedSecond[] = $message->getPayload(), - ], new NullLogger(), - new Injector($container), - $container, new ConsumeMiddlewareDispatcher($this->createMock(ConsumeMiddlewareFactoryInterface::class)), new FailureMiddlewareDispatcher($this->createMock(FailureMiddlewareFactoryInterface::class), []), - $callableFactory, + new HandlerResolver( + [ + 'test' => fn(MessageInterface $message): mixed => $this->messagesProcessed[] = $message->getPayload(), + 'test2' => fn(MessageInterface $message): mixed => $this->messagesProcessedSecond[] = $message->getPayload(), + ], + $container, + $callableFactory, + ), ); $messages = [1, 'foo', 'bar-baz']; @@ -61,13 +62,10 @@ public function testMessagesConsumedByHandlerClass(): void $container->method('has')->with(TestHandler::class)->willReturn(true); $callableFactory = new CallableFactory($container); $worker = new Worker( - [], new NullLogger(), - new Injector($container), - $container, new ConsumeMiddlewareDispatcher($this->createMock(ConsumeMiddlewareFactoryInterface::class)), new FailureMiddlewareDispatcher($this->createMock(FailureMiddlewareFactoryInterface::class), []), - $callableFactory, + new HandlerResolver([], $container, $callableFactory), ); $messages = [1, 'foo', 'bar-baz']; diff --git a/tests/Integration/MiddlewareTest.php b/tests/Integration/MiddlewareTest.php index 9118d0e5..bf1eadd6 100644 --- a/tests/Integration/MiddlewareTest.php +++ b/tests/Integration/MiddlewareTest.php @@ -8,7 +8,6 @@ use PHPUnit\Framework\TestCase; use Psr\Container\ContainerInterface; use Psr\Log\LoggerInterface; -use Yiisoft\Injector\Injector; use Yiisoft\Test\Support\Container\SimpleContainer; use Yiisoft\Test\Support\Log\SimpleLogger; use Yiisoft\Queue\Message\GenericMessage; @@ -26,6 +25,7 @@ use Yiisoft\Queue\Middleware\Push\PushMiddlewareFactory; use Yiisoft\Queue\SyncQueueProducer; use Yiisoft\Queue\QueueProducerInterface; +use Yiisoft\Queue\Message\Handler\Resolver\HandlerResolver; use Yiisoft\Queue\Tests\Integration\Support\TestMiddleware; use Yiisoft\Queue\Worker\Worker; use Yiisoft\Queue\Worker\WorkerInterface; @@ -104,13 +104,10 @@ public function testFullStackConsume(): void ); $worker = new Worker( - ['test' => static fn() => true], new SimpleLogger(), - new Injector($container), - $container, $consumeMiddlewareDispatcher, $failureMiddlewareDispatcher, - $callableFactory, + new HandlerResolver(['test' => static fn() => true], $container, $callableFactory), ); $message = new GenericMessage('test', ['initial']); diff --git a/tests/Integration/Support/TestHandler.php b/tests/Integration/Support/TestHandler.php index 34cd004d..422e3480 100644 --- a/tests/Integration/Support/TestHandler.php +++ b/tests/Integration/Support/TestHandler.php @@ -4,10 +4,10 @@ namespace Yiisoft\Queue\Tests\Integration\Support; -use Yiisoft\Queue\Message\MessageHandlerInterface; +use Yiisoft\Queue\Message\Handler\HandlerInterface; use Yiisoft\Queue\Message\MessageInterface; -final class TestHandler implements MessageHandlerInterface +final class TestHandler implements HandlerInterface { public function __construct(public array $messagesProcessed = []) {} diff --git a/tests/TestCase.php b/tests/TestCase.php index 5ed6f752..cb81a762 100644 --- a/tests/TestCase.php +++ b/tests/TestCase.php @@ -9,7 +9,6 @@ use Psr\Container\ContainerInterface; use Psr\Log\NullLogger; use RuntimeException; -use Yiisoft\Injector\Injector; use Yiisoft\Test\Support\Container\SimpleContainer; use Yiisoft\Queue\Adapter\AdapterInterface; use Yiisoft\Queue\Cli\LoopInterface; @@ -24,6 +23,7 @@ use Yiisoft\Queue\Middleware\Push\PushMiddlewareFactory; use Yiisoft\Queue\AsyncQueueProducer; use Yiisoft\Queue\QueueProducerInterface; +use Yiisoft\Queue\Message\Handler\Resolver\HandlerResolver; use Yiisoft\Queue\SyncQueueProducer; use Yiisoft\Queue\Worker\Worker; use Yiisoft\Queue\Worker\WorkerInterface; @@ -118,13 +118,14 @@ protected function createLoop(): LoopInterface protected function createWorker(): WorkerInterface { return new Worker( - $this->getMessageHandlers(), new NullLogger(), - new Injector($this->getContainer()), - $this->getContainer(), $this->getConsumeMiddlewareDispatcher(), $this->getFailureMiddlewareDispatcher(), - new CallableFactory($this->getContainer()), + new HandlerResolver( + $this->getMessageHandlers(), + $this->getContainer(), + new CallableFactory($this->getContainer()), + ), ); } diff --git a/tests/Unit/Message/Handler/Resolver/HandlerResolverTest.php b/tests/Unit/Message/Handler/Resolver/HandlerResolverTest.php new file mode 100644 index 00000000..94bd8911 --- /dev/null +++ b/tests/Unit/Message/Handler/Resolver/HandlerResolverTest.php @@ -0,0 +1,167 @@ + $handler], $container, new CallableFactory($container)); + + $resolvedHandler = $resolver->resolve($message->getType()); + $resolvedHandler->handle($message); + + $processedMessages = FakeHandler::$processedMessages; + FakeHandler::$processedMessages = []; + + $this->assertSame([$message], $processedMessages); + } + + public static function handlerDefinitionDataProvider(): iterable + { + yield 'definition' => [ + FakeHandler::class, + [FakeHandler::class => new FakeHandler()], + ]; + yield 'definition-object' => [ + [new FakeHandler(), 'handle'], + [], + ]; + yield 'definition-class' => [ + [FakeHandler::class, 'handle'], + [FakeHandler::class => new FakeHandler()], + ]; + yield 'definition-not-found-class-but-exist-in-container' => [ + ['not-found-class-name', 'handle'], + ['not-found-class-name' => new FakeHandler()], + ]; + yield 'callable' => [ + function (MessageInterface $message) { + FakeHandler::$processedMessages[] = $message; + }, + [], + ]; + } + + public function testResolveCachesResolvedHandler(): void + { + $container = new SimpleContainer([FakeHandler::class => new FakeHandler()]); + $resolver = new HandlerResolver(['simple' => FakeHandler::class], $container, new CallableFactory($container)); + + $this->assertSame($resolver->resolve('simple'), $resolver->resolve('simple')); + } + + public function testResolveStaticMethodHandler(): void + { + $container = new SimpleContainer(); + $resolver = new HandlerResolver( + ['static-handler' => StaticMessageHandler::handle(...)], + $container, + new CallableFactory($container), + ); + + StaticMessageHandler::$wasHandled = false; + $resolvedHandler = $resolver->resolve('static-handler'); + $resolvedHandler->handle(new GenericMessage('static-handler', null)); + + $this->assertTrue(StaticMessageHandler::$wasHandled); + } + + public function testResolveThrowsWhenDefinitionMethodUndefined(): void + { + $this->expectException(HandlerNotFoundException::class); + $this->expectExceptionMessage('Queue handler for message type "simple" does not exist'); + + $container = new SimpleContainer([FakeHandler::class => new FakeHandler()]); + $resolver = new HandlerResolver( + ['simple' => [FakeHandler::class, 'undefinedMethod']], + $container, + new CallableFactory($container), + ); + + $resolver->resolve('simple'); + } + + public function testResolveThrowsWhenDefinitionClassUndefined(): void + { + $this->expectException(HandlerNotFoundException::class); + $this->expectExceptionMessage('Queue handler for message type "simple" does not exist'); + + $container = new SimpleContainer([FakeHandler::class => new FakeHandler()]); + $resolver = new HandlerResolver( + ['simple' => ['UndefinedClass', 'handle']], + $container, + new CallableFactory($container), + ); + + $resolver->resolve('simple'); + } + + public function testResolveThrowsWhenDefinitionClassNotFoundInContainer(): void + { + $this->expectException(HandlerNotFoundException::class); + $this->expectExceptionMessage('Queue handler for message type "simple" does not exist'); + + $container = new SimpleContainer(); + $resolver = new HandlerResolver( + ['simple' => [FakeHandler::class, 'handle']], + $container, + new CallableFactory($container), + ); + + $resolver->resolve('simple'); + } + + public function testResolveThrowsWhenHandlerNotFoundInContainer(): void + { + $this->expectException(HandlerNotFoundException::class); + $this->expectExceptionMessage('Queue handler for message type "nonexistent" does not exist'); + + $container = new SimpleContainer(); + $resolver = new HandlerResolver([], $container, new CallableFactory($container)); + + $resolver->resolve('nonexistent'); + } + + public function testResolveThrowsWhenHandlerInContainerNotImplementingInterface(): void + { + $this->expectException(HandlerNotFoundException::class); + $this->expectExceptionMessage('Queue handler for message type "invalid" does not exist'); + + $container = new SimpleContainer([ + 'invalid' => new class { + public function handle(): void {} + }, + ]); + $resolver = new HandlerResolver([], $container, new CallableFactory($container)); + + $resolver->resolve('invalid'); + } + + public function testResolveThrowsWhenMessageTypeIsEmpty(): void + { + $this->expectException(HandlerNotFoundException::class); + $this->expectExceptionMessage('Queue handler for message type "" does not exist'); + + $container = new SimpleContainer(); + $resolver = new HandlerResolver([], $container, new CallableFactory($container)); + + $resolver->resolve(''); + } +} diff --git a/tests/Unit/WorkerTest.php b/tests/Unit/WorkerTest.php index 6840f828..1ef46e84 100644 --- a/tests/Unit/WorkerTest.php +++ b/tests/Unit/WorkerTest.php @@ -4,45 +4,39 @@ namespace Yiisoft\Queue\Tests\Unit; -use PHPUnit\Framework\Attributes\DataProvider; -use Psr\Container\ContainerInterface; use Psr\Log\LoggerInterface; use Psr\Log\NullLogger; use RuntimeException; -use Yiisoft\Injector\Injector; -use Yiisoft\Test\Support\Container\SimpleContainer; use Yiisoft\Test\Support\Log\SimpleLogger; use Yiisoft\Queue\Exception\MessageFailureException; +use Yiisoft\Queue\Message\Handler\CallableHandler; +use Yiisoft\Queue\Message\Handler\Resolver\HandlerResolverInterface; use Yiisoft\Queue\Message\GenericMessage; use Yiisoft\Queue\Message\MessageInterface; use Yiisoft\Queue\Middleware\Consume\ConsumeMiddlewareDispatcher; use Yiisoft\Queue\Middleware\Consume\ConsumeMiddlewareFactoryInterface; -use Yiisoft\Queue\Middleware\FailureHandling\FailureMiddlewareDispatcher; use Yiisoft\Queue\Middleware\Consume\ConsumeMiddlewareInterface; use Yiisoft\Queue\Middleware\FailureHandling\FailureHandlingRequest; -use Yiisoft\Queue\Middleware\FailureHandling\FailureMiddlewareInterface; +use Yiisoft\Queue\Middleware\FailureHandling\FailureMiddlewareDispatcher; use Yiisoft\Queue\Middleware\FailureHandling\FailureMiddlewareFactoryInterface; -use Yiisoft\Queue\Middleware\CallableFactory; +use Yiisoft\Queue\Middleware\FailureHandling\FailureMiddlewareInterface; use Yiisoft\Queue\Tests\App\FakeHandler; -use Yiisoft\Queue\Tests\App\StaticMessageHandler; use Yiisoft\Queue\Tests\TestCase; use Yiisoft\Queue\Worker\Worker; use PHPUnit\Framework\MockObject\MockObject; final class WorkerTest extends TestCase { - #[DataProvider('messageHandledDataProvider')] - public function testMessageHandled(mixed $handler, array $containerServices): void + public function testMessageHandled(): void { $message = new GenericMessage('simple', ['test-data']); $logger = new SimpleLogger(); - $container = new SimpleContainer($containerServices); - $handlers = ['simple' => $handler]; + $handlerResolver = $this->createHandlerResolver($message, static function (MessageInterface $message): void { + FakeHandler::$processedMessages[] = $message; + }); - $queueName = 'test-queue'; - $worker = $this->createWorkerByParams($handlers, $container, $logger); - - $worker->process($message, $queueName); + $worker = $this->createWorkerByParams($handlerResolver, $logger); + $worker->process($message, 'test-queue'); $processedMessages = FakeHandler::$processedMessages; FakeHandler::$processedMessages = []; @@ -54,93 +48,19 @@ public function testMessageHandled(mixed $handler, array $containerServices): vo $this->assertStringContainsString('Processing message without ID.', $messages[0]['message']); } - public static function messageHandledDataProvider(): iterable - { - yield 'definition' => [ - FakeHandler::class, - [FakeHandler::class => new FakeHandler()], - ]; - yield 'definition-object' => [ - [new FakeHandler(), 'handle'], - [], - ]; - yield 'definition-class' => [ - [FakeHandler::class, 'handle'], - [FakeHandler::class => new FakeHandler()], - ]; - yield 'definition-not-found-class-but-exist-in-container' => [ - ['not-found-class-name', 'handle'], - ['not-found-class-name' => new FakeHandler()], - ]; - yield 'static-definition' => [ - FakeHandler::staticHandle(...), - [FakeHandler::class => new FakeHandler()], - ]; - yield 'callable' => [ - function (MessageInterface $message) { - FakeHandler::$processedMessages[] = $message; - }, - [], - ]; - } - - public function testMessageFailWithDefinitionUndefinedMethodHandler(): void - { - $this->expectExceptionMessage('Queue handler for message type "simple" does not exist'); - - $message = new GenericMessage('simple', ['test-data']); - $handler = new FakeHandler(); - $container = new SimpleContainer([FakeHandler::class => $handler]); - $handlers = ['simple' => [FakeHandler::class, 'undefinedMethod']]; - - $queueName = 'test-queue'; - $worker = $this->createWorkerByParams($handlers, $container); - - $worker->process($message, $queueName); - } - - public function testMessageFailWithDefinitionUndefinedClassHandler(): void - { - $this->expectExceptionMessage('Queue handler for message type "simple" does not exist'); - - $message = new GenericMessage('simple', ['test-data']); - $logger = new SimpleLogger(); - $handler = new FakeHandler(); - $container = new SimpleContainer([FakeHandler::class => $handler]); - $handlers = ['simple' => ['UndefinedClass', 'handle']]; - - $queueName = 'test-queue'; - $worker = $this->createWorkerByParams($handlers, $container, $logger); - - $worker->process($message, $queueName); - } - - public function testMessageFailWithDefinitionClassNotFoundInContainerHandler(): void - { - $this->expectExceptionMessage('Queue handler for message type "simple" does not exist'); - $message = new GenericMessage('simple', ['test-data']); - $container = new SimpleContainer(); - $handlers = ['simple' => [FakeHandler::class, 'handle']]; - - $queueName = 'test-queue'; - $worker = $this->createWorkerByParams($handlers, $container); - - $worker->process($message, $queueName); - } - public function testMessageFailWithDefinitionHandlerException(): void { $message = new GenericMessage('simple', ['test-data']); $logger = new SimpleLogger(); - $handler = new FakeHandler(); - $container = new SimpleContainer([FakeHandler::class => $handler]); - $handlers = ['simple' => [FakeHandler::class, 'handleWithException']]; + $handlerResolver = $this->createHandlerResolver($message, static function (): never { + throw new RuntimeException('Test exception'); + }); - $queueName = 'test-queue'; - $worker = $this->createWorkerByParams($handlers, $container, $logger); + $worker = $this->createWorkerByParams($handlerResolver, $logger); try { - $worker->process($message, $queueName); + $worker->process($message, 'test-queue'); + self::fail('Exception was not thrown.'); } catch (MessageFailureException $exception) { self::assertSame($exception::class, MessageFailureException::class); self::assertSame($exception->getMessage(), "Processing of message without ID is stopped because of an exception:\nTest exception."); @@ -155,38 +75,6 @@ public function testMessageFailWithDefinitionHandlerException(): void } } - public function testHandlerNotFoundInContainer(): void - { - $message = new GenericMessage('nonexistent', ['test-data']); - $container = new SimpleContainer(); - $handlers = []; - - $queueName = 'test-queue'; - $worker = $this->createWorkerByParams($handlers, $container); - - $this->expectException(RuntimeException::class); - $this->expectExceptionMessage('Queue handler for message type "nonexistent" does not exist'); - $worker->process($message, $queueName); - } - - public function testHandlerInContainerNotImplementingInterface(): void - { - $message = new GenericMessage('invalid', ['test-data']); - $container = new SimpleContainer([ - 'invalid' => new class { - public function handle(): void {} - }, - ]); - $handlers = []; - - $queueName = 'test-queue'; - $worker = $this->createWorkerByParams($handlers, $container); - - $this->expectException(RuntimeException::class); - $this->expectExceptionMessage('Queue handler for message type "invalid" does not exist'); - $worker->process($message, $queueName); - } - public function testMessageFailureIsHandledSuccessfully(): void { $message = new GenericMessage('simple', null); @@ -212,42 +100,28 @@ public function testMessageFailureIsHandledSuccessfully(): void $failureMiddlewareFactory->method('createFailureMiddleware')->willReturn($failureMiddleware); $failureDispatcher = new FailureMiddlewareDispatcher($failureMiddlewareFactory, ['test-queue' => ['simple']]); - $container = new SimpleContainer(); - $worker = new Worker( - ['simple' => fn() => null], - new NullLogger(), - new Injector($container), - $container, - $consumeDispatcher, - $failureDispatcher, - new CallableFactory($container), - ); + $handlerResolver = $this->createHandlerResolver($message, static fn() => null); + $worker = $this->createWorkerByParams($handlerResolver, new NullLogger(), $consumeDispatcher, $failureDispatcher); $result = $worker->process($message, $queueName); self::assertSame($finalMessage, $result); } - public function testStaticMethodHandler(): void + private function createHandlerResolver(MessageInterface $message, callable $handler): HandlerResolverInterface { - $message = new GenericMessage('static-handler', ['test-data']); - $container = new SimpleContainer(); - $handlers = [ - 'static-handler' => StaticMessageHandler::handle(...), - ]; - - $queueName = 'test-queue'; - $worker = $this->createWorkerByParams($handlers, $container); + /** @var HandlerResolverInterface&MockObject $handlerResolver */ + $handlerResolver = $this->createMock(HandlerResolverInterface::class); + $handlerResolver->method('resolve')->with($message->getType())->willReturn(new CallableHandler($handler)); - StaticMessageHandler::$wasHandled = false; - $worker->process($message, $queueName); - $this->assertTrue(StaticMessageHandler::$wasHandled); + return $handlerResolver; } private function createWorkerByParams( - array $handlers, - ContainerInterface $container, + HandlerResolverInterface $handlerResolver, ?LoggerInterface $logger = null, + ?ConsumeMiddlewareDispatcher $consumeMiddlewareDispatcher = null, + ?FailureMiddlewareDispatcher $failureMiddlewareDispatcher = null, ): Worker { /** @var ConsumeMiddlewareFactoryInterface&MockObject $consumeMiddlewareFactory */ $consumeMiddlewareFactory = $this->createMock(ConsumeMiddlewareFactoryInterface::class); @@ -255,13 +129,10 @@ private function createWorkerByParams( $failureMiddlewareFactory = $this->createMock(FailureMiddlewareFactoryInterface::class); return new Worker( - $handlers, $logger ?? new NullLogger(), - new Injector($container), - $container, - new ConsumeMiddlewareDispatcher($consumeMiddlewareFactory), - new FailureMiddlewareDispatcher($failureMiddlewareFactory, []), - new CallableFactory($container), + $consumeMiddlewareDispatcher ?? new ConsumeMiddlewareDispatcher($consumeMiddlewareFactory), + $failureMiddlewareDispatcher ?? new FailureMiddlewareDispatcher($failureMiddlewareFactory, []), + $handlerResolver, ); } }