Skip to content

Commit eed4d36

Browse files
committed
[3.x] Add Io\Poll Event Loop
1 parent 0b45df3 commit eed4d36

4 files changed

Lines changed: 261 additions & 15 deletions

File tree

‎.github/workflows/ci.yml‎

Lines changed: 10 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -11,6 +11,7 @@ jobs:
1111
strategy:
1212
matrix:
1313
php:
14+
- 8.6
1415
- 8.5
1516
- 8.4
1617
- 8.3
@@ -29,7 +30,8 @@ jobs:
2930
coverage: ${{ matrix.php < 8.0 && 'xdebug' || 'pcov' }}
3031
ini-file: development
3132
ini-values: disable_functions='' # do not disable PCNTL functions on PHP < 8.1
32-
extensions: sockets, pcntl, event, ${{ matrix.php < 8.0 && 'ev-1.1.5' || 'ev' }}
33+
extensions: sockets, pcntl
34+
# extensions: sockets, pcntl, event, ${{ matrix.php < 8.0 && 'ev-1.1.5' || 'ev' }}
3335
env:
3436
fail-fast: true # fail step if any extension can not be installed
3537
- run: composer install
@@ -45,6 +47,7 @@ jobs:
4547
strategy:
4648
matrix:
4749
php:
50+
- 8.6
4851
- 8.5
4952
- 8.4
5053
- 8.3
@@ -63,11 +66,11 @@ jobs:
6366
coverage: ${{ matrix.php < 8.0 && 'xdebug' || 'pcov' }}
6467
ini-file: development
6568
extensions: sockets, pcntl
66-
- name: Install ext-uv
67-
run: |
68-
sudo apt-get update -q && sudo apt-get install libuv1-dev
69-
echo "yes" | sudo pecl install ${{ matrix.php >= 8.0 && 'uv-0.3.0' || 'uv-0.2.4' }}
70-
php -m | grep -q uv || echo "extension=uv.so" >> "$(php -r 'echo php_ini_loaded_file();')"
69+
# - name: Install ext-uv
70+
# run: |
71+
# sudo apt-get update -q && sudo apt-get install libuv1-dev
72+
# echo "yes" | sudo pecl install ${{ matrix.php >= 8.0 && 'uv-0.3.0' || 'uv-0.2.4' }}
73+
# php -m | grep -q uv || echo "extension=uv.so" >> "$(php -r 'echo php_ini_loaded_file();')"
7174
- run: composer install
7275
- run: vendor/bin/phpunit --coverage-text
7376
if: ${{ matrix.php >= 7.3 }}
@@ -81,6 +84,7 @@ jobs:
8184
strategy:
8285
matrix:
8386
php:
87+
- 8.6
8488
- 8.5
8589
- 8.4
8690
- 8.3

‎src/IoPollLoop.php‎

