-
Notifications
You must be signed in to change notification settings - Fork 4
/
AmqpDelayedMessageProducer.php
92 lines (75 loc) · 2.54 KB
/
AmqpDelayedMessageProducer.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
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
<?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 Humus\Amqp\Constants;
use Humus\Amqp\Producer;
use Prooph\Common\Messaging\Message;
use Prooph\Common\Messaging\MessageConverter;
use Prooph\Common\Messaging\MessageDataAssertion;
use Prooph\ServiceBus\Async\MessageProducer;
use Prooph\ServiceBus\Exception\RuntimeException;
use React\Promise\Deferred;
final class AmqpDelayedMessageProducer implements MessageProducer
{
/**
* @var Producer
*/
private $producer;
/**
* @var MessageConverter
*/
private $messageConverter;
/**
* @var string
*/
private $appId;
public function __construct(Producer $producer, MessageConverter $messageConverter, string $appId = '')
{
$this->producer = $producer;
$this->messageConverter = $messageConverter;
$this->appId = $appId;
}
/**
* @param Message $message
* @param Deferred|null $deferred
* @return void
* @throws RuntimeException
*/
public function __invoke(Message $message, Deferred $deferred = null): void
{
if (! $message instanceof DelayedMessage) {
throw new RuntimeException(\sprintf(
'Message is not a delayed message (instance of %s)',
DelayedMessage::class
));
}
if (null !== $deferred) {
throw new RuntimeException(__CLASS__ . ' cannot handle query messages which require future responses.');
}
$data = $this->arrayFromMessage($message);
$attributes = [
'headers' => [
'x-delay' => $message->delay(),
],
'app_id' => $this->appId,
'timestamp' => $message->createdAt()->getTimestamp(),
'type' => $message->messageName(),
];
$this->producer->publish($data, $message->messageName(), Constants::AMQP_NOPARAM, $attributes);
}
private function arrayFromMessage(Message $message): array
{
$messageData = $this->messageConverter->convertToArray($message);
MessageDataAssertion::assert($messageData);
$messageData['created_at'] = $message->createdAt()->format('Y-m-d\TH:i:s.u');
return $messageData;
}
}