forked from sroze/messenger-enqueue-transport
-
Notifications
You must be signed in to change notification settings - Fork 0
/
QueueInteropTransportFactory.php
109 lines (94 loc) · 3.24 KB
/
QueueInteropTransportFactory.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
<?php
/*
* This file is part of the Symfony package.
*
* (c) Fabien Potencier <[email protected]>
*
* For the full copyright and license information, please view the LICENSE
* file that was distributed with this source code.
*/
namespace Enqueue\MessengerAdapter;
use Interop\Queue\Context;
use Psr\Container\ContainerExceptionInterface;
use Psr\Container\ContainerInterface;
use Psr\Container\NotFoundExceptionInterface;
use RuntimeException;
use Symfony\Component\Messenger\Transport\Serialization\SerializerInterface;
use Symfony\Component\Messenger\Transport\TransportFactoryInterface;
use Symfony\Component\Messenger\Transport\TransportInterface;
/**
* Symfony Messenger transport factory.
*
* @author Samuel Roze <[email protected]>
*/
readonly class QueueInteropTransportFactory implements TransportFactoryInterface
{
public function __construct(
private SerializerInterface $serializer,
private ContainerInterface $container,
private bool $debug = false,
) {
}
// BC layer for Symfony 4.1 beta1
public function createReceiver(string $dsn, array $options): TransportInterface
{
return $this->createTransport($dsn, $options);
}
// BC layer for Symfony 4.1 beta1
public function createSender(string $dsn, array $options): TransportInterface
{
return $this->createTransport($dsn, $options);
}
public function createTransport(
string $dsn,
array $options,
SerializerInterface $serializer = null,
): TransportInterface {
[$contextManager, $dsnOptions] = $this->parseDsn($dsn);
$options = array_merge($dsnOptions, $options);
return new QueueInteropTransport(
$serializer ?? $this->serializer,
$contextManager,
$options,
$this->debug,
);
}
public function supports(string $dsn, array $options): bool
{
return str_starts_with($dsn, 'enqueue://');
}
/**
* @throws ContainerExceptionInterface
* @throws NotFoundExceptionInterface
*/
private function parseDsn(string $dsn): array
{
$parsedDsn = parse_url($dsn);
$enqueueContextName = $parsedDsn['host'];
$amqpOptions = [];
if (isset($parsedDsn['query'])) {
parse_str($parsedDsn['query'], $parsedQuery);
$parsedQuery = array_map(function ($e) {
return is_numeric($e) ? (int)$e : $e;
}, $parsedQuery);
$amqpOptions = array_replace_recursive($amqpOptions, $parsedQuery);
}
if (!$this->container->has($contextService = 'enqueue.transport.' . $enqueueContextName . '.context')) {
throw new RuntimeException(
sprintf(
'Can\'t find Enqueue\'s transport named "%s": Service "%s" is not found.',
$enqueueContextName,
$contextService,
),
);
}
$context = $this->container->get($contextService);
if (!$context instanceof Context) {
throw new RuntimeException(sprintf('Service "%s" not instanceof context', $contextService));
}
return [
new AmqpContextManager($context),
$amqpOptions,
];
}
}