diff --git a/composer.json b/composer.json index d34811d..b52eae0 100644 --- a/composer.json +++ b/composer.json @@ -26,12 +26,14 @@ "php": "^8.2" }, "require-dev": { + "amphp/amp": "^3.1", "guzzlehttp/promises": "^1.5.0 || ^2.0.0", "phpunit/phpunit": "^10.3", "react/promise": "^2.8 || ^3.0", "webonyx/graphql-php": "^15.0" }, "suggest": { + "amphp/amp": "To use with amphp/amp v3 futures (fiber-based)", "guzzlehttp/promises": "To use with Guzzle promise", "react/promise": "To use with ReactPhp promise", "webonyx/graphql-php": "To use with Webonyx GraphQL native promise" diff --git a/lib/promise-adapter/docs/usage.md b/lib/promise-adapter/docs/usage.md index 82bc6ec..808ab56 100644 --- a/lib/promise-adapter/docs/usage.md +++ b/lib/promise-adapter/docs/usage.md @@ -14,10 +14,21 @@ Optional to use ReactPhp: composer require "react/promise" ``` +Optional to use Amp v3: + +```sh +composer require "amphp/amp" +``` + ## Supported Adapter *Guzzle*: `Overblog\PromiseAdapter\Adapter\GuzzleHttpPromiseAdapter` *ReactPhp*: `Overblog\PromiseAdapter\Adapter\ReactPromiseAdapter` +*Amp v3*: `Overblog\PromiseAdapter\Adapter\AmpFutureAdapter` + To use a custom Promise lib you can implement `Overblog\PromiseAdapter\PromiseAdapterInterface` + +Adapters for promises that require event-loop scheduling and do not expose a +`then()` method can implement `Overblog\PromiseAdapter\AsyncPromiseAdapterInterface`. diff --git a/lib/promise-adapter/src/Adapter/AmpFutureAdapter.php b/lib/promise-adapter/src/Adapter/AmpFutureAdapter.php new file mode 100644 index 0000000..3a6d4a7 --- /dev/null +++ b/lib/promise-adapter/src/Adapter/AmpFutureAdapter.php @@ -0,0 +1,257 @@ + + * + * For the full copyright and license information, please view the LICENSE + * file that was distributed with this source code. + */ + +namespace Overblog\PromiseAdapter\Adapter; + +use Amp\CancelledException; +use Amp\DeferredFuture; +use Amp\Future; +use Overblog\PromiseAdapter\AsyncPromiseAdapterInterface; +use function Amp\async; +use function Amp\Future\await; +use function Amp\Future\awaitAll; + +/** + * Promise adapter backed by amphp/amp v3 fiber-based futures. + * + * Unlike the Guzzle/React adapters, amp v3 futures do not expose a `then()` + * method and only settle once the event loop advances (inside a fiber). + * The DataLoader core is patched to cooperate with this model via `Amp\async` + * and `Future::await()` instead of `->then()`. + * + * @implements AsyncPromiseAdapterInterface> + */ +class AmpFutureAdapter implements AsyncPromiseAdapterInterface +{ + /** @var \WeakMap, array{deferred: DeferredFuture|null, canceller: callable|null}>|null */ + private ?\WeakMap $cancellations = null; + + /** @var array> */ + private array $pending = []; + + /** @var array */ + private array $running = []; + + private int $nextPendingId = 0; + + /** + * @return Future + */ + public function create(&$resolve = null, &$reject = null, ?callable $canceller = null): Future + { + $deferred = new DeferredFuture(); + $future = $deferred->getFuture(); + $this->cancellations ??= new \WeakMap(); + $this->cancellations[$future] = [ + 'deferred' => $deferred, + 'canceller' => $canceller, + ]; + + $resolve = function ($value) use ($deferred, $future): void { + if ($deferred->isComplete()) { + return; + } + + $deferred->complete($value); + $this->markSettled($future); + }; + $reject = function (\Throwable $reason) use ($deferred, $future): void { + if ($deferred->isComplete()) { + return; + } + + $deferred->error($reason); + $this->markSettled($future); + }; + + return $future; + } + + /** + * @return Future + */ + public function createFulfilled($promiseOrValue = null): Future + { + if ($promiseOrValue instanceof Future) { + return $promiseOrValue; + } + + return Future::complete($promiseOrValue); + } + + /** + * @return Future + */ + public function createRejected($reason): Future + { + if ($reason instanceof Future) { + return $reason; + } + + return Future::error($reason); + } + + /** + * @return Future + */ + public function createAll($promisesOrValues): Future + { + $futures = []; + foreach ($promisesOrValues as $key => $value) { + $futures[$key] = $value instanceof Future ? $value : Future::complete($value); + } + + return async(static fn () => await($futures)); + } + + public function isPromise($value, $strict = false): bool + { + return $value instanceof Future; + } + + public function await($promise = null, $unwrap = false): mixed + { + if (null === $promise) { + $firstError = null; + + while ([] !== $this->pending) { + $pending = $this->pending; + $currentFiber = \Fiber::getCurrent(); + if (null !== $currentFiber) { + foreach ($this->running as $id => $fiber) { + if ($fiber === $currentFiber) { + unset($pending[$id]); + } + } + } + if ([] === $pending) { + break; + } + [$errors] = awaitAll($pending); + + if (null === $firstError && [] !== $errors) { + $firstError = reset($errors); + } + } + + if (null !== $firstError) { + throw $firstError; + } + + return null; + } + + if (!$promise instanceof Future) { + throw new \InvalidArgumentException(sprintf('The "%s" method must be called with an amp Future.', __METHOD__)); + } + + try { + return $promise->await(); + } catch (\Throwable $reason) { + if (!$unwrap) { + return $reason; + } + + throw $reason; + } + } + + public function cancel($promise): void + { + if (!$promise instanceof Future || null === $this->cancellations || !$this->cancellations->offsetExists($promise)) { + throw new \InvalidArgumentException(sprintf('The "%s" method must be called with a compatible Future.', __METHOD__)); + } + + $cancellation = $this->cancellations[$promise]; + $deferred = $cancellation['deferred']; + if (null === $deferred) { + return; + } + + $this->markSettled($promise); + $promise->ignore(); + + try { + if (null !== $cancellation['canceller']) { + ($cancellation['canceller'])(); + } + } catch (\Throwable $reason) { + if (!$deferred->isComplete()) { + $deferred->error($reason); + } + + return; + } + + if (!$deferred->isComplete()) { + $deferred->error(new CancelledException()); + } + } + + public function enqueue(callable $callback): void + { + $id = ++$this->nextPendingId; + $this->track(async(function () use ($callback, $id): void { + $this->running[$id] = \Fiber::getCurrent(); + + try { + $callback(); + } finally { + unset($this->running[$id]); + } + }), $id); + } + + public function observe($promise, callable $onFulfilled, callable $onRejected): void + { + if (!$promise instanceof Future) { + throw new \InvalidArgumentException(sprintf('The "%s" method must be called with a compatible Future.', __METHOD__)); + } + + $id = ++$this->nextPendingId; + $observer = async(function () use ($promise, $onFulfilled, $onRejected, $id): void { + $this->running[$id] = \Fiber::getCurrent(); + + try { + try { + $onFulfilled($promise->await()); + } catch (\Throwable $error) { + $onRejected($error); + } + } finally { + unset($this->running[$id]); + } + }); + + $this->track($observer, $id); + } + + private function markSettled(Future $future): void + { + if (null === $this->cancellations || !$this->cancellations->offsetExists($future)) { + return; + } + + $this->cancellations[$future] = [ + 'deferred' => null, + 'canceller' => null, + ]; + } + + private function track(Future $future, int $id): void + { + $tracked = $future->finally(function () use ($id): void { + unset($this->pending[$id]); + }); + $tracked->ignore(); + $this->pending[$id] = $tracked; + } +} diff --git a/lib/promise-adapter/src/AsyncPromiseAdapterInterface.php b/lib/promise-adapter/src/AsyncPromiseAdapterInterface.php new file mode 100644 index 0000000..8edcd09 --- /dev/null +++ b/lib/promise-adapter/src/AsyncPromiseAdapterInterface.php @@ -0,0 +1,38 @@ + + * + * For the full copyright and license information, please view the LICENSE + * file that was distributed with this source code. + */ + +namespace Overblog\PromiseAdapter; + +/** + * Supports promises that require event-loop scheduling and do not expose then(). + * + * @template TPromise + * + * @extends PromiseAdapterInterface + */ +interface AsyncPromiseAdapterInterface extends PromiseAdapterInterface +{ + /** + * Queue work for the next event-loop turn. + * + * A no-argument await() call must drain this work, including work that it + * enqueues or observes. cancel() must settle adapter-created promises + * without requiring an event-loop turn. + */ + public function enqueue(callable $callback): void; + + /** + * @param TPromise $promise + * + * A no-argument await() call must wait for the observer callbacks to finish. + */ + public function observe($promise, callable $onFulfilled, callable $onRejected): void; +} diff --git a/lib/promise-adapter/tests/AmpFutureAdapterTest.php b/lib/promise-adapter/tests/AmpFutureAdapterTest.php new file mode 100644 index 0000000..92b8458 --- /dev/null +++ b/lib/promise-adapter/tests/AmpFutureAdapterTest.php @@ -0,0 +1,125 @@ + + * + * For the full copyright and license information, please view the LICENSE + * file that was distributed with this source code. + */ + +namespace Overblog\PromiseAdapter\Tests; + +use Amp\CancelledException; +use Amp\Future; +use Overblog\PromiseAdapter\Adapter\AmpFutureAdapter; + +use function Amp\async; +use function Amp\delay; + +class AmpFutureAdapterTest extends \PHPUnit\Framework\TestCase +{ + public function testAwaitReturnsRejectionWithoutUnwrapByDefault(): void + { + $adapter = new AmpFutureAdapter(); + $error = new \RuntimeException('failed'); + + self::assertSame($error, $adapter->await(Future::error($error))); + } + + public function testCreateRejectedPreservesFuture(): void + { + $adapter = new AmpFutureAdapter(); + $future = Future::complete('value'); + + self::assertSame($future, $adapter->createRejected($future)); + } + + public function testCancelInvokesCancellerAndRejectsFuture(): void + { + $adapter = new AmpFutureAdapter(); + $error = new \RuntimeException('cancelled'); + $future = $adapter->create( + $resolve, + $reject, + static function () use ($error): void { + throw $error; + }, + ); + + $adapter->cancel($future); + + self::assertSame($error, $adapter->await($future)); + } + + public function testCancelRejectsInvalidFuture(): void + { + $adapter = new AmpFutureAdapter(); + + $this->expectException(\InvalidArgumentException::class); + $this->expectExceptionMessage('::cancel" method must be called with a compatible Future.'); + + $adapter->cancel(Future::complete(null)); + } + + public function testCancelWithoutCancellerRejectsFuture(): void + { + $adapter = new AmpFutureAdapter(); + $future = $adapter->create($resolve, $reject); + + $adapter->cancel($future); + + self::assertInstanceOf(CancelledException::class, $adapter->await($future)); + } + + public function testAwaitDrainsWorkEnqueuedByPendingWork(): void + { + $adapter = new AmpFutureAdapter(); + $calls = []; + $adapter->enqueue(function () use ($adapter, &$calls): void { + $calls[] = 'first'; + $adapter->enqueue(function () use (&$calls): void { + $calls[] = 'second'; + }); + }); + + $adapter->await(); + + self::assertSame(['first', 'second'], $calls); + } + + public function testFailedResolutionLeavesFutureCancellable(): void + { + $adapter = new AmpFutureAdapter(); + $future = $adapter->create($resolve, $reject); + + try { + $resolve(Future::complete('invalid nested future')); + self::fail('Expected resolving with a Future to fail.'); + } catch (\Error $error) { + self::assertSame('Cannot complete with an instance of Amp\Future', $error->getMessage()); + } + + $adapter->cancel($future); + + self::assertInstanceOf(CancelledException::class, $adapter->await($future)); + } + + public function testAwaitWaitsForWorkRunningInAnotherFiber(): void + { + $adapter = new AmpFutureAdapter(); + $completed = false; + $adapter->enqueue(function () use (&$completed): void { + delay(0.01); + $completed = true; + }); + + async(function () use ($adapter, &$completed): void { + delay(0); + $adapter->await(); + + self::assertTrue($completed); + })->await(); + } +} diff --git a/src/DataLoader.php b/src/DataLoader.php index f2236e2..3cf9d3b 100644 --- a/src/DataLoader.php +++ b/src/DataLoader.php @@ -11,6 +11,7 @@ namespace Overblog\DataLoader; +use Overblog\PromiseAdapter\AsyncPromiseAdapterInterface; use Overblog\PromiseAdapter\PromiseAdapterInterface; /** @@ -122,6 +123,12 @@ static function () { if (!$shouldBatch) { // Otherwise dispatch the (queue of one) immediately. $this->dispatchQueue(); + } elseif ($this->getPromiseAdapter() instanceof AsyncPromiseAdapterInterface) { + $this->getPromiseAdapter()->enqueue(function (): void { + if ($this->needProcess()) { + $this->dispatchQueue(); + } + }); } } @@ -205,7 +212,9 @@ public function __destruct() } } - $this->getPromiseAdapter()->await(); + if (!$this->getPromiseAdapter() instanceof AsyncPromiseAdapterInterface) { + $this->getPromiseAdapter()->await(); + } } } @@ -218,7 +227,9 @@ protected function process() { if ($this->needProcess()) { $this->getPromiseAdapter()->await(); - $this->dispatchQueue(); + if ($this->needProcess()) { + $this->dispatchQueue(); + } $this->getPromiseAdapter()->await(); } } @@ -265,36 +276,41 @@ public static function await($promise = null, $unwrap = true) return null; } - if (is_callable([$promise, 'then'])) { - $isPromiseCompleted = false; - $resolvedValue = null; - $rejectedReason = null; + if (!is_callable([$promise, 'then'])) { + if (is_object($promise) && null !== self::$promiseAdapters && isset(self::$promiseAdapters[$promise])) { + return self::$promiseAdapters[$promise]->await($promise, $unwrap); + } - $promise->then( - function ($value) use (&$isPromiseCompleted, &$resolvedValue) { - $isPromiseCompleted = true; - $resolvedValue = $value; - }, - function ($reason) use (&$isPromiseCompleted, &$rejectedReason) { - $isPromiseCompleted = true; - $rejectedReason = $reason; - } - ); + throw new \InvalidArgumentException(sprintf('The "%s" method must be called with a Promise ("then" method).', __METHOD__)); + } - //Promise is completed? - if ($isPromiseCompleted) { - // rejected ? - if ($rejectedReason instanceof \Throwable) { - if (!$unwrap) { - return $rejectedReason; - } - throw $rejectedReason; + $isPromiseCompleted = false; + $resolvedValue = null; + $rejectedReason = null; + + $promise->then( + function ($value) use (&$isPromiseCompleted, &$resolvedValue) { + $isPromiseCompleted = true; + $resolvedValue = $value; + }, + function ($reason) use (&$isPromiseCompleted, &$rejectedReason) { + $isPromiseCompleted = true; + $rejectedReason = $reason; + } + ); + + //Promise is completed? + if ($isPromiseCompleted) { + // rejected ? + if ($rejectedReason instanceof \Throwable) { + if (!$unwrap) { + return $rejectedReason; } - return $resolvedValue; + throw $rejectedReason; } - } else { - throw new \InvalidArgumentException(sprintf('The "%s" method must be called with a Promise ("then" method).', __METHOD__)); + + return $resolvedValue; } if (null !== self::$promiseAdapters && isset(self::$promiseAdapters[$promise])) { @@ -322,7 +338,7 @@ function ($reason) use (&$isPromiseCompleted, &$rejectedReason) { private static function awaitInstances() { - if (empty(self::$activeInstances)) { + if ([] === self::$activeInstances) { return; } @@ -409,6 +425,21 @@ private function dispatchQueueBatch(array $queue) return; } + $promiseAdapter = $this->getPromiseAdapter(); + if ($promiseAdapter instanceof AsyncPromiseAdapterInterface && $promiseAdapter->isPromise($batchPromise, true)) { + $promiseAdapter->observe( + $batchPromise, + function ($values) use ($queue): void { + $this->resolveDispatchedBatch($values, $queue); + }, + function (\Throwable $error) use ($queue): void { + $this->failedDispatch($queue, $error); + }, + ); + + return; + } + // Assert the expected response from batchLoadFn if (!$batchPromise || !is_callable([$batchPromise, 'then'])) { $this->failedDispatch($queue, new \RuntimeException( @@ -422,39 +453,50 @@ private function dispatchQueueBatch(array $queue) // Await the resolution of the call to batchLoadFn. $batchPromise->then( - function ($values) use ($keys, $queue) { - // Assert the expected resolution from batchLoadFn. - if (!is_array($values) && !$values instanceof \Traversable) { - throw new \RuntimeException( - 'DataLoader must be constructed with a function which accepts ' . - 'Array and returns Promise>, but the function did ' . - sprintf('not return a Promise of an Array: %s.', gettype($values)) - ); - } - if (count($values) !== count($keys)) { - throw new \RuntimeException( - 'DataLoader must be constructed with a function which accepts ' . - 'Array and returns Promise>, but the function did ' . - 'not return a Promise of an Array of the same length as the Array of keys.' - ); - } - - // Step through the values, resolving or rejecting each Promise in the - // loaded queue. - foreach ($queue as $index => $data) { - $value = $values[$index]; - if ($value instanceof \Throwable) { - $data['reject']($value); - } else { - $data['resolve']($value); - } - }; + function ($values) use ($queue) { + $this->resolveDispatchedBatch($values, $queue); } )->then(null, function ($error) use ($queue) { $this->failedDispatch($queue, $error); }); } + /** + * Fan a resolved batch result out to the individual queued promises. + * + * @param mixed $values + * @param array $queue + */ + private function resolveDispatchedBatch($values, $queue) + { + // Assert the expected response from batchLoadFn. + if (!is_array($values) && !$values instanceof \Traversable) { + throw new \RuntimeException( + 'DataLoader must be constructed with a function which accepts ' . + 'Array and returns Promise>, but the function did ' . + sprintf('not return a Promise of an Array: %s.', gettype($values)) + ); + } + if (count($values) !== count($queue)) { + throw new \RuntimeException( + 'DataLoader must be constructed with a function which accepts ' . + 'Array and returns Promise>, but the function did ' . + 'not return a Promise of an Array of the same length as the Array of keys.' + ); + } + + // Step through the values, resolving or rejecting each Promise in the + // loaded queue. + foreach ($queue as $index => $data) { + $value = $values[$index]; + if ($value instanceof \Throwable) { + $data['reject']($value); + } else { + $data['resolve']($value); + } + } + } + /** * Do not cache individual loads if the entire batch dispatch fails, * but still reject each request so they do not hang. diff --git a/tests/AmpFutureDataLoaderTest.php b/tests/AmpFutureDataLoaderTest.php new file mode 100644 index 0000000..f860927 --- /dev/null +++ b/tests/AmpFutureDataLoaderTest.php @@ -0,0 +1,204 @@ + + * + * For the full copyright and license information, please view the LICENSE + * file that was distributed with this source code. + */ + +namespace Overblog\DataLoader\Test; + +use Amp\Future; +use Error; +use Overblog\DataLoader\DataLoader; +use Overblog\PromiseAdapter\Adapter\AmpFutureAdapter; +use Overblog\PromiseAdapter\AsyncPromiseAdapterInterface; +use Overblog\PromiseAdapter\PromiseAdapterInterface; + +use function Amp\async; + +class AmpFutureDataLoaderTest extends TestCase +{ + protected function createPromiseAdapter(): PromiseAdapterInterface + { + return new AmpFutureAdapter(); + } + + public function testAwaitUsesTheRegisteredAmpFutureAdapter() + { + $loader = new DataLoader( + static fn (array $keys): Future => Future::complete(['value', 'value']), + $this->createPromiseAdapter(), + ); + + self::assertSame(['value', 'value'], DataLoader::await($loader->loadMany(['first', 'second']))); + } + + public function testRejectsEveryQueuedLoadWhenBatchFutureFailsWithAnError() + { + $loader = new DataLoader( + static fn (array $keys): Future => Future::error(new Error('batch load failed')), + new AmpFutureAdapter(), + ); + + $first = $loader->load('first'); + $second = $loader->load('second'); + + foreach ([$first, $second] as $future) { + try { + async(static fn () => $future->await())->await(); + self::fail('Expected the queued load to fail.'); + } catch (Error $error) { + self::assertSame('batch load failed', $error->getMessage()); + } + } + } + + public function testAwaitWithoutPromiseCompletesQueuedLoads(): void + { + $loader = new DataLoader( + static fn (array $keys): Future => Future::complete($keys), + new AmpFutureAdapter(), + ); + $future = $loader->load('value'); + + DataLoader::await(); + + self::assertTrue($future->isComplete()); + self::assertSame('value', $future->await()); + } + + public function testDestructionRejectsQueuedLoads(): void + { + $loader = new DataLoader( + static fn (array $keys): Future => Future::complete($keys), + new AmpFutureAdapter(), + ); + $future = $loader->load('value'); + + $loader->__destruct(); + unset($loader); + + $error = DataLoader::await($future, false); + self::assertInstanceOf(\RuntimeException::class, $error); + self::assertSame('DataLoader destroyed before promise complete.', $error->getMessage()); + } + + public function testSupportsDelegatingAsyncPromiseAdapter(): void + { + $adapter = new DelegatingAsyncPromiseAdapter(); + $loader = new DataLoader( + static fn (array $keys): Future => Future::complete($keys), + $adapter, + ); + + self::assertSame(['first', 'second'], DataLoader::await($loader->loadMany(['first', 'second']))); + } + + public function testDestructionDoesNotDispatchUnrelatedLoader(): void + { + $adapter = new AmpFutureAdapter(); + $loader = new DataLoader( + static fn (array $keys): Future => Future::complete($keys), + $adapter, + ); + $loader->load('cancelled')->ignore(); + + $unrelatedLoadCalls = []; + $unrelatedLoader = new DataLoader( + static function (array $keys) use (&$unrelatedLoadCalls): Future { + $unrelatedLoadCalls[] = $keys; + + return Future::complete($keys); + }, + $adapter, + ); + $unrelatedFuture = $unrelatedLoader->load('unrelated'); + + $loader->__destruct(); + unset($loader); + + self::assertSame([], $unrelatedLoadCalls); + self::assertFalse($unrelatedFuture->isComplete()); + self::assertSame('unrelated', DataLoader::await($unrelatedFuture)); + } + + public function testBatchLoaderCanAwaitAnotherLoader(): void + { + $adapter = new AmpFutureAdapter(); + $innerLoader = new DataLoader( + static fn (array $keys): Future => Future::complete($keys), + $adapter, + ); + $outerLoader = new DataLoader( + static function (array $keys) use ($innerLoader): Future { + return Future::complete(array_map( + static fn ($key) => DataLoader::await($innerLoader->load($key)), + $keys, + )); + }, + $adapter, + ); + + self::assertSame('value', DataLoader::await($outerLoader->load('value'))); + } +} + +/** @implements AsyncPromiseAdapterInterface> */ +final class DelegatingAsyncPromiseAdapter implements AsyncPromiseAdapterInterface +{ + private AmpFutureAdapter $adapter; + + public function __construct() + { + $this->adapter = new AmpFutureAdapter(); + } + + public function create(&$resolve = null, &$reject = null, ?callable $canceller = null): Future + { + return $this->adapter->create($resolve, $reject, $canceller); + } + + public function createFulfilled($promiseOrValue = null): Future + { + return $this->adapter->createFulfilled($promiseOrValue); + } + + public function createRejected($reason): Future + { + return $this->adapter->createRejected($reason); + } + + public function createAll($promisesOrValues): Future + { + return $this->adapter->createAll($promisesOrValues); + } + + public function isPromise($value, $strict = false): bool + { + return $this->adapter->isPromise($value, $strict); + } + + public function await($promise = null, $unwrap = false): mixed + { + return $this->adapter->await($promise, $unwrap); + } + + public function cancel($promise): void + { + $this->adapter->cancel($promise); + } + + public function enqueue(callable $callback): void + { + $this->adapter->enqueue($callback); + } + + public function observe($promise, callable $onFulfilled, callable $onRejected): void + { + $this->adapter->observe($promise, $onFulfilled, $onRejected); + } +}