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
4 changes: 2 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -5,13 +5,13 @@
# Worker

A small library to build stable, long-running workers that terminate gracefully when limits are exceeded
or a SIGTERM signal is received. Perfect for daemonized console commands running under
or a SIGTERM/SIGINT signal is received. Perfect for daemonized console commands running under
Docker, Kubernetes, supervisor or systemd, where the process manager restarts the worker after it exits.

## Features

* Configurable run, memory and time [limits](https://patchlevel.dev/docs/worker/latest/getting-started#limits)
* [Graceful shutdown](https://patchlevel.dev/docs/worker/latest/getting-started#graceful-shutdown-on-sigterm) on SIGTERM
* [Graceful shutdown](https://patchlevel.dev/docs/worker/latest/getting-started#graceful-shutdown) on SIGTERM and SIGINT
* Extensible via [events and custom listeners](https://patchlevel.dev/docs/worker/latest/events)
* [PSR-3 logging](https://patchlevel.dev/docs/worker/latest/getting-started#logging) of the worker lifecycle
* Plays well with [Symfony and Laravel console commands](https://patchlevel.dev/docs/worker/latest/integration)
Expand Down
2 changes: 1 addition & 1 deletion docs/events.md
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@ each carrying the worker instance:

All [limits](getting-started.md#limits) are implemented as event subscribers
(`StopWorkerOnIterationLimitListener`, `StopWorkerOnMemoryLimitListener`,
`StopWorkerOnTimeLimitListener`, `StopWorkerOnSigtermSignalListener`),
`StopWorkerOnTimeLimitListener`, `StopWorkerOnSignalListener`),
so you can add your own stop conditions the same way:

```php
Expand Down
15 changes: 9 additions & 6 deletions docs/getting-started.md
Original file line number Diff line number Diff line change
Expand Up @@ -32,7 +32,7 @@ stopping always happens *between* iterations, so your job is never interrupted h
## Limits

All options are optional. Without limits the worker runs until it is stopped
via `$stop()`, `$worker->stop()` or a SIGTERM signal.
via `$stop()`, `$worker->stop()` or a SIGTERM/SIGINT signal.

| Option | Type | Description |
|---------------|----------|---------------------------------------------------------------------------------------------------------------------------------|
Expand All @@ -51,11 +51,11 @@ Internally every limit is implemented as an event listener.
You can add your own stop conditions the same way, see [events & listeners](events.md).
:::

## Graceful shutdown on SIGTERM
## Graceful shutdown

If the `pcntl` extension is available, the worker automatically registers a SIGTERM handler.
When the process receives SIGTERM (e.g. from `docker stop`, a Kubernetes pod shutdown or supervisor),
the worker finishes the current iteration and then exits cleanly.
If the `pcntl` extension is available, the worker automatically registers a handler for SIGTERM and SIGINT.
When the process receives SIGTERM (e.g. from `docker stop`, a Kubernetes pod shutdown or supervisor)
or SIGINT (e.g. pressing `Ctrl+C`), the worker finishes the current iteration and then exits cleanly.

This makes the worker a good fit for process managers that send SIGTERM and restart the process,
e.g. to roll out a new version or to keep long-running processes fresh.
Expand All @@ -64,6 +64,9 @@ e.g. to roll out a new version or to keep long-running processes fresh.
Without `ext-pcntl` this feature is not available.
:::

If you need to react to other signals, register the `StopWorkerOnSignalListener` with your own list of signals
on a custom event dispatcher, see [events & listeners](events.md).

## Sleep

`run()` takes a sleep timer in milliseconds (default: `1000`):
Expand All @@ -78,7 +81,7 @@ the next iteration starts immediately. Pass `0` to disable sleeping entirely.
## Logging

The worker logs its lifecycle (start, iteration timings, sleep, stop reason) to the given PSR-3 logger.
Iteration details use the `debug` level; stop reasons (limit exceeded, SIGTERM received) use `info`.
Iteration details use the `debug` level; stop reasons (limit exceeded, signal received) use `info`.

With the `ConsoleLogger` from the [Symfony command example](integration.md#symfony), run the command with `-v`
to see stop reasons or `-vvv` to see everything.
4 changes: 2 additions & 2 deletions docs/introduction.md
Original file line number Diff line number Diff line change
@@ -1,15 +1,15 @@
# Worker

A small library to build stable, long-running workers that terminate gracefully when limits are exceeded
or a SIGTERM signal is received. Perfect for daemonized console commands running under
or a SIGTERM/SIGINT signal is received. Perfect for daemonized console commands running under
Docker, Kubernetes, supervisor or systemd, where the process manager restarts the worker after it exits.

It was extracted from the [event-sourcing](https://github.com/patchlevel/event-sourcing) library into a separate package.

## Features

* Configurable run, memory and time [limits](getting-started.md#limits)
* [Graceful shutdown](getting-started.md#graceful-shutdown-on-sigterm) on SIGTERM
* [Graceful shutdown](getting-started.md#graceful-shutdown) on SIGTERM and SIGINT
* Extensible via [events and custom listeners](events.md)
* [PSR-3 logging](getting-started.md#logging) of the worker lifecycle
* Plays well with [Symfony and Laravel console commands](integration.md)
Expand Down
4 changes: 2 additions & 2 deletions src/DefaultWorker.php
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,7 @@
use Patchlevel\Worker\Event\WorkerStoppedEvent;
use Patchlevel\Worker\Listener\StopWorkerOnIterationLimitListener;
use Patchlevel\Worker\Listener\StopWorkerOnMemoryLimitListener;
use Patchlevel\Worker\Listener\StopWorkerOnSigtermSignalListener;
use Patchlevel\Worker\Listener\StopWorkerOnSignalListener;
use Patchlevel\Worker\Listener\StopWorkerOnTimeLimitListener;
use Psr\Log\LoggerInterface;
use Psr\Log\NullLogger;
Expand All @@ -35,11 +35,11 @@
private readonly EventDispatcherInterface $eventDispatcher,
private readonly LoggerInterface|null $logger = null,
) {
$this->timeMeasure = static fn () => (int)round(microtime(true) * 1000);

Check warning on line 38 in src/DefaultWorker.php

View workflow job for this annotation

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

Escaped Mutant for Mutator "RoundingFamily": @@ @@ private readonly EventDispatcherInterface $eventDispatcher, private readonly LoggerInterface|null $logger = null, ) { - $this->timeMeasure = static fn () => (int)round(microtime(true) * 1000); + $this->timeMeasure = static fn () => (int)ceil(microtime(true) * 1000); } /** @PARAM positive-int|0 $sleepTimer in milliseconds */

Check warning on line 38 in src/DefaultWorker.php

View workflow job for this annotation

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

Escaped Mutant for Mutator "RoundingFamily": @@ @@ private readonly EventDispatcherInterface $eventDispatcher, private readonly LoggerInterface|null $logger = null, ) { - $this->timeMeasure = static fn () => (int)round(microtime(true) * 1000); + $this->timeMeasure = static fn () => (int)floor(microtime(true) * 1000); } /** @PARAM positive-int|0 $sleepTimer in milliseconds */

Check warning on line 38 in src/DefaultWorker.php

View workflow job for this annotation

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

Escaped Mutant for Mutator "IncrementInteger": @@ @@ private readonly EventDispatcherInterface $eventDispatcher, private readonly LoggerInterface|null $logger = null, ) { - $this->timeMeasure = static fn () => (int)round(microtime(true) * 1000); + $this->timeMeasure = static fn () => (int)round(microtime(true) * 1001); } /** @PARAM positive-int|0 $sleepTimer in milliseconds */

Check warning on line 38 in src/DefaultWorker.php

View workflow job for this annotation

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

Escaped Mutant for Mutator "Multiplication": @@ @@ private readonly EventDispatcherInterface $eventDispatcher, private readonly LoggerInterface|null $logger = null, ) { - $this->timeMeasure = static fn () => (int)round(microtime(true) * 1000); + $this->timeMeasure = static fn () => (int)round(microtime(true) / 1000); } /** @PARAM positive-int|0 $sleepTimer in milliseconds */

Check warning on line 38 in src/DefaultWorker.php

View workflow job for this annotation

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

Escaped Mutant for Mutator "DecrementInteger": @@ @@ private readonly EventDispatcherInterface $eventDispatcher, private readonly LoggerInterface|null $logger = null, ) { - $this->timeMeasure = static fn () => (int)round(microtime(true) * 1000); + $this->timeMeasure = static fn () => (int)round(microtime(true) * 999); } /** @PARAM positive-int|0 $sleepTimer in milliseconds */
}

/** @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');

Expand All @@ -60,10 +60,10 @@
$this->eventDispatcher->dispatch(new WorkerRunningEvent($this));

if ($this->shouldStop) {
break;

Check warning on line 63 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);
}

$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 "DecrementInteger": @@ @@ break; } - $sleepFor = max($sleepTimer - $ranTime, 0); + $sleepFor = max($sleepTimer - $ranTime, -1); if ($sleepFor <= 0) { continue;

if ($sleepFor <= 0) {
continue;
Expand Down Expand Up @@ -100,7 +100,7 @@
$eventDispatcher = new EventDispatcher();
}

$eventDispatcher->addSubscriber(new StopWorkerOnSigtermSignalListener($logger));
$eventDispatcher->addSubscriber(new StopWorkerOnSignalListener(logger: $logger));

if (isset($options['runLimit'])) {
$eventDispatcher->addSubscriber(
Expand Down
48 changes: 48 additions & 0 deletions src/Listener/StopWorkerOnSignalListener.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,48 @@
<?php

declare(strict_types=1);

namespace Patchlevel\Worker\Listener;

use Patchlevel\Worker\Event\WorkerStartedEvent;
use Psr\Log\LoggerInterface;
use Symfony\Component\EventDispatcher\EventSubscriberInterface;

use function function_exists;
use function pcntl_async_signals;
use function pcntl_signal;

use const SIGINT;
use const SIGTERM;

final class StopWorkerOnSignalListener implements EventSubscriberInterface
{
/** @param list<int>|null $signals defaults to SIGTERM and SIGINT */
public function __construct(
private readonly array|null $signals = null,
private readonly LoggerInterface|null $logger = null,
) {
}

public function onWorkerStarted(WorkerStartedEvent $event): void
{
pcntl_async_signals(true);

foreach ($this->signals ?? [SIGTERM, SIGINT] as $signal) {
pcntl_signal($signal, function (int $signal) use ($event): void {
$this->logger?->info('Worker received signal {signal}', ['signal' => $signal]);
$event->worker->stop();
});
}
}

/** @return array<class-string, string> */
public static function getSubscribedEvents(): array
{
if (!function_exists('pcntl_signal')) {
return [];
}

return [WorkerStartedEvent::class => 'onWorkerStarted'];
}
}
1 change: 1 addition & 0 deletions src/Listener/StopWorkerOnSigtermSignalListener.php
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@

use const SIGTERM;

/** @deprecated since 1.6, use StopWorkerOnSignalListener instead */
final class StopWorkerOnSigtermSignalListener implements EventSubscriberInterface
{
public function __construct(
Expand Down
4 changes: 2 additions & 2 deletions tests/Unit/DefaultWorkerTest.php
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,7 @@
use Patchlevel\Worker\Event\WorkerStartedEvent;
use Patchlevel\Worker\Listener\StopWorkerOnIterationLimitListener;
use Patchlevel\Worker\Listener\StopWorkerOnMemoryLimitListener;
use Patchlevel\Worker\Listener\StopWorkerOnSigtermSignalListener;
use Patchlevel\Worker\Listener\StopWorkerOnSignalListener;
use Patchlevel\Worker\Listener\StopWorkerOnTimeLimitListener;
use Patchlevel\Worker\Tests\ReturnCallback;
use PHPUnit\Framework\Attributes\CoversClass;
Expand Down Expand Up @@ -142,7 +142,7 @@ public function testOptions(): void

$invokationCount = $this->exactly(4);
$invokationParameters = [
[new StopWorkerOnSigtermSignalListener($logger)],
[new StopWorkerOnSignalListener(logger: $logger)],
[new StopWorkerOnIterationLimitListener(10, $logger)],
[new StopWorkerOnMemoryLimitListener(Bytes::parseFromString('10KB'), $logger)],
[new StopWorkerOnTimeLimitListener(20, $logger)],
Expand Down
107 changes: 107 additions & 0 deletions tests/Unit/Listener/StopWorkerOnSignalListenerTest.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,107 @@
<?php

declare(strict_types=1);

namespace Patchlevel\Worker\Tests\Unit\Listener;

use Patchlevel\Worker\DefaultWorker;
use Patchlevel\Worker\Listener\StopWorkerOnSignalListener;
use PHPUnit\Framework\Attributes\CoversClass;
use PHPUnit\Framework\Attributes\RequiresFunction;
use PHPUnit\Framework\TestCase;
use Psr\Log\LoggerInterface;
use Symfony\Component\EventDispatcher\EventDispatcher;

use function getmypid;
use function pcntl_signal;
use function posix_kill;

use const SIG_DFL;
use const SIGINT;
use const SIGTERM;
use const SIGUSR1;

#[CoversClass(StopWorkerOnSignalListener::class)]
#[RequiresFunction('pcntl_signal')]
#[RequiresFunction('posix_kill')]
final class StopWorkerOnSignalListenerTest extends TestCase
{
protected function tearDown(): void
{
pcntl_signal(SIGTERM, SIG_DFL);
pcntl_signal(SIGINT, SIG_DFL);
pcntl_signal(SIGUSR1, SIG_DFL);
}

public function testStopOnSigtermByDefault(): void
{
$logger = $this->createMock(LoggerInterface::class);
$logger
->expects($this->once())
->method('info')
->with('Worker received signal {signal}', ['signal' => SIGTERM]);

$calls = $this->runWorkerAndSendSignal(new StopWorkerOnSignalListener(logger: $logger), SIGTERM);

self::assertSame(1, $calls);
}

public function testStopOnSigintByDefault(): void
{
$logger = $this->createMock(LoggerInterface::class);
$logger
->expects($this->once())
->method('info')
->with('Worker received signal {signal}', ['signal' => SIGINT]);

$calls = $this->runWorkerAndSendSignal(new StopWorkerOnSignalListener(logger: $logger), SIGINT);

self::assertSame(1, $calls);
}

public function testStopOnCustomSignal(): void
{
$logger = $this->createMock(LoggerInterface::class);
$logger
->expects($this->once())
->method('info')
->with('Worker received signal {signal}', ['signal' => SIGUSR1]);

$calls = $this->runWorkerAndSendSignal(new StopWorkerOnSignalListener([SIGUSR1], $logger), SIGUSR1);

self::assertSame(1, $calls);
}

/**
* Sends the signal during the first job run. The job stops the worker
* itself after 5 runs, so more than 1 run means the signal was ignored.
*/
private function runWorkerAndSendSignal(StopWorkerOnSignalListener $listener, int $signal): int
{
$eventDispatcher = new EventDispatcher();
$eventDispatcher->addSubscriber($listener);

$calls = 0;

$worker = new DefaultWorker(
static function (callable $stop) use (&$calls, $signal): void {
$calls++;

if ($calls === 1) {
posix_kill((int)getmypid(), $signal);
}

if ($calls < 5) {
return;
}

$stop();
},
$eventDispatcher,
);

$worker->run(0);

return $calls;
}
}
Loading