Skip to content

Commit ba6bca9

Browse files
committed
Allow only loop consume strategy for Spool
1 parent dcf6049 commit ba6bca9

8 files changed

Lines changed: 50 additions & 73 deletions

File tree

src/Adapter/Amqp/Queue/AmqpQueueFactory.php

Lines changed: 1 addition & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -29,11 +29,6 @@ public function __construct(
2929
) {
3030
}
3131

32-
public function withPassive(): QueueFactoryInterface
33-
{
34-
return new AmqpQueueFactory($this->channelFactory, $this->definition->withPassive(true));
35-
}
36-
3732
public function create(): QueueInterface
3833
{
3934
if ($this->queue) {
@@ -59,14 +54,7 @@ public function create(): QueueInterface
5954

6055
$queue->declareQueue();
6156

62-
$amqpQueue = new AmqpQueue($channel, $queue);
63-
64-
if ($this->definition->passive) {
65-
return $amqpQueue;
66-
}
67-
68-
$this->queue = $amqpQueue;
69-
57+
$this->queue = new AmqpQueue($channel, $queue);
7058
$channel->getConnection()->attach($this);
7159

7260
foreach ($this->definition->bindings as $binding) {

src/Adapter/AmqpLib/Queue/AmqpQueue.php

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -52,13 +52,14 @@ public function consume(\Closure $handler, string $tag = ''): void
5252
false,
5353
$this->definition->exclusive,
5454
false,
55-
function (AMQPMessage $message) use ($handler, $queueName, &$stopConsuming): void {
55+
function (AMQPMessage $message) use ($handler, $queueName, &$stopConsuming, &$amqplibChannel): void {
5656
$receivedMessage = new AmqpReceivedMessage($message, $queueName);
5757

5858
$result = $handler($receivedMessage);
5959

6060
if (false === $result) {
6161
$stopConsuming = true;
62+
$amqplibChannel->basic_cancel($message->getConsumerTag(), false, true);
6263
}
6364
}
6465
);

src/Adapter/AmqpLib/Queue/AmqpQueueFactory.php

Lines changed: 0 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -30,11 +30,6 @@ public function __construct(
3030
) {
3131
}
3232

33-
public function withPassive(): QueueFactoryInterface
34-
{
35-
return new AmqpQueueFactory($this->channelFactory, $this->definition->withPassive(true));
36-
}
37-
3833
public function create(): QueueInterface
3934
{
4035
if ($this->queue) {
@@ -55,10 +50,6 @@ public function create(): QueueInterface
5550
$queue = new AmqpQueue($channel, $this->definition);
5651
$queue->declare();
5752

58-
if ($this->definition->passive) {
59-
return $queue;
60-
}
61-
6253
$connection->attach($this);
6354

6455
foreach ($this->definition->bindings as $binding) {

src/Consumer/AbstractConsumer.php

Lines changed: 0 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -124,18 +124,4 @@ protected function doRun(bool $autoAck = true, ?\Closure $beforeCallback = null,
124124
}
125125
}, $this->configuration->tagGenerator->generate());
126126
}
127-
128-
/**
129-
* Performs a graceful AMQP disconnect.
130-
*
131-
* A passively declared queue is created first to force a synchronous AMQP
132-
* round-trip on the current channel. This guarantees that all previously
133-
* sent frames (e.g. acknowledgements) were processed by the broker before
134-
* the connection is closed.
135-
*/
136-
protected function gracefulDisconnect(): void
137-
{
138-
$queue = $this->queueFactory->withPassive()->create();
139-
$queue->getChannel()->getConnection()->disconnect();
140-
}
141127
}

src/Consumer/Spool/SpoolConsumer.php

Lines changed: 15 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -16,12 +16,18 @@
1616
use FiveLab\Component\Amqp\AmqpEvents;
1717
use FiveLab\Component\Amqp\Channel\ChannelInterface;
1818
use FiveLab\Component\Amqp\Consumer\AbstractConsumer;
19+
use FiveLab\Component\Amqp\Consumer\ConsumerConfiguration;
1920
use FiveLab\Component\Amqp\Consumer\ConsumerStoppedReason;
21+
use FiveLab\Component\Amqp\Consumer\Handler\MessageHandlerInterface;
22+
use FiveLab\Component\Amqp\Consumer\Strategy\ConsumeStrategyInterface;
23+
use FiveLab\Component\Amqp\Consumer\Strategy\DefaultConsumeStrategy;
24+
use FiveLab\Component\Amqp\Consumer\Strategy\LoopConsumeStrategy;
2025
use FiveLab\Component\Amqp\Event\ConsumerStartedEvent;
2126
use FiveLab\Component\Amqp\Event\ConsumerStoppedEvent;
2227
use FiveLab\Component\Amqp\Exception\ConsumerTimeoutExceedException;
2328
use FiveLab\Component\Amqp\Message\MutableReceivedMessages;
2429
use FiveLab\Component\Amqp\Message\ReceivedMessage;
30+
use FiveLab\Component\Amqp\Queue\QueueFactoryInterface;
2531
use FiveLab\Component\Amqp\Queue\QueueInterface;
2632

2733
/**
@@ -35,6 +41,13 @@
3541
{
3642
public function run(): void
3743
{
44+
if (!$this->strategy instanceof LoopConsumeStrategy) {
45+
throw new \LogicException(\sprintf(
46+
'The "%s" consume strategy is not supported for spool consumer (only loop supported).',
47+
\get_class($this->strategy)
48+
));
49+
}
50+
3851
$this->allowConsuming();
3952

4053
$this->getEventDispatcher()?->dispatch(new ConsumerStartedEvent($this), AmqpEvents::CONSUMER_STARTED);
@@ -49,16 +62,13 @@ public function run(): void
4962
$receivedMessages = new MutableReceivedMessages();
5063
$endTime = \microtime(true) + $this->configuration->timeout;
5164

52-
$countOfProcessedMessages = 0;
53-
5465
try {
5566
$this->doRun(
5667
false,
57-
function (ReceivedMessage $message) use ($receivedMessages, &$countOfProcessedMessages, &$endTime): void {
68+
function (ReceivedMessage $message) use ($receivedMessages, &$endTime): void {
5869
$receivedMessages->push($message);
59-
$countOfProcessedMessages++;
6070

61-
if ($countOfProcessedMessages >= $this->configuration->prefetchCount) {
71+
if (\count($receivedMessages) >= $this->configuration->prefetchCount) {
6272
$this->strategy->stopConsume();
6373
}
6474

@@ -77,8 +87,6 @@ static function (ReceivedMessage $message): void {
7787
} catch (ConsumerTimeoutExceedException) {
7888
$this->flushMessages($receivedMessages);
7989

80-
$this->gracefulDisconnect();
81-
8290
$this->getEventDispatcher()?->dispatch(new ConsumerStoppedEvent($this, ConsumerStoppedReason::Timeout), AmqpEvents::CONSUMER_STOPPED);
8391

8492
continue;
@@ -91,18 +99,13 @@ static function (ReceivedMessage $message): void {
9199

92100
$receivedMessages->clear();
93101

94-
$this->gracefulDisconnect();
95-
96102
throw $error;
97103
}
98104

99105
$this->flushMessages($receivedMessages);
100-
// @todo: critical case, we can't reconnect per each batch block.
101-
$this->gracefulDisconnect();
102106
}
103107

104108
$this->flushMessages($receivedMessages);
105-
$this->gracefulDisconnect();
106109
}
107110

108111
private function flushMessages(MutableReceivedMessages $messages): void
@@ -123,8 +126,6 @@ private function flushMessages(MutableReceivedMessages $messages): void
123126

124127
$messages->clear();
125128

126-
$this->gracefulDisconnect();
127-
128129
throw $e;
129130
}
130131

src/Queue/QueueFactoryInterface.php

Lines changed: 0 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -21,11 +21,4 @@ interface QueueFactoryInterface
2121
* @return QueueInterface
2222
*/
2323
public function create(): QueueInterface;
24-
25-
/**
26-
* Create the queue factory with passive mode
27-
*
28-
* @return QueueFactoryInterface
29-
*/
30-
public function withPassive(): QueueFactoryInterface;
3124
}

tests/Functional/Adapter/SpoolConsumerTestCase.php

Lines changed: 29 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -71,7 +71,8 @@ public function shouldSuccessConsume(): void
7171
$consumer = new SpoolConsumer(
7272
$this->queueFactory,
7373
$this->messageHandler,
74-
new SpoolConsumerConfiguration(10, 1)
74+
new SpoolConsumerConfiguration(10, 1),
75+
new LoopConsumeStrategy(),
7576
);
7677

7778
$this->runConsumer($consumer);
@@ -99,7 +100,8 @@ public function shouldSuccessReturnMessagesToBrokerIfSpoolFailed(): void
99100
$consumer = new SpoolConsumer(
100101
$this->queueFactory,
101102
$this->messageHandler,
102-
new SpoolConsumerConfiguration(10, 1, 0, true)
103+
new SpoolConsumerConfiguration(10, 1, 0, true),
104+
new LoopConsumeStrategy(),
103105
);
104106

105107
try {
@@ -131,7 +133,8 @@ public function shouldNotReturnMessagesToBrokerIfSpoolFailedIfRequeueIsFalse():
131133
$consumer = new SpoolConsumer(
132134
$this->queueFactory,
133135
$this->messageHandler,
134-
new SpoolConsumerConfiguration(10, 1, 0, false)
136+
new SpoolConsumerConfiguration(10, 1, 0, false),
137+
new LoopConsumeStrategy(),
135138
);
136139

137140
try {
@@ -158,7 +161,8 @@ public function shouldReturnMessagesToBrokerIfFlushFailed(): void
158161
$consumer = new SpoolConsumer(
159162
$this->queueFactory,
160163
$this->messageHandler,
161-
new SpoolConsumerConfiguration(5, 1, 0, true)
164+
new SpoolConsumerConfiguration(5, 1, 0, true),
165+
new LoopConsumeStrategy(),
162166
);
163167

164168
try {
@@ -185,7 +189,8 @@ public function shouldNotReturnMessagesToBrokerIfFlushFailedAndRequeueIsFalse():
185189
$consumer = new SpoolConsumer(
186190
$this->queueFactory,
187191
$this->messageHandler,
188-
new SpoolConsumerConfiguration(5, 1, 0, false)
192+
new SpoolConsumerConfiguration(5, 1, 0, false),
193+
new LoopConsumeStrategy(),
189194
);
190195

191196
try {
@@ -221,7 +226,8 @@ public function shouldReturnMessagesToBrokerOnlyNotAckedMessagesIfFlushFalied():
221226
$consumer = new SpoolConsumer(
222227
$this->queueFactory,
223228
$this->messageHandler,
224-
new SpoolConsumerConfiguration(5, 1)
229+
new SpoolConsumerConfiguration(5, 1),
230+
new LoopConsumeStrategy(),
225231
);
226232

227233
try {
@@ -260,7 +266,8 @@ public function shouldThrowExceptionIfMessageHandlerTryAnsweringToBroker(): void
260266
$consumer = new SpoolConsumer(
261267
$this->queueFactory,
262268
$this->messageHandler,
263-
new SpoolConsumerConfiguration(10, 1)
269+
new SpoolConsumerConfiguration(10, 1),
270+
new LoopConsumeStrategy(),
264271
);
265272

266273
$this->expectException(\LogicException::class);
@@ -284,7 +291,8 @@ public function shouldNotThrowOrphanedEnvelope(): void
284291
$consumer = new SpoolConsumer(
285292
$this->queueFactory,
286293
$this->messageHandler,
287-
new SpoolConsumerConfiguration(10, 1)
294+
new SpoolConsumerConfiguration(10, 1),
295+
new LoopConsumeStrategy()
288296
);
289297

290298
$this->messageHandler->setHandlerCallback(function (ReceivedMessage $message) use (&$handledMessages, $consumer) {
@@ -329,7 +337,7 @@ public function shouldSuccessProcessOnStopAfterNExecutes(): void
329337
{
330338
$this->publishMessages(12);
331339

332-
$consumer = new SpoolConsumer($this->queueFactory, $this->messageHandler, new SpoolConsumerConfiguration(100, 1));
340+
$consumer = new SpoolConsumer($this->queueFactory, $this->messageHandler, new SpoolConsumerConfiguration(100, 1), new LoopConsumeStrategy());
333341
$consumer->setEventDispatcher($eventDispatcher = new EventDispatcher());
334342
$eventDispatcher->addListener(AmqpEvents::PROCESSED_MESSAGE, (new StopAfterNExecutesListener(5))->onProcessedMessage(...));
335343

@@ -363,7 +371,8 @@ public function shouldFlushWithZeroReadTimeout(int $prefetchCount, int $messageC
363371
$consumer = new SpoolConsumer(
364372
$this->queueFactory,
365373
$this->messageHandler,
366-
new SpoolConsumerConfiguration($prefetchCount, 1, 1, true)
374+
new SpoolConsumerConfiguration($prefetchCount, 1, 1, true),
375+
new LoopConsumeStrategy(),
367376
);
368377

369378
$this->queueFactory->create()->getChannel()->getConnection()->setReadTimeout(0);
@@ -388,7 +397,7 @@ public function shouldSuccessReRunAfterStop(): void
388397
{
389398
$this->publishMessages(10);
390399

391-
$consumer = new SpoolConsumer($this->queueFactory, $this->messageHandler, new SpoolConsumerConfiguration(10, 1, 1, true));
400+
$consumer = new SpoolConsumer($this->queueFactory, $this->messageHandler, new SpoolConsumerConfiguration(10, 1, 1, true), new LoopConsumeStrategy());
392401
$consumer->setEventDispatcher($eventDispatcher = new EventDispatcher());
393402

394403
$eventDispatcher->addListener(AmqpEvents::PROCESSED_MESSAGE, static function (ProcessedMessageEvent $event): void {
@@ -461,7 +470,8 @@ public function shouldSuccessAck2KMessagesForQuorumQueue(): void
461470
$consumer = new SpoolConsumer(
462471
$quorumQueueFactory,
463472
$this->messageHandler,
464-
new SpoolConsumerConfiguration(2000, 2)
473+
new SpoolConsumerConfiguration(2000, 2),
474+
new LoopConsumeStrategy(),
465475
);
466476

467477
$this->runConsumer($consumer);
@@ -483,7 +493,13 @@ private function runConsumer(SpoolConsumer $consumer, bool $changeReadTimeout =
483493
$consumer->getQueue()->getChannel()->getConnection()->setReadTimeout(0.2);
484494
}
485495

486-
$consumer->run();
496+
try {
497+
$consumer->run();
498+
} finally {
499+
\usleep(10_000);
500+
$consumer->getQueue()->getChannel()->getConnection()->disconnect();
501+
\usleep(10_000);
502+
}
487503
}
488504

489505
private function registerTimeoutListenerForStop(SpoolConsumer $consumer): void

tests/Functional/Signals/ConsumerSignalsInConsumeStrategyTest.php

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -93,7 +93,7 @@ public function shouldSuccessCatchSignalsInSpoolConsumer(): void
9393
{
9494
$consumerFile = \realpath(__DIR__.'/../../Consumers/spool-with-signals.php');
9595

96-
$process = new PhpSubprocess([$consumerFile, 'consume']);
96+
$process = new PhpSubprocess([$consumerFile, 'loop']);
9797
$process->setTimeout(10);
9898

9999
$sendSignal = false;
@@ -108,8 +108,9 @@ public function shouldSuccessCatchSignalsInSpoolConsumer(): void
108108
$output = $process->getOutput();
109109

110110
$expectedOutput = <<<OUTPUT
111-
bla 0
111+
tick 1
112112
handle signal: 15
113+
bla 0
113114
flush messages
114115
115116
OUTPUT;

0 commit comments

Comments
 (0)