diff --git a/README.md b/README.md index 836bc494..8f39156e 100644 --- a/README.md +++ b/README.md @@ -87,10 +87,10 @@ final class DownloadFileMessage extends Message Then create a handler that processes it: ```php +use Yiisoft\Queue\Message\Handler\HandlerInterface; use Yiisoft\Queue\Message\MessageInterface; -use Yiisoft\Queue\Message\MessageHandlerInterface; -final readonly class RemoteFileHandler implements MessageHandlerInterface +final readonly class RemoteFileHandler implements HandlerInterface { public function __construct( private FileDownloader $downloader, diff --git a/config/di.php b/config/di.php index 86224d7c..e8ae8707 100644 --- a/config/di.php +++ b/config/di.php @@ -21,14 +21,14 @@ use Yiisoft\Queue\Middleware\Push\PushMiddlewareConfig; use Yiisoft\Queue\Middleware\Push\PushMiddlewareFactory; use Yiisoft\Queue\Middleware\Push\PushMiddlewareFactoryInterface; +use Yiisoft\Queue\Message\Handler\HandlerResolver; use Yiisoft\Queue\Worker\Worker as QueueWorker; use Yiisoft\Queue\Worker\WorkerInterface; /* @var array $params */ return [ - QueueWorker::class => [ - 'class' => QueueWorker::class, + HandlerResolver::class => [ '__construct()' => [$params['yiisoft/queue']['handlers']], ], WorkerInterface::class => QueueWorker::class, 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/docs/guide/en/best-practices.md b/docs/guide/en/best-practices.md index 76d17124..d3cfaafe 100644 --- a/docs/guide/en/best-practices.md +++ b/docs/guide/en/best-practices.md @@ -9,7 +9,7 @@ This guide covers recommended practices for building reliable and maintainable q #### Bad ```php -final class ProcessPaymentHandler implements MessageHandlerInterface +final class ProcessPaymentHandler implements HandlerInterface { public function handle(MessageInterface $message): void { @@ -24,7 +24,7 @@ final class ProcessPaymentHandler implements MessageHandlerInterface #### Good ```php -final class ProcessPaymentHandler implements MessageHandlerInterface +final class ProcessPaymentHandler implements HandlerInterface { public function handle(MessageInterface $message): void { @@ -58,7 +58,7 @@ Avoid storing per-message state in handler properties. The container may return #### Bad ```php -final class ProcessPaymentHandler implements MessageHandlerInterface +final class ProcessPaymentHandler implements HandlerInterface { private array $processedIds = []; @@ -80,7 +80,7 @@ final class ProcessPaymentHandler implements MessageHandlerInterface #### Good ```php -final class ProcessPaymentHandler implements MessageHandlerInterface +final class ProcessPaymentHandler implements HandlerInterface { public function handle(MessageInterface $message): void { @@ -273,7 +273,7 @@ See [Message handler](message-handler.md) for details. ```php // Metrics collection in every handler -final class EmailHandler implements MessageHandlerInterface +final class EmailHandler implements HandlerInterface { public function handle(MessageInterface $message): void { diff --git a/docs/guide/en/configuration-manual.md b/docs/guide/en/configuration-manual.md index 5ad61cc7..63207220 100644 --- a/docs/guide/en/configuration-manual.md +++ b/docs/guide/en/configuration-manual.md @@ -15,8 +15,8 @@ To use the queue, you need to create instances of the following classes: ```php use Psr\Container\ContainerInterface; use Psr\Log\NullLogger; -use Yiisoft\Injector\Injector; use Yiisoft\Queue\Cli\SimpleLoop; +use Yiisoft\Queue\Message\Handler\HandlerResolver; use Yiisoft\Queue\Middleware\CallableFactory; use Yiisoft\Queue\Middleware\Consume\ConsumeMiddlewareDispatcher; use Yiisoft\Queue\Middleware\Consume\ConsumeMiddlewareFactory; @@ -57,13 +57,10 @@ $pushMiddlewareConfig = new PushMiddlewareConfig( // Create worker $worker = new Worker( - $handlers, $logger, - new Injector($container), - $container, $consumeMiddlewareDispatcher, $failureMiddlewareDispatcher, - $callableFactory, + new HandlerResolver($handlers, $container), ); // Create loop (SignalLoop requires ext-pcntl; SimpleLoop works without it) diff --git a/docs/guide/en/configuration-with-config.md b/docs/guide/en/configuration-with-config.md index 603dcbfb..7adbbbd3 100644 --- a/docs/guide/en/configuration-with-config.md +++ b/docs/guide/en/configuration-with-config.md @@ -7,7 +7,7 @@ If you are using [yiisoft/config](https://github.com/yiisoft/config) (i.e. insta In [yiisoft/app](https://github.com/yiisoft/app) / [yiisoft/app-api](https://github.com/yiisoft/app-api) templates you typically add or adjust configuration in `config/params.php`. If your project structure differs, put configuration into any params config file that is loaded by [yiisoft/config](https://github.com/yiisoft/config). -When your message type equals the FQCN of a handler class that implements `Yiisoft\Queue\Message\MessageHandlerInterface`, nothing else has to be configured: the DI container resolves the class automatically. See [Message handler](message-handler.md) for details and trade-offs. +When your message type equals the FQCN of a handler class that implements `Yiisoft\Queue\Message\Handler\HandlerInterface`, nothing else has to be configured: the DI container resolves the class automatically. See [Message handler](message-handler.md) for details and trade-offs. Advanced applications eventually need the following tweaks: diff --git a/docs/guide/en/message-handler-advanced.md b/docs/guide/en/message-handler-advanced.md index f6015b61..706b3345 100644 --- a/docs/guide/en/message-handler-advanced.md +++ b/docs/guide/en/message-handler-advanced.md @@ -8,7 +8,7 @@ For a conceptual overview of what messages and handlers are, see [Messages and h Handler definitions are configured in: - `$params['yiisoft/queue']['handlers']` when using [yiisoft/config](https://github.com/yiisoft/config), or -- the `$handlers` argument of `Yiisoft\Queue\Worker\Worker` when creating it manually. +- the `$handlers` argument of `Yiisoft\Queue\Message\Handler\HandlerResolver` when creating it manually. ## Supported handler definition formats @@ -76,7 +76,7 @@ return [ ]; ``` -Handler definition should be either an [extended callable definition](./callable-definitions-extended.md) or a container identifier that resolves to a `MessageHandlerInterface` instance. +Handler definition should be either an [extended callable definition](./callable-definitions-extended.md) or a container identifier that resolves to a `HandlerInterface` instance. ## When mapping by short names is a better idea @@ -111,7 +111,7 @@ This way external producers never need to know your internal PHP class names. The worker recognises three callable signatures: -- `MessageHandlerInterface` — implement the interface; the worker calls `handle(MessageInterface $message): void` directly (covered in [Message handler](message-handler.md)). +- `HandlerInterface` — implement the interface; the worker calls `handle(MessageInterface $message): void` directly (covered in [Message handler](message-handler.md)). - Invokable class — add `__invoke(MessageInterface $message): void`. - Explicit method — reference as `[HandlerClass::class, 'handle']` with `handle(MessageInterface $message): void` as the entry point. @@ -129,4 +129,4 @@ return [ ]; ``` -This config is consumed by the DI definitions from [`config/di.php`](../../../config/di.php) where the `Worker` is constructed with `$params['yiisoft/queue']['handlers']`. +This config is consumed by the DI definitions from [`config/di.php`](../../../config/di.php) where the `HandlerResolver` is constructed with `$params['yiisoft/queue']['handlers']`. diff --git a/docs/guide/en/message-handler.md b/docs/guide/en/message-handler.md index ec22d344..ad7de070 100644 --- a/docs/guide/en/message-handler.md +++ b/docs/guide/en/message-handler.md @@ -2,11 +2,11 @@ > If you are new to the concept of messages and handlers, read [Messages and handlers: concepts](messages-and-handlers.md) first. -The simplest setup requires no configuration at all: create a dedicated class implementing `Yiisoft\Queue\Message\MessageHandlerInterface` and use its FQCN as the message type when pushing a message. +The simplest setup requires no configuration at all: create a dedicated class implementing `Yiisoft\Queue\Message\Handler\HandlerInterface` and use its FQCN as the message type when pushing a message. ## HandlerInterface implementation (without type mapping) -If your handler implements `Yiisoft\Queue\Message\MessageHandlerInterface`, you can use the class FQCN as the message type. The DI container resolves the handler automatically. +If your handler implements `Yiisoft\Queue\Message\Handler\HandlerInterface`, you can use the class FQCN as the message type. The DI container resolves the handler automatically. > By default the [yiisoft/di](https://github.com/yiisoft/di) container resolves all FQCNs into corresponding class objects. @@ -46,7 +46,7 @@ new RemoteFileMessage('https://...'); **Handler**: ```php -final class RemoteFileHandler implements \Yiisoft\Queue\Message\MessageHandlerInterface +final class RemoteFileHandler implements \Yiisoft\Queue\Message\Handler\HandlerInterface { public function handle(\Yiisoft\Queue\Message\MessageInterface $message): void { diff --git a/docs/guide/en/messages-and-handlers.md b/docs/guide/en/messages-and-handlers.md index 60473d83..7c08654b 100644 --- a/docs/guide/en/messages-and-handlers.md +++ b/docs/guide/en/messages-and-handlers.md @@ -99,7 +99,7 @@ The message has no business logic, no dependencies. It is a value object — a t The handler receives the message and acts on it: ```php -final class SendEmailHandler implements \Yiisoft\Queue\Message\MessageHandlerInterface +final class SendEmailHandler implements \Yiisoft\Queue\Message\Handler\HandlerInterface { public function __construct(private Mailer $mailer) {} diff --git a/docs/guide/en/performance-tuning.md b/docs/guide/en/performance-tuning.md index 63602e9a..b983939b 100644 --- a/docs/guide/en/performance-tuning.md +++ b/docs/guide/en/performance-tuning.md @@ -93,7 +93,7 @@ public function handle(MessageInterface $message): void ```php // Bad - accumulates in memory -class Handler implements MessageHandlerInterface +class Handler implements HandlerInterface { private static array $cache = []; @@ -105,7 +105,7 @@ class Handler implements MessageHandlerInterface } // Good - use external cache -class Handler implements MessageHandlerInterface +class Handler implements HandlerInterface { public function __construct(private CacheInterface $cache) {} @@ -328,7 +328,7 @@ public function handle(MessageInterface $message): void If your message handler only reads data, use read replicas: ```php -final class GenerateReportHandler implements MessageHandlerInterface +final class GenerateReportHandler implements HandlerInterface { public function __construct( private ConnectionInterface $readDb, // Read replica diff --git a/src/Message/Handler/CallableHandler.php b/src/Message/Handler/CallableHandler.php new file mode 100644 index 00000000..35cab228 --- /dev/null +++ b/src/Message/Handler/CallableHandler.php @@ -0,0 +1,29 @@ +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 @@ + + */ + private array $cache = []; + + private readonly Injector $injector; + private readonly CallableFactory $callableFactory; + + /** + * @param (array|callable|HandlerInterface|string)[] $handlers Handler definitions indexed by message type. + * @param ContainerInterface $container Container used to resolve handlers. + * @param ContainerInterface|null $callableDependencyContainer Container used to resolve callable handler + * dependencies. If not set, the main container is used. + * + * @psalm-param array $handlers + */ + public function __construct( + private readonly array $handlers, + private readonly ContainerInterface $container, + ?ContainerInterface $callableDependencyContainer = null, + ) { + $this->injector = new Injector($callableDependencyContainer ?? $this->container); + $this->callableFactory = new CallableFactory($this->container); + } + + /** + * Get a handler for the given message type. + * + * @param string $messageType Message type. + * + * @throws HandlerNotFoundException If no handler exists for the message type. + * @throws InvalidHandlerConfigurationException If the handler definition is configured incorrectly. + * @throws ContainerExceptionInterface Error while retrieving the entry from container. + */ + public function resolve(string $messageType): HandlerInterface + { + if ($messageType === '') { + throw new LogicException('Message type cannot be empty.'); + } + + if (array_key_exists($messageType, $this->cache)) { + return $this->cache[$messageType]; + } + + $this->cache[$messageType] = $this->internalResolve($messageType); + + return $this->cache[$messageType]; + } + + /** + * @throws HandlerNotFoundException + * @throws InvalidHandlerConfigurationException + * @throws ContainerExceptionInterface + */ + private function internalResolve(string $messageType): HandlerInterface + { + $definition = $this->handlers[$messageType] ?? $messageType; + + if ($definition instanceof HandlerInterface) { + return $definition; + } + + if (is_callable($definition)) { + return $this->createCallableHandler($messageType, $definition); + } + + if (is_string($definition)) { + return $this->getHandlerFromContainer($messageType, $definition); + } + + return $this->createCallableHandler($messageType, $definition); + } + + /** + * @throws HandlerNotFoundException + * @throws InvalidHandlerConfigurationException + * @throws ContainerExceptionInterface + */ + private function getHandlerFromContainer(string $messageType, string $id): HandlerInterface + { + if (!$this->container->has($id)) { + throw new HandlerNotFoundException($messageType); + } + + $handler = $this->container->get($id); + + if ($handler instanceof HandlerInterface) { + return $handler; + } + + if (is_callable($handler)) { + return $this->createCallableHandler($messageType, $handler); + } + + throw new InvalidHandlerConfigurationException( + $messageType, + sprintf( + 'Resolved from container handler should be an instance of "%s" or callable, got "%s".', + HandlerInterface::class, + get_debug_type($handler), + ), + ); + } + + /** + * @throws InvalidHandlerConfigurationException + * @throws ContainerExceptionInterface + */ + private function createCallableHandler(string $messageType, mixed $definition): CallableHandler + { + try { + $callable = $this->callableFactory->create($definition); + } catch (InvalidCallableConfigurationException $exception) { + throw new InvalidHandlerConfigurationException($messageType, $exception->getMessage(), $exception); + } + + $callable = function (MessageInterface $message) use ($callable): void { + $this->injector->invoke($callable, [$message]); + }; + + return new CallableHandler($callable); + } +} diff --git a/src/Message/Handler/InvalidHandlerConfigurationException.php b/src/Message/Handler/InvalidHandlerConfigurationException.php new file mode 100644 index 00000000..dd14bb2c --- /dev/null +++ b/src/Message/Handler/InvalidHandlerConfigurationException.php @@ -0,0 +1,27 @@ +middleware === null) { - $this->middleware = ($this->middlewareFactory)(); - } + $this->middleware ??= ($this->middlewareFactory)(); return $this->middleware->processConsume($request, $this->handler); } diff --git a/src/Middleware/FailureHandling/FailureMiddlewareDispatcher.php b/src/Middleware/FailureHandling/FailureMiddlewareDispatcher.php index 3647e6ad..ebc1a6a5 100644 --- a/src/Middleware/FailureHandling/FailureMiddlewareDispatcher.php +++ b/src/Middleware/FailureHandling/FailureMiddlewareDispatcher.php @@ -43,9 +43,7 @@ public function dispatch( } $definitions = array_reverse($this->middlewareDefinitions[$queueName]); - if (!isset($this->stack[$queueName])) { - $this->stack[$queueName] = new FailureMiddlewareStack($this->buildMiddlewares(...$definitions), $finishHandler); - } + $this->stack[$queueName] ??= new FailureMiddlewareStack($this->buildMiddlewares(...$definitions), $finishHandler); return $this->stack[$queueName]->handleFailure($request); } @@ -83,9 +81,7 @@ public function withMiddlewares(array $middlewareDefinitions): self private function init(): void { - if (!isset($this->middlewareDefinitions[self::DEFAULT_PIPELINE])) { - $this->middlewareDefinitions[self::DEFAULT_PIPELINE] = []; - } + $this->middlewareDefinitions[self::DEFAULT_PIPELINE] ??= []; } /** diff --git a/src/Middleware/FailureHandling/FailureMiddlewareStack.php b/src/Middleware/FailureHandling/FailureMiddlewareStack.php index a9d7b21d..6040ba7d 100644 --- a/src/Middleware/FailureHandling/FailureMiddlewareStack.php +++ b/src/Middleware/FailureHandling/FailureMiddlewareStack.php @@ -64,9 +64,7 @@ public function __construct( public function handleFailure(FailureHandlingRequest $request): FailureHandlingRequest { - if ($this->middleware === null) { - $this->middleware = ($this->middlewareFactory)(); - } + $this->middleware ??= ($this->middlewareFactory)(); return $this->middleware->processFailure($request, $this->handler); } diff --git a/src/Middleware/Push/PushMiddlewareDispatcher.php b/src/Middleware/Push/PushMiddlewareDispatcher.php index 4347d8a5..54e3abb5 100644 --- a/src/Middleware/Push/PushMiddlewareDispatcher.php +++ b/src/Middleware/Push/PushMiddlewareDispatcher.php @@ -37,9 +37,7 @@ public function __construct( */ public function dispatch(MessageInterface $message): MessageInterface { - if ($this->stack === null) { - $this->stack = new PushMiddlewareStack($this->buildMiddlewares(), $this->finishHandler); - } + $this->stack ??= new PushMiddlewareStack($this->buildMiddlewares(), $this->finishHandler); return $this->stack->handlePush($message); } diff --git a/src/Middleware/Push/PushMiddlewareFactory.php b/src/Middleware/Push/PushMiddlewareFactory.php index 4326e749..967b274c 100644 --- a/src/Middleware/Push/PushMiddlewareFactory.php +++ b/src/Middleware/Push/PushMiddlewareFactory.php @@ -27,7 +27,7 @@ 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, MessageHandlerPushInterface $handler): + * - A callable with `function(MessageInterface $message, PushHandlerInterface $handler): * MessageInterface` signature. * - A controller handler action in format `[TestController::class, 'index']`. `TestController` instance will * be created and `index()` method will be executed. diff --git a/src/Middleware/Push/PushMiddlewareStack.php b/src/Middleware/Push/PushMiddlewareStack.php index de249b03..4ea47696 100644 --- a/src/Middleware/Push/PushMiddlewareStack.php +++ b/src/Middleware/Push/PushMiddlewareStack.php @@ -68,9 +68,7 @@ public function __construct( public function handlePush(MessageInterface $message): MessageInterface { - if ($this->middleware === null) { - $this->middleware = ($this->middlewareFactory)(); - } + $this->middleware ??= ($this->middlewareFactory)(); return $this->middleware->processPush($message, $this->handler); } diff --git a/src/Worker/Worker.php b/src/Worker/Worker.php index 65f7c8b9..8460a3ee 100644 --- a/src/Worker/Worker.php +++ b/src/Worker/Worker.php @@ -4,46 +4,27 @@ namespace Yiisoft\Queue\Worker; -use Closure; -use Psr\Container\ContainerInterface; use Psr\Log\LoggerInterface; -use RuntimeException; use Throwable; -use Yiisoft\Injector\Injector; use Yiisoft\Queue\Exception\MessageFailureException; +use Yiisoft\Queue\Message\Handler\HandlerResolver; use Yiisoft\Queue\Message\MessageInterface; -use Yiisoft\Queue\Message\MessageHandlerInterface; -use Yiisoft\Queue\Middleware\CallableFactory; -use Yiisoft\Queue\Middleware\InvalidCallableConfigurationException; use Yiisoft\Queue\Middleware\Consume\ConsumeFinalHandler; use Yiisoft\Queue\Middleware\Consume\ConsumeMiddlewareDispatcher; use Yiisoft\Queue\Middleware\Consume\ConsumeRequest; -use Yiisoft\Queue\Middleware\Consume\ConsumeHandlerInterface; use Yiisoft\Queue\Middleware\FailureHandling\FailureFinalHandler; use Yiisoft\Queue\Middleware\FailureHandling\FailureHandlingRequest; use Yiisoft\Queue\Middleware\FailureHandling\FailureMiddlewareDispatcher; -use Yiisoft\Queue\Middleware\FailureHandling\FailureHandlerInterface; use Yiisoft\Queue\QueueProducerInterface; use Yiisoft\Queue\Message\IdEnvelope; -use function array_key_exists; -use function is_string; -use function sprintf; - final class Worker implements WorkerInterface { - /** @var array 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 HandlerResolver $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..d3a75ca4 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\HandlerResolver; use Yiisoft\Queue\Tests\Benchmark\Support\VoidAdapter; use Yiisoft\Queue\Worker\Worker; use Yiisoft\Test\Support\Container\SimpleContainer; @@ -43,18 +43,18 @@ 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, + ), ); $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..e978602f 100644 --- a/tests/Integration/MessageConsumingTest.php +++ b/tests/Integration/MessageConsumingTest.php @@ -6,14 +6,13 @@ 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; use Yiisoft\Queue\Middleware\Consume\ConsumeMiddlewareDispatcher; use Yiisoft\Queue\Middleware\Consume\ConsumeMiddlewareFactoryInterface; use Yiisoft\Queue\Middleware\FailureHandling\FailureMiddlewareDispatcher; use Yiisoft\Queue\Middleware\FailureHandling\FailureMiddlewareFactoryInterface; +use Yiisoft\Queue\Message\Handler\HandlerResolver; use Yiisoft\Queue\Tests\Integration\Support\TestHandler; use Yiisoft\Queue\Tests\TestCase; use Yiisoft\Queue\Worker\Worker; @@ -29,18 +28,17 @@ public function testMessagesConsumed(): void $this->messagesProcessedSecond = []; $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, + ), ); $messages = [1, 'foo', 'bar-baz']; @@ -59,15 +57,11 @@ public function testMessagesConsumedByHandlerClass(): void $container = $this->createMock(ContainerInterface::class); $container->method('get')->with(TestHandler::class)->willReturn($handler); $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), ); $messages = [1, 'foo', 'bar-baz']; diff --git a/tests/Integration/MiddlewareTest.php b/tests/Integration/MiddlewareTest.php index 9118d0e5..50d4e7bb 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\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), ); $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..eef5914d 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\HandlerResolver; use Yiisoft\Queue\SyncQueueProducer; use Yiisoft\Queue\Worker\Worker; use Yiisoft\Queue\Worker\WorkerInterface; @@ -57,36 +57,28 @@ protected function setUp(): void */ protected function getQueue(): QueueProducerInterface { - if ($this->queue === null) { - $this->queue = $this->createQueue(); - } + $this->queue ??= $this->createQueue(); return $this->queue; } protected function getLoop(): LoopInterface { - if ($this->loop === null) { - $this->loop = $this->createLoop(); - } + $this->loop ??= $this->createLoop(); return $this->loop; } protected function getWorker(): WorkerInterface { - if ($this->worker === null) { - $this->worker = $this->createWorker(); - } + $this->worker ??= $this->createWorker(); return $this->worker; } protected function getContainer(): ContainerInterface { - if ($this->container === null) { - $this->container = $this->createContainer(); - } + $this->container ??= $this->createContainer(); return $this->container; } @@ -118,13 +110,13 @@ 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(), + ), ); } diff --git a/tests/Unit/Message/Handler/Resolver/HandlerResolverTest.php b/tests/Unit/Message/Handler/Resolver/HandlerResolverTest.php new file mode 100644 index 00000000..e741ecf1 --- /dev/null +++ b/tests/Unit/Message/Handler/Resolver/HandlerResolverTest.php @@ -0,0 +1,201 @@ + $handler], $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); + + $this->assertSame($resolver->resolve('simple'), $resolver->resolve('simple')); + } + + public function testResolveStaticMethodHandler(): void + { + $container = new SimpleContainer(); + $resolver = new HandlerResolver( + ['static-handler' => StaticMessageHandler::handle(...)], + $container, + ); + + StaticMessageHandler::$wasHandled = false; + $resolvedHandler = $resolver->resolve('static-handler'); + $resolvedHandler->handle(new GenericMessage('static-handler', null)); + + $this->assertTrue(StaticMessageHandler::$wasHandled); + } + + public function testResolveStaticMethodStringHandler(): void + { + $container = new SimpleContainer(); + $resolver = new HandlerResolver( + ['static-handler' => StaticMessageHandler::class . '::handle'], + $container, + ); + + StaticMessageHandler::$wasHandled = false; + $resolvedHandler = $resolver->resolve('static-handler'); + $resolvedHandler->handle(new GenericMessage('static-handler', null)); + + $this->assertTrue(StaticMessageHandler::$wasHandled); + } + + public function testResolveNamedFunctionHandler(): void + { + $message = new GenericMessage('named-function-handler', null); + $resolver = new HandlerResolver( + ['named-function-handler' => __NAMESPACE__ . '\\namedFunctionHandler'], + new SimpleContainer(), + ); + + try { + $resolver->resolve('named-function-handler')->handle($message); + + $this->assertSame([$message], FakeHandler::$processedMessages); + } finally { + FakeHandler::$processedMessages = []; + } + } + + public function testResolveThrowsWhenDefinitionMethodUndefined(): void + { + $this->expectException(InvalidHandlerConfigurationException::class); + $this->expectExceptionMessage('Queue handler for message type "simple" is configured incorrectly'); + + $container = new SimpleContainer([FakeHandler::class => new FakeHandler()]); + $resolver = new HandlerResolver( + ['simple' => [FakeHandler::class, 'undefinedMethod']], + $container, + ); + + $resolver->resolve('simple'); + } + + public function testResolveThrowsWhenDefinitionClassUndefined(): void + { + $this->expectException(InvalidHandlerConfigurationException::class); + $this->expectExceptionMessage('Queue handler for message type "simple" is configured incorrectly'); + + $container = new SimpleContainer([FakeHandler::class => new FakeHandler()]); + $resolver = new HandlerResolver( + ['simple' => ['UndefinedClass', 'handle']], + $container, + ); + + $resolver->resolve('simple'); + } + + public function testResolveThrowsWhenDefinitionClassNotFoundInContainer(): void + { + $this->expectException(InvalidHandlerConfigurationException::class); + $this->expectExceptionMessage('Queue handler for message type "simple" is configured incorrectly'); + + $container = new SimpleContainer(); + $resolver = new HandlerResolver( + ['simple' => [FakeHandler::class, 'handle']], + $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); + + $resolver->resolve('nonexistent'); + } + + public function testResolveThrowsWhenHandlerInContainerNotImplementingInterface(): void + { + $this->expectException(InvalidHandlerConfigurationException::class); + $this->expectExceptionMessage('Queue handler for message type "invalid" is configured incorrectly'); + + $container = new SimpleContainer([ + 'invalid' => new class { + public function handle(): void {} + }, + ]); + $resolver = new HandlerResolver([], $container); + + $resolver->resolve('invalid'); + } + + public function testResolveThrowsWhenMessageTypeIsEmpty(): void + { + $this->expectException(LogicException::class); + $this->expectExceptionMessage('Message type cannot be empty.'); + + $container = new SimpleContainer(); + $resolver = new HandlerResolver([], $container); + + $resolver->resolve(''); + } +} + +function namedFunctionHandler(MessageInterface $message): void +{ + FakeHandler::$processedMessages[] = $message; +} diff --git a/tests/Unit/WorkerTest.php b/tests/Unit/WorkerTest.php index 6840f828..df648666 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\HandlerResolver; 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,29 @@ 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): HandlerResolver { - $message = new GenericMessage('static-handler', ['test-data']); $container = new SimpleContainer(); - $handlers = [ - 'static-handler' => StaticMessageHandler::handle(...), - ]; - - $queueName = 'test-queue'; - $worker = $this->createWorkerByParams($handlers, $container); - StaticMessageHandler::$wasHandled = false; - $worker->process($message, $queueName); - $this->assertTrue(StaticMessageHandler::$wasHandled); + return new HandlerResolver( + [$message->getType() => $handler], + $container, + ); } private function createWorkerByParams( - array $handlers, - ContainerInterface $container, + HandlerResolver $handlerResolver, ?LoggerInterface $logger = null, + ?ConsumeMiddlewareDispatcher $consumeMiddlewareDispatcher = null, + ?FailureMiddlewareDispatcher $failureMiddlewareDispatcher = null, ): Worker { /** @var ConsumeMiddlewareFactoryInterface&MockObject $consumeMiddlewareFactory */ $consumeMiddlewareFactory = $this->createMock(ConsumeMiddlewareFactoryInterface::class); @@ -255,13 +130,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, ); } }