Skip to content

Commit f22cb2e

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

5 files changed

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

‎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();

‎src/Watcher/StreamWatcher.php‎

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,15 @@
1+
<?php
2+
3+
namespace React\EventLoop\Watcher;
4+
5+
final class StreamWatcher
6+
{
7+
public function __construct(
8+
public readonly int $key,
9+
/** @var ?callable */
10+
public $readListener = null,
11+
/** @var ?callable */
12+
public $writeListener = null,
13+
) {
14+
}
15+
}

‎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)