|
7 | 7 | use Interop\Queue\PsrMessage; |
8 | 8 | use Interop\Queue\PsrProcessor; |
9 | 9 |
|
10 | | -class RouterProcessor implements PsrProcessor |
| 10 | +final class RouterProcessor implements PsrProcessor |
11 | 11 | { |
12 | 12 | /** |
13 | 13 | * @var DriverInterface |
14 | 14 | */ |
15 | 15 | private $driver; |
16 | 16 |
|
17 | 17 | /** |
18 | | - * @var array |
| 18 | + * @var RouteCollection |
19 | 19 | */ |
20 | | - private $eventRoutes; |
| 20 | + private $routeCollection; |
21 | 21 |
|
22 | | - /** |
23 | | - * @var array |
24 | | - */ |
25 | | - private $commandRoutes; |
26 | | - |
27 | | - /** |
28 | | - * @param DriverInterface $driver |
29 | | - * @param array $eventRoutes |
30 | | - * @param array $commandRoutes |
31 | | - */ |
32 | | - public function __construct(DriverInterface $driver, array $eventRoutes = [], array $commandRoutes = []) |
| 22 | + public function __construct(DriverInterface $driver, RouteCollection $routeCollection) |
33 | 23 | { |
34 | 24 | $this->driver = $driver; |
35 | | - |
36 | | - $this->eventRoutes = $eventRoutes; |
37 | | - $this->commandRoutes = $commandRoutes; |
| 25 | + $this->routeCollection = $routeCollection; |
38 | 26 | } |
39 | 27 |
|
40 | | - /** |
41 | | - * @param string $topicName |
42 | | - * @param string $queueName |
43 | | - * @param string $processorName |
44 | | - */ |
45 | | - public function add($topicName, $queueName, $processorName) |
| 28 | + public function process(PsrMessage $message, PsrContext $context): Result |
46 | 29 | { |
47 | | - if (Config::COMMAND_TOPIC === $topicName) { |
48 | | - $this->commandRoutes[$processorName] = $queueName; |
49 | | - } else { |
50 | | - $this->eventRoutes[$topicName][] = [$processorName, $queueName]; |
| 30 | + $topic = $message->getProperty(Config::PARAMETER_TOPIC_NAME); |
| 31 | + if ($topic) { |
| 32 | + return $this->routeEvent($topic, $message); |
51 | 33 | } |
52 | | - } |
53 | 34 |
|
54 | | - /** |
55 | | - * {@inheritdoc} |
56 | | - */ |
57 | | - public function process(PsrMessage $message, PsrContext $context) |
58 | | - { |
59 | | - $topicName = $message->getProperty(Config::PARAMETER_TOPIC_NAME); |
60 | | - if (false == $topicName) { |
61 | | - return Result::reject(sprintf( |
62 | | - 'Got message without required parameter: "%s"', |
63 | | - Config::PARAMETER_TOPIC_NAME |
64 | | - )); |
| 35 | + $command = $message->getProperty(Config::PARAMETER_COMMAND_NAME); |
| 36 | + if ($command) { |
| 37 | + return $this->routeCommand($command, $message); |
65 | 38 | } |
66 | 39 |
|
67 | | - if (Config::COMMAND_TOPIC === $topicName) { |
68 | | - return $this->routeCommand($message); |
69 | | - } |
70 | | - |
71 | | - return $this->routeEvent($message); |
| 40 | + return Result::reject(sprintf( |
| 41 | + 'Got message without required parameters. Either "%s" or "%s" property should be set', |
| 42 | + Config::PARAMETER_TOPIC_NAME, |
| 43 | + Config::PARAMETER_COMMAND_NAME |
| 44 | + )); |
72 | 45 | } |
73 | 46 |
|
74 | | - /** |
75 | | - * @param PsrMessage $message |
76 | | - * |
77 | | - * @return string|Result |
78 | | - */ |
79 | | - private function routeEvent(PsrMessage $message) |
| 47 | + private function routeEvent(string $topic, PsrMessage $message): Result |
80 | 48 | { |
81 | | - $topicName = $message->getProperty(Config::PARAMETER_TOPIC_NAME); |
| 49 | + $count = 0; |
| 50 | + foreach ($this->routeCollection->topicRoutes($topic) as $route) { |
| 51 | + $processorMessage = clone $message; |
| 52 | + $processorMessage->setProperty(Config::PARAMETER_PROCESSOR_NAME, $route->getProcessor()); |
| 53 | + $processorMessage->setProperty(Config::PARAMETER_PROCESSOR_QUEUE_NAME, $route->getQueue()); |
82 | 54 |
|
83 | | - if (array_key_exists($topicName, $this->eventRoutes)) { |
84 | | - foreach ($this->eventRoutes[$topicName] as $route) { |
85 | | - $processorMessage = clone $message; |
86 | | - $processorMessage->setProperty(Config::PARAMETER_PROCESSOR_NAME, $route[0]); |
87 | | - $processorMessage->setProperty(Config::PARAMETER_PROCESSOR_QUEUE_NAME, $route[1]); |
| 55 | + $this->driver->sendToProcessor($this->driver->createClientMessage($processorMessage)); |
88 | 56 |
|
89 | | - $this->driver->sendToProcessor($this->driver->createClientMessage($processorMessage)); |
90 | | - } |
| 57 | + ++$count; |
91 | 58 | } |
92 | 59 |
|
93 | | - return self::ACK; |
| 60 | + return Result::ack(sprintf('Routed to "%d" event subscribers', $count)); |
94 | 61 | } |
95 | 62 |
|
96 | | - /** |
97 | | - * @param PsrMessage $message |
98 | | - * |
99 | | - * @return string|Result |
100 | | - */ |
101 | | - private function routeCommand(PsrMessage $message) |
| 63 | + private function routeCommand(string $command, PsrMessage $message): Result |
102 | 64 | { |
103 | | - $commandName = $message->getProperty(Config::PARAMETER_COMMAND_NAME); |
104 | | - if (false == $commandName) { |
105 | | - return Result::reject(sprintf( |
106 | | - 'Got message without required parameter: "%s"', |
107 | | - Config::PARAMETER_COMMAND_NAME |
108 | | - )); |
| 65 | + $route = $this->routeCollection->commandRoute($command); |
| 66 | + if (false == $route) { |
| 67 | + throw new \LogicException(sprintf('The command "%s" processor not found', $command)); |
109 | 68 | } |
110 | 69 |
|
111 | | - if (isset($this->commandRoutes[$commandName])) { |
112 | | - $processorMessage = clone $message; |
113 | | - $processorMessage->setProperty(Config::PARAMETER_PROCESSOR_QUEUE_NAME, $this->commandRoutes[$commandName]); |
114 | | - $processorMessage->setProperty(Config::PARAMETER_PROCESSOR_NAME, $commandName); |
| 70 | + $processorMessage = clone $message; |
| 71 | + $processorMessage->setProperty(Config::PARAMETER_PROCESSOR_NAME, $route->getProcessor()); |
| 72 | + $processorMessage->setProperty(Config::PARAMETER_PROCESSOR_QUEUE_NAME, $route->getQueue()); |
115 | 73 |
|
116 | | - $this->driver->sendToProcessor($this->driver->createClientMessage($processorMessage)); |
117 | | - } |
| 74 | + $this->driver->sendToProcessor($this->driver->createClientMessage($processorMessage)); |
118 | 75 |
|
119 | | - return self::ACK; |
| 76 | + return Result::ack('Routed to the command processor'); |
120 | 77 | } |
121 | 78 | } |
0 commit comments