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
27 changes: 27 additions & 0 deletions docs/getting-started.md
Original file line number Diff line number Diff line change
Expand Up @@ -81,6 +81,33 @@ The job's own run time is subtracted from the sleep: if the job took 300ms and t
is 500ms, the worker only sleeps 200ms. If the job took longer than the sleep timer,
the next iteration starts immediately. Pass `0` to disable sleeping entirely.

### Skip the sleep when there is work

If the job returns `true`, the worker skips the sleep and starts the next iteration immediately.
This is useful for queue consumers: as long as there are messages, they are processed without a pause,
and the worker only sleeps once the queue is empty.

```php
use Patchlevel\Worker\DefaultWorker;

$worker = DefaultWorker::create(
static function (callable $stop) use ($queue): bool {
$message = $queue->pop();

if ($message === null) {
return false; // nothing to do, sleep
}

handle($message);

return true; // there may be more, continue immediately
},
);

$worker->run(1000);
```
Any other return value, including no return value at all, keeps the regular sleep behaviour.

## Logging

The worker logs its lifecycle (start, iteration timings, sleep, stop reason) to the given PSR-3 logger.
Expand Down
10 changes: 7 additions & 3 deletions src/DefaultWorker.php
Original file line number Diff line number Diff line change
Expand Up @@ -28,7 +28,7 @@

private readonly ClockInterface $clock;

/** @param Closure(Closure):void $job */
/** @param Closure(Closure):(bool|void) $job return true if the job did work to skip the sleep */

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

isnt it more true|void ? :D

@DavidBadura DavidBadura Oct 2, 2026 •

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Only true changes the behaviour, that's right. I'd still keep bool though, for two reasons.

false reads as an explicit "nothing to do, sleep", which makes the intent clearer than mixing return; and return true;:

$worker = DefaultWorker::create(
    static function (callable $stop) use ($queue): bool {
        $message = $queue->pop();

        if ($message === null) {
            return false; // queue is empty, sleep
        }

        handle($message);

        return true; // there may be more, continue immediately
    },
);

And it lets a job pass a bool result straight through, which is common for consumer APIs:

// with bool|void
static fn (callable $stop): bool => $consumer->consumeOne();

// with true|void
static function (callable $stop) use ($consumer) {
    if ($consumer->consumeOne()) {
        return true;
    }
};

With true|void, PHPStan rejects the first closure because false is not allowed, so the result has to be converted by hand. That needs an extra branch just to drop false, and the closure can't have a native return type anymore: true|void isn't valid PHP, and true|null would throw a TypeError on the path without a return, unless you add an explicit return null.

public function __construct(
private readonly Closure $job,
private readonly EventDispatcherInterface $eventDispatcher,
Expand All @@ -39,7 +39,7 @@
}

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

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');
{
$this->logger?->debug('Worker starting');

Expand All @@ -53,7 +53,7 @@

$startTime = $this->milliseconds();

($this->job)($this->stop(...));
$didWork = ($this->job)($this->stop(...)) === true;

$endTime = $this->milliseconds();
$ranTime = $endTime - $startTime;
Expand All @@ -63,17 +63,21 @@
$this->eventDispatcher->dispatch(new WorkerRunningEvent($this));

if ($this->shouldStop) {
break;

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; } if ($didWork) {
}

if ($didWork) {
continue;
}

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

Check warning on line 73 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 80 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 80 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 80 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;
Expand All @@ -95,7 +99,7 @@
}

/**
* @param Closure(Closure):void $job
* @param Closure(Closure):(bool|void) $job
* @param array{runLimit?: (positive-int|null), memoryLimit?: (string|null), timeLimit?: (positive-int|null)} $options
*/
public static function create(
Expand Down
41 changes: 41 additions & 0 deletions tests/Unit/DefaultWorkerTest.php
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,8 @@
use Symfony\Component\EventDispatcher\EventDispatcher;
use Symfony\Component\EventDispatcher\EventDispatcherInterface;

use function array_shift;

#[CoversClass(DefaultWorker::class)]
final class DefaultWorkerTest extends TestCase
{
Expand Down Expand Up @@ -243,6 +245,45 @@ public function testRunWorkerNotSleeping(): void
$worker->run(5);
}

public function testRunWorkerSkipsSleepWhenJobDidWork(): void
{
$logger = $this->createMock(LoggerInterface::class);
$logger
->expects($this->exactly(11))
->method('debug')
->willReturnCallback(
new ReturnCallback([
[['Worker starting', []]],
[['Worker starting job run', []]],
[['Worker finished job run ({ranTime}ms)', ['ranTime' => 10]]],
[['Worker starting job run', []]],
[['Worker finished job run ({ranTime}ms)', ['ranTime' => 10]]],
[['Worker sleep for {sleepTimer}ms', ['sleepTimer' => 190]]],
[['Worker starting job run', []]],
[['Worker finished job run ({ranTime}ms)', ['ranTime' => 10]]],
[['Worker received stop signal', []]],
[['Worker stopped', []]],
[['Worker terminated', []]],
]),
);

$eventDispatcher = new EventDispatcher();
$eventDispatcher->addSubscriber(new StopWorkerOnIterationLimitListener(3));

$jobResults = [true, false, true];

$worker = new DefaultWorker(
static function () use (&$jobResults): bool {
return (bool)array_shift($jobResults);
},
$eventDispatcher,
$logger,
new TestClock(tick: 10),
);

$worker->run(200);
}

public function testDefaultCreate(): void
{
$calls = 0;
Expand Down
Loading