1
0
mirror of https://github.com/danog/amp.git synced 2025-01-22 13:21:16 +01:00
amp/lib/functions.php

1082 lines
33 KiB
PHP
Raw Normal View History

<?php
2016-08-15 23:46:26 -05:00
2018-06-18 20:00:01 +02:00
namespace Amp
{
use React\Promise\PromiseInterface as ReactPromise;
/**
* Returns a new function that wraps $callback in a promise/coroutine-aware function that automatically runs
2017-05-03 15:21:49 +02:00
* Generators as coroutines. The returned function always returns a promise when invoked. Errors have to be handled
* by the callback caller or they will go unnoticed.
*
2017-05-03 15:21:49 +02:00
* Use this function to create a coroutine-aware callable for a promise-aware callback caller.
*
* @template TReturn
* @template TPromise
2020-04-19 15:38:22 +02:00
* @template TGeneratorReturn
* @template TGeneratorPromise
*
* @template TGenerator as TGeneratorReturn|Promise<TGeneratorPromise>
* @template T as TReturn|Promise<TPromise>|\Generator<mixed, mixed, mixed, TGenerator>
*
* @formatter:off
*
* @param callable(...mixed): T $callback
*
* @return callable
* @psalm-return (T is Promise ? (callable(mixed...): Promise<TPromise>) : (T is \Generator ? (TGenerator is Promise ? (callable(mixed...): Promise<TGeneratorPromise>) : (callable(mixed...): Promise<TGeneratorReturn>)) : (callable(mixed...): Promise<TReturn>)))
2020-04-19 15:38:22 +02:00
*
* @formatter:on
*
2017-05-03 15:21:49 +02:00
* @see asyncCoroutine()
*
* @psalm-suppress InvalidReturnType
*/
2018-06-18 20:00:01 +02:00
function coroutine(callable $callback): callable
{
/** @psalm-suppress InvalidReturnStatement */
return static function (...$args) use ($callback): Promise {
2017-05-03 15:21:49 +02:00
return call($callback, ...$args);
};
}
/**
* Returns a new function that wraps $callback in a promise/coroutine-aware function that automatically runs
2017-05-03 15:21:49 +02:00
* Generators as coroutines. The returned function always returns void when invoked. Errors are forwarded to the
* loop's error handler using `Amp\Promise\rethrow()`.
*
2017-05-03 15:21:49 +02:00
* Use this function to create a coroutine-aware callable for a non-promise-aware callback caller.
*
2020-04-28 22:34:37 +02:00
* @param callable(...mixed): mixed $callback
*
* @return callable
2020-04-28 22:34:37 +02:00
* @psalm-return callable(mixed...): void
*
2017-05-03 15:21:49 +02:00
* @see coroutine()
*/
2018-06-18 20:00:01 +02:00
function asyncCoroutine(callable $callback): callable
{
return static function (...$args) use ($callback) {
2017-05-03 15:21:49 +02:00
Promise\rethrow(call($callback, ...$args));
};
}
/**
* Calls the given function, always returning a promise. If the function returns a Generator, it will be run as a
* coroutine. If the function throws, a failed promise will be returned.
*
* @template TReturn
* @template TPromise
* @template TGeneratorReturn
2020-04-19 15:38:22 +02:00
* @template TGeneratorPromise
*
* @template TGenerator as TGeneratorReturn|Promise<TGeneratorPromise>
* @template T as TReturn|Promise<TPromise>|\Generator<mixed, mixed, mixed, TGenerator>
*
2020-04-19 15:38:22 +02:00
* @formatter:off
*
* @param callable(...mixed): T $callback
* @param mixed ...$args Arguments to pass to the function.
*
* @return Promise
* @psalm-return (T is Promise ? Promise<TPromise> : (T is \Generator ? (TGenerator is Promise ? Promise<TGeneratorPromise> : Promise<TGeneratorReturn>) : Promise<TReturn>))
2020-04-19 15:38:22 +02:00
*
* @formatter:on
*/
2018-06-18 20:00:01 +02:00
function call(callable $callback, ...$args): Promise
{
try {
$result = $callback(...$args);
} catch (\Throwable $exception) {
return new Failure($exception);
}
2016-05-21 12:19:48 -05:00
if ($result instanceof \Generator) {
return new Coroutine($result);
}
2016-12-11 16:17:51 +01:00
2017-03-14 11:56:36 -05:00
if ($result instanceof Promise) {
return $result;
2016-05-21 12:19:48 -05:00
}
2016-05-21 09:44:52 -05:00
2017-03-14 11:56:36 -05:00
if ($result instanceof ReactPromise) {
return Promise\adapt($result);
2017-03-14 11:56:36 -05:00
}
return new Success($result);
2017-02-22 15:52:30 -06:00
}
2017-05-03 15:21:49 +02:00
/**
* Calls the given function. If the function returns a Generator, it will be run as a coroutine. If the function
* throws or returns a failing promise, the failure is forwarded to the loop error handler.
*
2020-04-28 22:34:37 +02:00
* @param callable(...mixed): mixed $callback
* @param mixed ...$args Arguments to pass to the function.
2017-05-03 15:21:49 +02:00
*
* @return void
2017-05-03 15:21:49 +02:00
*/
2018-06-18 20:00:01 +02:00
function asyncCall(callable $callback, ...$args)
{
2017-05-03 15:21:49 +02:00
Promise\rethrow(call($callback, ...$args));
}
2019-08-02 22:37:42 +02:00
/**
* Sleeps for the specified number of milliseconds.
*
* @param int $milliseconds
*
* @return Delayed
*/
function delay(int $milliseconds): Delayed
2019-08-02 22:37:42 +02:00
{
return new Delayed($milliseconds);
}
2019-11-11 20:02:09 +01:00
/**
* Returns the current time relative to an arbitrary point in time.
*
* @return int Time in milliseconds.
*/
function getCurrentTime(): int
{
return Internal\getCurrentTime();
}
2017-02-22 15:52:30 -06:00
}
2018-06-18 20:00:01 +02:00
namespace Amp\Promise
{
use Amp\Deferred;
2017-04-23 14:39:19 +02:00
use Amp\Loop;
use Amp\MultiReasonException;
use Amp\Promise;
use Amp\Success;
use Amp\TimeoutException;
use React\Promise\PromiseInterface as ReactPromise;
2018-11-25 07:56:42 -09:00
use function Amp\call;
use function Amp\Internal\createTypeError;
2016-05-21 09:44:52 -05:00
/**
2017-04-23 15:47:52 +02:00
* Registers a callback that will forward the failure reason to the event loop's error handler if the promise fails.
*
2017-04-23 15:47:52 +02:00
* Use this function if you neither return the promise nor handle a possible error yourself to prevent errors from
* going entirely unnoticed.
*
2019-11-11 20:02:09 +01:00
* @param Promise|ReactPromise $promise Promise to register the handler on.
*
2020-03-28 22:20:44 +01:00
* @return void
* @throws \TypeError If $promise is not an instance of \Amp\Promise or \React\Promise\PromiseInterface.
2020-03-28 22:20:44 +01:00
*
*/
2018-06-18 20:00:01 +02:00
function rethrow($promise)
{
if (!$promise instanceof Promise) {
if ($promise instanceof ReactPromise) {
$promise = adapt($promise);
} else {
throw createTypeError([Promise::class, ReactPromise::class], $promise);
}
}
$promise->onResolve(static function ($exception) {
if ($exception) {
throw $exception;
}
2016-05-21 12:19:48 -05:00
});
2016-05-21 09:44:52 -05:00
}
/**
* Runs the event loop until the promise is resolved. Should not be called within a running event loop.
*
2017-04-23 15:47:52 +02:00
* Use this function only in synchronous contexts to wait for an asynchronous operation. Use coroutines and yield to
* await promise resolution in a fully asynchronous application instead.
*
2020-04-28 22:34:37 +02:00
* @template TPromise
* @template T as Promise<TPromise>|ReactPromise
*
2019-11-11 20:02:09 +01:00
* @param Promise|ReactPromise $promise Promise to wait for.
*
* @return mixed Promise success value.
*
2020-04-28 22:34:37 +02:00
* @psalm-param T $promise
* @psalm-return (T is Promise ? TPromise : mixed)
*
* @throws \TypeError If $promise is not an instance of \Amp\Promise or \React\Promise\PromiseInterface.
2017-04-23 15:47:52 +02:00
* @throws \Error If the event loop stopped without the $promise being resolved.
* @throws \Throwable Promise failure reason.
*/
2018-06-18 20:00:01 +02:00
function wait($promise)
{
if (!$promise instanceof Promise) {
if ($promise instanceof ReactPromise) {
$promise = adapt($promise);
} else {
throw createTypeError([Promise::class, ReactPromise::class], $promise);
}
}
$resolved = false;
try {
Loop::run(function () use (&$resolved, &$value, &$exception, $promise) {
$promise->onResolve(function ($e, $v) use (&$resolved, &$value, &$exception) {
Loop::stop();
$resolved = true;
$exception = $e;
$value = $v;
});
});
} catch (\Throwable $throwable) {
throw new \Error("Loop exceptionally stopped without resolving the promise", 0, $throwable);
}
2016-05-21 12:19:48 -05:00
if (!$resolved) {
throw new \Error("Loop stopped without resolving the promise");
2016-05-21 12:19:48 -05:00
}
2016-05-21 09:44:52 -05:00
if ($exception) {
throw $exception;
2016-05-21 12:19:48 -05:00
}
2016-05-21 09:44:52 -05:00
return $value;
}
2016-05-21 09:44:52 -05:00
/**
2017-04-23 15:47:52 +02:00
* Creates an artificial timeout for any `Promise`.
*
* If the timeout expires before the promise is resolved, the returned promise fails with an instance of
2017-04-23 15:47:52 +02:00
* `Amp\TimeoutException`.
*
* @template TReturn
2016-05-21 09:44:52 -05:00
*
* @param Promise<TReturn>|ReactPromise $promise Promise to which the timeout is applied.
* @param int $timeout Timeout in milliseconds.
*
* @return Promise<TReturn>
*
* @throws \TypeError If $promise is not an instance of \Amp\Promise or \React\Promise\PromiseInterface.
2016-05-21 09:44:52 -05:00
*/
2018-06-18 20:00:01 +02:00
function timeout($promise, int $timeout): Promise
{
if (!$promise instanceof Promise) {
if ($promise instanceof ReactPromise) {
$promise = adapt($promise);
} else {
throw createTypeError([Promise::class, ReactPromise::class], $promise);
}
}
$deferred = new Deferred;
2016-05-21 09:44:52 -05:00
$watcher = Loop::delay($timeout, static function () use (&$deferred) {
2017-04-07 18:47:44 +02:00
$temp = $deferred; // prevent double resolve
$deferred = null;
2017-04-07 18:47:44 +02:00
$temp->fail(new TimeoutException);
2016-05-21 09:44:52 -05:00
});
Loop::unreference($watcher);
2016-05-21 09:44:52 -05:00
$promise->onResolve(function () use (&$deferred, $promise, $watcher) {
if ($deferred !== null) {
Loop::cancel($watcher);
$deferred->resolve($promise);
}
});
2016-05-21 09:44:52 -05:00
2017-04-07 12:19:37 -05:00
return $deferred->promise();
2016-05-21 09:44:52 -05:00
}
2018-11-25 07:56:42 -09:00
/**
* Creates an artificial timeout for any `Promise`.
*
* If the promise is resolved before the timeout expires, the result is returned
*
* If the timeout expires before the promise is resolved, a default value is returned
*
* @template TReturn
2018-11-25 07:56:42 -09:00
*
* @param Promise<TReturn>|ReactPromise $promise Promise to which the timeout is applied.
* @param int $timeout Timeout in milliseconds.
* @param TReturn $default
*
* @return Promise<TReturn>
2018-11-25 07:56:42 -09:00
*
* @throws \TypeError If $promise is not an instance of \Amp\Promise or \React\Promise\PromiseInterface.
*/
function timeoutWithDefault($promise, int $timeout, $default = null): Promise
{
$promise = timeout($promise, $timeout);
return call(static function () use ($promise, $default) {
2018-11-25 07:56:42 -09:00
try {
return yield $promise;
} catch (TimeoutException $exception) {
return $default;
}
});
}
/**
* Adapts any object with a done(callable $onFulfilled, callable $onRejected) or then(callable $onFulfilled,
* callable $onRejected) method to a promise usable by components depending on placeholders implementing
* \AsyncInterop\Promise.
*
* @param object $promise Object with a done() or then() method.
*
2019-11-11 20:02:09 +01:00
* @return Promise Promise resolved by the $thenable object.
*
* @throws \Error If the provided object does not have a then() method.
*/
2018-06-18 20:00:01 +02:00
function adapt($promise): Promise
{
$deferred = new Deferred;
2016-05-21 09:44:52 -05:00
if (\method_exists($promise, 'done')) {
$promise->done([$deferred, 'resolve'], [$deferred, 'fail']);
} elseif (\method_exists($promise, 'then')) {
$promise->then([$deferred, 'resolve'], [$deferred, 'fail']);
} else {
throw new \Error("Object must have a 'then' or 'done' method");
}
return $deferred->promise();
}
2016-05-21 09:44:52 -05:00
/**
* Returns a promise that is resolved when all promises are resolved. The returned promise will not fail.
* Returned promise succeeds with a two-item array delineating successful and failed promise results,
* with keys identical and corresponding to the original given array.
*
* This function is the same as some() with the notable exception that it will never fail even
* if all promises in the array resolve unsuccessfully.
*
2019-11-11 20:02:09 +01:00
* @param Promise[]|ReactPromise[] $promises
*
2019-11-11 20:02:09 +01:00
* @return Promise
*
* @throws \Error If a non-Promise is in the array.
*/
2018-06-18 20:00:01 +02:00
function any(array $promises): Promise
{
return some($promises, 0);
2016-05-21 09:44:52 -05:00
}
/**
* Returns a promise that succeeds when all promises succeed, and fails if any promise fails. Returned
* promise succeeds with an array of values used to succeed each contained promise, with keys corresponding to
* the array of promises.
*
2019-11-11 20:02:09 +01:00
* @param Promise[]|ReactPromise[] $promises Array of only promises.
*
2019-11-11 20:02:09 +01:00
* @return Promise
*
* @throws \Error If a non-Promise is in the array.
*
2020-04-19 15:38:22 +02:00
* @template TValue
*
* @psalm-param array<array-key, Promise<TValue>|ReactPromise> $promises
* @psalm-assert array<array-key, Promise<TValue>|ReactPromise> $promises $promises
* @psalm-return Promise<array<array-key, TValue>>
*/
2018-06-18 20:00:01 +02:00
function all(array $promises): Promise
{
if (empty($promises)) {
return new Success([]);
}
2016-05-21 09:44:52 -05:00
$deferred = new Deferred;
$result = $deferred->promise();
2016-05-21 09:44:52 -05:00
$pending = \count($promises);
$values = [];
foreach ($promises as $key => $promise) {
if ($promise instanceof ReactPromise) {
$promise = adapt($promise);
} elseif (!$promise instanceof Promise) {
throw createTypeError([Promise::class, ReactPromise::class], $promise);
2016-05-23 21:32:41 -05:00
}
$values[$key] = null; // add entry to array to preserve order
2017-04-07 12:19:37 -05:00
$promise->onResolve(function ($exception, $value) use (&$deferred, &$values, &$pending, $key) {
2017-04-07 18:47:44 +02:00
if ($pending === 0) {
return;
}
2016-05-21 09:44:52 -05:00
if ($exception) {
2017-04-07 18:47:44 +02:00
$pending = 0;
$deferred->fail($exception);
2017-04-07 12:19:37 -05:00
$deferred = null;
return;
}
2016-05-21 09:44:52 -05:00
$values[$key] = $value;
if (0 === --$pending) {
$deferred->resolve($values);
}
});
}
2016-05-21 12:19:48 -05:00
return $result;
2016-05-21 09:44:52 -05:00
}
/**
* Returns a promise that succeeds when the first promise succeeds, and fails only if all promises fail.
*
2019-11-11 20:02:09 +01:00
* @param Promise[]|ReactPromise[] $promises Array of only promises.
*
2019-11-11 20:02:09 +01:00
* @return Promise
*
* @throws \Error If the array is empty or a non-Promise is in the array.
*/
2018-06-18 20:00:01 +02:00
function first(array $promises): Promise
{
if (empty($promises)) {
throw new \Error("No promises provided");
}
2016-05-21 09:44:52 -05:00
$deferred = new Deferred;
$result = $deferred->promise();
2016-05-21 12:19:48 -05:00
$pending = \count($promises);
$exceptions = [];
foreach ($promises as $key => $promise) {
if ($promise instanceof ReactPromise) {
$promise = adapt($promise);
} elseif (!$promise instanceof Promise) {
throw createTypeError([Promise::class, ReactPromise::class], $promise);
2016-07-31 00:31:04 -05:00
}
$exceptions[$key] = null; // add entry to array to preserve order
2017-12-05 08:48:36 +01:00
$promise->onResolve(function ($error, $value) use (&$deferred, &$exceptions, &$pending, $key) {
2017-04-07 18:47:44 +02:00
if ($pending === 0) {
2016-07-31 00:31:04 -05:00
return;
2016-05-21 09:44:52 -05:00
}
2017-04-13 18:49:32 +02:00
if (!$error) {
2017-04-07 18:47:44 +02:00
$pending = 0;
$deferred->resolve($value);
2017-04-07 12:19:37 -05:00
$deferred = null;
return;
}
2017-04-13 18:49:32 +02:00
$exceptions[$key] = $error;
if (0 === --$pending) {
$deferred->fail(new MultiReasonException($exceptions));
}
});
}
return $result;
2016-05-21 09:44:52 -05:00
}
/**
* Resolves with a two-item array delineating successful and failed Promise results.
*
* The returned promise will only fail if the given number of required promises fail.
*
2019-11-11 20:02:09 +01:00
* @param Promise[]|ReactPromise[] $promises Array of only promises.
* @param int $required Number of promises that must succeed for the
* returned promise to succeed.
*
2019-11-11 20:02:09 +01:00
* @return Promise
*
* @throws \Error If a non-Promise is in the array.
*/
2018-06-18 20:00:01 +02:00
function some(array $promises, int $required = 1): Promise
{
2017-03-27 11:42:11 -05:00
if ($required < 0) {
throw new \Error("Number of promises required must be non-negative");
}
$pending = \count($promises);
2016-05-21 09:44:52 -05:00
2017-03-27 11:42:11 -05:00
if ($required > $pending) {
throw new \Error("Too few promises provided");
}
if (empty($promises)) {
return new Success([[], []]);
}
$deferred = new Deferred;
$result = $deferred->promise();
$values = [];
$exceptions = [];
foreach ($promises as $key => $promise) {
if ($promise instanceof ReactPromise) {
$promise = adapt($promise);
} elseif (!$promise instanceof Promise) {
throw createTypeError([Promise::class, ReactPromise::class], $promise);
}
$values[$key] = $exceptions[$key] = null; // add entry to arrays to preserve order
$promise->onResolve(static function ($exception, $value) use (
2018-06-18 20:00:01 +02:00
&$values,
&$exceptions,
&$pending,
$key,
$required,
$deferred
) {
if ($exception) {
$exceptions[$key] = $exception;
unset($values[$key]);
} else {
$values[$key] = $value;
unset($exceptions[$key]);
}
2016-07-18 23:29:19 -05:00
if (0 === --$pending) {
if (\count($values) < $required) {
$deferred->fail(new MultiReasonException($exceptions));
2017-04-07 18:47:44 +02:00
} else {
$deferred->resolve([$exceptions, $values]);
}
}
});
}
return $result;
}
2018-11-26 09:36:46 -09:00
/**
* Wraps a promise into another promise, altering the exception or result.
*
2019-11-11 20:02:09 +01:00
* @param Promise|ReactPromise $promise
* @param callable $callback
*
2018-11-26 09:36:46 -09:00
* @return Promise
*/
function wrap($promise, callable $callback): Promise
{
if ($promise instanceof ReactPromise) {
$promise = adapt($promise);
} elseif (!$promise instanceof Promise) {
throw createTypeError([Promise::class, ReactPromise::class], $promise);
}
$deferred = new Deferred();
$promise->onResolve(static function (\Throwable $exception = null, $result) use ($deferred, $callback) {
2018-11-26 09:36:46 -09:00
try {
$result = $callback($exception, $result);
} catch (\Throwable $exception) {
$deferred->fail($exception);
return;
}
$deferred->resolve($result);
});
return $deferred->promise();
}
2016-07-18 23:23:25 -05:00
}
2018-06-18 20:00:01 +02:00
namespace Amp\Iterator
{
use Amp\Delayed;
use Amp\Emitter;
2017-04-27 10:51:06 -05:00
use Amp\Iterator;
use Amp\Producer;
use Amp\Promise;
2020-05-13 10:15:21 -05:00
use Amp\Stream;
2018-10-05 21:01:57 +02:00
use function Amp\call;
2017-04-26 13:06:41 -05:00
use function Amp\coroutine;
use function Amp\Internal\createTypeError;
/**
2017-04-27 10:51:06 -05:00
* Creates an iterator from the given iterable, emitting the each value. The iterable may contain promises. If any
* promise fails, the iterator will fail with the same reason.
*
* @param array|\Traversable $iterable Elements to emit.
2018-06-18 20:00:01 +02:00
* @param int $delay Delay between element emissions in milliseconds.
*
* @return Iterator
*
* @throws \TypeError If the argument is not an array or instance of \Traversable.
*/
2018-06-18 20:00:01 +02:00
function fromIterable(/* iterable */
$iterable,
int $delay = 0
): Iterator {
if (!$iterable instanceof \Traversable && !\is_array($iterable)) {
throw createTypeError(["array", "Traversable"], $iterable);
2016-05-24 11:47:14 -05:00
}
if ($delay) {
return new Producer(static function (callable $emit) use ($iterable, $delay) {
foreach ($iterable as $value) {
yield new Delayed($delay);
yield $emit($value);
}
});
}
return new Producer(static function (callable $emit) use ($iterable) {
foreach ($iterable as $value) {
yield $emit($value);
}
});
2016-05-24 11:47:14 -05:00
}
2016-12-11 16:17:51 +01:00
/**
* @template TValue
* @template TReturn
*
* @param Iterator<TValue> $iterator
* @param callable (TValue $value): TReturn $onEmit
*
* @return Iterator<TReturn>
*/
2018-06-18 20:00:01 +02:00
function map(Iterator $iterator, callable $onEmit): Iterator
{
return new Producer(static function (callable $emit) use ($iterator, $onEmit) {
2017-04-27 10:51:06 -05:00
while (yield $iterator->advance()) {
yield $emit($onEmit($iterator->getCurrent()));
}
});
}
/**
* @template TValue
*
* @param Iterator<TValue> $iterator
* @param callable(TValue $value):bool $filter
*
* @return Iterator<TValue>
*/
2018-06-18 20:00:01 +02:00
function filter(Iterator $iterator, callable $filter): Iterator
{
return new Producer(static function (callable $emit) use ($iterator, $filter) {
2017-04-27 10:51:06 -05:00
while (yield $iterator->advance()) {
if ($filter($iterator->getCurrent())) {
yield $emit($iterator->getCurrent());
}
}
});
}
/**
2017-04-27 10:51:06 -05:00
* Creates an iterator that emits values emitted from any iterator in the array of iterators.
*
* @param Iterator[] $iterators
*
* @return Iterator
*/
2018-06-18 20:00:01 +02:00
function merge(array $iterators): Iterator
{
$emitter = new Emitter;
2017-05-01 00:29:23 -05:00
$result = $emitter->iterate();
$coroutine = coroutine(static function (Iterator $iterator) use (&$emitter) {
2017-04-27 10:51:06 -05:00
while ((yield $iterator->advance()) && $emitter !== null) {
yield $emitter->emit($iterator->getCurrent());
2017-04-26 13:06:41 -05:00
}
});
$coroutines = [];
2017-04-27 10:51:06 -05:00
foreach ($iterators as $iterator) {
if (!$iterator instanceof Iterator) {
throw createTypeError([Iterator::class], $iterator);
}
2020-03-28 22:20:44 +01:00
2017-04-27 10:51:06 -05:00
$coroutines[] = $coroutine($iterator);
2016-05-24 11:47:14 -05:00
}
2016-12-11 16:17:51 +01:00
Promise\all($coroutines)->onResolve(static function ($exception) use (&$emitter) {
if ($exception) {
$emitter->fail($exception);
$emitter = null;
} else {
2017-04-26 13:14:10 -05:00
$emitter->complete();
}
});
return $result;
2016-08-01 11:10:59 -05:00
}
2016-12-11 16:17:51 +01:00
/**
2017-04-27 10:51:06 -05:00
* Concatenates the given iterators into a single iterator, emitting values from a single iterator at a time. The
* prior iterator must complete before values are emitted from any subsequent iterators. Iterators are concatenated
* in the order given (iteration order of the array).
*
* @param Iterator[] $iterators
*
* @return Iterator
*/
2018-06-18 20:00:01 +02:00
function concat(array $iterators): Iterator
{
2017-04-27 10:51:06 -05:00
foreach ($iterators as $iterator) {
if (!$iterator instanceof Iterator) {
throw createTypeError([Iterator::class], $iterator);
}
}
2016-12-11 16:17:51 +01:00
$emitter = new Emitter;
$previous = [];
$promise = Promise\all($previous);
$coroutine = coroutine(static function (Iterator $iterator, callable $emit) {
2017-04-27 10:51:06 -05:00
while (yield $iterator->advance()) {
yield $emit($iterator->getCurrent());
2017-04-26 13:06:41 -05:00
}
});
2017-04-27 10:51:06 -05:00
foreach ($iterators as $iterator) {
$emit = coroutine(static function ($value) use ($emitter, $promise) {
static $pending = true, $failed = false;
if ($failed) {
return;
}
if ($pending) {
try {
yield $promise;
$pending = false;
} catch (\Throwable $exception) {
$failed = true;
2017-04-27 10:51:06 -05:00
return; // Prior iterator failed.
}
}
2016-12-11 16:17:51 +01:00
yield $emitter->emit($value);
});
2017-04-27 10:51:06 -05:00
$previous[] = $coroutine($iterator, $emit);
$promise = Promise\all($previous);
2016-08-01 11:10:59 -05:00
}
2016-12-11 16:17:51 +01:00
$promise->onResolve(static function ($exception) use ($emitter) {
if ($exception) {
$emitter->fail($exception);
return;
}
2016-12-11 16:17:51 +01:00
2017-04-26 13:14:10 -05:00
$emitter->complete();
});
2016-08-01 11:10:59 -05:00
2017-05-01 00:29:23 -05:00
return $emitter->iterate();
2016-05-24 11:47:14 -05:00
}
2018-10-05 21:01:57 +02:00
2020-05-06 18:57:29 +02:00
/**
* Discards all remaining items and returns the number of discarded items.
*
* @template TValue
*
* @param Iterator $iterator
*
* @return Promise
*
* @psalm-param Iterator<TValue> $iterator
* @psalm-return Promise<int>
*/
function discard(Iterator $iterator): Promise
{
return call(static function () use ($iterator): \Generator {
$count = 0;
while (yield $iterator->advance()) {
$count++;
}
return $count;
});
}
2018-10-05 21:01:57 +02:00
/**
* Collects all items from an iterator into an array.
*
* @template TValue
*
2018-10-05 21:01:57 +02:00
* @param Iterator $iterator
*
* @psalm-param Iterator<TValue> $iterator
*
* @return Promise
2020-05-13 10:15:21 -05:00
* @psalm-return Promise<array<int, TValue>>
2018-10-05 21:01:57 +02:00
*/
function toArray(Iterator $iterator): Promise
2018-10-05 21:01:57 +02:00
{
2020-05-13 10:15:21 -05:00
return call(static function () use ($iterator): \Generator {
/** @psalm-var list $array */
2018-10-05 21:01:57 +02:00
$array = [];
2018-10-05 21:01:57 +02:00
while (yield $iterator->advance()) {
$array[] = $iterator->getCurrent();
}
return $array;
});
}
2020-05-13 10:15:21 -05:00
function fromStream(Stream $stream): Iterator
{
return new Producer(function (callable $emit) use ($stream): \Generator {
while (null !== $value = $stream->continue()) {
yield $emit($value);
}
});
}
}
namespace Amp\Stream
{
use Amp\AsyncGenerator;
use Amp\Delayed;
use Amp\Iterator;
use Amp\Promise;
use Amp\Stream;
use Amp\StreamSource;
use function Amp\call;
use function Amp\coroutine;
use function Amp\Internal\createTypeError;
/**
* Creates a stream from the given iterable, emitting the each value. The iterable may contain promises. If any
* promise fails, the returned stream will fail with the same reason.
*
* @param array|\Traversable $iterable Elements to yield.
* @param int $delay Delay between elements yielded in milliseconds.
*
* @return Stream
*
* @throws \TypeError If the argument is not an array or instance of \Traversable.
*/
function fromIterable(/* iterable */
$iterable,
int $delay = 0
): Stream {
if (!$iterable instanceof \Traversable && !\is_array($iterable)) {
throw createTypeError(["array", "Traversable"], $iterable);
}
if ($delay) {
return new AsyncGenerator(static function (callable $yield) use ($iterable, $delay) {
foreach ($iterable as $value) {
yield new Delayed($delay);
yield $yield($value instanceof Promise ? yield $value : $value);
}
});
}
return new AsyncGenerator(static function (callable $yield) use ($iterable) {
foreach ($iterable as $value) {
yield $yield($value instanceof Promise ? yield $value : $value);
}
});
}
/**
* @template TValue
* @template TReturn
*
* @param Stream $stream
* @param callable (TValue $value, int $key): TReturn $onYield
*
* @psalm-param Stream<TValue> $stream
*
* @return Stream
*
* @psalm-return Stream<TReturn>
*/
function map(Stream $stream, callable $onYield): Stream
{
return new AsyncGenerator(static function (callable $yield) use ($stream, $onYield) {
while (list($value, $key) = yield $stream->continue()) {
yield $yield($onYield($value, $key));
}
});
}
/**
* @template TValue
*
* @param Stream $stream
* @param callable(TValue $value, int $key):bool $filter
*
* @psalm-param Stream<TValue> $stream
*
* @return Stream
*
* @psalm-return Stream<TValue>
*/
function filter(Stream $stream, callable $filter): Stream
{
return new AsyncGenerator(static function (callable $yield) use ($stream, $filter) {
while (list($value, $key) = yield $stream->continue()) {
if ($filter($value, $key)) {
yield $yield($value, $key);
}
}
});
}
/**
* Creates a stream that yields values emitted from any stream in the array of streams.
*
* @param Stream[] $streams
*
* @return Stream
*/
function merge(array $streams): Stream
{
$source = new StreamSource;
$result = $source->stream();
$coroutine = coroutine(static function (Stream $stream) use (&$source) {
while ((list($value) = yield $stream->continue()) && $source !== null) {
yield $source->yield($value);
}
});
$coroutines = [];
foreach ($streams as $stream) {
if (!$stream instanceof Stream) {
throw createTypeError([Stream::class], $stream);
}
$coroutines[] = $coroutine($stream);
}
Promise\all($coroutines)->onResolve(static function ($exception) use (&$source) {
$temp = $source;
$source = null;
if ($exception) {
$temp->fail($exception);
} else {
$temp->complete();
}
});
return $result;
}
/**
* Concatenates the given streams into a single stream, yielding values from a single stream at a time. The
* prior stream must complete before values are yielded from any subsequent streams. Streams are concatenated
* in the order given (iteration order of the array).
*
* @param Stream[] $streams
*
* @return Stream
*/
function concat(array $streams): Stream
{
foreach ($streams as $stream) {
if (!$stream instanceof Stream) {
throw createTypeError([Stream::class], $stream);
}
}
$source = new StreamSource;
$previous = [];
$promise = Promise\all($previous);
$coroutine = coroutine(static function (Stream $stream, callable $yield) {
while (list($value) = yield $stream->continue()) {
yield $yield($value);
}
});
foreach ($streams as $iterator) {
$emit = coroutine(static function ($value) use ($source, $promise) {
static $pending = true, $failed = false;
if ($failed) {
return;
}
if ($pending) {
try {
yield $promise;
$pending = false;
} catch (\Throwable $exception) {
$failed = true;
return; // Prior iterator failed.
}
}
yield $source->yield($value);
});
$previous[] = $coroutine($iterator, $emit);
$promise = Promise\all($previous);
}
$promise->onResolve(static function ($exception) use ($source) {
if ($exception) {
$source->fail($exception);
return;
}
$source->complete();
});
return $source->stream();
}
/**
* Discards all remaining items and returns the number of discarded items.
*
* @template TValue
*
* @param Stream $stream
*
* @psalm-param Stream<TValue> $stream
*
* @return Promise<int>
*
* @psalm-return Promise<int>
*/
function discard(Stream $stream): Promise
{
return call(static function () use ($stream): \Generator {
$count = 0;
while (yield $stream->continue()) {
$count++;
}
return $count;
});
}
/**
* Collects all items from a stream into an array.
*
* @template TValue
*
* @param Stream $stream
*
* @psalm-param Stream<TValue> $stream
*
* @return Promise
*
* @psalm-return Promise<array<int, TValue>>
*/
function toArray(Stream $stream): Promise
{
return call(static function () use ($stream): \Generator {
/** @psalm-var list $array */
$array = [];
while (list($value) = yield $stream->continue()) {
$array[] = $value;
}
return $array;
});
}
function fromIterator(Iterator $iterator): Stream
{
return new AsyncGenerator(function (callable $yield) use ($iterator): \Generator {
while (yield $iterator->advance()) {
yield $yield($iterator->getCurrent());
}
});
}
2016-05-24 11:47:14 -05:00
}