Skip to content
Open
Show file tree
Hide file tree
Changes from 12 commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,7 @@ jobs:

strategy:
matrix:
php: [ '8.1', '8.2', '8.3' ]
php: [ '8.1', '8.2', '8.3', '8.4' ]

steps:
- name: 'Init repository'
Expand Down
14 changes: 13 additions & 1 deletion .php-cs-fixer.dist.php
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,19 @@
]
);

$config = new M6Web\CS\Config\BedrockStreaming();
$baseConfig = new M6Web\CS\Config\BedrockStreaming();

$config = new PhpCsFixer\Config('Bedrock Streaming');
$config->setFinder($finder);

$override_rules = array_merge(
$baseConfig->getRules(),
[
// Adding strict_types should be part of another PR
'declare_strict_types' => false,
]
);

$config->setRules($override_rules);

return $config;
4 changes: 2 additions & 2 deletions composer.json
Original file line number Diff line number Diff line change
Expand Up @@ -25,13 +25,13 @@
},
"require-dev": {
"phpunit/phpunit": "^10.5",
"amphp/amp": "^2.0",
"amphp/amp": "^3.0",
"guzzlehttp/guzzle": "^7.4",
"m6web/php-cs-fixer-config": "^2.0",
"ext-curl": "^8.1",
"react/event-loop": "^1.0",
"react/promise": "^2.7",
"phpstan/phpstan": "^1.0",
"phpstan/phpstan": "^1.10",
"symfony/http-client": "^6.4",
"psr/http-factory": "^1.0",
"http-interop/http-factory-guzzle": "^1.0"
Expand Down
153 changes: 81 additions & 72 deletions src/Adapter/Amp/EventLoop.php
Original file line number Diff line number Diff line change
Expand Up @@ -4,37 +4,45 @@

namespace M6Web\Tornado\Adapter\Amp;

use Amp\DeferredFuture;
use Amp\Future;
use M6Web\Tornado\Adapter\Common;
use M6Web\Tornado\Deferred;
use M6Web\Tornado\Promise;

class EventLoop implements \M6Web\Tornado\EventLoop
{
private Common\Internal\FailingPromiseCollection $unhandledFailingPromises;

public function __construct()
{
$this->unhandledFailingPromises = new Common\Internal\FailingPromiseCollection();
}

/**
* {@inheritdoc}
*/
public function wait(Promise $promise)
{
try {
$result = \Amp\Promise\wait(
Internal\PromiseWrapper::toHandledPromise($promise, $this->unhandledFailingPromises)->getAmpPromise()
);
$result = \Amp\Future\await([Internal\PromiseWrapper::toHandledPromise($promise, $this->unhandledFailingPromises)->ampFuture]);
$this->unhandledFailingPromises->throwIfWatchedFailingPromiseExists();

return $result;
return $result[0] ?? null;
} catch (\Error $error) {
// Modify exceptions sent by Amp itself
if ($error->getCode() !== 0) {
throw $error;
}
switch ($error->getMessage()) {
case 'Loop stopped without resolving the promise':
throw new \Error('Impossible to resolve the promise, no more task to execute.', 0, $error);
case 'Loop exceptionally stopped without resolving the promise':
throw $error->getPrevious() ?? $error;
default:
throw $error;

if (str_starts_with($error->getMessage(), 'Event loop terminated without resuming the current suspension')) {
throw new \Error('Impossible to resolve the promise, no more task to execute.', 0, $error);
}

throw match ($error->getMessage()) {
'Loop exceptionally stopped without resolving the promise' => $error->getPrevious() ?? $error,
default => $error,
};
}
}

Expand All @@ -43,7 +51,7 @@ public function wait(Promise $promise)
*/
public function async(\Generator $generator): Promise
{
$wrapper = function (\Generator $generator, \Amp\Deferred $deferred): \Generator {
$wrapper = function (\Generator $generator, DeferredFuture $deferred) {
try {
while ($generator->valid()) {
$blockingPromise = $generator->current();
Expand All @@ -53,13 +61,13 @@ public function async(\Generator $generator): Promise
$blockingPromise = Internal\PromiseWrapper::toHandledPromise(
$blockingPromise,
$this->unhandledFailingPromises
)->getAmpPromise();
)->ampFuture;

// Forwards promise value/exception to underlying generator
$blockingPromiseValue = null;
$blockingPromiseException = null;
try {
$blockingPromiseValue = yield $blockingPromise;
$blockingPromiseValue = $blockingPromise->await();
} catch (\Throwable $throwable) {
$blockingPromiseException = $throwable;
}
Expand All @@ -70,37 +78,48 @@ public function async(\Generator $generator): Promise
}
}
} catch (\Throwable $throwable) {
$deferred->fail($throwable);
$deferred->error($throwable);

return;
}

$deferred->resolve($generator->getReturn());
$deferred->complete($generator->getReturn());
};

$deferred = new \Amp\Deferred();
\Amp\Promise\rethrow(new \Amp\Coroutine($wrapper($generator, $deferred)));
$deferred = new DeferredFuture();
\Amp\async(fn () => $wrapper($generator, $deferred));

return Internal\PromiseWrapper::createUnhandled($deferred->promise(), $this->unhandledFailingPromises);
return Internal\PromiseWrapper::createUnhandled($deferred->getFuture(), $this->unhandledFailingPromises);
}

/**
* {@inheritdoc}
*/
public function promiseAll(Promise ...$promises): Promise
{
return Internal\PromiseWrapper::createUnhandled(
\Amp\Promise\all(
array_map(
fn (Promise $promise) => Internal\PromiseWrapper::toHandledPromise(
$promise,
$this->unhandledFailingPromises
)->getAmpPromise(),
$promises
)
),
$this->unhandledFailingPromises
$orderedResults = \array_fill_keys(\array_keys($promises), null);
$futures = array_map(
fn (Promise $promise): Future => Internal\PromiseWrapper::toHandledPromise($promise, $this->unhandledFailingPromises)->ampFuture,
$promises
);

$future = \Amp\async(function () use (&$orderedResults, $futures): array {
[$errors, $values] = \Amp\Future\awaitAll($futures);

ksort($errors);

if (count($errors) > 0) {
throw reset($errors);
}
Comment thread
e-zannelli marked this conversation as resolved.
Outdated

foreach ($values as $index => $value) {
$orderedResults[$index] = $value;
}

return $orderedResults;
});

return Internal\PromiseWrapper::createUnhandled($future, $this->unhandledFailingPromises);
}

/**
Expand All @@ -125,82 +144,81 @@ public function promiseRace(Promise ...$promises): Promise
return $this->promiseFulfilled(null);
}

$deferred = new \Amp\Deferred();
$deferred = new DeferredFuture();
$isFirstPromise = true;

$wrapPromise = function (\Amp\Promise $promise) use ($deferred, &$isFirstPromise): \Generator {
$wrapPromise = function (Future $future) use ($deferred, &$isFirstPromise) {
try {
$result = yield $promise;
$result = $future->await();
if ($isFirstPromise) {
$isFirstPromise = false;
$deferred->resolve($result);
$deferred->complete($result);
}
} catch (\Throwable $throwable) {
if ($isFirstPromise) {
$isFirstPromise = false;
$deferred->fail($throwable);
$deferred->error($throwable);
}
}
};

$promises = array_map(
$futures = array_map(
fn (Promise $promise) => Internal\PromiseWrapper::toHandledPromise(
$promise,
$this->unhandledFailingPromises
)->getAmpPromise(),
)->ampFuture,
$promises
);

foreach ($promises as $index => $promise) {
\Amp\Promise\rethrow(new \Amp\Coroutine($wrapPromise($promise)));
foreach ($futures as $future) {
\Amp\async(fn () => $wrapPromise($future));
}

return Internal\PromiseWrapper::createUnhandled($deferred->promise(), $this->unhandledFailingPromises);
return Internal\PromiseWrapper::createUnhandled($deferred->getFuture(), $this->unhandledFailingPromises);
}

/**
* {@inheritdoc}
*/
public function promiseFulfilled($value): Promise
{
return Internal\PromiseWrapper::createHandled(new \Amp\Success($value));
return Internal\PromiseWrapper::createHandled(Future::complete($value));
}

/**
* {@inheritdoc}
*/
public function promiseRejected(\Throwable $throwable): Promise
{
// Manually created promises are considered as handled.
return Internal\PromiseWrapper::createHandled(new \Amp\Failure($throwable));
return Internal\PromiseWrapper::createHandled(Future::error($throwable));
}

/**
* {@inheritdoc}
*/
public function idle(): Promise
{
$deferred = new \Amp\Deferred();
$deferred = new DeferredFuture();

\Amp\Loop::defer(function () use ($deferred): void {
$deferred->resolve();
\Revolt\EventLoop::defer(function () use ($deferred): void {
$deferred->complete();
});

return Internal\PromiseWrapper::createUnhandled($deferred->promise(), $this->unhandledFailingPromises);
return Internal\PromiseWrapper::createUnhandled($deferred->getFuture(), $this->unhandledFailingPromises);
}

/**
* {@inheritdoc}
*/
public function delay(int $milliseconds): Promise
{
$deferred = new \Amp\Deferred();
$deferred = new DeferredFuture();

\Amp\Loop::delay($milliseconds, function () use ($deferred): void {
$deferred->resolve();
\Revolt\EventLoop::delay($milliseconds / 1000, function () use ($deferred): void {
$deferred->complete();
});

return Internal\PromiseWrapper::createUnhandled($deferred->promise(), $this->unhandledFailingPromises);
return Internal\PromiseWrapper::createUnhandled($deferred->getFuture(), $this->unhandledFailingPromises);
}

/**
Expand All @@ -209,9 +227,8 @@ public function delay(int $milliseconds): Promise
public function deferred(): Deferred
{
return new Internal\Deferred(
$deferred = new \Amp\Deferred(),
// Manually created promises are considered as handled.
Internal\PromiseWrapper::createHandled($deferred->promise())
$deferred = new DeferredFuture(),
Internal\PromiseWrapper::createHandled($deferred->getFuture())
);
}

Expand All @@ -220,42 +237,34 @@ public function deferred(): Deferred
*/
public function readable($stream): Promise
{
$deferred = new \Amp\Deferred();
$deferred = new DeferredFuture();

\Amp\Loop::onReadable(
\Revolt\EventLoop::onReadable(
$stream,
function ($watcherId, $stream) use ($deferred): void {
\Amp\Loop::cancel($watcherId);
$deferred->resolve($stream);
\Revolt\EventLoop::cancel($watcherId);
$deferred->complete($stream);
}
);

return Internal\PromiseWrapper::createUnhandled($deferred->promise(), $this->unhandledFailingPromises);
return Internal\PromiseWrapper::createUnhandled($deferred->getFuture(), $this->unhandledFailingPromises);
}

/**
* {@inheritdoc}
*/
public function writable($stream): Promise
{
$deferred = new \Amp\Deferred();
$deferred = new DeferredFuture();

\Amp\Loop::onWritable(
\Revolt\EventLoop::onWritable(
$stream,
function ($watcherId, $stream) use ($deferred): void {
\Amp\Loop::cancel($watcherId);
$deferred->resolve($stream);
\Revolt\EventLoop::cancel($watcherId);
$deferred->complete($stream);
}
);

return Internal\PromiseWrapper::createUnhandled($deferred->promise(), $this->unhandledFailingPromises);
return Internal\PromiseWrapper::createUnhandled($deferred->getFuture(), $this->unhandledFailingPromises);
}

public function __construct()
{
$this->unhandledFailingPromises = new Common\Internal\FailingPromiseCollection();
}

/** @var Common\Internal\FailingPromiseCollection */
private $unhandledFailingPromises;
}
Loading
Loading