Skip to content

Commit 94b23e6

Browse files
committed
Update amqp package. Add listeners to configuration. Remove middlewares/event_handlers configurations.
1 parent 42880fc commit 94b23e6

19 files changed

Lines changed: 635 additions & 328 deletions

.github/workflows/tests.yaml

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -57,6 +57,9 @@ jobs:
5757
- php: 8.4
5858
symfony: '~7.0'
5959

60+
- php: 8.5
61+
symfony: '~7.0'
62+
6063

6164
steps:
6265
- name: Checkout

Dockerfile

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -12,6 +12,7 @@ RUN \
1212
RUN \
1313
apt-get install -y --no-install-recommends \
1414
yes | pecl install xdebug && \
15+
docker-php-ext-install pcntl && \
1516
docker-php-ext-enable xdebug
1617

1718
# Install composer

composer.json

Lines changed: 7 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -20,18 +20,20 @@
2020

2121
"require": {
2222
"php": "^8.2",
23-
"fivelab/amqp": "~2.2.0",
24-
"symfony/dependency-injection": "~6.4 | ~7.0",
25-
"symfony/framework-bundle": "~6.4 | ~7.0"
23+
"fivelab/amqp": "dev-master",
24+
"symfony/dependency-injection": "~6.4 | ~7.0 | ~8.0",
25+
"symfony/framework-bundle": "~6.4 | ~7.0 | ~8.0"
2626
},
2727

2828
"require-dev": {
2929
"phpunit/phpunit": "~11.5",
3030
"phpmetrics/phpmetrics": "^3.0@rc",
3131
"escapestudios/symfony2-coding-standard": "~3.5",
3232
"matthiasnoback/symfony-dependency-injection-test": "~6.0",
33-
"symfony/console": "~6.4 | ~7.0",
34-
"symfony/expression-language": "~6.4 | ~7.0",
33+
"symfony/console": "~6.4 | ~7.0 | 8.0",
34+
"symfony/expression-language": "~6.4 | ~7.0 | 8.0",
35+
"doctrine/persistence": "~4.0",
36+
"doctrine/dbal": "~4.0",
3537
"fivelab/ci-rules": "dev-master"
3638
},
3739

src/DependencyInjection/AmqpExtension.php

Lines changed: 29 additions & 36 deletions
Original file line numberDiff line numberDiff line change
@@ -14,6 +14,8 @@
1414
namespace FiveLab\Bundle\AmqpBundle\DependencyInjection;
1515

1616
use FiveLab\Bundle\AmqpBundle\Factory\DriverFactory;
17+
use FiveLab\Bundle\AmqpBundle\Listener\PingDbalConnectionsListener;
18+
use FiveLab\Bundle\AmqpBundle\Listener\ReleaseMemoryListener;
1719
use FiveLab\Component\Amqp\Channel\ChannelFactoryInterface;
1820
use FiveLab\Component\Amqp\Connection\ConnectionFactoryInterface;
1921
use FiveLab\Component\Amqp\Connection\Dsn;
@@ -33,7 +35,6 @@
3335
use FiveLab\Component\Amqp\Queue\Definition\Arguments\SingleActiveCustomerArgument;
3436
use FiveLab\Component\Amqp\Queue\QueueFactoryInterface;
3537
use Symfony\Component\Config\FileLocator;
36-
use Symfony\Component\DependencyInjection\Argument\ServiceClosureArgument;
3738
use Symfony\Component\DependencyInjection\ChildDefinition;
3839
use Symfony\Component\DependencyInjection\Compiler\ServiceLocatorTagPass;
3940
use Symfony\Component\DependencyInjection\ContainerBuilder;
@@ -71,14 +72,16 @@ public function load(array $configs, ContainerBuilder $container): void
7172
$loader->load('definitions.php');
7273
$loader->load('consumers.php');
7374
$loader->load('publishers.php');
75+
$loader->load('listeners.php');
7476
$loader->load('services.php');
7577

78+
$this->configureListeners($container, $config['listeners']);
7679
$this->configureConnections($container, $config['connections']);
7780
$this->configureChannels($container, $config['channels']);
7881
$this->configureExchanges($container, $config['exchanges']);
7982
$this->configureQueues($container, $config['queues'], $config['queue_default_arguments'] ?? []);
8083
$this->configurePublishers($container, $config['publishers'], $config['publisher_middleware']);
81-
$this->configureConsumers($container, $config['consumers'], $config['consumer_middleware'], $config['consumer_event_handlers'], $config['consumer_defaults']);
84+
$this->configureConsumers($container, $config['consumers'], $config['consumer_defaults']);
8285

