Skip to content
Draft
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
18 changes: 18 additions & 0 deletions config/di.php
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@
declare(strict_types=1);

use Psr\Container\ContainerInterface;
use Yiisoft\Definitions\Reference;
use Yiisoft\Queue\Cli\LoopInterface;
use Yiisoft\Queue\Cli\SignalLoop;
use Yiisoft\Queue\Cli\SimpleLoop;
Expand All @@ -21,6 +22,10 @@
use Yiisoft\Queue\Middleware\Push\PushMiddlewareConfig;
use Yiisoft\Queue\Middleware\Push\PushMiddlewareFactory;
use Yiisoft\Queue\Middleware\Push\PushMiddlewareFactoryInterface;
use Yiisoft\Queue\Middleware\Worker\WorkerFinalHandler;
use Yiisoft\Queue\Middleware\Worker\WorkerMiddlewareDispatcher;
use Yiisoft\Queue\Middleware\Worker\WorkerMiddlewareFactory;
use Yiisoft\Queue\Middleware\Worker\WorkerMiddlewareFactoryInterface;
use Yiisoft\Queue\Message\Handler\HandlerResolver;
use Yiisoft\Queue\Worker\Worker as QueueWorker;
use Yiisoft\Queue\Worker\WorkerInterface;
Expand All @@ -40,6 +45,7 @@
PushMiddlewareFactoryInterface::class => PushMiddlewareFactory::class,
ConsumeMiddlewareFactoryInterface::class => ConsumeMiddlewareFactory::class,
FailureMiddlewareFactoryInterface::class => FailureMiddlewareFactory::class,
WorkerMiddlewareFactoryInterface::class => WorkerMiddlewareFactory::class,
PushMiddlewareConfig::class => [
'__construct()' => ['commonMiddlewareDefinitions' => $params['yiisoft/queue']['middlewares-push']],
],
Expand All @@ -49,6 +55,18 @@
FailureMiddlewareDispatcher::class => [
'__construct()' => ['middlewareDefinitions' => $params['yiisoft/queue']['middlewares-fail']],
],
WorkerMiddlewareDispatcher::class => [
'__construct()' => [
'middlewareDefinitions' => $params['yiisoft/queue']['middlewares-worker'],
'finalHandler' => Reference::to(WorkerFinalHandler::class),
],
],
WorkerFinalHandler::class => [
'__construct()' => [
'handlerResolver' => Reference::to(HandlerResolver::class),
'consumeMiddlewareDispatcher' => Reference::to(ConsumeMiddlewareDispatcher::class),
],
],
MessageEncoderInterface::class => JsonMessageEncoder::class,
MessageSerializerInterface::class => MessageSerializer::class,
MessageClassResolverInterface::class => [
Expand Down
1 change: 1 addition & 0 deletions config/params.php
Original file line number Diff line number Diff line change
Expand Up @@ -46,6 +46,7 @@
'middlewares-push' => [],
'middlewares-consume' => [],
'middlewares-fail' => [],
'middlewares-worker' => [],
],
'yiisoft/yii-debug' => [
'collectors' => [
Expand Down
3 changes: 2 additions & 1 deletion src/AsyncQueueProducer.php
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@
use Yiisoft\Queue\Message\MessageInterface;
use Yiisoft\Queue\Middleware\Push\AdapterPushHandler;
use Yiisoft\Queue\Middleware\Push\PushMiddlewareConfig;
use Yiisoft\Queue\Middleware\Push\PushRequest;
use Yiisoft\Queue\Middleware\Push\PushMiddlewareDispatcher;

/**
Expand Down Expand Up @@ -50,7 +51,7 @@ public function push(MessageInterface $message): MessageInterface
'Preparing to push message with message type "{messageType}".',
['messageType' => $message->getType()],
);
$message = $this->dispatcher->dispatch($message);
$message = $this->dispatcher->dispatch(new PushRequest($message, $this->queueName))->getMessage();
$id = IdEnvelope::fromMessage($message)->getId();
$this->logger->info(
$id === null
Expand Down
5 changes: 2 additions & 3 deletions src/Middleware/Push/AdapterPushHandler.php
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,6 @@
namespace Yiisoft\Queue\Middleware\Push;

use Yiisoft\Queue\Adapter\AdapterInterface;
use Yiisoft\Queue\Message\MessageInterface;

/**
* @internal
Expand All @@ -16,8 +15,8 @@ public function __construct(
private readonly AdapterInterface $adapter,
) {}

public function handlePush(MessageInterface $message): MessageInterface
public function handlePush(PushRequest $request): PushRequest
{
return $this->adapter->push($message);
return $request->withMessage($this->adapter->push($request->getMessage()));
}
}
9 changes: 5 additions & 4 deletions src/Middleware/Push/Implementation/IdMiddleware.php
Original file line number Diff line number Diff line change
Expand Up @@ -5,25 +5,26 @@
namespace Yiisoft\Queue\Middleware\Push\Implementation;

use Yiisoft\Queue\Message\IdEnvelope;
use Yiisoft\Queue\Message\MessageInterface;
use Yiisoft\Queue\Middleware\Push\PushHandlerInterface;
use Yiisoft\Queue\Middleware\Push\PushMiddlewareInterface;
use Yiisoft\Queue\Middleware\Push\PushRequest;

/**
* A middleware for message ID setting.
*/
final class IdMiddleware implements PushMiddlewareInterface
{
public function processPush(MessageInterface $message, PushHandlerInterface $handler): MessageInterface
public function processPush(PushRequest $request, PushHandlerInterface $handler): PushRequest
{
$message = $request->getMessage();
$envelope = IdEnvelope::fromMessage($message);

if ($envelope->getId() === null) {
return $handler->handlePush(
new IdEnvelope($message, uniqid('yii3-message-', true)),
$request->withMessage(new IdEnvelope($message, uniqid('yii3-message-', true))),
);
}

return $handler->handlePush($message);
return $handler->handlePush($request);
}
}
4 changes: 1 addition & 3 deletions src/Middleware/Push/PushHandlerInterface.php
Original file line number Diff line number Diff line change
Expand Up @@ -4,9 +4,7 @@

namespace Yiisoft\Queue\Middleware\Push;

use Yiisoft\Queue\Message\MessageInterface;

interface PushHandlerInterface
{
public function handlePush(MessageInterface $message): MessageInterface;
public function handlePush(PushRequest $request): PushRequest;
}
7 changes: 3 additions & 4 deletions src/Middleware/Push/PushMiddlewareDispatcher.php
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,6 @@
namespace Yiisoft\Queue\Middleware\Push;

use Closure;
use Yiisoft\Queue\Message\MessageInterface;

/**
* @internal Used internally by {@see SyncQueueProducer} and {@see AsyncQueueProducer}.
Expand Down Expand Up @@ -33,13 +32,13 @@ public function __construct(
/**
* Dispatch message through middleware to get response.
*
* @param MessageInterface $message Message to pass to middleware.
* @param PushRequest $request Request to pass to middleware.
*/
public function dispatch(MessageInterface $message): MessageInterface
public function dispatch(PushRequest $request): PushRequest
{
$this->stack ??= new PushMiddlewareStack($this->buildMiddlewares(), $this->finalHandler);

return $this->stack->handlePush($message);
return $this->stack->handlePush($request);
}

/**
Expand Down
14 changes: 6 additions & 8 deletions src/Middleware/Push/PushMiddlewareFactory.php
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,6 @@

namespace Yiisoft\Queue\Middleware\Push;

use Yiisoft\Queue\Message\MessageInterface;
use Yiisoft\Queue\Middleware\InvalidMiddlewareDefinitionException;
use Yiisoft\Queue\Middleware\MiddlewareFactory;

Expand All @@ -25,15 +24,14 @@ final class PushMiddlewareFactory extends MiddlewareFactory implements PushMiddl
*
* - A middleware object.
* - A name of a middleware class. The middleware instance will be obtained from container and executed.
* - A callable with `function(MessageInterface $message, PushHandlerInterface $handler):
* MessageInterface` signature.
* - A callable with `function(PushRequest $request, PushHandlerInterface $handler): PushRequest` signature.
* - A controller handler action in format `[TestController::class, 'index']`. `TestController` instance will
* be created and `index()` method will be executed.
* - A function returning a middleware. The middleware returned will be executed.
*
* For handler action and callable
* typed parameters are automatically injected using dependency injection container.
* Current message and handler could be obtained by type-hinting for {@see MessageInterface}
* Current request and handler could be obtained by type-hinting for {@see PushRequest}
* and {@see PushHandlerInterface}.
*
* @throws InvalidMiddlewareDefinitionException
Expand Down Expand Up @@ -74,15 +72,15 @@ public function __construct(callable $callback)
$this->callback = $callback;
}

public function processPush(MessageInterface $message, PushHandlerInterface $handler): MessageInterface
public function processPush(PushRequest $request, PushHandlerInterface $handler): PushRequest
{
$response = ($this->callback)($message, $handler);
if ($response instanceof MessageInterface) {
$response = ($this->callback)($request, $handler);
if ($response instanceof PushRequest) {
return $response;
}

if ($response instanceof PushMiddlewareInterface) {
return $response->processPush($message, $handler);
return $response->processPush($request, $handler);
}

throw new InvalidMiddlewareDefinitionException($this->callback);
Expand Down
4 changes: 1 addition & 3 deletions src/Middleware/Push/PushMiddlewareInterface.php
Original file line number Diff line number Diff line change
Expand Up @@ -4,9 +4,7 @@

namespace Yiisoft\Queue\Middleware\Push;

use Yiisoft\Queue\Message\MessageInterface;

interface PushMiddlewareInterface
{
public function processPush(MessageInterface $message, PushHandlerInterface $handler): MessageInterface;
public function processPush(PushRequest $request, PushHandlerInterface $handler): PushRequest;
}
9 changes: 4 additions & 5 deletions src/Middleware/Push/PushMiddlewareStack.php
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,6 @@
namespace Yiisoft\Queue\Middleware\Push;

use Closure;
use Yiisoft\Queue\Message\MessageInterface;

/**
* @internal
Expand All @@ -31,10 +30,10 @@ public function __construct(
private readonly PushHandlerInterface $finalHandler,
) {}

public function handlePush(MessageInterface $message): MessageInterface
public function handlePush(PushRequest $request): PushRequest
{
$this->stack ??= $this->build();
return $this->stack->handlePush($message);
return $this->stack->handlePush($request);
}

private function build(): PushHandlerInterface
Expand Down Expand Up @@ -66,11 +65,11 @@ public function __construct(
private readonly PushHandlerInterface $handler,
) {}

public function handlePush(MessageInterface $message): MessageInterface
public function handlePush(PushRequest $request): PushRequest
{
$this->middleware ??= ($this->middlewareFactory)();

return $this->middleware->processPush($message, $this->handler);
return $this->middleware->processPush($request, $this->handler);
}
};
}
Expand Down
38 changes: 38 additions & 0 deletions src/Middleware/Push/PushRequest.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,38 @@
<?php

declare(strict_types=1);

namespace Yiisoft\Queue\Middleware\Push;

use Yiisoft\Queue\Message\MessageInterface;

final class PushRequest
{
public function __construct(private MessageInterface $message, private string $queueName) {}

public function getMessage(): MessageInterface
{
return $this->message;
}

public function getQueueName(): string
{
return $this->queueName;
}

public function withMessage(MessageInterface $message): self
{
$instance = clone $this;
$instance->message = $message;

return $instance;
}

public function withQueueName(string $queueName): self
{
$instance = clone $this;
$instance->queueName = $queueName;

return $instance;
}
}
7 changes: 3 additions & 4 deletions src/Middleware/Push/SynchronousPushHandler.php
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,6 @@

namespace Yiisoft\Queue\Middleware\Push;

use Yiisoft\Queue\Message\MessageInterface;
use Yiisoft\Queue\Worker\WorkerInterface;

/**
Expand All @@ -17,10 +16,10 @@ public function __construct(
private readonly string $queueName,
) {}

public function handlePush(MessageInterface $message): MessageInterface
public function handlePush(PushRequest $request): PushRequest
{
$this->worker->process($message, $this->queueName);
$this->worker->process($request->getMessage(), $request->getQueueName());

return $message;
return $request;
}
}
27 changes: 27 additions & 0 deletions src/Middleware/Worker/WorkerFinalHandler.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,27 @@
<?php

declare(strict_types=1);

namespace Yiisoft\Queue\Middleware\Worker;

use Yiisoft\Queue\Message\Handler\HandlerResolver;
use Yiisoft\Queue\Middleware\Consume\ConsumeFinalHandler;
use Yiisoft\Queue\Middleware\Consume\ConsumeMiddlewareDispatcher;
use Yiisoft\Queue\Middleware\Consume\ConsumeRequest;

final class WorkerFinalHandler implements WorkerHandlerInterface
{
public function __construct(
private readonly HandlerResolver $handlerResolver,
private readonly ConsumeMiddlewareDispatcher $consumeMiddlewareDispatcher,
) {}

public function handleWorker(WorkerRequest $request): WorkerRequest
{
$handler = $this->handlerResolver->resolve($request->getMessage()->getType());
$consumeRequest = new ConsumeRequest($request->getMessage(), $request->getQueueName());
$this->consumeMiddlewareDispatcher->dispatch($consumeRequest, new ConsumeFinalHandler($handler->handle(...)));

return $request;
}
}
10 changes: 10 additions & 0 deletions src/Middleware/Worker/WorkerHandlerInterface.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,10 @@
<?php

declare(strict_types=1);

namespace Yiisoft\Queue\Middleware\Worker;

interface WorkerHandlerInterface
{
public function handleWorker(WorkerRequest $request): WorkerRequest;
}
46 changes: 46 additions & 0 deletions src/Middleware/Worker/WorkerMiddlewareDispatcher.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,46 @@
<?php

declare(strict_types=1);

namespace Yiisoft\Queue\Middleware\Worker;

use Closure;

final class WorkerMiddlewareDispatcher
{
private ?WorkerMiddlewareStack $stack = null;

/**
* @param mixed[] $middlewareDefinitions Middleware definitions.
*/
public function __construct(
private readonly WorkerMiddlewareFactoryInterface $middlewareFactory,
private readonly array $middlewareDefinitions,
private readonly WorkerHandlerInterface $finalHandler,
) {}

public function dispatch(WorkerRequest $request): WorkerRequest
{
$this->stack ??= new WorkerMiddlewareStack($this->buildMiddlewares(), $this->finalHandler);

return $this->stack->handleWorker($request);
}

/**
* @psalm-return list<Closure():WorkerMiddlewareInterface>
*/
private function buildMiddlewares(): array
{
/** @var list<Closure():WorkerMiddlewareInterface> $middlewares */
$middlewares = [];
$factory = $this->middlewareFactory;

foreach ($this->middlewareDefinitions as $middlewareDefinition) {
$middlewares[] = static fn(): WorkerMiddlewareInterface => $factory->createWorkerMiddleware(
$middlewareDefinition,
);
}

return $middlewares;
}
}
Loading
Loading