From de7212c9a754175ba6393ee36abc0184dfedf96a Mon Sep 17 00:00:00 2001 From: David Badura Date: Fri, 2 Oct 2026 17:20:39 +0200 Subject: [PATCH] Dispatch WorkerStoppedEvent when the job throws If the job threw an exception, run() left without dispatching WorkerStoppedEvent, so cleanup listeners never ran in exactly the case they are most needed. The event is now always dispatched and carries the exception, which is rethrown afterwards. --- docs/events.md | 16 +++++++ docs/getting-started.md | 3 ++ src/DefaultWorker.php | 51 ++++++++++++---------- src/Event/WorkerStoppedEvent.php | 2 + tests/Unit/DefaultWorkerTest.php | 73 ++++++++++++++++++++++++++++++++ 5 files changed, 123 insertions(+), 22 deletions(-) diff --git a/docs/events.md b/docs/events.md index fbee5d7..44d289c 100644 --- a/docs/events.md +++ b/docs/events.md @@ -9,6 +9,22 @@ each carrying the worker instance: | `WorkerRunningEvent` | after every iteration | | `WorkerStoppedEvent` | once, after the worker has stopped | +`WorkerStoppedEvent` is also dispatched if the job throws an exception. +In that case the exception is available via `$event->exception` and is rethrown afterwards, +so listeners can clean up resources and still distinguish a crash from a regular stop: + +```php +use Patchlevel\Worker\Event\WorkerStoppedEvent; + +$eventDispatcher->addListener( + WorkerStoppedEvent::class, + static function (WorkerStoppedEvent $event): void { + if ($event->exception !== null) { + // the worker crashed + } + }, +); +``` All [limits](getting-started.md#limits) are implemented as event subscribers (`StopWorkerOnIterationLimitListener`, `StopWorkerOnMemoryLimitListener`, `StopWorkerOnTimeLimitListener`, `StopWorkerOnSignalListener`), diff --git a/docs/getting-started.md b/docs/getting-started.md index d84e8f2..1685ec1 100644 --- a/docs/getting-started.md +++ b/docs/getting-started.md @@ -29,6 +29,9 @@ The job receives a `$stop` callback: calling it tells the worker to exit the loo after the current iteration has finished. The worker never aborts a running job — stopping always happens *between* iterations, so your job is never interrupted halfway. +If the job throws an exception, the worker stops, dispatches the `WorkerStoppedEvent` +and rethrows the exception, so the process exits with an error. + ## Limits All options are optional. Without limits the worker runs until it is stopped diff --git a/src/DefaultWorker.php b/src/DefaultWorker.php index 42d5cc3..82c090b 100644 --- a/src/DefaultWorker.php +++ b/src/DefaultWorker.php @@ -17,6 +17,7 @@ use Psr\Log\NullLogger; use Symfony\Component\EventDispatcher\EventDispatcher; use Symfony\Component\EventDispatcher\EventDispatcherInterface; +use Throwable; use function max; use function usleep; @@ -44,41 +45,47 @@ public function run(int $sleepTimer = 1000): void $this->eventDispatcher->dispatch(new WorkerStartedEvent($this)); - while (!$this->shouldStop) { - $this->logger?->debug('Worker starting job run'); + $exception = null; - $startTime = $this->milliseconds(); + try { + while (!$this->shouldStop) { + $this->logger?->debug('Worker starting job run'); - ($this->job)($this->stop(...)); + $startTime = $this->milliseconds(); - $endTime = $this->milliseconds(); - $ranTime = $endTime - $startTime; + ($this->job)($this->stop(...)); - $this->logger?->debug('Worker finished job run ({ranTime}ms)', ['ranTime' => $ranTime]); + $endTime = $this->milliseconds(); + $ranTime = $endTime - $startTime; - $this->eventDispatcher->dispatch(new WorkerRunningEvent($this)); + $this->logger?->debug('Worker finished job run ({ranTime}ms)', ['ranTime' => $ranTime]); - if ($this->shouldStop) { - break; - } + $this->eventDispatcher->dispatch(new WorkerRunningEvent($this)); - $sleepFor = max($sleepTimer - $ranTime, 0); + if ($this->shouldStop) { + break; + } - if ($sleepFor <= 0) { - continue; - } + $sleepFor = max($sleepTimer - $ranTime, 0); - $this->logger?->debug('Worker sleep for {sleepTimer}ms', ['sleepTimer' => $sleepFor]); - usleep($sleepFor * 1000); - } + if ($sleepFor <= 0) { + continue; + } - $this->shouldStop = false; + $this->logger?->debug('Worker sleep for {sleepTimer}ms', ['sleepTimer' => $sleepFor]); + usleep($sleepFor * 1000); + } + } catch (Throwable $exception) { + throw $exception; + } finally { + $this->shouldStop = false; - $this->logger?->debug('Worker stopped'); + $this->logger?->debug('Worker stopped'); - $this->eventDispatcher->dispatch(new WorkerStoppedEvent($this)); + $this->eventDispatcher->dispatch(new WorkerStoppedEvent($this, $exception)); - $this->logger?->debug('Worker terminated'); + $this->logger?->debug('Worker terminated'); + } } public function stop(): void diff --git a/src/Event/WorkerStoppedEvent.php b/src/Event/WorkerStoppedEvent.php index b33da86..5cc38fa 100644 --- a/src/Event/WorkerStoppedEvent.php +++ b/src/Event/WorkerStoppedEvent.php @@ -5,11 +5,13 @@ namespace Patchlevel\Worker\Event; use Patchlevel\Worker\Worker; +use Throwable; final class WorkerStoppedEvent { public function __construct( public readonly Worker $worker, + public readonly Throwable|null $exception = null, ) { } } diff --git a/tests/Unit/DefaultWorkerTest.php b/tests/Unit/DefaultWorkerTest.php index 2340024..53b518c 100644 --- a/tests/Unit/DefaultWorkerTest.php +++ b/tests/Unit/DefaultWorkerTest.php @@ -8,6 +8,7 @@ use Patchlevel\Worker\DefaultWorker; use Patchlevel\Worker\Event\WorkerRunningEvent; use Patchlevel\Worker\Event\WorkerStartedEvent; +use Patchlevel\Worker\Event\WorkerStoppedEvent; use Patchlevel\Worker\Listener\StopWorkerOnIterationLimitListener; use Patchlevel\Worker\Listener\StopWorkerOnMemoryLimitListener; use Patchlevel\Worker\Listener\StopWorkerOnSignalListener; @@ -17,6 +18,7 @@ use PHPUnit\Framework\Attributes\CoversClass; use PHPUnit\Framework\TestCase; use Psr\Log\LoggerInterface; +use RuntimeException; use Symfony\Component\EventDispatcher\EventDispatcher; use Symfony\Component\EventDispatcher\EventDispatcherInterface; @@ -289,4 +291,75 @@ static function () use (&$calls): void { self::assertSame(0, $calls); } + + public function testStoppedEventWithoutException(): void + { + $stoppedEvent = null; + + $eventDispatcher = new EventDispatcher(); + $eventDispatcher->addListener( + WorkerStoppedEvent::class, + static function (WorkerStoppedEvent $event) use (&$stoppedEvent): void { + $stoppedEvent = $event; + }, + ); + + $worker = new DefaultWorker( + static function (callable $stop): void { + $stop(); + }, + $eventDispatcher, + ); + + $worker->run(0); + + self::assertInstanceOf(WorkerStoppedEvent::class, $stoppedEvent); + self::assertSame($worker, $stoppedEvent->worker); + self::assertNull($stoppedEvent->exception); + } + + public function testJobExceptionDispatchesStoppedEvent(): void + { + $exception = new RuntimeException('job failed'); + $stoppedEvent = null; + + $eventDispatcher = new EventDispatcher(); + $eventDispatcher->addListener( + WorkerStoppedEvent::class, + static function (WorkerStoppedEvent $event) use (&$stoppedEvent): void { + $stoppedEvent = $event; + }, + ); + + $logger = $this->createMock(LoggerInterface::class); + $logger + ->expects($this->exactly(4)) + ->method('debug') + ->willReturnCallback( + new ReturnCallback([ + [['Worker starting', []]], + [['Worker starting job run', []]], + [['Worker stopped', []]], + [['Worker terminated', []]], + ]), + ); + + $worker = new DefaultWorker( + static function () use ($exception): void { + throw $exception; + }, + $eventDispatcher, + $logger, + ); + + try { + $worker->run(0); + self::fail('Expected exception was not thrown'); + } catch (RuntimeException $e) { + self::assertSame($exception, $e); + } + + self::assertInstanceOf(WorkerStoppedEvent::class, $stoppedEvent); + self::assertSame($exception, $stoppedEvent->exception); + } }