8386
if ($this->isConfigEnabled($container, $config['round_robin'])) {
8487
$loader->load('round-robin.php');
@@ -93,8 +96,6 @@ public function load(array $configs, ContainerBuilder $container): void
9396
$container,
9497
$config['delay'],
9598
$config['publisher_middleware'],
96-
$config['consumer_middleware'],
97-
$config['consumer_event_handlers'],
9899
$config['queue_default_arguments'] ?? [],
99100
$config['consumer_defaults'] ?? []
100101
);
@@ -132,6 +133,21 @@ private function setDefaultStrategy(array $config): array
132133
return $config;
133134
}
134135

136+
private function configureListeners(ContainerBuilder $container, array $listeners): void
137+
{
138+
if (null !== $listeners['release_memory']) {
139+
$container->getDefinition(ReleaseMemoryListener::class)
140+
->setAbstract(false)
141+
->replaceArgument(1, $listeners['release_memory']);
142+
}
143+
144+
if (null !== $listeners['ping_dbal_connections']) {
145+
$container->getDefinition(PingDbalConnectionsListener::class)
146+
->setAbstract(false)
147+
->replaceArgument(1, $listeners['ping_dbal_connections']);
148+
}
149+
}
150+
135151
private function configureConnections(ContainerBuilder $container, array $connections): void
136152
{
137153
$registryDef = $container->getDefinition('fivelab.amqp.connection_factory_registry');
@@ -517,7 +533,7 @@ private function configureQueues(ContainerBuilder $container, array $queues, arr
517533
$container->setParameter('fivelab.amqp.queue_factories', \array_keys($this->queueFactories));
518534
}
519535

520-
private function configureConsumers(ContainerBuilder $container, array $consumers, array $globalMiddlewares, array $eventHandlers, array $defaults): void
536+
private function configureConsumers(ContainerBuilder $container, array $consumers, array $defaults): void
521537
{
522538
$consumerRegistryDef = $container->getDefinition('fivelab.amqp.consumer_registry');
523539
$checkConsumerRegistryDef = $container->getDefinition('fivelab.amqp.consumer_checker_registry');
@@ -574,19 +590,6 @@ private function configureConsumers(ContainerBuilder $container, array $consumer
574590
$queueFactoryServiceId = $this->queueFactories[$consumer['queue']];
575591
}
576592

577-
// Configure middleware for consumer
578-
$consumerMiddlewares = \array_merge($globalMiddlewares, $consumer['middleware']);
579-
580-
$middlewareList = \array_map(static function (string $serviceId) {
581-
return new Reference($serviceId);
582-
}, $consumerMiddlewares);
583-
584-
$middlewareServiceId = \sprintf('fivelab.amqp.consumer.%s.middlewares', $key);
585-
$middlewareServiceDef = new ChildDefinition('fivelab.amqp.consumer.middlewares.abstract');
586-
$middlewareServiceDef->setArguments($middlewareList);
587-
588-
$container->setDefinition($middlewareServiceId, $middlewareServiceDef);
589-
590593
// Configure message handler for consumer
591594
$messageHandlersServiceId = \sprintf('fivelab.amqp.consumer.%s.message_handler', $key);
592595
$messageHandlersServiceDef = new ChildDefinition('fivelab.amqp.consumer.message_handler.abstract');
@@ -642,9 +645,8 @@ private function configureConsumers(ContainerBuilder $container, array $consumer
642645
$consumerServiceDef
643646
->replaceArgument(0, new Reference($queueFactoryServiceId))
644647
->replaceArgument(1, new Reference($messageHandlersServiceId))
645-
->replaceArgument(2, new Reference($middlewareServiceId))
646-
->replaceArgument(3, new Reference($consumerConfigurationServiceId))
647-
->replaceArgument(4, new Reference($strategyServiceId));
648+
->replaceArgument(2, new Reference($consumerConfigurationServiceId))
649+
->replaceArgument(3, new Reference($strategyServiceId));
648650
} elseif ('spool' === $consumer['mode']) {
649651
// Configure spool consumer
650652
$consumerConfigurationServiceId = \sprintf('fivelab.amqp.consumer.%s.configuration', $key);
@@ -662,9 +664,8 @@ private function configureConsumers(ContainerBuilder $container, array $consumer
662664
$consumerServiceDef
663665
->replaceArgument(0, new Reference($queueFactoryServiceId))
664666
->replaceArgument(1, new Reference($messageHandlersServiceId))
665-
->replaceArgument(2, new Reference($middlewareServiceId))
666-
->replaceArgument(3, new Reference($consumerConfigurationServiceId))
667-
->replaceArgument(4, new Reference($strategyServiceId));
667+
->replaceArgument(2, new Reference($consumerConfigurationServiceId))
668+
->replaceArgument(3, new Reference($strategyServiceId));
668669
} elseif ('loop' === $consumer['mode']) {
669670
$consumerConfigurationServiceId = \sprintf('fivelab.amqp.consumer.%s.configuration', $key);
670671
$consumerConfigurationServiceDef = new ChildDefinition('fivelab.amqp.consumer_loop.configuration.abstract');
@@ -680,23 +681,15 @@ private function configureConsumers(ContainerBuilder $container, array $consumer
680681
$consumerServiceDef
681682
->replaceArgument(0, new Reference($queueFactoryServiceId))
682683
->replaceArgument(1, new Reference($messageHandlersServiceId))
683-
->replaceArgument(2, new Reference($middlewareServiceId))
684-
->replaceArgument(3, new Reference($consumerConfigurationServiceId))
685-
->replaceArgument(4, new Reference($strategyServiceId));
684+
->replaceArgument(2, new Reference($consumerConfigurationServiceId))
685+
->replaceArgument(3, new Reference($strategyServiceId));
686686
} else {
687687
throw new \InvalidArgumentException(\sprintf(
688688
'Unknown mode "%s".',
689689
$consumer['mode']
690690
));
691691
}
692692

693-
foreach ($eventHandlers as $eventHandler) {
694-
$consumerServiceDef->addMethodCall(
695-
'addEventHandler',
696-
[new ServiceClosureArgument(new Reference($eventHandler)), true]
697-
);
698-
}
699-
700693
$container->setDefinition($consumerConfigurationServiceId, $consumerConfigurationServiceDef);
701694
$container->setDefinition($consumerServiceId, $consumerServiceDef);
702695

@@ -830,7 +823,7 @@ private function configureRoundRobin(ContainerBuilder $container, array $config)
830823
->replaceArgument(2, \array_keys($this->consumers));
831824
}
832825

833-
private function configureDelay(ContainerBuilder $container, array $config, array $globalPublisherMiddlewares, array $globalConsumerMiddlewares, array $consumerEventHandlers, array $queueDefaultArguments, array $consumerDefaults): void
826+
private function configureDelay(ContainerBuilder $container, array $config, array $globalPublisherMiddlewares, array $queueDefaultArguments, array $consumerDefaults): void
834827
{
835828
// Configure exchange
836829
$this->configureExchanges($container, [
@@ -949,7 +942,7 @@ private function configureDelay(ContainerBuilder $container, array $config, arra
949942
'idle_timeout' => 100000,
950943
],
951944
],
952-
], $globalConsumerMiddlewares, $consumerEventHandlers, $consumerDefaults);
945+
], $consumerDefaults);
953946
}
954947

