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
16 changes: 16 additions & 0 deletions docs/events.md
Original file line number Diff line number Diff line change
Expand Up @@ -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`),
Expand Down
3 changes: 3 additions & 0 deletions docs/getting-started.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
51 changes: 29 additions & 22 deletions src/DefaultWorker.php
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -38,47 +39,53 @@
}

/** @param positive-int|0 $sleepTimer in milliseconds */
public function run(int $sleepTimer = 1000): void

Check warning on line 42 in src/DefaultWorker.php

View workflow job for this annotation

GitHub Actions / Mutation tests (locked, 8.4, ubuntu-latest)

Escaped Mutant for Mutator "IncrementInteger": @@ @@ } /** @PARAM positive-int|0 $sleepTimer in milliseconds */ - public function run(int $sleepTimer = 1000): void + public function run(int $sleepTimer = 1001): void { $this->logger?->debug('Worker starting');

Check warning on line 42 in src/DefaultWorker.php

View workflow job for this annotation

GitHub Actions / Mutation tests (locked, 8.4, ubuntu-latest)

Escaped Mutant for Mutator "DecrementInteger": @@ @@ } /** @PARAM positive-int|0 $sleepTimer in milliseconds */ - public function run(int $sleepTimer = 1000): void + public function run(int $sleepTimer = 999): void { $this->logger?->debug('Worker starting');
{
$this->logger?->debug('Worker starting');

$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;

Check warning on line 66 in src/DefaultWorker.php

View workflow job for this annotation

GitHub Actions / Mutation tests on diff (locked, 8.4, ubuntu-latest)

Escaped Mutant for Mutator "Break_": @@ @@ $this->eventDispatcher->dispatch(new WorkerRunningEvent($this)); if ($this->shouldStop) { - break; + continue; } $sleepFor = max($sleepTimer - $ranTime, 0);

Check warning on line 66 in src/DefaultWorker.php

View workflow job for this annotation

GitHub Actions / Mutation tests (locked, 8.4, ubuntu-latest)

Escaped Mutant for Mutator "Break_": @@ @@ $this->eventDispatcher->dispatch(new WorkerRunningEvent($this)); if ($this->shouldStop) { - break; + continue; } $sleepFor = max($sleepTimer - $ranTime, 0);
}

if ($sleepFor <= 0) {
continue;
}
$sleepFor = max($sleepTimer - $ranTime, 0);

Check warning on line 69 in src/DefaultWorker.php

View workflow job for this annotation

GitHub Actions / Mutation tests on diff (locked, 8.4, ubuntu-latest)

Escaped Mutant for Mutator "DecrementInteger": @@ @@ break; } - $sleepFor = max($sleepTimer - $ranTime, 0); + $sleepFor = max($sleepTimer - $ranTime, -1); if ($sleepFor <= 0) { continue;

Check warning on line 69 in src/DefaultWorker.php

View workflow job for this annotation

GitHub Actions / Mutation tests (locked, 8.4, ubuntu-latest)

Escaped Mutant for Mutator "DecrementInteger": @@ @@ break; } - $sleepFor = max($sleepTimer - $ranTime, 0); + $sleepFor = max($sleepTimer - $ranTime, -1); if ($sleepFor <= 0) { continue;

$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);

Check warning on line 76 in src/DefaultWorker.php

View workflow job for this annotation

GitHub Actions / Mutation tests on diff (locked, 8.4, ubuntu-latest)

Escaped Mutant for Mutator "FunctionCallRemoval": @@ @@ } $this->logger?->debug('Worker sleep for {sleepTimer}ms', ['sleepTimer' => $sleepFor]); - usleep($sleepFor * 1000); + } } catch (Throwable $exception) { throw $exception;

Check warning on line 76 in src/DefaultWorker.php

View workflow job for this annotation

GitHub Actions / Mutation tests on diff (locked, 8.4, ubuntu-latest)

Escaped Mutant for Mutator "DecrementInteger": @@ @@ } $this->logger?->debug('Worker sleep for {sleepTimer}ms', ['sleepTimer' => $sleepFor]); - usleep($sleepFor * 1000); + usleep($sleepFor * 999); } } catch (Throwable $exception) { throw $exception;

Check warning on line 76 in src/DefaultWorker.php

View workflow job for this annotation

GitHub Actions / Mutation tests on diff (locked, 8.4, ubuntu-latest)

Escaped Mutant for Mutator "IncrementInteger": @@ @@ } $this->logger?->debug('Worker sleep for {sleepTimer}ms', ['sleepTimer' => $sleepFor]); - usleep($sleepFor * 1000); + usleep($sleepFor * 1001); } } catch (Throwable $exception) { throw $exception;

Check warning on line 76 in src/DefaultWorker.php

View workflow job for this annotation

GitHub Actions / Mutation tests (locked, 8.4, ubuntu-latest)

Escaped Mutant for Mutator "FunctionCallRemoval": @@ @@ } $this->logger?->debug('Worker sleep for {sleepTimer}ms', ['sleepTimer' => $sleepFor]); - usleep($sleepFor * 1000); + } } catch (Throwable $exception) { throw $exception;

Check warning on line 76 in src/DefaultWorker.php

View workflow job for this annotation

GitHub Actions / Mutation tests (locked, 8.4, ubuntu-latest)

Escaped Mutant for Mutator "IncrementInteger": @@ @@ } $this->logger?->debug('Worker sleep for {sleepTimer}ms', ['sleepTimer' => $sleepFor]); - usleep($sleepFor * 1000); + usleep($sleepFor * 1001); } } catch (Throwable $exception) { throw $exception;

Check warning on line 76 in src/DefaultWorker.php

View workflow job for this annotation

GitHub Actions / Mutation tests (locked, 8.4, ubuntu-latest)

Escaped Mutant for Mutator "DecrementInteger": @@ @@ } $this->logger?->debug('Worker sleep for {sleepTimer}ms', ['sleepTimer' => $sleepFor]); - usleep($sleepFor * 1000); + usleep($sleepFor * 999); } } catch (Throwable $exception) { throw $exception;
}
} 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
Expand Down
2 changes: 2 additions & 0 deletions src/Event/WorkerStoppedEvent.php
Original file line number Diff line number Diff line change
Expand Up @@ -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,
) {
}
}
73 changes: 73 additions & 0 deletions tests/Unit/DefaultWorkerTest.php
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;

Expand Down Expand Up @@ -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);
}
}
Loading