Lines changed: 221 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,221 @@
1+
<?php
2+
3+
namespace React\EventLoop;
4+
5+
use React\EventLoop\Tick\FutureTickQueue;
6+
use React\EventLoop\Timer\Timer;
7+
use React\EventLoop\Timer\Timers;
8+
use SplObjectStorage;
9+
10+
final class IoPollLoop implements LoopInterface
11+
{
12+
/**
13+
* @internal
14+
* This is about 22 years, and the max \Time\Duration` takes for seconds
15+
*/
16+
const MAX_DURATION_SECONDS = 9_223_372_035;
17+
18+
private $running = false;
19+
private $context;
20+
private $futureTickQueue;
21+
private $timers;
22+
private $pcntl = false;
23+
private $pcntlPoll = false;
24+
private $signals;
25+
private $watchers = [];
26+
private $readListeners = [];
27+
private $writeListeners = [];
28+
29+
public function __construct()
30+
{
31+
$this->context = new \Io\Poll\Context();
32+
$this->futureTickQueue = new FutureTickQueue();
33+
$this->timers = new Timers();
34+
$this->pcntl = \function_exists('pcntl_signal') && \function_exists('pcntl_signal_dispatch');
35+
$this->pcntlPoll = $this->pcntl && !\function_exists('pcntl_async_signals');
36+
$this->signals = new SignalsHandler();
37+
38+
// prefer async signals if available (PHP 7.1+) or fall back to dispatching on each tick
39+
if ($this->pcntl && !$this->pcntlPoll) {
40+
\pcntl_async_signals(true);
41+
}
42+
}
43+
44+
public function addReadStream($stream, $listener)
45+
{
46+
$key = (int) $stream;
47+
if (!isset($this->readListeners[$key])) {
48+
$this->readListeners[$key] = $listener;
49+
}
50+
$this->manageStream($key, $stream, \Io\Poll\Event::Read, true);
51+
}
52+
53+
public function addWriteStream($stream, $listener)
54+
{
55+
$key = (int) $stream;
56+
if (!isset($this->writeListeners[$key])) {
57+
$this->writeListeners[$key] = $listener;
58+
}
59+
$this->manageStream($key, $stream, \Io\Poll\Event::Write, true);
60+
}
61+
62+
public function removeReadStream($stream)
63+
{
64+
$key = (int) $stream;
65+
unset($this->readListeners[$key]);
66+
$this->manageStream($key, $stream, \Io\Poll\Event::Read, false);
67+
}
68+
69+
public function removeWriteStream($stream)
70+
{
71+
$key = (int) $stream;
72+
unset($this->writeListeners[$key]);
73+
$this->manageStream($key, $stream, \Io\Poll\Event::Write, false);
74+
}
75+
76+
public function addTimer($interval, $callback)
77+
{
78+
$timer = new Timer($interval, $callback, false);
79+
80+
$this->timers->add($timer);
81+
82+
return $timer;
83+
}
84+
85+
public function addPeriodicTimer($interval, $callback)
86+
{
87+
$timer = new Timer($interval, $callback, true);
88+
89+
$this->timers->add($timer);
90+
91+
return $timer;
92+
}
93+
94+
public function cancelTimer(TimerInterface $timer)
95+
{
96+
$this->timers->cancel($timer);
97+
}
98+
99+
public function futureTick($listener)
100+
{
101+
$this->futureTickQueue->add($listener);
102+
}
103+
104+
public function addSignal($signal, $listener)
105+
{
106+
if ($this->pcntl === false) {
107+
throw new \BadMethodCallException('Event loop feature "signals" isn\'t supported by the "StreamSelectLoop"');
108+
}
109+
110+
$first = $this->signals->count($signal) === 0;
111+
$this->signals->add($signal, $listener);
112+
113+
if ($first) {
114+
\pcntl_signal($signal, [$this->signals, 'call']);
115+
}
116+
}
117+
118+
public function removeSignal($signal, $listener)
119+
{
120+
if (!$this->signals->count($signal)) {
121+
return;
122+
}
123+
124+
$this->signals->remove($signal, $listener);
125+
126+
if ($this->signals->count($signal) === 0) {
127+
\pcntl_signal($signal, \SIG_DFL);
128+
}
129+
}
130+
131+
public function run()
132+
{
133+
$this->running = true;
134+
135+
while ($this->running) {
136+
$this->futureTickQueue->tick();
137+
138+
$this->timers->tick();
139+
140+
// Future-tick queue has pending callbacks ...
141+
if (!$this->futureTickQueue->isEmpty()) {
142+
$timeout = 0;
143+
144+
// There is a pending timer, only block until it is due ...
145+
} elseif ($scheduledAt = $this->timers->getFirst()) {
146+
$timeout = ($scheduledAt - $this->timers->getTime()) / 1_000_000_000;
147+
if ($timeout < 0) {
148+
$timeout = 0;
149+
} else {
150+
$timeout = $timeout > self::MAX_DURATION_SECONDS ? self::MAX_DURATION_SECONDS : $timeout;
151+
}
152+
153+
// The only possible event is stream or signal activity, so wait forever ...
154+
} elseif ($this->readListeners || $this->writeListeners || !$this->signals->isEmpty()) {
155+
$timeout = self::MAX_DURATION_SECONDS;
156+
157+
// There's nothing left to do ...
158+
} else {
159+
break;
160+
}
161+
162+
$seconds = (int)$timeout;
163+
$nanoseconds = (int)(($timeout - $seconds) * 1_000_000_000);
164+
$waitTimeout = \Time\Duration::fromSeconds($seconds, $nanoseconds);
165+
166+
foreach ($this->context->wait($waitTimeout) as $watcher) {
167+
$stream = $watcher->getHandle()->getStream();
168+
$key = $watcher->getData();
169+
$triggeredEvents = $watcher->getTriggeredEvents();
170+
171+
if (array_key_exists($key, $this->readListeners) && in_array(\Io\Poll\Event::Read, $triggeredEvents)) {
172+
\call_user_func($this->readListeners[$key], $stream);
173+
}
174+
175+
if (array_key_exists($key, $this->writeListeners) && in_array(\Io\Poll\Event::Write, $triggeredEvents)) {
176+
\call_user_func($this->writeListeners[$key], $stream);
177+
}
178+
}
179+
}
180+
}
181+
182+
public function stop()
183+
{
184+
$this->running = false;
185+
}
186+
187+
private function manageStream($key, $stream, \Io\Poll\Event $event, bool $add)
188+
{
189+
if (!array_key_exists($key, $this->watchers)) {
190+
if (!$add) {
191+
return;
192+
}
193+
194+
$handle = new \StreamPollHandle($stream);
195+
$this->watchers[$key] = $this->context->add($handle, [$event], $key);
196+
197+
return;
198+
}
199+
200+
if (!$this->watchers[$key]->isActive()) {
201+
unset($this->watchers[$key]);
202+
203+
return;
204+
}
205+
206+
$events = $this->watchers[$key]->getWatchedEvents();
207+
$events = array_filter($events, static fn (\Io\Poll\Event $watchedEvent)=> $watchedEvent === $event);
208+
if ($add) {
209+
$events[] = $event;
210+
}
211+
212+
if (count($events) > 0) {
213+
$this->watchers[$key]->modifyEvents($events);
214+
215+
return;
216+
}
217+
218+
$this->watchers[$key]->remove();
219+
unset($this->watchers[$key]);
220+
}
221+
}

