.gitattributes000064400000000054151360554300007440 0ustar00/tests export-ignore /.github export-ignore LICENSE000064400000002054151360554300005554 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.json000064400000002014151360554300007265 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.php000064400000003155151360554300010624 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.php000064400000002446151360554300010776 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.php000064400000001342151360554300010327 0ustar00isEmpty() ? new Channel(1) : $this->pop(); } public function release(Channel $channel) { $channel->errCode = 0; $this->push($channel); } } src/Concurrent.php000064400000004420151360554300010170 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.php000064400000007414151360554300010023 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.php000064400000000547151360554300014373 0ustar00throwable; } } src/Exception/InvalidArgumentException.php000064400000000533151360554300014755 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.php000064400000000541151360554300013311 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.php000064400000001646151360554300007274 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/Mutex.php000064400000002603151360554300007151 0ustar00push(1, $timeout); if ($channel->isTimeout() || $channel->isClosing()) { return false; } return true; } public static function unlock(string $key, float $timeout = 5): bool { if (isset(static::$channels[$key])) { $channel = static::$channels[$key]; $channel->pop($timeout); if ($channel->isTimeout()) { // unlock more than once return false; } } return true; } public static function clear(string $key): void { if (isset(static::$channels[$key])) { $channel = static::$channels[$key]; static::$channels[$key] = null; $channel->close(); } } } src/Parallel.php000064400000006202151360554300007602 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.php000064400000003572151360554300007776 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.php000064400000003226151360554300007304 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; } }