Low-budget Event Loop In PHP using Fibers.

<?php
declare(strict_types=1);
final class Loop
{
private static ?self $instance = null;
private array $ready = [];
private array $timers = [];
private array $readers = [];
private array $writers = [];
private int $nextTimerId = 1;
private ?Closure $exceptionHandler = null;
private function __construct() {}
private static function get(): self
{
return self::$instance ??= new self;
}
public static function spawn(Closure $callback): void
{
$loop = self::get();
$loop->ready[] = [
'fiber' => new Fiber($callback),
'value' => null,
'exception' => null,
];
}
public static function timer(Closure $callback, float $seconds = 0, bool $periodic = false): Closure
{
if ($seconds < 0) {
throw new InvalidArgumentException('Timer seconds cannot be negative.');
}
$loop = self::get();
$id = $loop->nextTimerId++;
$delayNs = (int) round($seconds * 1_000_000_000);
$loop->timers[$id] = [
'time' => hrtime(true) + $delayNs,
'callback' => $callback,
'interval' => $periodic ? $delayNs : null,
];
return function () use ($id, $loop): void {
unset($loop->timers[$id]);
};
}
public static function sleep(float $seconds): void
{
if ($seconds < 0) {
throw new InvalidArgumentException('Sleep seconds cannot be negative.');
}
Fiber::suspend([
'type' => 'sleep',
'payload' => [
'delay' => (int) round($seconds * 1_000_000_000),
],
]);
}
/**
* @param resource $stream
* @return resource|false
*/
public static function accept($stream, ?float $timeout = null)
{
return Fiber::suspend([
'type' => 'reader',
'payload' => [
'operation' => 'accept',
'stream' => $stream,
'timeout' => $timeout,
],
]);
}
/** @param resource $stream */
public static function read($stream, int $length): string|false
{
return Fiber::suspend([
'type' => 'reader',
'payload' => [
'operation' => 'read',
'stream' => $stream,
'length' => $length,
],
]);
}
/** @param resource $stream */
public static function write($stream, string $data, ?int $length = null): int|false
{
return Fiber::suspend([
'type' => 'writer',
'payload' => [
'stream' => $stream,
'data' => $data,
'length' => $length,
],
]);
}
public static function onException(Closure $handler): void
{
$loop = self::get();
$loop->exceptionHandler = $handler;
}
public static function run(): void
{
$loop = self::get();
while (true) {
/*
* Expire timers before running ready Fibers.
*/
$now = hrtime(true);
foreach ($loop->timers as $id => $timer) {
if ($timer['time'] > $now) {
continue;
}
$loop->ready[] = [
'fiber' => new Fiber($timer['callback']),
'value' => null,
'exception' => null,
];
if ($timer['interval'] === null) {
unset($loop->timers[$id]);
continue;
}
if ($timer['interval'] === 0) {
$timer['time'] = $now;
} else {
do {
$timer['time'] += $timer['interval'];
} while ($timer['time'] <= $now);
}
$loop->timers[$id] = $timer;
}
/*
* Run ready Fibers.
*/
while ($loop->ready) {
/** @var array{fiber: Fiber, value: mixed, exception: ?Throwable} $task */
$task = array_shift($loop->ready);
$fiber = $task['fiber'];
try {
if (! $fiber->isStarted()) {
$syscall = $fiber->start();
} elseif ($task['exception'] !== null) {
$syscall = $fiber->throw($task['exception']);
} else {
$syscall = $fiber->resume($task['value']);
}
} catch (Throwable $e) {
if ($loop->exceptionHandler !== null) {
($loop->exceptionHandler)($e, $fiber);
continue;
}
throw $e;
}
if ($fiber->isTerminated()) {
continue;
}
if (! is_array($syscall) || ! isset($syscall['type'], $syscall['payload']) || ! is_string($syscall['type']) || ! is_array($syscall['payload'])) {
$loop->ready[] = [
'fiber' => $fiber,
'value' => null,
'exception' => new RuntimeException('Invalid syscall.'),
];
continue;
}
$payload = $syscall['payload'];
switch ($syscall['type']) {
case 'sleep':
$delay = $payload['delay'];
$id = $loop->nextTimerId++;
$loop->timers[$id] = [
'time' => hrtime(true) + $delay,
'callback' => function () use ($loop, $fiber): void {
$loop->ready[] = [
'fiber' => $fiber,
'value' => null,
'exception' => null,
];
},
'interval' => null,
];
break;
case 'reader':
$stream = $payload['stream'];
$id = (int) $stream;
if (isset($loop->readers[$id])) {
$loop->ready[] = [
'fiber' => $fiber,
'value' => null,
'exception' => new RuntimeException('A read operation is already pending for this stream.'),
];
break;
}
$loop->readers[$id] = [
'stream' => $stream,
'fiber' => $fiber,
'operation' => $payload['operation'],
'length' => $payload['length'] ?? null,
'timeout' => $payload['timeout'] ?? null,
];
break;
case 'writer':
$stream = $payload['stream'];
$id = (int) $stream;
if (isset($loop->writers[$id])) {
$loop->ready[] = [
'fiber' => $fiber,
'value' => null,
'exception' => new RuntimeException('A write operation is already pending for this stream.'),
];
break;
}
$loop->writers[$id] = [
'stream' => $stream,
'fiber' => $fiber,
'data' => $payload['data'],
'length' => $payload['length'],
];
break;
default:
$loop->ready[] = [
'fiber' => $fiber,
'value' => null,
'exception' => new RuntimeException("Unknown syscall: {$syscall['type']}."),
];
break;
}
}
/*
* Nothing is waiting.
*/
if (! $loop->timers && ! $loop->readers && ! $loop->writers) {
break;
}
$read = [];
$write = [];
$except = null;
$timeout = null;
foreach ($loop->readers as $watcher) {
$read[] = $watcher['stream'];
}
foreach ($loop->writers as $watcher) {
$write[] = $watcher['stream'];
}
/*
* Find the nearest timer.
*/
if ($loop->timers) {
$now = hrtime(true);
$next = PHP_INT_MAX;
foreach ($loop->timers as $timer) {
$next = min($next, $timer['time']);
}
$remaining = max(0, $next - $now);
$timeout = max(1, intdiv($remaining, 1_000));
}
/*
* If only timers are pending.
*/
if (! $read && ! $write) {
if ($timeout !== null && $timeout > 0) {
$seconds = intdiv($timeout, 1_000_000);
$microseconds = $timeout % 1_000_000;
usleep(($seconds * 1_000_000) + $microseconds);
}
continue;
}
/*
* Wait for I/O or the next timer.
*/
if ($timeout === null) {
$seconds = null;
$microseconds = null;
} else {
$seconds = intdiv($timeout, 1_000_000);
$microseconds = $timeout % 1_000_000;
}
$result = @stream_select($read, $write, $except, $seconds, $microseconds);
if ($result === false) {
continue;
}
/*
* Handle readable streams.
*/
foreach ($read as $stream) {
$id = (int) $stream;
if (! isset($loop->readers[$id])) {
continue;
}
$watcher = $loop->readers[$id];
unset($loop->readers[$id]);
switch ($watcher['operation']) {
case 'accept':
$client = @stream_socket_accept($stream, $watcher['timeout']);
if ($client !== false) {
stream_set_blocking($client, false);
}
$loop->ready[] = [
'fiber' => $watcher['fiber'],
'value' => $client,
'exception' => null,
];
break;
case 'read':
$data = @fread($stream, $watcher['length']);
$loop->ready[] = [
'fiber' => $watcher['fiber'],
'value' => $data,
'exception' => null,
];
break;
default:
$loop->ready[] = [
'fiber' => $watcher['fiber'],
'value' => null,
'exception' => new RuntimeException("Unknown reader operation: {$watcher['operation']}."),
];
break;
}
}
/*
* Handle writable streams.
*/
foreach ($write as $stream) {
$id = (int) $stream;
if (! isset($loop->writers[$id])) {
continue;
}
$watcher = $loop->writers[$id];
unset($loop->writers[$id]);
$written = @fwrite($stream, $watcher['data'], $watcher['length']);
$loop->ready[] = [
'fiber' => $watcher['fiber'],
'value' => $written,
'exception' => null,
];
}
}
}
}
Loop::timer(function () {
static $count = 0;
echo "Tick {$count}!\n";
$count++;
}, 0, true);
Loop::timer(function () {
echo "5 seconds passed.\n";
}, 5);
Loop::spawn(function () {
$server = stream_socket_server('tcp://127.0.0.1:8888');
stream_set_blocking($server, false);
echo "Server listening on http://127.0.0.1:8888\n";
while (true) {
$client = Loop::accept($server);
if (! $client) {
continue;
}
Loop::spawn(function () use ($client) {
$request = Loop::read($client, 8192);
echo "Received: {$request}\n";
$body = str_repeat("Hello, World!\n", 10);
Loop::write(
$client,
"HTTP/1.1 200 OK\r\n".
'Content-Length: '.mb_strlen($body)."\r\n".
"Content-Type: text/plain; charset=utf-8\r\n".
"Connection: close\r\n".
"\r\n".
$body
);
fclose($client);
});
}
});
Loop::run();
添加评论
点赞收藏
点踩分享查看原文
评论
?
参与讨论