diff --git a/docs/getting-started.md b/docs/getting-started.md index e6948a6..b319e7f 100644 --- a/docs/getting-started.md +++ b/docs/getting-started.md @@ -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. @@ -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. diff --git a/src/DefaultWorker.php b/src/DefaultWorker.php index 00374fe..44e565e 100644 --- a/src/DefaultWorker.php +++ b/src/DefaultWorker.php @@ -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; @@ -99,8 +100,8 @@ public function stop(): void } /** - * @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, @@ -133,6 +134,12 @@ public static function create( ); } + if (isset($options['heartbeatFile'])) { + $eventDispatcher->addSubscriber( + new HeartbeatListener($options['heartbeatFile'], $logger), + ); + } + return new self( $job, $eventDispatcher, diff --git a/src/Listener/HeartbeatListener.php b/src/Listener/HeartbeatListener.php new file mode 100644 index 0000000..25bd256 --- /dev/null +++ b/src/Listener/HeartbeatListener.php @@ -0,0 +1,63 @@ +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 */ + public static function getSubscribedEvents(): array + { + return [ + WorkerStartedEvent::class => 'onWorkerStarted', + WorkerRunningEvent::class => 'onWorkerRunning', + WorkerStoppedEvent::class => 'onWorkerStopped', + ]; + } +} diff --git a/tests/Unit/DefaultWorkerTest.php b/tests/Unit/DefaultWorkerTest.php index 3c1f317..a3c4498 100644 --- a/tests/Unit/DefaultWorkerTest.php +++ b/tests/Unit/DefaultWorkerTest.php @@ -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; @@ -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); } @@ -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); @@ -170,6 +172,7 @@ static function ($stop): void { 'runLimit' => 10, 'memoryLimit' => '10KB', 'timeLimit' => 20, + 'heartbeatFile' => '/tmp/worker-heartbeat', ], $logger, $eventDispatcher, diff --git a/tests/Unit/Listener/HeartbeatListenerTest.php b/tests/Unit/Listener/HeartbeatListenerTest.php new file mode 100644 index 0000000..9298a21 --- /dev/null +++ b/tests/Unit/Listener/HeartbeatListenerTest.php @@ -0,0 +1,118 @@ +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(), + ); + } +}