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 src/Debug/QueueWorkerInterfaceProxy.php
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,9 @@
use Yiisoft\Queue\QueueProducerInterface;
use Yiisoft\Queue\Worker\WorkerInterface;

/**
* Debug proxy for {@see WorkerInterface} that records processed messages into the {@see QueueCollector}.
*/
final class QueueWorkerInterfaceProxy implements WorkerInterface
{
public function __construct(
Expand All @@ -19,8 +22,8 @@ public function process(
MessageInterface $message,
string $queueName,
?QueueProducerInterface $retryProducer = null,
): MessageInterface {
): void {
$this->collector->collectWorkerProcessing($message, $queueName);
return $this->worker->process($message, $queueName, $retryProducer);
$this->worker->process($message, $queueName, $retryProducer);
}
}
8 changes: 3 additions & 5 deletions src/Worker/Worker.php
Original file line number Diff line number Diff line change
Expand Up @@ -34,7 +34,7 @@ public function process(
MessageInterface $message,
string $queueName,
?QueueProducerInterface $retryProducer = null,
): MessageInterface {
): void {
$messageId = IdEnvelope::fromMessage($message)->getId();
if ($messageId === null) {
$this->logger->info('Processing message without ID.');
Expand All @@ -46,15 +46,13 @@ public function process(
try {
$handler = $this->handlerResolver->resolve($message->getType());
$finishHandler = new ConsumeFinalHandler($handler->handle(...));
return $this->consumeMiddlewareDispatcher->dispatch($request, $finishHandler)->getMessage();
$this->consumeMiddlewareDispatcher->dispatch($request, $finishHandler);
} catch (Throwable $exception) {
$request = new FailureHandlingRequest($request->getMessage(), $exception, $request->getQueueName(), $retryProducer);

try {
$result = $this->failureMiddlewareDispatcher->dispatch($request, new FailureFinalHandler());
$this->failureMiddlewareDispatcher->dispatch($request, new FailureFinalHandler());
$this->logger->info($exception->getMessage());

return $result->getMessage();
} catch (Throwable $exception) {
$exception = new MessageFailureException($message, $exception);
$this->logger->error($exception->getMessage());
Expand Down
9 changes: 7 additions & 2 deletions src/Worker/WorkerInterface.php
Original file line number Diff line number Diff line change
Expand Up @@ -7,12 +7,17 @@
use Yiisoft\Queue\Message\MessageInterface;
use Yiisoft\Queue\QueueProducerInterface;

/**
* Executes a message: runs it through the consume pipeline and, on failure, through the failure pipeline.
*/
interface WorkerInterface
{
/** @param string $queueName Logical execution queue name. */
/**
* @param string $queueName Logical execution queue name.
*/
public function process(
MessageInterface $message,
string $queueName,
?QueueProducerInterface $retryProducer = null,
): MessageInterface;
): void;
}
4 changes: 1 addition & 3 deletions stubs/StubWorker.php
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,5 @@ public function process(
MessageInterface $message,
string $queueName,
?QueueProducerInterface $retryProducer = null,
): MessageInterface {
return $message;
}
): void {}
}
13 changes: 10 additions & 3 deletions tests/Integration/MiddlewareTest.php
Original file line number Diff line number Diff line change
Expand Up @@ -95,17 +95,24 @@ public function testFullStackConsume(): void
[],
);

$handledMessage = null;
$worker = new Worker(
new SimpleLogger(),
$consumeMiddlewareDispatcher,
$failureMiddlewareDispatcher,
new HandlerResolver(['test' => static fn() => true], $container),
new HandlerResolver(
['test' => static function (MessageInterface $message) use (&$handledMessage): void {
$handledMessage = $message;
}],
$container,
),
);

$message = new GenericMessage('test', ['initial']);
$messageConsumed = $worker->process($message, 'test-queue');
$worker->process($message, 'test-queue');

self::assertEquals($stack, $messageConsumed->getPayload());
self::assertInstanceOf(MessageInterface::class, $handledMessage);
self::assertEquals($stack, $handledMessage->getPayload());
}

public function testFullStackFailure(): void
Expand Down
4 changes: 1 addition & 3 deletions tests/Unit/Debug/QueueWorkerInterfaceProxyTest.php
Original file line number Diff line number Diff line change
Expand Up @@ -19,9 +19,7 @@ public function testProcessDelegatesToWorker(): void
$collector->startup();
$proxy = new QueueWorkerInterfaceProxy(new StubWorker(), $collector);

$result = $proxy->process($message, 'chan');

self::assertSame($message, $result);
$proxy->process($message, 'chan');

$collected = $collector->getCollected();
self::assertArrayHasKey('processingMessages', $collected);
Expand Down
12 changes: 3 additions & 9 deletions tests/Unit/Stubs/StubWorkerTest.php
Original file line number Diff line number Diff line change
Expand Up @@ -12,15 +12,9 @@ final class StubWorkerTest extends TestCase
{
public function testBase(): void
{
$worker = new StubWorker();

$sourceMessage = new GenericMessage('test', 42);
$this->expectNotToPerformAssertions();

$message = $worker->process($sourceMessage, 'test-queue');

$this->assertSame($sourceMessage, $message);
$this->assertSame('test', $message->getType());
$this->assertSame(42, $message->getPayload());
$this->assertSame([], $message->getMeta());
$worker = new StubWorker();
$worker->process(new GenericMessage('test', 42), 'test-queue');
}
}
16 changes: 9 additions & 7 deletions tests/Unit/WorkerTest.php
Original file line number Diff line number Diff line change
Expand Up @@ -94,7 +94,13 @@ public function testMessageFailureIsHandledSuccessfully(): void
$finalMessage = new GenericMessage('final', null);
/** @var FailureMiddlewareInterface&MockObject $failureMiddleware */
$failureMiddleware = $this->createMock(FailureMiddlewareInterface::class);
$failureMiddleware->method('processFailure')->willReturn(new FailureHandlingRequest($finalMessage, $originalException, $queueName));
$failureMiddleware
->expects(self::once())
->method('processFailure')
->with(self::callback(
static fn(FailureHandlingRequest $request): bool => $request->getException() === $originalException,
))
->willReturn(new FailureHandlingRequest($finalMessage, $originalException, $queueName));

/** @var FailureMiddlewareFactoryInterface&MockObject $failureMiddlewareFactory */
$failureMiddlewareFactory = $this->createMock(FailureMiddlewareFactoryInterface::class);
Expand All @@ -104,9 +110,7 @@ public function testMessageFailureIsHandledSuccessfully(): void
$handlerResolver = $this->createHandlerResolver($message, static fn() => null);
$worker = $this->createWorkerByParams($handlerResolver, new NullLogger(), $consumeDispatcher, $failureDispatcher);

$result = $worker->process($message, $queueName);

self::assertSame($finalMessage, $result);
$worker->process($message, $queueName);
}

public function testUnresolvableHandlerIsHandledByFailurePipeline(): void
Expand All @@ -133,9 +137,7 @@ public function testUnresolvableHandlerIsHandledByFailurePipeline(): void

$worker = $this->createWorkerByParams($handlerResolver, failureMiddlewareDispatcher: $failureDispatcher);

$result = $worker->process($message, $queueName);

self::assertSame($finalMessage, $result);
$worker->process($message, $queueName);
}

private function createHandlerResolver(MessageInterface $message, callable $handler): HandlerResolver
Expand Down
Loading