diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index f65e97c..84b2eb4 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -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' diff --git a/.php-cs-fixer.dist.php b/.php-cs-fixer.dist.php index de8fa1d..807ccdd 100644 --- a/.php-cs-fixer.dist.php +++ b/.php-cs-fixer.dist.php @@ -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; diff --git a/composer.json b/composer.json index 516a5f6..4ba7b97 100644 --- a/composer.json +++ b/composer.json @@ -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" diff --git a/src/Adapter/Amp/EventLoop.php b/src/Adapter/Amp/EventLoop.php index 6337156..c75aeb6 100644 --- a/src/Adapter/Amp/EventLoop.php +++ b/src/Adapter/Amp/EventLoop.php @@ -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, + }; } } @@ -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(); @@ -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; } @@ -70,18 +78,18 @@ 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); } /** @@ -89,18 +97,23 @@ public function async(\Generator $generator): Promise */ 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 { + $values = \Amp\Future\await($futures); + + foreach ($values as $index => $value) { + $orderedResults[$index] = $value; + } + + return $orderedResults; + }); + + return Internal\PromiseWrapper::createUnhandled($future, $this->unhandledFailingPromises); } /** @@ -125,37 +138,37 @@ 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); } /** @@ -163,7 +176,7 @@ public function promiseRace(Promise ...$promises): Promise */ public function promiseFulfilled($value): Promise { - return Internal\PromiseWrapper::createHandled(new \Amp\Success($value)); + return Internal\PromiseWrapper::createHandled(Future::complete($value)); } /** @@ -171,8 +184,7 @@ public function promiseFulfilled($value): Promise */ 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)); } /** @@ -180,13 +192,13 @@ public function promiseRejected(\Throwable $throwable): Promise */ 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); } /** @@ -194,13 +206,13 @@ public function idle(): Promise */ 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); } /** @@ -209,9 +221,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()) ); } @@ -220,17 +231,17 @@ 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); } /** @@ -238,24 +249,16 @@ function ($watcherId, $stream) use ($deferred): void { */ 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; } diff --git a/src/Adapter/Amp/Internal/Deferred.php b/src/Adapter/Amp/Internal/Deferred.php index 028a5e1..4851cd6 100644 --- a/src/Adapter/Amp/Internal/Deferred.php +++ b/src/Adapter/Amp/Internal/Deferred.php @@ -4,6 +4,7 @@ namespace M6Web\Tornado\Adapter\Amp\Internal; +use Amp\DeferredFuture; use M6Web\Tornado\Promise; /** @@ -15,11 +16,13 @@ class Deferred implements \M6Web\Tornado\Deferred { /** - * @param \Amp\Deferred $ampDeferred + * @param DeferredFuture $ampDeferred * @param PromiseWrapper $promise */ - public function __construct(private readonly \Amp\Deferred $ampDeferred, private readonly PromiseWrapper $promise) - { + public function __construct( + private readonly DeferredFuture $ampDeferred, + private readonly PromiseWrapper $promise, + ) { } /** @@ -30,20 +33,12 @@ public function getPromise(): Promise return $this->promise; } - /** - * @return PromiseWrapper - */ - public function getPromiseWrapper(): PromiseWrapper - { - return $this->promise; - } - /** * @param TValue $value */ public function resolve($value): void { - $this->ampDeferred->resolve($value); + $this->ampDeferred->complete($value); } /** @@ -51,6 +46,6 @@ public function resolve($value): void */ public function reject(\Throwable $throwable): void { - $this->ampDeferred->fail($throwable); + $this->ampDeferred->error($throwable); } } diff --git a/src/Adapter/Amp/Internal/PromiseWrapper.php b/src/Adapter/Amp/Internal/PromiseWrapper.php index 0f2914d..6e56d07 100644 --- a/src/Adapter/Amp/Internal/PromiseWrapper.php +++ b/src/Adapter/Amp/Internal/PromiseWrapper.php @@ -4,6 +4,7 @@ namespace M6Web\Tornado\Adapter\Amp\Internal; +use Amp\Future; use M6Web\Tornado\Adapter\Common\Internal\FailingPromiseCollection; use M6Web\Tornado\Promise; @@ -20,24 +21,22 @@ class PromiseWrapper implements Promise /** * Use named (static) constructor instead * - * @param \Amp\Promise $ampPromise + * @param Future $ampFuture */ private function __construct( - private readonly \Amp\Promise $ampPromise, + public readonly Future $ampFuture, private bool $isHandled, ) { } /** - * @param \Amp\Promise $ampPromise - * * @return self */ - public static function createUnhandled(\Amp\Promise $ampPromise, FailingPromiseCollection $failingPromiseCollection): self + public static function createUnhandled(Future $ampPromise, FailingPromiseCollection $failingPromiseCollection): self { $promiseWrapper = new self($ampPromise, false); - $promiseWrapper->ampPromise->onResolve( - function (?\Throwable $reason, $value) use ($promiseWrapper, $failingPromiseCollection): void { + $promiseWrapper->ampFuture->catch( + function (?\Throwable $reason) use ($promiseWrapper, $failingPromiseCollection): void { if ($reason !== null && !$promiseWrapper->isHandled) { $failingPromiseCollection->watchFailingPromise($promiseWrapper, $reason); } @@ -48,21 +47,13 @@ function (?\Throwable $reason, $value) use ($promiseWrapper, $failingPromiseColl } /** - * @param \Amp\Promise $ampPromise - * * @return self */ - public static function createHandled(\Amp\Promise $ampPromise): self + public static function createHandled(Future $ampPromise): self { - return new self($ampPromise, true); - } + $ampPromise->ignore(); - /** - * @return \Amp\Promise - */ - public function getAmpPromise(): \Amp\Promise - { - return $this->ampPromise; + return new self($ampPromise, true); } /** @@ -75,6 +66,7 @@ public static function toHandledPromise(Promise $promise, FailingPromiseCollecti assert($promise instanceof self, new \Error('Input promise was not created by this adapter.')); $promise->isHandled = true; + $promise->ampFuture->ignore(); $failingPromiseCollection->unwatchPromise($promise); return $promise; diff --git a/src/Adapter/Tornado/SynchronousEventLoop.php b/src/Adapter/Tornado/SynchronousEventLoop.php index f42fd0a..d61e70f 100644 --- a/src/Adapter/Tornado/SynchronousEventLoop.php +++ b/src/Adapter/Tornado/SynchronousEventLoop.php @@ -156,10 +156,8 @@ public function delay(int $milliseconds): Promise public function deferred(): Deferred { $deferred = new class implements Deferred { - /** @var SynchronousEventLoop */ - public $eventLoop; - /** @var ?Promise */ - private $promise; + public SynchronousEventLoop $eventLoop; + private ?Promise $promise = null; public function getPromise(): Promise { diff --git a/tests/Adapter/Amp/EventLoopTest.php b/tests/Adapter/Amp/EventLoopTest.php index f8f2a5f..ad6ba75 100644 --- a/tests/Adapter/Amp/EventLoopTest.php +++ b/tests/Adapter/Amp/EventLoopTest.php @@ -6,8 +6,9 @@ use M6Web\Tornado\Adapter\Amp; use M6Web\Tornado\EventLoop; +use M6WebTest\Tornado\EventLoopTestCase; -class EventLoopTest extends \M6WebTest\Tornado\EventLoopTestCase +class EventLoopTest extends EventLoopTestCase { protected function createEventLoop(): EventLoop { @@ -22,7 +23,6 @@ public function testStreamShouldReadFromWritable(string $expectedSequence = ''): protected function tearDown(): void { - \Amp\Loop::set((new \Amp\Loop\DriverFactory())->create()); gc_collect_cycles(); // extensions using an event loop may otherwise leak the file descriptors to the loop } } diff --git a/tests/EventLoopTest/PromiseAllTestTrait.php b/tests/EventLoopTest/PromiseAllTestTrait.php index 85d77da..e82be42 100644 --- a/tests/EventLoopTest/PromiseAllTestTrait.php +++ b/tests/EventLoopTest/PromiseAllTestTrait.php @@ -101,4 +101,18 @@ public function testPromiseAllCatchableException(): void $eventLoop->wait($eventLoop->async($createGenerator())) ); } + + public function testPromiseAllShouldPreserveTheOrderOfArrayWithStringKeys(): void + { + $eventLoop = $this->createEventLoop(); + $expectedValues = [ + 'c' => 1, + 'a' => 2, + 'b' => 3, + ]; + $promises = \array_map([$eventLoop, 'promiseFulfilled'], $expectedValues); + $promise = $eventLoop->promiseAll(...$promises); + + $this->assertSame($expectedValues, $eventLoop->wait($promise)); + } } diff --git a/tests/EventLoopTest/PromiseForeachTestTrait.php b/tests/EventLoopTest/PromiseForeachTestTrait.php index e960ba3..1954e52 100644 --- a/tests/EventLoopTest/PromiseForeachTestTrait.php +++ b/tests/EventLoopTest/PromiseForeachTestTrait.php @@ -43,7 +43,6 @@ public function testPromiseForeachShouldThrowIfCallbackDoesNotReturnGenerator(): { $eventLoop = $this->createEventLoop(); $callback = function (): void { - return; }; $this->expectException(\TypeError::class); diff --git a/tests/EventLoopTest/StreamsTestTrait.php b/tests/EventLoopTest/StreamsTestTrait.php index 789ee7c..4190d5f 100644 --- a/tests/EventLoopTest/StreamsTestTrait.php +++ b/tests/EventLoopTest/StreamsTestTrait.php @@ -48,7 +48,6 @@ public function testStreamShouldReadFromWritable(string $expectedSequence = 'W0R $sequence .= "W$token"; // Write twice slower yield $eventLoop->idle(); - // yield $eventLoop->idle(); } fclose($stream); };