Repository navigation
Expand file tree
/
Copy pathMessageFactory.php
More file actions
120 lines (98 loc) · 3.24 KB
/
Copy pathMessageFactory.php
File metadata and controls
120 lines (98 loc) · 3.24 KB
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
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
<?php
namespace Codibly\QueuesBundle;
use AppBundle\Entity\Message\Message;
use Codibly\QueuesBundle\HashGenerator\HashGeneratorInterface;
use Bernard\Message\PlainMessage;
use Psr\Log\LoggerInterface;
class MessageFactory
{
/**
* @var LoggerInterface
*/
private $logger;
private $internalMessagesBinding = [];
/**
* @var HashGeneratorInterface
*/
private $hashGenerator;
/**
* MessageFactory constructor.
* @param LoggerInterface $logger
* @param HashGeneratorInterface $hashGenerator
*/
public function __construct(LoggerInterface $logger, HashGeneratorInterface $hashGenerator)
{
$this->logger = $logger;
$this->hashGenerator = $hashGenerator;
}
public function addInternalMessageBinding(string $name, string $internalMessageClass)
{
if (!class_exists($internalMessageClass)) {
throw new \InvalidArgumentException('Internal message class must be full-namespace version.');
}
if (array_key_exists($name, $this->internalMessagesBinding)) {
return;
}
$this->internalMessagesBinding[$name] = $internalMessageClass;
}
/**
* @param PlainMessage $message
* @return Message
*/
public function createFromBernardBundle(PlainMessage $message): Message
{
$name = $message->getName();
if (!$this->binded($name)) {
throw new \InvalidArgumentException('Message binding don\'t exist');
}
return $this->createFromArray($name, $message->all());
}
public function createNew(string $name, array $data): Message
{
if (!$this->binded($name)) {
throw new \InvalidArgumentException('Message binding don\'t exist');
}
$hash = $this->hashGenerator->generateHash(json_encode($data));
return $this->createFromArray(
$name,
array_merge(
[
'messageId' => $hash,
],
$data
)
);
}
public function createFromArray(string $name, array $data): Message
{
try {
$className = $this->internalMessagesBinding[$name];
$internalMessage = new $className($data);
return $internalMessage;
} catch (\InvalidArgumentException $e) {
$this->logger->error(
sprintf(
'Internal message %s cannot be created from data consumed from queue: %s because of: %s',
$name,
json_encode($data),
$e->getMessage()
)
);
throw new \InvalidArgumentException('Unable to create internal message from consumed data', 0, $e);
}
}
private function binded(string $name): bool
{
if (!array_key_exists($name, $this->internalMessagesBinding)) {
$this->logger->error(
sprintf(
'Message consumed from queue with name: %s don\'t have internal binding. Binded messages: %s',
$name,
json_encode(array_keys($this->internalMessagesBinding))
)
);
return false;
}
return true;
}
}