.gitattributes000064400000000054151360554200007437 0ustar00/tests export-ignore /.github export-ignore LICENSE000064400000002054151360554200005553 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.json000064400000002014151360554200007264 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.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.0-dev" } } } src/Channel/Caller.php000064400000003154151360554200010622 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.php000064400000002445151360554200010774 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.php000064400000001341151360554200010325 0ustar00isEmpty() ? new Channel(1) : $this->pop(); } public function release(Channel $channel) { $channel->errCode = 0; $this->push($channel); } } src/Concurrent.php000064400000004417151360554200010175 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.php000064400000006026151360554200010020 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/ChannelClosedException.php000064400000000546151360554200014371 0ustar00throwable; } } src/Exception/InvalidArgumentException.php000064400000000532151360554200014753 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.php000064400000000540151360554200013307 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.php000064400000001645151360554200007272 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.php000064400000006201151360554200007600 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.php000064400000002506151360554200011240 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.php000064400000003065151360554200007304 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; } }