Skip to content

Commit 9fc3d90

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

5 files changed

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

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