-
Notifications
You must be signed in to change notification settings - Fork 439
/
Copy pathSpoolProducer.php
65 lines (50 loc) · 1.42 KB
/
SpoolProducer.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
<?php
namespace Enqueue\Client;
use Enqueue\Rpc\Promise;
class SpoolProducer implements ProducerInterface
{
/**
* @var ProducerInterface
*/
private $realProducer;
/**
* @var array
*/
private $events;
/**
* @var array
*/
private $commands;
public function __construct(ProducerInterface $realProducer)
{
$this->realProducer = $realProducer;
$this->events = new \SplQueue();
$this->commands = new \SplQueue();
}
public function sendCommand(string $command, $message, bool $needReply = false): ?Promise
{
if ($needReply) {
return $this->realProducer->sendCommand($command, $message, $needReply);
}
$this->commands->enqueue([$command, $message]);
return null;
}
public function sendEvent(string $topic, $message): void
{
$this->events->enqueue([$topic, $message]);
}
/**
* When it is called it sends all previously queued messages.
*/
public function flush(): void
{
while (false == $this->events->isEmpty()) {
list($topic, $message) = $this->events->dequeue();
$this->realProducer->sendEvent($topic, $message);
}
while (false == $this->commands->isEmpty()) {
list($command, $message) = $this->commands->dequeue();
$this->realProducer->sendCommand($command, $message);
}
}
}