955948
private function createArgumentDefinition(ContainerBuilder $container, string $serviceId, string $class, mixed ...$values): Reference

src/DependencyInjection/Configuration.php

Lines changed: 30 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,7 @@ public function getConfigTreeBuilder(): TreeBuilder
2828

2929
$rootNode
3030
->children()
31+
->append($this->getListenersDefinition())
3132
->append($this->getDelayDefinition())
3233
->append($this->getRoundRobinDefinition())
3334
->append($this->getConnectionsNodeDefinition())
@@ -37,15 +38,42 @@ public function getConfigTreeBuilder(): TreeBuilder
3738
->append($this->getQueueArgumentsNodeDefinition('queue_default_arguments'))
3839
->append($this->getConsumersNodeDefinition())
3940
->append($this->getConsumerDefaults())
40-
->append($this->getConsumerEventHandlersNodeDefinition())
4141
->append($this->getPublishersNodeDefinition())
42-
->append($this->getMiddlewareNodeDefinition('consumer_'))
4342
->append($this->getMiddlewareNodeDefinition('publisher_'))
4443
->end();
4544

4645
return $treeBuilder;
4746
}
4847

48+
private function getListenersDefinition(): ArrayNodeDefinition
49+
{
50+
$node = new ArrayNodeDefinition('listeners');
51+
52+
$node
53+
->addDefaultsIfNotSet()
54+
->children()
55+
->scalarNode('release_memory')
56+
->defaultValue(null)
57+
->info('Enable release memory listener (true - clear before handle, false - after handle, null - disable listener).')
58+
->validate()
59+
->ifFalse(static fn (mixed $value) => null === $value || \is_bool($value))
60+
->thenInvalid('Invalid value for "release_memory". Must be bool or null.')
61+
->end()
62+
->end()
63+
64+
->scalarNode('ping_dbal_connections')
65+
->defaultValue(null)
66+
->info('Enable ping DBAL connections listeners (number - seconds for ping interval, null - disable listener)')
67+
->validate()
68+
->ifFalse(static fn (mixed $value) => null === $value || \is_int($value) || $value < 1)
69+
->thenInvalid('Invalid value for "ping_dbal_connections". Must be integer or null.')
70+
->end()
71+
->end()
72+
->end();
73+
74+
return $node;
75+
}
76+
4977
private function getDelayDefinition(): NodeDefinition
5078
{
5179
$node = new ArrayNodeDefinition('delay');
@@ -407,8 +435,6 @@ private function getConsumersNodeDefinition(): NodeDefinition
407435
->end()
408436
->end()
409437

410-
->append($this->getMiddlewareNodeDefinition())
411-
412438
->arrayNode('options')
413439
->addDefaultsIfNotSet()
414440
->children()
@@ -662,18 +688,6 @@ private function getMiddlewareNodeDefinition(string $nodePrefix = ''): NodeDefin
662688
return $node;
663689
}
664690

