Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 4 additions & 4 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
2 changes: 1 addition & 1 deletion composer.json
Original file line number Diff line number Diff line change
Expand Up @@ -65,7 +65,7 @@
},
"extra": {
"branch-alias": {
"dev-master": "0.3.x-dev"
"dev-master": "0.4.x-dev"
}
}
}
2 changes: 1 addition & 1 deletion src/DI/Pass/EventPass.php
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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
{
Expand Down
2 changes: 1 addition & 1 deletion src/DI/Pass/HandlerPass.php
Original file line number Diff line number Diff line change
Expand Up @@ -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,
];
Expand Down
3 changes: 1 addition & 2 deletions src/DI/Pass/RoutingPass.php
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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();
Expand Down
15 changes: 11 additions & 4 deletions src/DI/Utils/BuilderMan.php
Original file line number Diff line number Diff line change
Expand Up @@ -64,6 +64,7 @@ public function getTransportToFailureTransportsServiceMapping(): array
$builder = $this->pass->getContainerBuilder();

$transports = $this->getTransports();
/** @var array<string, string> $definitions */
$definitions = $builder->findByTag(MessengerExtension::FAILURE_TRANSPORT_TAG);

$transportsMapping = [];
Expand Down Expand Up @@ -253,9 +254,12 @@ private function getServiceDefinitionsByTag(string $tag): array
{
$builder = $this->pass->getContainerBuilder();

/** @var array<string, string> $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;
Expand All @@ -268,9 +272,12 @@ private function getServiceNamesByTag(string $tag): array
{
$builder = $this->pass->getContainerBuilder();

/** @var array<string, string> $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;
Expand Down
57 changes: 57 additions & 0 deletions src/EventListener/StopWorkerOnTimeLimitListener.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,57 @@
<?php declare(strict_types = 1);

namespace Contributte\Messenger\EventListener;

use Psr\Log\LoggerInterface;
use Symfony\Component\EventDispatcher\EventSubscriberInterface;
use Symfony\Component\Messenger\Event\WorkerRunningEvent;
use Symfony\Component\Messenger\Event\WorkerStartedEvent;
use Symfony\Component\Messenger\Exception\InvalidArgumentException;

/**
* Stops the worker after the configured time limit (messenger.worker.timeLimit).
*
* Symfony deprecated its own StopWorkerOnTimeLimitListener in 8.1 in favour of
* the "time_limit" worker option, which can only be passed per consume command.
* This listener keeps the global config option working on Symfony 7.4 and 8.x.
*/
class StopWorkerOnTimeLimitListener implements EventSubscriberInterface
{

private float $endTime = 0;

public function __construct(
private int $timeLimitInSeconds,
private ?LoggerInterface $logger = null,
)
{
if ($timeLimitInSeconds <= 0) {
throw new InvalidArgumentException('Time limit must be greater than zero.');
}
}

/**
* @return array<class-string, string>
*/
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]);
}
}

}
2 changes: 1 addition & 1 deletion tests/Cases/DI/MessengerExtension.worker.phpt
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;
Expand Down
97 changes: 97 additions & 0 deletions tests/Cases/EventListener/StopWorkerOnTimeLimitListener.phpt
Original file line number Diff line number Diff line change
@@ -0,0 +1,97 @@
<?php declare(strict_types = 1);

namespace Tests\Cases\EventListener;

use Contributte\Messenger\EventListener\StopWorkerOnTimeLimitListener;
use Contributte\Messenger\Logger\BufferLogger;
use Contributte\Tester\Toolkit;
use Symfony\Component\Messenger\Event\WorkerRunningEvent;
use Symfony\Component\Messenger\Event\WorkerStartedEvent;
use Symfony\Component\Messenger\Exception\InvalidArgumentException;
use Symfony\Component\Messenger\MessageBus;
use Symfony\Component\Messenger\Worker;
use Tester\Assert;

require_once __DIR__ . '/../../bootstrap.php';

function createWorker(): Worker
{
return new class ([], new MessageBus()) extends Worker {

public int $stopped = 0;

public function stop(): void
{
$this->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);
});
Loading