Skip to content

Commit 6c7696d

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

4 files changed

Lines changed: 254 additions & 9 deletions

File tree

‎.github/workflows/ci.yml‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -9,8 +9,10 @@ jobs:
99
name: PHPUnit (PHP ${{ matrix.php }})
1010
runs-on: ubuntu-24.04
1111
strategy:
12+
fail-fast: false
1213
matrix:
1314
php:
15+
- 8.6
1416
- 8.5
1517
- 8.4
1618
- 8.3
@@ -43,8 +45,10 @@ jobs:
4345
runs-on: ubuntu-24.04
4446
continue-on-error: true
4547
strategy:
48+
fail-fast: false
4649
matrix:
4750
php:
51+
- 8.6
4852
- 8.5
4953
- 8.4
5054
- 8.3
@@ -81,6 +85,7 @@ jobs:
8185
strategy:
8286
matrix:
8387
php:
88+
- 8.6
8489
- 8.5
8590
- 8.4
8691
- 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)