-
Notifications
You must be signed in to change notification settings - Fork 55
/
QueueInteropTransportFactory.php
97 lines (81 loc) · 3.02 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
<?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\ContainerInterface;
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]>
*/
class QueueInteropTransportFactory implements TransportFactoryInterface
{
private $serializer;
private $debug;
private $container;
public function __construct(SerializerInterface $serializer, ContainerInterface $container, bool $debug = false)
{
$this->serializer = $serializer;
$this->container = $container;
$this->debug = $debug;
}
// 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 0 === strpos($dsn, 'enqueue://');
}
private function parseDsn(string $dsn): array
{
$parsedDsn = parse_url($dsn);
$enqueueContextName = $parsedDsn['host'];
$amqpOptions = array();
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 array(
new AmqpContextManager($context),
$amqpOptions,
);
}
}