.gitattributes000064400000000054151360552510007441 0ustar00/tests export-ignore /.github export-ignore LICENSE000064400000002054151360552510005555 0ustar00The MIT License (MIT) Copyright (c) Hyperf Permission is hereby granted, free of charge, to any person obtaining a copy of this software and associated documentation files (the "Software"), to deal in the Software without restriction, including without limitation the rights to use, copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the Software, and to permit persons to whom the Software is furnished to do so, subject to the following conditions: The above copyright notice and this permission notice shall be included in all copies or substantial portions of the Software. THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE. composer.json000064400000002021151360552510007264 0ustar00{ "name": "hyperf/coroutine", "description": "Hyperf Coroutine", "license": "MIT", "keywords": [ "php", "swoole", "hyperf", "coroutine" ], "homepage": "https://hyperf.io", "support": { "docs": "https://hyperf.wiki", "issues": "https://github.com/hyperf/hyperf/issues", "pull-request": "https://github.com/hyperf/hyperf/pulls", "source": "https://github.com/hyperf/hyperf" }, "require": { "php": ">=8.0", "hyperf/context": "~3.0.0", "hyperf/contract": "~3.0.0", "hyperf/engine": "^1.2|^2.0" }, "autoload": { "psr-4": { "Hyperf\\Coroutine\\": "src/" }, "files": [ "src/Functions.php" ] }, "autoload-dev": { "psr-4": { "HyperfTest\\Coroutine\\": "tests/" } }, "config": { "sort-packages": true }, "extra": { "branch-alias": { "dev-master": "3.0-dev" } } } src/Concurrent.php000064400000004417151360552510010177 0ustar00channel = new Channel($limit); } public function __call($name, $arguments) { if (in_array($name, ['isFull', 'isEmpty'])) { return $this->channel->{$name}(...$arguments); } throw new InvalidArgumentException(sprintf('The method %s is not supported.', $name)); } public function getLimit(): int { return $this->limit; } public function length(): int { return $this->channel->getLength(); } public function getLength(): int { return $this->channel->getLength(); } public function getRunningCoroutineCount(): int { return $this->getLength(); } public function getChannel(): Channel { return $this->channel; } public function create(callable $callable): void { $this->channel->push(true); Coroutine::create(function () use ($callable) { try { $callable(); } catch (Throwable $exception) { if (ApplicationContext::hasContainer()) { $container = ApplicationContext::getContainer(); if ($container->has(StdoutLoggerInterface::class) && $container->has(FormatterInterface::class)) { $logger = $container->get(StdoutLoggerInterface::class); $formatter = $container->get(FormatterInterface::class); $logger->error($formatter->format($exception)); } } } finally { $this->channel->pop(); } }); } } src/Coroutine.php000064400000006026151360552510010022 0ustar00getId(); } catch (\Throwable) { return -1; } } public static function inCoroutine(): bool { return Co::id() > 0; } public static function stats(): array { return Co::stats(); } public static function exists(int $id): bool { return Co::exists($id); } private static function printLog(Throwable $throwable): void { if (ApplicationContext::hasContainer()) { $container = ApplicationContext::getContainer(); if ($container->has(StdoutLoggerInterface::class)) { $logger = $container->get(StdoutLoggerInterface::class); if ($container->has(FormatterInterface::class)) { $formatter = $container->get(FormatterInterface::class); $logger->warning($formatter->format($throwable)); } else { $logger->warning((string) $throwable); } } } } } src/Exception/ExceptionThrower.php000064400000000754151360552510013324 0ustar00throwable; } } src/Exception/InvalidArgumentException.php000064400000000532151360552510014755 0ustar00results; } public function setResults(array $results) { $this->results = $results; } public function getThrowables(): array { return $this->throwables; } public function setThrowables(array $throwables) { return $this->throwables = $throwables; } } src/Exception/TimeoutException.php000064400000000540151360552510013311 0ustar00 $callable) { $parallel->add($callable, $key); } return $parallel->wait(); } function wait(Closure $closure, ?float $timeout = null) { if (ApplicationContext::hasContainer()) { $waiter = ApplicationContext::getContainer()->get(Waiter::class); return $waiter->wait($closure, $timeout); } return (new Waiter())->wait($closure, $timeout); } function co(callable $callable): bool|int { $id = Coroutine::create($callable); return $id > 0 ? $id : false; } function defer(callable $callable): void { Coroutine::defer($callable); } function go(callable $callable): bool|int { $id = Coroutine::create($callable); return $id > 0 ? $id : false; } /** * Run callable in non-coroutine environment, all hook functions by Swoole only available in the callable. * * @param array|callable $callbacks */ function run($callbacks, int $flags = SWOOLE_HOOK_ALL): bool { if (Coroutine::inCoroutine()) { throw new RuntimeException('Function \'run\' only execute in non-coroutine environment.'); } \Swoole\Runtime::enableCoroutine($flags); /* @phpstan-ignore-next-line */ $result = \Swoole\Coroutine\run(...(array) $callbacks); \Swoole\Runtime::enableCoroutine(false); return $result; } src/Locker.php000064400000003363151360552510007273 0ustar00= 2.0, use `Co::yield()` instead. match (Constant::ENGINE) { 'Swoole' => SwooleCoroutine::yield(), /* @phpstan-ignore-next-line */ default => Co::yield(), }; return false; } public static function unlock(string $key): void { if (self::has($key)) { $ids = self::get($key); foreach ($ids as $id) { if ($id > 0) { // TODO: When the verion of `hyperf/engine` >= 2.0, use `Co::resumeById()` instead. match (Constant::ENGINE) { 'Swoole' => SwooleCoroutine::resume($id), /* @phpstan-ignore-next-line */ default => Co::resumeById($id), }; } } self::clear($key); } } } src/Parallel.php000064400000006201151360552510007602 0ustar00 0) { $this->concurrentChannel = new Channel($concurrent); } } public function add(callable $callable, $key = null) { if (is_null($key)) { $this->callbacks[] = $callable; } else { $this->callbacks[$key] = $callable; } } public function wait(bool $throw = true): array { $wg = new WaitGroup(); $wg->add(count($this->callbacks)); foreach ($this->callbacks as $key => $callback) { $this->concurrentChannel && $this->concurrentChannel->push(true); $this->results[$key] = null; Coroutine::create(function () use ($callback, $key, $wg) { try { $this->results[$key] = $callback(); } catch (Throwable $throwable) { $this->throwables[$key] = $throwable; unset($this->results[$key]); } finally { $this->concurrentChannel && $this->concurrentChannel->pop(); $wg->done(); } }); } $wg->wait(); if ($throw && ($throwableCount = count($this->throwables)) > 0) { $message = 'Detecting ' . $throwableCount . ' throwable occurred during parallel execution:' . PHP_EOL . $this->formatThrowables($this->throwables); $executionException = new ParallelExecutionException($message); $executionException->setResults($this->results); $executionException->setThrowables($this->throwables); unset($this->results, $this->throwables); throw $executionException; } return $this->results; } public function count(): int { return count($this->callbacks); } public function clear(): void { $this->callbacks = []; $this->results = []; $this->throwables = []; } /** * Format throwables into a nice list. * * @param Throwable[] $throwables */ private function formatThrowables(array $throwables): string { $output = ''; foreach ($throwables as $key => $value) { $output .= sprintf('(%s) %s: %s' . PHP_EOL . '%s' . PHP_EOL, $key, get_class($value), $value->getMessage(), $value->getTraceAsString()); } return $output; } } src/Traits/Container.php000064400000002506151360552510011242 0ustar00chan = new Channel(1); if ($delta > 0) { $this->add($delta); } } public function add(int $delta = 1): void { if ($this->waiting) { throw new BadMethodCallException('WaitGroup misuse: add called concurrently with wait'); } $count = $this->count + $delta; if ($count < 0) { throw new InvalidArgumentException('WaitGroup misuse: negative counter'); } $this->count = $count; } public function done(): void { $count = $this->count - 1; if ($count < 0) { throw new BadMethodCallException('WaitGroup misuse: negative counter'); } $this->count = $count; if ($count === 0 && $this->waiting) { $this->chan->push(true); } } public function wait(float $timeout = -1): bool { if ($this->waiting) { throw new BadMethodCallException('WaitGroup misuse: reused before previous wait has returned'); } if ($this->count > 0) { $this->waiting = true; $done = $this->chan->pop($timeout); $this->waiting = false; return $done; } return true; } public function count(): int { return $this->count; } } src/Waiter.php000064400000003065151360552510007306 0ustar00popTimeout = $timeout; } /** * @param null|float $timeout seconds */ public function wait(Closure $closure, ?float $timeout = null) { if ($timeout === null) { $timeout = $this->popTimeout; } $channel = new Channel(1); Coroutine::create(function () use ($channel, $closure) { try { $result = $closure(); } catch (Throwable $exception) { $result = new ExceptionThrower($exception); } finally { $channel->push($result ?? null, $this->pushTimeout); } }); $result = $channel->pop($timeout); if ($result === false && $channel->isTimeout()) { throw new WaitTimeoutException(sprintf('Channel wait failed, reason: Timed out for %s s', $timeout)); } if ($result instanceof ExceptionThrower) { throw $result->getThrowable(); } return $result; } }