Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 5 additions & 2 deletions config/di.php
Original file line number Diff line number Diff line change
Expand Up @@ -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')
Expand Down
4 changes: 2 additions & 2 deletions config/params.php
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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,
Expand Down
30 changes: 30 additions & 0 deletions src/Message/Handler/CallableHandler.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,30 @@
<?php

declare(strict_types=1);

namespace Yiisoft\Queue\Message\Handler;

use Yiisoft\Queue\Message\MessageInterface;

/**
* Handles a message by invoking the given callable.
*
* @psalm-type MessageHandlerCallable = callable(MessageInterface $message): void
*/
final class CallableHandler implements HandlerInterface
{
/**
* @param callable $handler Callable invoked to handle a message.
* Format: `function (MessageInterface $message): void`.
*
* @psalm-param MessageHandlerCallable $handler
*/
public function __construct(
private readonly mixed $handler,
) {}

public function handle(MessageInterface $message): void
{
($this->handler)($message);
}
}
18 changes: 18 additions & 0 deletions src/Message/Handler/HandlerInterface.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,18 @@
<?php

declare(strict_types=1);

namespace Yiisoft\Queue\Message\Handler;

use Yiisoft\Queue\Message\MessageInterface;

/**
* Handles a message.
*/
interface HandlerInterface
{
/**
* Handle the given message.
*/
public function handle(MessageInterface $message): void;
}
25 changes: 25 additions & 0 deletions src/Message/Handler/Resolver/HandlerNotFoundException.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,25 @@
<?php

declare(strict_types=1);

namespace Yiisoft\Queue\Message\Handler\Resolver;

use LogicException;
use Throwable;

use function sprintf;

/**
* Thrown when a handler for the given message type is not found.
*/
final class HandlerNotFoundException extends LogicException
{
public function __construct(string $messageType, int $code = 0, ?Throwable $previous = null)
{
parent::__construct(
sprintf('Queue handler for message type "%s" does not exist.', $messageType),
$code,
$previous,
);
}
}
61 changes: 61 additions & 0 deletions src/Message/Handler/Resolver/HandlerResolver.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,61 @@
<?php

declare(strict_types=1);

namespace Yiisoft\Queue\Message\Handler\Resolver;

use Psr\Container\ContainerInterface;
use Yiisoft\Queue\Message\Handler\CallableHandler;
use Yiisoft\Queue\Message\Handler\HandlerInterface;
use Yiisoft\Queue\Message\MessageInterface;
use Yiisoft\Queue\Middleware\CallableFactory;
use Yiisoft\Queue\Middleware\InvalidCallableConfigurationException;

use function array_key_exists;
use function is_string;

/**
* Resolves message handlers from configuration, a DI container, or a callable factory.
*/
final class HandlerResolver implements HandlerResolverInterface
{
/** @var array<non-empty-string, HandlerInterface> Cache of resolved handlers */
private array $cache = [];

public function __construct(
/** @var array<non-empty-string, array|callable|object|string|null> */
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);
}
}
}
22 changes: 22 additions & 0 deletions src/Message/Handler/Resolver/HandlerResolverInterface.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,22 @@
<?php

declare(strict_types=1);

namespace Yiisoft\Queue\Message\Handler\Resolver;

use Yiisoft\Queue\Message\Handler\HandlerInterface;

/**
* Resolves a message type to a handler.
*/
interface HandlerResolverInterface
{
/**
* Get a handler for the given message type.
*
* @param string $messageType Message type.
*
* @throws HandlerNotFoundException If no handler exists for the message type.
*/
public function resolve(string $messageType): HandlerInterface;
}
10 changes: 0 additions & 10 deletions src/Message/MessageHandlerInterface.php

This file was deleted.

75 changes: 6 additions & 69 deletions src/Worker/Worker.php
Original file line number Diff line number Diff line change
Expand Up @@ -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\Resolver\HandlerResolverInterface;
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<non-empty-string, callable|null> Cache of resolved handlers */
private array $handlersCached = [];

public function __construct(
/** @var array<non-empty-string, array|callable|object|string|null> */
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,
) {}

/**
Expand All @@ -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();
Expand All @@ -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();
}
}
15 changes: 8 additions & 7 deletions tests/Benchmark/QueueBench.php
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;
Expand All @@ -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);
Expand Down
Loading
Loading