From f22cb2ee126966132ac41c9b4684a4e3d1614170 Mon Sep 17 00:00:00 2001 From: Cees-Jan Kiewiet Date: Fri, 25 Sep 2026 14:22:58 +0200 Subject: [PATCH] [3.x] Add Io\Poll Event Loop --- .github/workflows/ci.yml | 16 ++- src/IoPollLoop.php | 231 ++++++++++++++++++++++++++++++++++ src/Loop.php | 22 ++-- src/Watcher/StreamWatcher.php | 15 +++ tests/IoPollLoopTest.php | 17 +++ 5 files changed, 286 insertions(+), 15 deletions(-) create mode 100644 src/IoPollLoop.php create mode 100644 src/Watcher/StreamWatcher.php create mode 100644 tests/IoPollLoopTest.php diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 7bde850f..9efb80f7 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -11,6 +11,7 @@ jobs: strategy: matrix: php: + - 8.6 - 8.5 - 8.4 - 8.3 @@ -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 @@ -45,6 +47,7 @@ jobs: strategy: matrix: php: + - 8.6 - 8.5 - 8.4 - 8.3 @@ -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 }} @@ -81,6 +84,7 @@ jobs: strategy: matrix: php: + - 8.6 - 8.5 - 8.4 - 8.3 diff --git a/src/IoPollLoop.php b/src/IoPollLoop.php new file mode 100644 index 00000000..62ae070a --- /dev/null +++ b/src/IoPollLoop.php @@ -0,0 +1,231 @@ + */ + private $watchers = []; + + 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) + { + $this->manageStream($stream, \Io\Poll\Event::Read, $listener); + } + + public function addWriteStream($stream, $listener) + { + $this->manageStream($stream, \Io\Poll\Event::Write, $listener); + } + + public function removeReadStream($stream) + { + $this->manageStream($stream, \Io\Poll\Event::Read); + } + + public function removeWriteStream($stream) + { + $this->manageStream($stream, \Io\Poll\Event::Write); + } + + 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()) { + $duration = \Time\Duration::fromSeconds(0); + + // There is a pending timer, only block until it is due ... + } elseif ($scheduledAt = $this->timers->getFirst()) { + $timeout = $scheduledAt - $this->timers->getTime(); + if ($timeout < 0) { + $timeout = 0; + } + + $seconds = (int)$timeout; + $nanoseconds = (int)(($timeout - $seconds) * 1_000_000_000); + + $duration = \Time\Duration::fromSeconds($seconds, $nanoseconds); + + // The only possible event is stream or signal activity, so wait forever ... + } elseif ($this->watchers || !$this->signals->isEmpty()) { + $duration = null; + + // There's nothing left to do ... + } else { + break; + } + + if ($this->pcntlPoll) { + \pcntl_signal_dispatch(); + } + + foreach ($this->context->wait($duration) as $watcher) { + $stream = $watcher->getHandle()->getStream(); + $streamWatcher = $watcher->getData(); + $triggeredEvents = $watcher->getTriggeredEvents(); + + if ($streamWatcher->readListener !== null && in_array(\Io\Poll\Event::Read, $triggeredEvents)) { + \call_user_func($streamWatcher->readListener, $stream); + } + + if ($streamWatcher->writeListener !== null && in_array(\Io\Poll\Event::Write, $triggeredEvents)) { + \call_user_func($streamWatcher->writeListener, $stream); + } + } + } + } + + public function stop() + { + $this->running = false; + } + + private function manageStream($stream, \Io\Poll\Event $event, ?callable $listener = null) + { + $key = (int) $stream; + if (!array_key_exists($key, $this->watchers)) { + if ($listener === null) { + return; + } + + $streamWatcher = new StreamWatcher($key); + $this->updateStreamWatcherListener($streamWatcher, $event, $listener); + $handle = new \StreamPollHandle($stream); + $watcher = $this->context->add($handle, [$event], $streamWatcher); + $this->watchers[$key] = $watcher; + + return; + } + + if (!isset($watcher)) { + $watcher = $this->watchers[$key]; + } + + if (!$watcher->isActive()) { + $watcher->remove(); + unset($this->watchers[$key]); + + return; + } + + $events = $watcher->getWatchedEvents(); + $events = array_filter($events, static fn (\Io\Poll\Event $watchedEvent)=> $watchedEvent === $event); + if ($listener !== null) { + $events[] = $event; + } + $this->updateStreamWatcherListener($watcher->getData(), $event, $listener); + + if (count($events) > 0 && ($watcher->getData()->readListener !== null || $watcher->getData()->writeListener !== null)) { + $watcher->modifyEvents($events); + + return; + } + + $watcher->remove(); + unset($this->watchers[$key]); + } + + private function updateStreamWatcherListener(StreamWatcher $streamWatcher, \Io\Poll\Event $event, ?callable $listener = null): void + { + if ($event === \Io\Poll\Event::Read) { + $streamWatcher->readListener = $listener; + } elseif ($event === \Io\Poll\Event::Write) { + $streamWatcher->writeListener = $listener; + } + } +} diff --git a/src/Loop.php b/src/Loop.php index 732c5d5e..da50f442 100644 --- a/src/Loop.php +++ b/src/Loop.php @@ -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(); diff --git a/src/Watcher/StreamWatcher.php b/src/Watcher/StreamWatcher.php new file mode 100644 index 00000000..c38a2f1e --- /dev/null +++ b/src/Watcher/StreamWatcher.php @@ -0,0 +1,15 @@ +markTestSkipped('IOPollLoop tests skipped because IO Poll is not available.'); + } + + return new IoPollLoop(); + } +}