From eed4d36fc32dbd10354d32b2efd8a626fb94b86d 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 | 221 +++++++++++++++++++++++++++++++++++++++ src/Loop.php | 22 ++-- tests/IoPollLoopTest.php | 17 +++ 4 files changed, 261 insertions(+), 15 deletions(-) create mode 100644 src/IoPollLoop.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..45ca9227 --- /dev/null +++ b/src/IoPollLoop.php @@ -0,0 +1,221 @@ +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]); + } +} 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/tests/IoPollLoopTest.php b/tests/IoPollLoopTest.php new file mode 100644 index 00000000..9593da32 --- /dev/null +++ b/tests/IoPollLoopTest.php @@ -0,0 +1,17 @@ +markTestSkipped('IOPollLoop tests skipped because IO Poll is not available.'); + } + + return new IoPollLoop(); + } +}