-
Notifications
You must be signed in to change notification settings - Fork 7
/
AbstractTransactionalProducer.php
175 lines (146 loc) · 3.46 KB
/
AbstractTransactionalProducer.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
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
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
<?php
namespace Skrz\Bundle\BunnyBundle;
use Bunny\Exception\BunnyException;
use Skrz\Meta\JSON\JsonMetaInterface;
use Skrz\Meta\MetaInterface;
use Skrz\Meta\Protobuf\ProtobufMetaInterface;
class AbstractTransactionalProducer
{
/** @var string */
private $exchange;
/** @var string */
private $routingKey;
/** @var boolean */
private $mandatory;
/** @var boolean */
private $immediate;
/** @var string */
private $metaClassName;
/** @var object */
private $meta;
/** @var string */
private $beforeMethod;
/** @var string */
private $contentType;
/** @var boolean */
private $autoCommit = false;
/** @var BunnyManager */
protected $manager;
public function __construct(
$exchange,
$routingKey,
$mandatory,
$immediate,
$metaClassName,
$beforeMethod,
$contentType,
BunnyManager $manager
) {
$this->exchange = $exchange;
$this->routingKey = $routingKey;
$this->mandatory = $mandatory;
$this->immediate = $immediate;
$this->metaClassName = $metaClassName;
$this->beforeMethod = $beforeMethod;
$this->contentType = $contentType;
$this->manager = $manager;
}
public function createMeta()
{
if ($this->metaClassName) {
/** @var MetaInterface $metaClassName */
$metaClassName = $this->metaClassName;
return $metaClassName::getInstance();
} else {
return null;
}
}
public function getMeta()
{
if ($this->meta === null) {
$this->meta = $this->createMeta();
}
return $this->meta;
}
/**
* @param object $message
* @param string $routingKey
* @param array $headers
* @throws BunnyException
*/
public function publish($message, $routingKey = null, array $headers = [])
{
if (!$this->getMeta()) {
throw new BunnyException("Could not create meta class {$this->metaClassName}.");
}
if (is_string($message)) {
$message = $this->meta->fromJson($message);
}
if ($this->beforeMethod) {
$this->{$this->beforeMethod}($message, $this->manager->getTransactionalChannel());
}
switch ($this->contentType) {
case ContentTypes::APPLICATION_JSON:
if ($this->meta instanceof JsonMetaInterface) {
$message = $this->meta->toJson($message);
} else {
throw new BunnyException("Cannot serialize message to JSON.");
}
break;
case ContentTypes::APPLICATION_PROTOBUF:
if ($this->meta instanceof ProtobufMetaInterface) {
$message = $this->meta->toProtobuf($message);
} else {
throw new BunnyException("Cannot serialize message to Protobuf.");
}
break;
default:
throw new BunnyException("Unhandled content type '{$this->contentType}'.");
}
if ($routingKey === null) {
$routingKey = $this->routingKey;
}
$headers["content-type"] = $this->contentType;
$this->manager->getTransactionalChannel()->publish(
$message,
$headers,
$this->exchange,
$routingKey,
$this->mandatory,
$this->immediate
);
if ($this->autoCommit) {
$this->commit();
}
}
/**
* turn on/off automatic commit
* @param bool $bool
*/
public function setAutoCommit($bool = true)
{
$this->autoCommit = $bool;
}
/**
* commit messages
*/
public function commit()
{
try {
$this->manager->getTransactionalChannel()->txCommit();
} catch (\Exception $e) {
throw new BunnyException("Cannot commit message.");
}
}
/**
* rollback messages
*/
public function rollback()
{
try {
$this->manager->getTransactionalChannel()->txRollback();
} catch (\Exception $e) {
throw new BunnyException("Cannot rollback message.");
}
}
}