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
1 change: 1 addition & 0 deletions composer.json
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@
],
"require": {
"php": "~8.2.0 || ~8.3.0 || ~8.4.0 || ~8.5.0",
"psr/clock": "^1.0",
"psr/log": "^1.1 || ^2.0 || ^3.0",
"symfony/event-dispatcher": "^5.4.26 || ^6.4.0 || ^7.0.0 || ^8.0.0"
},
Expand Down
98 changes: 49 additions & 49 deletions composer.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

21 changes: 21 additions & 0 deletions docs/getting-started.md
Original file line number Diff line number Diff line change
Expand Up @@ -85,3 +85,24 @@ Iteration details use the `debug` level; stop reasons (limit exceeded, signal re

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.

## Clock

The worker uses a [PSR-20](https://www.php-fig.org/psr/psr-20/) clock to measure the job's run time
and to check the `timeLimit`. By default it uses the system clock. You can pass your own clock to `create`,
e.g. a mock clock to test time-dependent behaviour without actually waiting:

```php
use Patchlevel\Worker\DefaultWorker;
use Symfony\Component\Clock\MockClock;

$worker = DefaultWorker::create(
$job,
['timeLimit' => 3600],
$logger,
clock: new MockClock(),
);
```
:::note
The clock is only used for measuring time. The sleep between iterations still waits for real.
:::
22 changes: 14 additions & 8 deletions src/DefaultWorker.php
Original file line number Diff line number Diff line change
Expand Up @@ -12,34 +12,33 @@
use Patchlevel\Worker\Listener\StopWorkerOnMemoryLimitListener;
use Patchlevel\Worker\Listener\StopWorkerOnSignalListener;
use Patchlevel\Worker\Listener\StopWorkerOnTimeLimitListener;
use Psr\Clock\ClockInterface;
use Psr\Log\LoggerInterface;
use Psr\Log\NullLogger;
use Symfony\Component\EventDispatcher\EventDispatcher;
use Symfony\Component\EventDispatcher\EventDispatcherInterface;

use function max;
use function microtime;
use function round;
use function usleep;

final class DefaultWorker implements Worker
{
private bool $shouldStop = false;

/** @var Closure():int */
private Closure $timeMeasure;
private readonly ClockInterface $clock;

/** @param Closure(Closure):void $job */
public function __construct(
private readonly Closure $job,
private readonly EventDispatcherInterface $eventDispatcher,
private readonly LoggerInterface|null $logger = null,
ClockInterface|null $clock = null,
) {
$this->timeMeasure = static fn () => (int)round(microtime(true) * 1000);
$this->clock = $clock ?? new SystemClock();
}

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

Check warning on line 41 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 41 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 @@ -48,11 +47,11 @@
while (!$this->shouldStop) {
$this->logger?->debug('Worker starting job run');

$startTime = ($this->timeMeasure)();
$startTime = $this->milliseconds();

($this->job)($this->stop(...));

$endTime = ($this->timeMeasure)();
$endTime = $this->milliseconds();
$ranTime = $endTime - $startTime;

$this->logger?->debug('Worker finished job run ({ranTime}ms)', ['ranTime' => $ranTime]);
Expand All @@ -60,17 +59,17 @@
$this->eventDispatcher->dispatch(new WorkerRunningEvent($this));

if ($this->shouldStop) {
break;

Check warning on line 62 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 65 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;
}

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

Check warning on line 72 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); + } $this->shouldStop = false;

Check warning on line 72 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); } $this->shouldStop = false;

Check warning on line 72 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); } $this->shouldStop = false;
}

$this->shouldStop = false;
Expand All @@ -97,6 +96,7 @@
array $options = [],
LoggerInterface $logger = new NullLogger(),
EventDispatcherInterface|null $eventDispatcher = null,
ClockInterface|null $clock = null,
): self {
if ($eventDispatcher === null) {
$eventDispatcher = new EventDispatcher();
Expand All @@ -118,14 +118,20 @@

if (isset($options['timeLimit'])) {
$eventDispatcher->addSubscriber(
new StopWorkerOnTimeLimitListener($options['timeLimit'], $logger),
new StopWorkerOnTimeLimitListener($options['timeLimit'], $logger, $clock),
);
}

return new self(
$job,
$eventDispatcher,
$logger,
$clock,
);
}

private function milliseconds(): int
{
return (int)$this->clock->now()->format('Uv');
}
}
16 changes: 11 additions & 5 deletions src/Listener/StopWorkerOnTimeLimitListener.php
Original file line number Diff line number Diff line change
Expand Up @@ -4,32 +4,38 @@

namespace Patchlevel\Worker\Listener;

use DateInterval;
use DateTimeImmutable;
use Patchlevel\Worker\Event\WorkerRunningEvent;
use Patchlevel\Worker\Event\WorkerStartedEvent;
use Patchlevel\Worker\SystemClock;
use Psr\Clock\ClockInterface;
use Psr\Log\LoggerInterface;
use Symfony\Component\EventDispatcher\EventSubscriberInterface;

use function time;

final class StopWorkerOnTimeLimitListener implements EventSubscriberInterface
{
private float $endTime = 0;
private DateTimeImmutable|null $endTime = null;

private readonly ClockInterface $clock;

/** @param positive-int $timeLimit in seconds */
public function __construct(
private readonly int $timeLimit,
private readonly LoggerInterface|null $logger = null,
ClockInterface|null $clock = null,
) {
$this->clock = $clock ?? new SystemClock();
}

public function onWorkerStarted(): void
{
$this->endTime = time() + $this->timeLimit;
$this->endTime = $this->clock->now()->add(new DateInterval('PT' . $this->timeLimit . 'S'));
}

public function onWorkerRunning(WorkerRunningEvent $event): void
{
if ($this->endTime >= time()) {
if ($this->endTime !== null && $this->clock->now() < $this->endTime) {
return;
}

Expand Down
16 changes: 16 additions & 0 deletions src/SystemClock.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,16 @@
<?php

declare(strict_types=1);

namespace Patchlevel\Worker;

use DateTimeImmutable;
use Psr\Clock\ClockInterface;

final class SystemClock implements ClockInterface
{
public function now(): DateTimeImmutable
{
return new DateTimeImmutable();
}
}
36 changes: 36 additions & 0 deletions tests/TestClock.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,36 @@
<?php

declare(strict_types=1);

namespace Patchlevel\Worker\Tests;

use DateTimeImmutable;
use Psr\Clock\ClockInterface;

use function intdiv;
use function sprintf;

final class TestClock implements ClockInterface
{
private int $milliseconds = 1_767_225_600_000;

/** @param int $tick milliseconds the clock moves forward on every now() call */
public function __construct(
private readonly int $tick = 0,
) {
}

public function now(): DateTimeImmutable
{
$now = new DateTimeImmutable(sprintf('@%d.%03d', intdiv($this->milliseconds, 1000), $this->milliseconds % 1000));

$this->milliseconds += $this->tick;

return $now;
}

public function advance(int $milliseconds): void
{
$this->milliseconds += $milliseconds;
}
}
Loading
Loading