1
0
mirror of https://github.com/danog/amp.git synced 2024-12-13 01:47:33 +01:00
amp/test/Pipeline/MergeTest.php

103 lines
2.8 KiB
PHP
Raw Normal View History

<?php
2020-08-23 16:18:28 +02:00
namespace Amp\Test\Pipeline;
2020-05-13 17:15:21 +02:00
use Amp\AsyncGenerator;
2017-05-02 07:22:53 +02:00
use Amp\Delayed;
2020-05-17 21:41:42 +02:00
use Amp\PHPUnit\AsyncTestCase;
use Amp\PHPUnit\TestException;
2020-08-23 16:18:28 +02:00
use Amp\Pipeline;
2020-09-28 05:19:52 +02:00
use function Amp\await;
2020-05-17 21:41:42 +02:00
class MergeTest extends AsyncTestCase
2018-06-18 20:00:01 +02:00
{
2020-05-17 21:41:42 +02:00
public function getArrays(): array
2018-06-18 20:00:01 +02:00
{
return [
2017-03-14 00:52:57 +01:00
[[\range(1, 3), \range(4, 6)], [1, 4, 2, 5, 3, 6]],
[[\range(1, 5), \range(6, 8)], [1, 6, 2, 7, 3, 8, 4, 5]],
[[\range(1, 4), \range(5, 10)], [1, 5, 2, 6, 3, 7, 4, 8, 9, 10]],
];
}
/**
2017-03-14 00:52:57 +01:00
* @dataProvider getArrays
*
2020-09-28 05:19:52 +02:00
* @param array $array
* @param array $expected
*/
2020-09-28 05:19:52 +02:00
public function testMerge(array $array, array $expected): void
2018-06-18 20:00:01 +02:00
{
2020-08-23 16:18:28 +02:00
$pipelines = \array_map(static function (array $iterator): Pipeline {
return Pipeline\fromIterable($iterator);
2020-09-28 05:19:52 +02:00
}, $array);
2020-08-23 16:18:28 +02:00
$pipeline = Pipeline\merge($pipelines);
2017-03-14 00:52:57 +01:00
2020-09-28 05:19:52 +02:00
while (null !== $value = $pipeline->continue()) {
2020-05-21 17:11:22 +02:00
$this->assertSame(\array_shift($expected), $value);
2020-05-17 21:41:42 +02:00
}
2017-04-27 17:32:53 +02:00
}
/**
* @depends testMerge
*/
2020-09-28 05:19:52 +02:00
public function testMergeWithDelayedYields(): void
2018-06-18 20:00:01 +02:00
{
2020-08-23 16:18:28 +02:00
$pipelines = [];
2020-05-17 21:41:42 +02:00
$values1 = [new Delayed(10, 1), new Delayed(50, 2), new Delayed(70, 3)];
$values2 = [new Delayed(20, 4), new Delayed(40, 5), new Delayed(60, 6)];
$expected = [1, 4, 5, 2, 6, 3];
2020-09-28 05:19:52 +02:00
$pipelines[] = new AsyncGenerator(function () use ($values1) {
2020-05-17 21:41:42 +02:00
foreach ($values1 as $value) {
2020-09-28 05:19:52 +02:00
yield await($value);
2020-05-17 21:41:42 +02:00
}
});
2020-09-28 05:19:52 +02:00
$pipelines[] = new AsyncGenerator(function () use ($values2) {
2020-05-17 21:41:42 +02:00
foreach ($values2 as $value) {
2020-09-28 05:19:52 +02:00
yield await($value);
2017-04-27 17:32:53 +02:00
}
});
2020-05-17 21:41:42 +02:00
2020-08-23 16:18:28 +02:00
$pipeline = Pipeline\merge($pipelines);
2020-05-17 21:41:42 +02:00
2020-09-28 05:19:52 +02:00
while (null !== $value = $pipeline->continue()) {
2020-05-21 17:11:22 +02:00
$this->assertSame(\array_shift($expected), $value);
2020-05-17 21:41:42 +02:00
}
}
/**
* @depends testMerge
*/
2020-09-28 05:19:52 +02:00
public function testMergeWithFailedPipeline(): void
2018-06-18 20:00:01 +02:00
{
2020-05-17 21:41:42 +02:00
$exception = new TestException;
2020-09-28 05:19:52 +02:00
$generator = new AsyncGenerator(static function () use ($exception) {
yield 1; // Emit once before failing.
2020-05-17 21:41:42 +02:00
throw $exception;
});
2020-05-17 21:41:42 +02:00
2020-08-23 16:18:28 +02:00
$pipeline = Pipeline\merge([$generator, Pipeline\fromIterable(\range(1, 5))]);
2020-05-17 21:41:42 +02:00
try {
/** @noinspection PhpStatementHasEmptyBodyInspection */
2020-09-28 05:19:52 +02:00
while ($pipeline->continue()) {
2020-05-17 21:41:42 +02:00
;
}
2020-08-23 16:18:28 +02:00
$this->fail("The exception used to fail the pipeline should be thrown from continue()");
2020-05-17 21:41:42 +02:00
} catch (TestException $reason) {
$this->assertSame($exception, $reason);
}
}
2020-09-28 05:19:52 +02:00
public function testNonPipeline(): void
2018-06-18 20:00:01 +02:00
{
2020-05-13 17:15:21 +02:00
$this->expectException(\TypeError::class);
2020-05-17 21:41:42 +02:00
/** @noinspection PhpParamsInspection */
2020-08-23 16:18:28 +02:00
Pipeline\merge([1]);
}
}