Skip to content
Draft
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
16 changes: 10 additions & 6 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@ jobs:
strategy:
matrix:
php:
- 8.6
- 8.5
- 8.4
- 8.3
Expand All @@ -29,7 +30,8 @@ jobs:
coverage: ${{ matrix.php < 8.0 && 'xdebug' || 'pcov' }}
ini-file: development
ini-values: disable_functions='' # do not disable PCNTL functions on PHP < 8.1
extensions: sockets, pcntl, event, ${{ matrix.php < 8.0 && 'ev-1.1.5' || 'ev' }}
extensions: sockets, pcntl
# extensions: sockets, pcntl, event, ${{ matrix.php < 8.0 && 'ev-1.1.5' || 'ev' }}
env:
fail-fast: true # fail step if any extension can not be installed
- run: composer install
Expand All @@ -45,6 +47,7 @@ jobs:
strategy:
matrix:
php:
- 8.6
- 8.5
- 8.4
- 8.3
Expand All @@ -63,11 +66,11 @@ jobs:
coverage: ${{ matrix.php < 8.0 && 'xdebug' || 'pcov' }}
ini-file: development
extensions: sockets, pcntl
- name: Install ext-uv
run: |
sudo apt-get update -q && sudo apt-get install libuv1-dev
echo "yes" | sudo pecl install ${{ matrix.php >= 8.0 && 'uv-0.3.0' || 'uv-0.2.4' }}
php -m | grep -q uv || echo "extension=uv.so" >> "$(php -r 'echo php_ini_loaded_file();')"
# - name: Install ext-uv
# run: |
# sudo apt-get update -q && sudo apt-get install libuv1-dev
# echo "yes" | sudo pecl install ${{ matrix.php >= 8.0 && 'uv-0.3.0' || 'uv-0.2.4' }}
# php -m | grep -q uv || echo "extension=uv.so" >> "$(php -r 'echo php_ini_loaded_file();')"
- run: composer install
- run: vendor/bin/phpunit --coverage-text
if: ${{ matrix.php >= 7.3 }}
Expand All @@ -81,6 +84,7 @@ jobs:
strategy:
matrix:
php:
- 8.6
- 8.5
- 8.4
- 8.3
Expand Down
221 changes: 221 additions & 0 deletions src/IoPollLoop.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,221 @@
<?php

namespace React\EventLoop;

use React\EventLoop\Tick\FutureTickQueue;
use React\EventLoop\Timer\Timer;
use React\EventLoop\Timer\Timers;
use SplObjectStorage;

