Skip to content
Open
Show file tree
Hide file tree
Changes from 2 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);
Comment thread
samdark marked this conversation as resolved.
}
}
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