Skip to content

Commit 054ec29

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

4 files changed

Lines changed: 279 additions & 35 deletions

File tree

‎.github/workflows/ci.yml‎

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

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