diff --git a/README.md b/README.md index 5bfe08e..0a16f8a 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.4` | `master` | 3.1+ | 7.4 / 8.x | `>=8.2` | +| stable | `^0.3` | `master` | 3.1+ | 7.4 / 8.x | `>=8.2` | ## Development 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" } } } 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/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..a08c109 100644 --- a/src/DI/Utils/BuilderMan.php +++ b/src/DI/Utils/BuilderMan.php @@ -64,6 +64,7 @@ public function getTransportToFailureTransportsServiceMapping(): array $builder = $this->pass->getContainerBuilder(); $transports = $this->getTransports(); + /** @var array $definitions */ $definitions = $builder->findByTag(MessengerExtension::FAILURE_TRANSPORT_TAG); $transportsMapping = []; @@ -253,9 +254,12 @@ private function getServiceDefinitionsByTag(string $tag): array { $builder = $this->pass->getContainerBuilder(); + /** @var array $tags */ + $tags = $builder->findByTag($tag); + $definitions = []; - foreach ($builder->findByTag($tag) as $serviceName => $tagValue) { - $definitions[(string) $tagValue] = $builder->getDefinition($serviceName); + foreach ($tags as $serviceName => $tagValue) { + $definitions[$tagValue] = $builder->getDefinition($serviceName); } return $definitions; @@ -268,9 +272,12 @@ private function getServiceNamesByTag(string $tag): array { $builder = $this->pass->getContainerBuilder(); + /** @var array $tags */ + $tags = $builder->findByTag($tag); + $definitions = []; - foreach ($builder->findByTag($tag) as $serviceName => $tagValue) { - $definitions[(string) $tagValue] = $serviceName; + foreach ($tags as $serviceName => $tagValue) { + $definitions[$tagValue] = $serviceName; } return $definitions; 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); +});