665-
private function getConsumerEventHandlersNodeDefinition(): NodeDefinition
666-
{
667-
$node = new ArrayNodeDefinition('consumer_event_handlers');
668-
669-
$node
670-
->defaultValue([])
671-
->prototype('scalar')
672-
->end();
673-
674-
return $node;
675-
}
676-
677691
private function getStrategyNodeDefinition(?string $defaultValue = null): NodeDefinition
678692
{
679693
return (new ScalarNodeDefinition('strategy'))
Lines changed: 80 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,80 @@
1+
<?php
2+
3+
/*
4+
* This file is part of the FiveLab AmqpBundle package
5+
*
6+
* (c) FiveLab
7+
*
8+
* For the full copyright and license information, please view the LICENSE
9+
* file that was distributed with this source code
10+
*/
11+
12+
declare(strict_types = 1);
13+
14+
namespace FiveLab\Bundle\AmqpBundle\Listener;
15+
16+
use Doctrine\DBAL\Connection;
17+
use Doctrine\Persistence\ConnectionRegistry;
18+
use FiveLab\Component\Amqp\Event\ConsumerTickEvent;
19+
use Symfony\Component\Console\ConsoleEvents;
20+
use Symfony\Component\Console\Event\ConsoleSignalEvent;
21+
use Symfony\Component\Console\Output\OutputInterface;
22+
use Symfony\Component\EventDispatcher\EventSubscriberInterface;
23+
24+
readonly class PingDbalConnectionsListener implements EventSubscriberInterface
25+
{
26+
private \ArrayObject $options;
27+
28+
public function __construct(private ConnectionRegistry $registry, private int $interval, private array $signals = [\SIGALRM])
29+
{
30+
$this->options = new \ArrayObject([
31+
'ping' => false,
32+
'last_ping' => \time(),
33+
'output' => null,
34+
]);
35+
}
36+
37+
public static function getSubscribedEvents(): array
38+
{
39+
return [
40+
ConsoleEvents::SIGNAL => ['onConsoleSignal', 0],
41+
ConsumerTickEvent::class => ['onConsumerTick', 0],
42+
];
43+
}
44+
45+
public function onConsoleSignal(ConsoleSignalEvent $event): void
46+
{
47+
if (!\in_array($event->getHandlingSignal(), $this->signals, true)) {
48+
return;
49+
}
50+
51+
$nextPing = $this->options['last_ping'] + $this->interval;
52+
53+
if ($nextPing < \time()) {
54+
$this->options->offsetSet('ping', true);
55+
$this->options->offsetSet('output', $event->getOutput());
56+
}
57+
}
58+
59+
public function onConsumerTick(): void
60+
{
61+
if (!$this->options['ping']) {
62+
return;
63+
}
64+
65+
$this->options->offsetSet('ping', false);
66+
$this->options->offsetSet('last_ping', \time());
67+
68+
/** @var Connection $connection */
69+
foreach ($this->registry->getConnections() as $key => $connection) {
70+
if ($connection->isConnected()) {
71+
$connection->executeQuery('SELECT 1');
72+
73+
$this->options['output']->writeln(\sprintf(
74+
'Ping <comment>%s</comment> database connection.',
75+
$key
76+
), OutputInterface::VERBOSITY_DEBUG);
77+
}
78+
}
79+
}
80+
}

0 commit comments

Comments
 (0)