‎src/Loop.php‎

Lines changed: 13 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -238,16 +238,20 @@ public static function stop()
238238
private static function create()
239239
{
240240
// @codeCoverageIgnoreStart
241-
if (\function_exists('uv_loop_new')) {
242-
return new ExtUvLoop();
243-
}
244-
245-
if (\class_exists('EvLoop', false)) {
246-
return new ExtEvLoop();
247-
}
241+
// if (\function_exists('uv_loop_new')) {
242+
// return new ExtUvLoop();
243+
// }
244+
//
245+
// if (\class_exists('EvLoop', false)) {
246+
// return new ExtEvLoop();
247+
// }
248+
//
249+
// if (\class_exists('EventBase', false)) {
250+
// return new ExtEventLoop();
251+
// }
248252

249-
if (\class_exists('EventBase', false)) {
250-
return new ExtEventLoop();
253+
if (\class_exists('Io\Poll\Context', false)) {
254+
return new IoPollLoop();
251255
}
252256

253257
return new StreamSelectLoop();

‎tests/IoPollLoopTest.php‎

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,17 @@
1+
<?php
2+
3+
namespace React\Tests\EventLoop;
4+
5+
use React\EventLoop\IoPollLoop;
6+
7+
class IoPollLoopTest extends \React\Tests\EventLoop\AbstractLoopTest
8+
{
9+
public function createLoop()
10+
{
11+
if (!\class_exists('Io\Poll\Context', false)) {
12+
$this->markTestSkipped('IOPollLoop tests skipped because IO Poll is not available.');
13+
}
14+
15+
return new IoPollLoop();
16+
}
17+
}

0 commit comments

Comments
 (0)