Final Class Yiisoft\Queue\Worker\Worker
| Inheritance | Yiisoft\ |
|---|---|
| Implements | Yiisoft\ |
Public Methods
| Method | Description | Defined By |
|---|---|---|
| __construct() | Yiisoft\ |
|
| process() | Yiisoft\ |
Method Details
| public mixed __construct ( array $handlers, \ | ||
| $handlers | array | |
| $logger | \ |
|
| $injector | \ |
|
| $container | \ |
|
| $consumeMiddlewareDispatcher | Yiisoft\ |
|
| $failureMiddlewareDispatcher | Yiisoft\ |
|
| $callableFactory | Yiisoft\ |
|
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,
) {}
| public Yiisoft\ | ||
| $message | Yiisoft\ |
|
| $queueName | string | |
| $retryProducer | ?\ |
|
| throws | Throwable | |
|---|---|---|
public function process(
MessageInterface $message,
string $queueName,
?QueueProducerInterface $retryProducer = null,
): MessageInterface {
$messageId = IdEnvelope::fromMessage($message)->getId();
if ($messageId === null) {
$this->logger->info('Processing message without ID.');
} else {
$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));
}
$request = new ConsumeRequest($message, $queueName);
$closure = fn(MessageInterface $message): mixed => $this->injector->invoke($handler, [$message]);
try {
return $this->consumeMiddlewareDispatcher->dispatch($request, $this->createConsumeHandler($closure))->getMessage();
} catch (Throwable $exception) {
$request = new FailureHandlingRequest($request->getMessage(), $exception, $request->getQueueName(), $retryProducer);
try {
$result = $this->failureMiddlewareDispatcher->dispatch($request, $this->createFailureHandler());
$this->logger->info($exception->getMessage());
return $result->getMessage();
} catch (Throwable $exception) {
$exception = new MessageFailureException($message, $exception);
$this->logger->error($exception->getMessage());
throw $exception;
}
}
}
User Contributed Notes
Leave a comment
Join the conversation to share a note.