.gitattributes000064400000000054151360551660007446 0ustar00/tests export-ignore /.github export-ignore LICENSE000064400000002054151360551660005562 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.json000064400000002014151360551660007273 0ustar00{ "name": "hyperf/coroutine", "description": "Hyperf Coroutine", "license": "MIT", "keywords": [ "php", "swoole", "hyperf", "coroutine" ], "homepage": "https://hyperf.io", "support": { "issues": "https://github.com/hyperf/hyperf/issues", "source": "https://github.com/hyperf/hyperf", "docs": "https://hyperf.wiki", "pull-request": "https://github.com/hyperf/hyperf/pulls" }, "require": { "php": ">=8.1", "hyperf/context": "~3.1.0", "hyperf/contract": "~3.1.0", "hyperf/engine": "^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.1-dev" } } } src/Channel/Caller.php000064400000003155151360551660010632 0ustar00initInstance(); } public function call(Closure $closure) { $release = true; $channel = $this->channel; try { $instance = $channel->pop($this->waitTimeout); if ($instance === false) { if ($channel->isClosing()) { throw new ChannelClosedException('The channel was closed.'); } if ($channel->isTimeout()) { throw new WaitTimeoutException('The instance pop from channel timeout.'); } } $result = $closure($instance); } catch (ChannelClosedException|WaitTimeoutException $exception) { $release = false; throw $exception; } finally { $release && $channel->push($instance ?? null); } return $result; } public function initInstance(): void { $this->channel?->close(); $this->channel = new Channel(1); $this->channel->push($this->closure->__invoke()); } } src/Channel/Manager.php000064400000002446151360551660011004 0ustar00channels[$id])) { return $this->channels[$id]; } if ($initialize) { return $this->channels[$id] = $this->make($this->size); } return null; } public function make(int $limit): Channel { return new Channel($limit); } public function close(int $id): void { if ($channel = $this->channels[$id] ?? null) { $channel->close(); } unset($this->channels[$id]); } public function getChannels(): array { return $this->channels; } public function flush(): void { $channels = $this->getChannels(); foreach ($channels as $id => $channel) { $this->close($id); } } } src/Channel/Pool.php000064400000001342151360551660010335 0ustar00isEmpty() ? new Channel(1) : $this->pop(); } public function release(Channel $channel) { $channel->errCode = 0; $this->push($channel); } } src/Concurrent.php000064400000004420151360551660010176 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.php000064400000007414151360551660010031 0ustar00getId(); } catch (Throwable) { return -1; } } /** * Create a coroutine with a copy of the parent coroutine context. */ public static function fork(callable $callable, array $keys = []): int { $cid = static::id(); $callable = static function () use ($callable, $cid, $keys) { Context::copy($cid, $keys); $callable(); }; return static::create($callable); } 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/ChannelClosedException.php000064400000000547151360551660014401 0ustar00throwable; } } src/Exception/InvalidArgumentException.php000064400000000533151360551660014763 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.php000064400000000541151360551660013317 0ustar00 $callable) { $parallel->add($callable, $key); } return $parallel->wait(); } /** * @template TReturn * * @param Closure():TReturn $closure * @return TReturn */ 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.'); } Runtime::enableCoroutine($flags); /* @phpstan-ignore-next-line */ $result = \Swoole\Coroutine\run(...(array) $callbacks); Runtime::enableCoroutine(0); return $result; } src/Locker.php000064400000001646151360551660007302 0ustar00pop(-1); return false; } public static function unlock(string $key): void { if (isset(static::$channels[$key])) { $channel = static::$channels[$key]; static::$channels[$key] = null; $channel->close(); } } } src/Parallel.php000064400000006202151360551660007610 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/WaitGroup.php000064400000003572151360551660010004 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.php000064400000003226151360551660007312 0ustar00popTimeout = $timeout; } /** * @template TReturn * * @param Closure():TReturn $closure * @param null|float $timeout seconds * @return TReturn */ 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; } }