This repository has been archived by the owner on Jul 15, 2021. It is now read-only.
-
Notifications
You must be signed in to change notification settings - Fork 2
/
AbstractPublisher.php
76 lines (59 loc) · 2.63 KB
/
AbstractPublisher.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
declare(strict_types=1);
namespace Wakeapp\Bundle\RabbitQueueBundle\Publisher;
use PhpAmqpLib\Message\AMQPMessage;
use PhpAmqpLib\Wire\AMQPTable;
use Wakeapp\Bundle\RabbitQueueBundle\Client\RabbitMqClient;
use Wakeapp\Bundle\RabbitQueueBundle\Definition\DefinitionInterface;
use Wakeapp\Bundle\RabbitQueueBundle\Enum\QueueTypeEnum;
use Wakeapp\Bundle\RabbitQueueBundle\Registry\HydratorRegistry;
abstract class AbstractPublisher implements PublisherInterface
{
public const QUEUE_TYPE = QueueTypeEnum::FIFO;
private const DEFAULT_NAME = '';
protected HydratorRegistry $hydratorRegistry;
protected string $hydratorName;
protected RabbitMqClient $client;
public function __construct(RabbitMqClient $client, HydratorRegistry $hydratorRegistry, string $hydratorName)
{
$this->hydratorRegistry = $hydratorRegistry;
$this->hydratorName = $hydratorName;
$this->client = $client;
}
abstract protected function prepareOptions(DefinitionInterface $definition, array $options): array;
public function publish(DefinitionInterface $definition, string $dataString, array $options = []): void
{
$exchangeName = $this->getDefinitionExchangeName($definition);
$queueName = $this->getDefinitionQueueName($definition);
$message = new AMQPMessage($dataString, [
'delivery_mode' => AMQPMessage::DELIVERY_MODE_PERSISTENT,
'content_type' => $this->hydratorRegistry->getHydrator($this->hydratorName)::getKey(),
]);
$amqpTableOptions = $this->prepareOptions($definition, $options);
if (!empty($amqpTableOptions)) {
$message->set('application_headers', new AMQPTable($amqpTableOptions));
}
$this->client->publish($message, $exchangeName, $queueName);
}
abstract public static function getQueueType(): string;
protected function getDefinitionExchangeName(DefinitionInterface $definition): string
{
if ($definition->getQueueType() === (QueueTypeEnum::FIFO | QueueTypeEnum::DEDUPLICATE)) {
return self::DEFAULT_NAME;
}
return $definition->getQueueType() === QueueTypeEnum::FIFO
? self::DEFAULT_NAME
: $definition->getEntryPointName()
;
}
protected function getDefinitionQueueName(DefinitionInterface $definition): string
{
if ($definition->getQueueType() === (QueueTypeEnum::FIFO | QueueTypeEnum::DEDUPLICATE)) {
return $definition::getQueueName();
}
return $definition->getQueueType() === QueueTypeEnum::FIFO
? $definition::getQueueName()
: self::DEFAULT_NAME
;
}
}