final class IoPollLoop implements LoopInterface
{
/**
* @internal
* This is about 22 years, and the max \Time\Duration` takes for seconds
*/
const MAX_DURATION_SECONDS = 9_223_372_035;

private $running = false;
private $context;
private $futureTickQueue;
private $timers;
private $pcntl = false;
private $pcntlPoll = false;
private $signals;
private $watchers = [];
private $readListeners = [];
private $writeListeners = [];

public function __construct()
{
$this->context = new \Io\Poll\Context();
$this->futureTickQueue = new FutureTickQueue();
$this->timers = new Timers();
$this->pcntl = \function_exists('pcntl_signal') && \function_exists('pcntl_signal_dispatch');
$this->pcntlPoll = $this->pcntl && !\function_exists('pcntl_async_signals');
$this->signals = new SignalsHandler();

// prefer async signals if available (PHP 7.1+) or fall back to dispatching on each tick
if ($this->pcntl && !$this->pcntlPoll) {
\pcntl_async_signals(true);
}
}

public function addReadStream($stream, $listener)
{
$key = (int) $stream;
if (!isset($this->readListeners[$key])) {
$this->readListeners[$key] = $listener;
}
$this->manageStream($key, $stream, \Io\Poll\Event::Read, true);
}

public function addWriteStream($stream, $listener)
{
$key = (int) $stream;
if (!isset($this->writeListeners[$key])) {
$this->writeListeners[$key] = $listener;
}
$this->manageStream($key, $stream, \Io\Poll\Event::Write, true);
}

public function removeReadStream($stream)
{
$key = (int) $stream;
unset($this->readListeners[$key]);
$this->manageStream($key, $stream, \Io\Poll\Event::Read, false);
}

public function removeWriteStream($stream)
{
$key = (int) $stream;
unset($this->writeListeners[$key]);
$this->manageStream($key, $stream, \Io\Poll\Event::Write, false);
}

public function addTimer($interval, $callback)
{
$timer = new Timer($interval, $callback, false);

$this->timers->add($timer);

return $timer;
}

public function addPeriodicTimer($interval, $callback)
{
$timer = new Timer($interval, $callback, true);

$this->timers->add($timer);

return $timer;
}

public function cancelTimer(TimerInterface $timer)
{
$this->timers->cancel($timer);
}

public function futureTick($listener)
{
$this->futureTickQueue->add($listener);
}

public function addSignal($signal, $listener)
{
if ($this->pcntl === false) {
throw new \BadMethodCallException('Event loop feature "signals" isn\'t supported by the "StreamSelectLoop"');
}

$first = $this->signals->count($signal) === 0;
$this->signals->add($signal, $listener);

if ($first) {
\pcntl_signal($signal, [$this->signals, 'call']);
}
}

public function removeSignal($signal, $listener)
{
if (!$this->signals->count($signal)) {
return;
}

$this->signals->remove($signal, $listener);

if ($this->signals->count($signal) === 0) {
\pcntl_signal($signal, \SIG_DFL);
}
}

public function run()
{
$this->running = true;

while ($this->running) {
$this->futureTickQueue->tick();

$this->timers->tick();

// Future-tick queue has pending callbacks ...
if (!$this->futureTickQueue->isEmpty()) {
$timeout = 0;

// There is a pending timer, only block until it is due ...
} elseif ($scheduledAt = $this->timers->getFirst()) {
$timeout = ($scheduledAt - $this->timers->getTime()) / 1_000_000_000;
if ($timeout < 0) {
$timeout = 0;
} else {
$timeout = $timeout > self::MAX_DURATION_SECONDS ? self::MAX_DURATION_SECONDS : $timeout;
}

// The only possible event is stream or signal activity, so wait forever ...
} elseif ($this->readListeners || $this->writeListeners || !$this->signals->isEmpty()) {
$timeout = self::MAX_DURATION_SECONDS;

// There's nothing left to do ...
} else {
break;
}

$seconds = (int)$timeout;
$nanoseconds = (int)(($timeout - $seconds) * 1_000_000_000);
$waitTimeout = \Time\Duration::fromSeconds($seconds, $nanoseconds);

foreach ($this->context->wait($waitTimeout) as $watcher) {
$stream = $watcher->getHandle()->getStream();
$key = $watcher->getData();
$triggeredEvents = $watcher->getTriggeredEvents();

if (array_key_exists($key, $this->readListeners) && in_array(\Io\Poll\Event::Read, $triggeredEvents)) {
\call_user_func($this->readListeners[$key], $stream);
}

if (array_key_exists($key, $this->writeListeners) && in_array(\Io\Poll\Event::Write, $triggeredEvents)) {
\call_user_func($this->writeListeners[$key], $stream);
}
}
}
}

public function stop()
{
$this->running = false;
}

private function manageStream($key, $stream, \Io\Poll\Event $event, bool $add)
{
if (!array_key_exists($key, $this->watchers)) {
if (!$add) {
return;
}

$handle = new \StreamPollHandle($stream);
$this->watchers[$key] = $this->context->add($handle, [$event], $key);

return;
}

if (!$this->watchers[$key]->isActive()) {
unset($this->watchers[$key]);

return;
}

$events = $this->watchers[$key]->getWatchedEvents();
$events = array_filter($events, static fn (\Io\Poll\Event $watchedEvent)=> $watchedEvent === $event);
if ($add) {
$events[] = $event;
}

if (count($events) > 0) {
$this->watchers[$key]->modifyEvents($events);

return;
}

$this->watchers[$key]->remove();
unset($this->watchers[$key]);
}
}
22 changes: 13 additions & 9 deletions src/Loop.php
Original file line number Diff line number Diff line change
Expand Up @@ -238,16 +238,20 @@ public static function stop()
private static function create()
{
// @codeCoverageIgnoreStart
if (\function_exists('uv_loop_new')) {
return new ExtUvLoop();
}

if (\class_exists('EvLoop', false)) {
return new ExtEvLoop();
}
// if (\function_exists('uv_loop_new')) {
// return new ExtUvLoop();
// }
//
// if (\class_exists('EvLoop', false)) {
// return new ExtEvLoop();
// }
//
// if (\class_exists('EventBase', false)) {
// return new ExtEventLoop();
// }

if (\class_exists('EventBase', false)) {
return new ExtEventLoop();
if (\class_exists('Io\Poll\Context', false)) {
return new IoPollLoop();
}

return new StreamSelectLoop();
Expand Down
17 changes: 17 additions & 0 deletions tests/IoPollLoopTest.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,17 @@
<?php

namespace React\Tests\EventLoop;

use React\EventLoop\IoPollLoop;

class IoPollLoopTest extends \React\Tests\EventLoop\AbstractLoopTest
{
public function createLoop()
{
if (!\class_exists('Io\Poll\Context', false)) {
$this->markTestSkipped('IOPollLoop tests skipped because IO Poll is not available.');
}

return new IoPollLoop();
}
}
Loading