From facaf0b1813a170592fbcf600f27593bc7bd8dd9 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Milan=20Felix=20=C5=A0ulc?= Date: Fri, 2 Oct 2026 23:24:03 +0200 Subject: [PATCH 1/6] Worker: own StopWorkerOnTimeLimitListener (Symfony 8.1 deprecated its listener) --- src/DI/Pass/EventPass.php | 2 +- .../StopWorkerOnTimeLimitListener.php | 57 +++++++++++ tests/Cases/DI/MessengerExtension.worker.phpt | 2 +- .../StopWorkerOnTimeLimitListener.phpt | 97 +++++++++++++++++++ 4 files changed, 156 insertions(+), 2 deletions(-) create mode 100644 src/EventListener/StopWorkerOnTimeLimitListener.php create mode 100644 tests/Cases/EventListener/StopWorkerOnTimeLimitListener.phpt diff --git a/src/DI/Pass/EventPass.php b/src/DI/Pass/EventPass.php index 76419c1..461454a 100644 --- a/src/DI/Pass/EventPass.php +++ b/src/DI/Pass/EventPass.php @@ -4,6 +4,7 @@ use Contributte\Messenger\Container\NetteContainer; use Contributte\Messenger\DI\Utils\BuilderMan; +use Contributte\Messenger\EventListener\StopWorkerOnTimeLimitListener; use Nette\DI\Definitions\ServiceDefinition; use Nette\DI\Definitions\Statement; use Symfony\Component\EventDispatcher\EventDispatcher; @@ -18,7 +19,6 @@ use Symfony\Component\Messenger\EventListener\StopWorkerOnMemoryLimitListener; use Symfony\Component\Messenger\EventListener\StopWorkerOnMessageLimitListener; use Symfony\Component\Messenger\EventListener\StopWorkerOnRestartSignalListener; -use Symfony\Component\Messenger\EventListener\StopWorkerOnTimeLimitListener; class EventPass extends AbstractPass { diff --git a/src/EventListener/StopWorkerOnTimeLimitListener.php b/src/EventListener/StopWorkerOnTimeLimitListener.php new file mode 100644 index 0000000..6c2c246 --- /dev/null +++ b/src/EventListener/StopWorkerOnTimeLimitListener.php @@ -0,0 +1,57 @@ + + */ + public static function getSubscribedEvents(): array + { + return [ + WorkerStartedEvent::class => 'onWorkerStarted', + WorkerRunningEvent::class => 'onWorkerRunning', + ]; + } + + public function onWorkerStarted(): void + { + $this->endTime = microtime(true) + $this->timeLimitInSeconds; + } + + public function onWorkerRunning(WorkerRunningEvent $event): void + { + if ($this->endTime < microtime(true)) { + $event->getWorker()->stop(); + $this->logger?->info('Worker stopped due to time limit of {timeLimit}s exceeded', ['timeLimit' => $this->timeLimitInSeconds]); + } + } + +} diff --git a/tests/Cases/DI/MessengerExtension.worker.phpt b/tests/Cases/DI/MessengerExtension.worker.phpt index d2aad57..b5f3ec7 100644 --- a/tests/Cases/DI/MessengerExtension.worker.phpt +++ b/tests/Cases/DI/MessengerExtension.worker.phpt @@ -3,6 +3,7 @@ namespace Tests\Cases\DI; use Contributte\EventDispatcher\DI\EventDispatcherExtension; +use Contributte\Messenger\EventListener\StopWorkerOnTimeLimitListener; use Contributte\Tester\Toolkit; use Nette\DI\Compiler; use Symfony\Component\EventDispatcher\EventDispatcher; @@ -11,7 +12,6 @@ use Symfony\Component\Messenger\EventListener\StopWorkerOnCustomStopExceptionLis use Symfony\Component\Messenger\EventListener\StopWorkerOnFailureLimitListener; use Symfony\Component\Messenger\EventListener\StopWorkerOnMemoryLimitListener; use Symfony\Component\Messenger\EventListener\StopWorkerOnMessageLimitListener; -use Symfony\Component\Messenger\EventListener\StopWorkerOnTimeLimitListener; use Tester\Assert; use Tests\Toolkit\Container; use Tests\Toolkit\Helpers; diff --git a/tests/Cases/EventListener/StopWorkerOnTimeLimitListener.phpt b/tests/Cases/EventListener/StopWorkerOnTimeLimitListener.phpt new file mode 100644 index 0000000..e92c0af --- /dev/null +++ b/tests/Cases/EventListener/StopWorkerOnTimeLimitListener.phpt @@ -0,0 +1,97 @@ +stopped++; + } + + }; +} + +// Subscribed events +Toolkit::test(function (): void { + Assert::equal([ + WorkerStartedEvent::class => 'onWorkerStarted', + WorkerRunningEvent::class => 'onWorkerRunning', + ], StopWorkerOnTimeLimitListener::getSubscribedEvents()); +}); + +// Invalid time limit +Toolkit::test(function (): void { + Assert::exception( + static fn () => new StopWorkerOnTimeLimitListener(0), + InvalidArgumentException::class, + 'Time limit must be greater than zero.' + ); + + Assert::exception( + static fn () => new StopWorkerOnTimeLimitListener(-1), + InvalidArgumentException::class, + 'Time limit must be greater than zero.' + ); +}); + +// Worker is not stopped before the time limit +Toolkit::test(function (): void { + $worker = createWorker(); + $logger = new BufferLogger(); + $listener = new StopWorkerOnTimeLimitListener(60, $logger); + + $listener->onWorkerStarted(); + $listener->onWorkerRunning(new WorkerRunningEvent($worker, false)); + + Assert::same(0, $worker->stopped); + Assert::count(0, $logger->obtain()); +}); + +// Worker is stopped once the time limit is exceeded +Toolkit::test(function (): void { + $worker = createWorker(); + $logger = new BufferLogger(); + $listener = new StopWorkerOnTimeLimitListener(1, $logger); + + $listener->onWorkerStarted(); + usleep(1100000); + $listener->onWorkerRunning(new WorkerRunningEvent($worker, true)); + + Assert::same(1, $worker->stopped); + Assert::equal([ + [ + 'level' => 'info', + 'message' => 'Worker stopped due to time limit of {timeLimit}s exceeded', + 'context' => ['timeLimit' => 1], + ], + ], $logger->obtain()); +}); + +// Works without logger +Toolkit::test(function (): void { + $worker = createWorker(); + $listener = new StopWorkerOnTimeLimitListener(1); + + $listener->onWorkerStarted(); + usleep(1100000); + $listener->onWorkerRunning(new WorkerRunningEvent($worker, false)); + + Assert::same(1, $worker->stopped); +}); From 183d5d493a6533a0436a772b08a94d97bba40244 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Milan=20Felix=20=C5=A0ulc?= Date: Fri, 2 Oct 2026 23:24:03 +0200 Subject: [PATCH 2/6] Code: validate string tags in BuilderMan, drop redundant priority fallback --- src/DI/Pass/HandlerPass.php | 2 +- src/DI/Pass/RoutingPass.php | 3 +-- src/DI/Utils/BuilderMan.php | 32 +++++++++++++++++++++++--------- 3 files changed, 25 insertions(+), 12 deletions(-) diff --git a/src/DI/Pass/HandlerPass.php b/src/DI/Pass/HandlerPass.php index 564d1d1..5028997 100644 --- a/src/DI/Pass/HandlerPass.php +++ b/src/DI/Pass/HandlerPass.php @@ -146,7 +146,7 @@ private function getAttributesOptions(string $serviceClass, string $serviceName, 'bus' => $attribute->bus ?? $defaultBusName, 'alias' => null, 'method' => $attribute->method ?? self::DEFAULT_METHOD_NAME, - 'priority' => $attribute->priority ?? self::DEFAULT_PRIORITY, + 'priority' => $attribute->priority, 'handles' => $attribute->handles ?? null, 'from_transport' => $attribute->fromTransport ?? null, ]; diff --git a/src/DI/Pass/RoutingPass.php b/src/DI/Pass/RoutingPass.php index c6f028f..c1f98be 100644 --- a/src/DI/Pass/RoutingPass.php +++ b/src/DI/Pass/RoutingPass.php @@ -2,7 +2,6 @@ namespace Contributte\Messenger\DI\Pass; -use Contributte\Messenger\DI\MessengerExtension; use Contributte\Messenger\DI\Utils\BuilderMan; use Contributte\Messenger\Exception\LogicalException; use Nette\DI\Definitions\ServiceDefinition; @@ -29,7 +28,7 @@ public function beforePassCompile(): void { $builder = $this->getContainerBuilder(); $config = $this->getConfig(); - $transports = array_values($builder->findByTag(MessengerExtension::TRANSPORT_TAG)); + $transports = array_keys(BuilderMan::of($this)->getTransports()); // Scan message classes for #[AsMessage] attribute routing $attributeRouting = BuilderMan::of($this)->getAttributeRouting(); diff --git a/src/DI/Utils/BuilderMan.php b/src/DI/Utils/BuilderMan.php index 7996527..88703b3 100644 --- a/src/DI/Utils/BuilderMan.php +++ b/src/DI/Utils/BuilderMan.php @@ -64,7 +64,7 @@ public function getTransportToFailureTransportsServiceMapping(): array $builder = $this->pass->getContainerBuilder(); $transports = $this->getTransports(); - $definitions = $builder->findByTag(MessengerExtension::FAILURE_TRANSPORT_TAG); + $definitions = $this->findStringTags(MessengerExtension::FAILURE_TRANSPORT_TAG); $transportsMapping = []; foreach ($definitions as $serviceName => $failureTransport) { @@ -90,8 +90,7 @@ public function getFailedTransports(): array $builder = $this->pass->getContainerBuilder(); $transports = $this->getTransports(); - /** @var array $definitions */ - $definitions = $builder->findByTag(MessengerExtension::FAILURE_TRANSPORT_TAG); + $definitions = $this->findStringTags(MessengerExtension::FAILURE_TRANSPORT_TAG); $transportsMapping = []; @@ -254,8 +253,8 @@ private function getServiceDefinitionsByTag(string $tag): array $builder = $this->pass->getContainerBuilder(); $definitions = []; - foreach ($builder->findByTag($tag) as $serviceName => $tagValue) { - $definitions[(string) $tagValue] = $builder->getDefinition($serviceName); + foreach ($this->findStringTags($tag) as $serviceName => $tagValue) { + $definitions[$tagValue] = $builder->getDefinition($serviceName); } return $definitions; @@ -266,14 +265,29 @@ private function getServiceDefinitionsByTag(string $tag): array */ private function getServiceNamesByTag(string $tag): array { - $builder = $this->pass->getContainerBuilder(); - $definitions = []; - foreach ($builder->findByTag($tag) as $serviceName => $tagValue) { - $definitions[(string) $tagValue] = $serviceName; + foreach ($this->findStringTags($tag) as $serviceName => $tagValue) { + $definitions[$tagValue] = $serviceName; } return $definitions; } + /** + * @return array service name => tag value + */ + private function findStringTags(string $tag): array + { + $tags = []; + foreach ($this->pass->getContainerBuilder()->findByTag($tag) as $serviceName => $tagValue) { + if (!is_string($tagValue)) { + throw new LogicalException(sprintf('Tag "%s" of service "%s" must be a string, %s given.', $tag, $serviceName, get_debug_type($tagValue))); + } + + $tags[$serviceName] = $tagValue; + } + + return $tags; + } + } From 0aa183fab6a5fde3b8d70e542df7e15a0d69d2dd Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Milan=20Felix=20=C5=A0ulc?= Date: Fri, 2 Oct 2026 23:24:03 +0200 Subject: [PATCH 3/6] Docs: update version table with Symfony support --- README.md | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/README.md b/README.md index 5bfe08e..dd8af58 100644 --- a/README.md +++ b/README.md @@ -32,10 +32,10 @@ For details on how to use this package, check out our [documentation](.docs). ## Version -| State | Version | Branch | Nette | PHP | -|--------|---------|----------|-------|---------| -| dev | `^0.3` | `master` | 3.2+ | `>=8.2` | -| stable | `^0.2` | `master` | 3.2+ | `>=8.2` | +| State | Version | Branch | Nette | Symfony | PHP | +|--------|---------|----------|-------|------------|---------| +| dev | `^0.3` | `master` | 3.1+ | 7.4 / 8.x | `>=8.2` | +| stable | `^0.2` | `master` | 3.1+ | 6.x | `>=8.0` | ## Development From 89d54dbdff485fb6c549a466812027040c37650b Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Milan=20Felix=20=C5=A0ulc?= Date: Fri, 2 Oct 2026 23:32:59 +0200 Subject: [PATCH 4/6] Composer: open v0.4.x --- composer.json | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/composer.json b/composer.json index a758e8e..a5890fb 100644 --- a/composer.json +++ b/composer.json @@ -65,7 +65,7 @@ }, "extra": { "branch-alias": { - "dev-master": "0.3.x-dev" + "dev-master": "0.4.x-dev" } } } From cc27e69e702740a30f8869b0747cec6ac05e48de Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Milan=20Felix=20=C5=A0ulc?= Date: Fri, 2 Oct 2026 23:32:59 +0200 Subject: [PATCH 5/6] Docs: update version table for v0.3.0 --- README.md | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/README.md b/README.md index dd8af58..0a16f8a 100644 --- a/README.md +++ b/README.md @@ -34,8 +34,8 @@ For details on how to use this package, check out our [documentation](.docs). | State | Version | Branch | Nette | Symfony | PHP | |--------|---------|----------|-------|------------|---------| -| dev | `^0.3` | `master` | 3.1+ | 7.4 / 8.x | `>=8.2` | -| stable | `^0.2` | `master` | 3.1+ | 6.x | `>=8.0` | +| dev | `^0.4` | `master` | 3.1+ | 7.4 / 8.x | `>=8.2` | +| stable | `^0.3` | `master` | 3.1+ | 7.4 / 8.x | `>=8.2` | ## Development From 3330325e411d40b65a4ac27b434d643da7f37f52 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Milan=20Felix=20=C5=A0ulc?= Date: Fri, 2 Oct 2026 23:40:16 +0200 Subject: [PATCH 6/6] DI: use findByTag directly in BuilderMan --- src/DI/Utils/BuilderMan.php | 35 ++++++++++++++--------------------- 1 file changed, 14 insertions(+), 21 deletions(-) diff --git a/src/DI/Utils/BuilderMan.php b/src/DI/Utils/BuilderMan.php index 88703b3..a08c109 100644 --- a/src/DI/Utils/BuilderMan.php +++ b/src/DI/Utils/BuilderMan.php @@ -64,7 +64,8 @@ public function getTransportToFailureTransportsServiceMapping(): array $builder = $this->pass->getContainerBuilder(); $transports = $this->getTransports(); - $definitions = $this->findStringTags(MessengerExtension::FAILURE_TRANSPORT_TAG); + /** @var array $definitions */ + $definitions = $builder->findByTag(MessengerExtension::FAILURE_TRANSPORT_TAG); $transportsMapping = []; foreach ($definitions as $serviceName => $failureTransport) { @@ -90,7 +91,8 @@ public function getFailedTransports(): array $builder = $this->pass->getContainerBuilder(); $transports = $this->getTransports(); - $definitions = $this->findStringTags(MessengerExtension::FAILURE_TRANSPORT_TAG); + /** @var array $definitions */ + $definitions = $builder->findByTag(MessengerExtension::FAILURE_TRANSPORT_TAG); $transportsMapping = []; @@ -252,8 +254,11 @@ private function getServiceDefinitionsByTag(string $tag): array { $builder = $this->pass->getContainerBuilder(); + /** @var array $tags */ + $tags = $builder->findByTag($tag); + $definitions = []; - foreach ($this->findStringTags($tag) as $serviceName => $tagValue) { + foreach ($tags as $serviceName => $tagValue) { $definitions[$tagValue] = $builder->getDefinition($serviceName); } @@ -265,29 +270,17 @@ private function getServiceDefinitionsByTag(string $tag): array */ private function getServiceNamesByTag(string $tag): array { + $builder = $this->pass->getContainerBuilder(); + + /** @var array $tags */ + $tags = $builder->findByTag($tag); + $definitions = []; - foreach ($this->findStringTags($tag) as $serviceName => $tagValue) { + foreach ($tags as $serviceName => $tagValue) { $definitions[$tagValue] = $serviceName; } return $definitions; } - /** - * @return array service name => tag value - */ - private function findStringTags(string $tag): array - { - $tags = []; - foreach ($this->pass->getContainerBuilder()->findByTag($tag) as $serviceName => $tagValue) { - if (!is_string($tagValue)) { - throw new LogicalException(sprintf('Tag "%s" of service "%s" must be a string, %s given.', $tag, $serviceName, get_debug_type($tagValue))); - } - - $tags[$serviceName] = $tagValue; - } - - return $tags; - } - }