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
47 changes: 42 additions & 5 deletions docs/getting-started.md
Original file line number Diff line number Diff line change
Expand Up @@ -37,11 +37,12 @@ and rethrows the exception, so the process exits with an error.
All options are optional. Without limits the worker runs until it is stopped
via `$stop()`, `$worker->stop()` or a SIGTERM/SIGINT signal.

| Option | Type | Description |
|---------------|----------|-------------------------------------------------------------------------------------------------------------------------------------------------------|
| `runLimit` | `int` | Stop after this number of iterations. |
| `memoryLimit` | `string` | Stop when memory usage exceeds this value, e.g. `128MB` or `128M`. Supported units: `B`, `K`/`KB`, `M`/`MB`, `G`/`GB` (case-insensitive, 1024-based). |
| `timeLimit` | `int` | Stop after this number of seconds. |
| Option | Type | Description |
|-----------------|----------|-------------------------------------------------------------------------------------------------------------------------------------------------------|
| `runLimit` | `int` | Stop after this number of iterations. |
| `memoryLimit` | `string` | Stop when memory usage exceeds this value, e.g. `128MB` or `128M`. Supported units: `B`, `K`/`KB`, `M`/`MB`, `G`/`GB` (case-insensitive, 1024-based). |
| `timeLimit` | `int` | Stop after this number of seconds. |
| `heartbeatFile` | `string` | Touch this file on start and after every iteration, see [heartbeat](#heartbeat). |

Limits are checked after each iteration. When a limit is exceeded, the worker logs the reason and stops gracefully.

Expand Down Expand Up @@ -108,6 +109,42 @@ $worker->run(1000);
```
Any other return value, including no return value at all, keeps the regular sleep behaviour.

## Heartbeat

A process manager only sees whether the worker process is running, not whether it is stuck in a job.
With the `heartbeatFile` option the worker touches a file when it starts and after every iteration,
and removes it when it stops. A liveness probe can then check how old the file is:

```php
use Patchlevel\Worker\DefaultWorker;

$worker = DefaultWorker::create(
$job,
['heartbeatFile' => '/tmp/worker-heartbeat'],
$logger,
);
```
For example as a Kubernetes liveness probe that restarts the pod if the file is older than 60 seconds:

```yaml
livenessProbe:
exec:
command:
- sh
- -c
- test $(( $(date +%s) - $(stat -c %Y /tmp/worker-heartbeat) )) -lt 60
initialDelaySeconds: 10
periodSeconds: 30
```
:::warning
The file is only updated between iterations. Choose the threshold larger than your longest job
plus the sleep timer, otherwise a slow but healthy worker is considered dead.
:::

:::note
If the file cannot be written, the worker logs a warning and keeps running.
:::

## Logging

The worker logs its lifecycle (start, iteration timings, sleep, stop reason) to the given PSR-3 logger.
Expand Down
11 changes: 9 additions & 2 deletions src/DefaultWorker.php
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@
use Patchlevel\Worker\Event\WorkerRunningEvent;
use Patchlevel\Worker\Event\WorkerStartedEvent;
use Patchlevel\Worker\Event\WorkerStoppedEvent;
use Patchlevel\Worker\Listener\HeartbeatListener;
use Patchlevel\Worker\Listener\StopWorkerOnIterationLimitListener;
use Patchlevel\Worker\Listener\StopWorkerOnMemoryLimitListener;
use Patchlevel\Worker\Listener\StopWorkerOnSignalListener;
Expand Down Expand Up @@ -39,7 +40,7 @@
}

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

Check warning on line 43 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 43 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 @@ -63,21 +64,21 @@
$this->eventDispatcher->dispatch(new WorkerRunningEvent($this));

if ($this->shouldStop) {
break;

Check warning on line 67 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; } if ($didWork) {
}

if ($didWork) {
continue;
}

$sleepFor = max($sleepTimer - $ranTime, 0);

Check warning on line 74 in src/DefaultWorker.php

View workflow job for this annotation

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

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

if ($sleepFor <= 0) {
continue;
}

$this->logger?->debug('Worker sleep for {sleepTimer}ms', ['sleepTimer' => $sleepFor]);
usleep($sleepFor * 1000);

Check warning on line 81 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 81 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;

Check warning on line 81 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;
}
} catch (Throwable $exception) {
throw $exception;
Expand All @@ -99,8 +100,8 @@
}

/**
* @param Closure(Closure):(bool|void) $job
* @param array{runLimit?: (positive-int|null), memoryLimit?: (string|null), timeLimit?: (positive-int|null)} $options
* @param Closure(Closure):(bool|void) $job
* @param array{runLimit?: (positive-int|null), memoryLimit?: (string|null), timeLimit?: (positive-int|null), heartbeatFile?: (string|null)} $options
*/
public static function create(
Closure $job,
Expand Down Expand Up @@ -133,6 +134,12 @@
);
}

if (isset($options['heartbeatFile'])) {
$eventDispatcher->addSubscriber(
new HeartbeatListener($options['heartbeatFile'], $logger),
);
}

return new self(
$job,
$eventDispatcher,
Expand Down
63 changes: 63 additions & 0 deletions src/Listener/HeartbeatListener.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,63 @@
<?php

declare(strict_types=1);

namespace Patchlevel\Worker\Listener;

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

use function is_file;
use function touch;
use function unlink;

final class HeartbeatListener implements EventSubscriberInterface
{
/** @param string $file touched on start and after every iteration, removed when the worker stops */
public function __construct(
private readonly string $file,
private readonly LoggerInterface|null $logger = null,
) {
}

public function onWorkerStarted(): void
{
$this->beat();
}

public function onWorkerRunning(): void
{
$this->beat();
}

public function onWorkerStopped(): void
{
if (!is_file($this->file)) {
return;
}

unlink($this->file);
}

private function beat(): void
{
if (@touch($this->file)) {
return;
}

$this->logger?->warning('Worker could not update heartbeat file {file}', ['file' => $this->file]);
}

/** @return array<class-string, string> */
public static function getSubscribedEvents(): array
{
return [
WorkerStartedEvent::class => 'onWorkerStarted',
WorkerRunningEvent::class => 'onWorkerRunning',
WorkerStoppedEvent::class => 'onWorkerStopped',
];
}
}
7 changes: 5 additions & 2 deletions tests/Unit/DefaultWorkerTest.php
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@
use Patchlevel\Worker\Event\WorkerRunningEvent;
use Patchlevel\Worker\Event\WorkerStartedEvent;
use Patchlevel\Worker\Event\WorkerStoppedEvent;
use Patchlevel\Worker\Listener\HeartbeatListener;
use Patchlevel\Worker\Listener\StopWorkerOnIterationLimitListener;
use Patchlevel\Worker\Listener\StopWorkerOnMemoryLimitListener;
use Patchlevel\Worker\Listener\StopWorkerOnSignalListener;
Expand Down Expand Up @@ -61,7 +62,7 @@ static function (object $event): object {
$this->assertSame($invokationParameters[$invokationCount->numberOfInvocations() - 1], $parameters);
});

$worker = new DefaultWorker(static fn () => null, $eventDispatcher, $logger);
$worker = new DefaultWorker(static fn () => null, $eventDispatcher, $logger, new TestClock());
$worker->run(200);
}

Expand Down Expand Up @@ -146,12 +147,13 @@ public function testOptions(): void

$clock = new TestClock();

$invokationCount = $this->exactly(4);
$invokationCount = $this->exactly(5);
$invokationParameters = [
[new StopWorkerOnSignalListener(logger: $logger)],
[new StopWorkerOnIterationLimitListener(10, $logger)],
[new StopWorkerOnMemoryLimitListener(Bytes::parseFromString('10KB'), $logger)],
[new StopWorkerOnTimeLimitListener(20, $logger, $clock)],
[new HeartbeatListener('/tmp/worker-heartbeat', $logger)],
];

$eventDispatcher = $this->createMock(EventDispatcherInterface::class);
Expand All @@ -170,6 +172,7 @@ static function ($stop): void {
'runLimit' => 10,
'memoryLimit' => '10KB',
'timeLimit' => 20,
'heartbeatFile' => '/tmp/worker-heartbeat',
],
$logger,
$eventDispatcher,
Expand Down
118 changes: 118 additions & 0 deletions tests/Unit/Listener/HeartbeatListenerTest.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,118 @@
<?php

declare(strict_types=1);

namespace Patchlevel\Worker\Tests\Unit\Listener;

use Patchlevel\Worker\Event\WorkerRunningEvent;
use Patchlevel\Worker\Event\WorkerStartedEvent;
use Patchlevel\Worker\Event\WorkerStoppedEvent;
use Patchlevel\Worker\Listener\HeartbeatListener;
use PHPUnit\Framework\Attributes\CoversClass;
use PHPUnit\Framework\TestCase;
use Psr\Log\LoggerInterface;

use function clearstatcache;
use function filemtime;
use function is_file;
use function sys_get_temp_dir;
use function time;
use function touch;
use function uniqid;
use function unlink;

#[CoversClass(HeartbeatListener::class)]
final class HeartbeatListenerTest extends TestCase
{
private string $file;

protected function setUp(): void
{
$this->file = sys_get_temp_dir() . '/' . uniqid('worker-heartbeat-', true);
}

protected function tearDown(): void
{
if (!is_file($this->file)) {
return;
}

unlink($this->file);
}

public function testCreateFileOnStart(): void
{
$logger = $this->createMock(LoggerInterface::class);
$logger
->expects($this->never())
->method('warning');

$listener = new HeartbeatListener($this->file, $logger);
$listener->onWorkerStarted();

self::assertFileExists($this->file);
}

public function testUpdateFileAfterIteration(): void
{
touch($this->file, time() - 60);

$logger = $this->createMock(LoggerInterface::class);
$logger
->expects($this->never())
->method('warning');

$listener = new HeartbeatListener($this->file, $logger);
$listener->onWorkerRunning();

clearstatcache(true, $this->file);

self::assertGreaterThanOrEqual(time() - 1, filemtime($this->file));
}

public function testRemoveFileOnStop(): void
{
touch($this->file);

$listener = new HeartbeatListener($this->file);
$listener->onWorkerStopped();

self::assertFileDoesNotExist($this->file);
}

public function testStopWithoutFile(): void
{
$listener = new HeartbeatListener($this->file);
$listener->onWorkerStopped();

self::assertFileDoesNotExist($this->file);
}

public function testLogWarningIfFileCannotBeWritten(): void
{
$file = sys_get_temp_dir() . '/' . uniqid('worker-heartbeat-', true) . '/missing-dir/heartbeat';

$logger = $this->createMock(LoggerInterface::class);
$logger
->expects($this->once())
->method('warning')
->with('Worker could not update heartbeat file {file}', ['file' => $file]);

$listener = new HeartbeatListener($file, $logger);
$listener->onWorkerRunning();

self::assertFileDoesNotExist($file);
}

public function testSubscribedEvents(): void
{
self::assertSame(
[
WorkerStartedEvent::class => 'onWorkerStarted',
WorkerRunningEvent::class => 'onWorkerRunning',
WorkerStoppedEvent::class => 'onWorkerStopped',
],
HeartbeatListener::getSubscribedEvents(),
);
}
}
Loading