-
Notifications
You must be signed in to change notification settings - Fork 4
/
AmqpCommandConsumerCallback.php
76 lines (62 loc) · 2.02 KB
/
AmqpCommandConsumerCallback.php
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
<?php
/**
* This file is part of the prooph/humus-amqp-producer.
* (c) 2016-2019 Alexander Miertsch <kontakt@codeliner.ws>
* (c) 2016-2019 Sascha-Oliver Prolic <saschaprolic@googlemail.com>
*
* For the full copyright and license information, please view the LICENSE
* file that was distributed with this source code.
*/
declare(strict_types=1);
namespace Prooph\ServiceBus\Message\HumusAmqp;
use DateTimeImmutable;
use DateTimeZone;
use Humus\Amqp\DeliveryResult;
use Humus\Amqp\Envelope;
use Humus\Amqp\Queue;
use Prooph\Common\Messaging\MessageFactory;
use Prooph\EventStore\Exception\ConcurrencyException;
use Prooph\ServiceBus\CommandBus;
final class AmqpCommandConsumerCallback
{
/**
* @var CommandBus
*/
private $commandBus;
/**
* @var MessageFactory
*/
private $messageFactory;
public function __construct(CommandBus $commandBus, MessageFactory $messageFactory)
{
$this->commandBus = $commandBus;
$this->messageFactory = $messageFactory;
}
public function __invoke(Envelope $envelope, Queue $queue): DeliveryResult
{
$data = \json_decode($envelope->getBody(), true);
if (! isset($data['created_at'])) {
return DeliveryResult::MSG_REJECT();
}
$data['created_at'] = DateTimeImmutable::createFromFormat(
'Y-m-d\TH:i:s.u',
$data['created_at'],
new DateTimeZone('UTC')
);
if (false === $data['created_at']) {
return DeliveryResult::MSG_REJECT();
}
try {
$command = $this->messageFactory->createMessageFromArray($envelope->getType(), $data);
$this->commandBus->dispatch($command);
} catch (\Throwable $e) {
while ($e = $e->getPrevious()) {
if ($e instanceof ConcurrencyException) {
return DeliveryResult::MSG_REJECT_REQUEUE();
}
}
return DeliveryResult::MSG_REJECT();
}
return DeliveryResult::MSG_ACK();
}
}