.gitattributes000064400000000105151360553570007445 0ustar00/.github export-ignore /examples export-ignore /tests export-ignore .gitignore000064400000000044151360553570006544 0ustar00/vendor/ composer.lock *.cache *.log.php-cs-fixer.php000064400000005225151360553570007656 0ustar00setRiskyAllowed(true) ->setRules([ '@PSR2' => true, '@Symfony' => true, '@DoctrineAnnotation' => true, '@PhpCsFixer' => true, 'header_comment' => [ 'comment_type' => 'PHPDoc', 'header' => $header, 'separate' => 'none', 'location' => 'after_declare_strict', ], 'array_syntax' => [ 'syntax' => 'short' ], 'list_syntax' => [ 'syntax' => 'short' ], 'concat_space' => [ 'spacing' => 'one' ], 'blank_line_before_statement' => [ 'statements' => [ 'declare', ], ], 'general_phpdoc_annotation_remove' => [ 'annotations' => [ 'author' ], ], 'ordered_imports' => [ 'imports_order' => [ 'class', 'function', 'const', ], 'sort_algorithm' => 'alpha', ], 'single_line_comment_style' => [ 'comment_types' => [ ], ], 'yoda_style' => [ 'always_move_variable' => false, 'equal' => false, 'identical' => false, ], 'phpdoc_align' => [ 'align' => 'left', ], 'multiline_whitespace_before_semicolons' => [ 'strategy' => 'no_multi_line', ], 'constant_case' => [ 'case' => 'lower', ], 'global_namespace_import' => [ 'import_classes' => true, 'import_constants' => true, 'import_functions' => true, ], 'class_attributes_separation' => true, 'combine_consecutive_unsets' => true, 'declare_strict_types' => true, 'linebreak_after_opening_tag' => true, 'lowercase_static_reference' => true, 'no_useless_else' => true, 'no_unused_imports' => true, 'not_operator_with_successor_space' => true, 'not_operator_with_space' => false, 'ordered_class_elements' => true, 'php_unit_strict' => false, 'phpdoc_separation' => false, 'single_quote' => true, 'standardize_not_equals' => true, 'multiline_comment_opening_closing' => true, ]) ->setFinder( PhpCsFixer\Finder::create() ->exclude('vendor') ->in(__DIR__) ) ->setUsingCache(false); .phpstorm.meta.php000064400000000164151360553570010147 0ustar00" version="1.0" license="MIT" app.name="Hyperf" ARG timezone ARG PHP_VERSION ENV TIMEZONE=${timezone:-"Asia/Shanghai"} ENV COMPOSER_ROOT_VERSION="v1.2.0" # update RUN set -ex \ # show php version and extensions && php -v \ && php -m \ && php --ri swoole \ # ---------- some config ---------- && cd "/etc/php${PHP_VERSION%\.*}" \ # - config PHP && { \ echo "upload_max_filesize=128M"; \ echo "post_max_size=128M"; \ echo "memory_limit=1G"; \ echo "date.timezone=${TIMEZONE}"; \ } | tee conf.d/99_overrides.ini \ # - config timezone && ln -sf /usr/share/zoneinfo/${TIMEZONE} /etc/localtime \ && echo "${TIMEZONE}" > /etc/timezone \ # ---------- clear works ---------- && rm -rf /var/cache/apk/* /tmp/* /usr/share/man \ && echo -e "\033[42;37m Build Completed :).\033[0m\n" WORKDIR /opt/www COPY . /opt/www RUN composer install -o README.md000064400000000234151360553570006034 0ustar00# Swoole Engine ![Swoole Engine Test](https://github.com/hyperf/engine/workflows/Swoole%20Engine%20Test/badge.svg) ``` composer require hyperf/engine ``` composer.json000064400000002310151360553570007274 0ustar00{ "name": "hyperf/engine", "type": "library", "license": "MIT", "keywords": [ "php", "hyperf" ], "description": "", "autoload": { "psr-4": { "Hyperf\\Engine\\": "src/" } }, "autoload-dev": { "psr-4": { "HyperfTest\\": "tests" } }, "require": { "php": ">=8.0" }, "require-dev": { "friendsofphp/php-cs-fixer": "^3.0", "hyperf/guzzle": "^2.2", "hyperf/http-message": " ^2.0", "phpstan/phpstan": "^1.0", "phpunit/phpunit": "^9.4", "swoole/ide-helper": "dev-master" }, "suggest": { "ext-swoole": ">=4.5" }, "minimum-stability": "dev", "prefer-stable": true, "config": { "optimize-autoloader": true, "sort-packages": true }, "extra": { "branch-alias": { "dev-master": "1.8-dev" }, "hyperf": { "config": "Hyperf\\Engine\\ConfigProvider" } }, "scripts": { "test": "phpunit -c phpunit.xml --colors=always", "analyse": "phpstan analyse --memory-limit 1024M -l 0 ./src", "cs-fix": "php-cs-fixer fix $1" } } phpunit.xml000064400000000755151360553570006776 0ustar00 ./tests/ src/Channel.php000064400000006522151360553570007433 0ustar00 80000 && SWOOLE_VERSION_ID >= 50000) { class Channel extends \Swoole\Coroutine\Channel implements ChannelInterface { protected bool $closed = false; public function push(mixed $data, float $timeout = -1): bool { return parent::push($data, $timeout); } public function pop(float $timeout = -1): mixed { return parent::pop($timeout); } public function getCapacity(): int { return $this->capacity; } public function getLength(): int { return $this->length(); } public function isAvailable(): bool { return ! $this->isClosing(); } public function close(): bool { $this->closed = true; return parent::close(); } public function hasProducers(): bool { throw new RuntimeException('Not supported.'); } public function hasConsumers(): bool { throw new RuntimeException('Not supported.'); } public function isReadable(): bool { throw new RuntimeException('Not supported.'); } public function isWritable(): bool { throw new RuntimeException('Not supported.'); } public function isClosing(): bool { return $this->closed || $this->errCode === SWOOLE_CHANNEL_CLOSED; } public function isTimeout(): bool { return ! $this->closed && $this->errCode === SWOOLE_CHANNEL_TIMEOUT; } } } else { class Channel extends \Swoole\Coroutine\Channel implements ChannelInterface { /** * @var bool */ protected $closed = false; public function getCapacity(): int { return $this->capacity; } public function getLength(): int { return $this->length(); } public function isAvailable(): bool { return ! $this->isClosing(); } public function close(): bool { $this->closed = true; return parent::close(); } public function hasProducers(): bool { throw new RuntimeException('Not supported.'); } public function hasConsumers(): bool { throw new RuntimeException('Not supported.'); } public function isReadable(): bool { throw new RuntimeException('Not supported.'); } public function isWritable(): bool { throw new RuntimeException('Not supported.'); } public function isClosing(): bool { return $this->closed || $this->errCode === SWOOLE_CHANNEL_CLOSED; } public function isTimeout(): bool { return ! $this->closed && $this->errCode === SWOOLE_CHANNEL_TIMEOUT; } } } src/ConfigProvider.php000064400000001110151360553570010767 0ustar00 [ SocketFactoryInterface::class => SocketFactory::class, ], ]; } } src/Constant.php000064400000001041151360553570007643 0ustar00 80000 && SWOOLE_VERSION_ID >= 50000) { interface ChannelInterface { /** * @param float|int $timeout [optional] = -1 */ public function push(mixed $data, float $timeout = -1): bool; /** * @param float $timeout seconds [optional] = -1 * @return mixed when pop failed, return false */ public function pop(float $timeout = -1): mixed; /** * Swow: When the channel is closed, all the data in it will be destroyed. * Swoole: When the channel is closed, the data in it can still be popped out, but push behavior will no longer succeed. */ public function close(): bool; public function getCapacity(): int; public function getLength(): int; public function isAvailable(): bool; public function hasProducers(): bool; public function hasConsumers(): bool; public function isEmpty(): bool; public function isFull(): bool; public function isReadable(): bool; public function isWritable(): bool; public function isClosing(): bool; public function isTimeout(): bool; } } else { interface ChannelInterface { /** * @param mixed $data [required] * @param float|int $timeout [optional] = -1 * @return bool */ public function push($data, $timeout = -1); /** * @param float $timeout seconds [optional] = -1 * @return mixed when pop failed, return false */ public function pop($timeout = -1); /** * Swow: When the channel is closed, all the data in it will be destroyed. * Swoole: When the channel is closed, the data in it can still be popped out, but push behavior will no longer succeed. * @return mixed */ public function close(): bool; /** * @return int */ public function getCapacity(); /** * @return int */ public function getLength(); /** * @return bool */ public function isAvailable(); /** * @return bool */ public function hasProducers(); /** * @return bool */ public function hasConsumers(); /** * @return bool */ public function isEmpty(); /** * @return bool */ public function isFull(); /** * @return bool */ public function isReadable(); /** * @return bool */ public function isWritable(); /** * @return bool */ public function isClosing(); /** * @return bool */ public function isTimeout(); } } src/Contract/CoroutineInterface.php000064400000003626151360553570013432 0ustar00 true, * 'package_max_length' => 1024 * 1024 * 2, * 'package_length_type' => 'N', * 'package_length_offset' => 0, * 'package_body_offset' => 4, * ] */ public function getProtocol(): array; } src/Contract/SocketInterface.php000064400000000546151360553570012711 0ustar00callable = $callable; } public static function create(callable $callable, ...$data) { $coroutine = new static($callable); $coroutine->execute(...$data); return $coroutine; } public function execute(...$data) { $this->id = SwooleCo::create($this->callable, ...$data); return $this; } public function getId() { if (is_null($this->id)) { throw new RuntimeException('Coroutine was not be executed.'); } return $this->id; } public static function id() { return SwooleCo::getCid(); } public static function pid(?int $id = null) { if ($id) { $cid = SwooleCo::getPcid($id); if ($cid === false) { throw new CoroutineDestroyedException(sprintf('Coroutine #%d has been destroyed.', $id)); } } else { $cid = SwooleCo::getPcid(); } if ($cid === false) { throw new RunningInNonCoroutineException('Non-Coroutine environment don\'t has parent coroutine id.'); } return max(0, $cid); } public static function set(array $config) { SwooleCo::set($config); } /** * @return null|ArrayObject */ public static function getContextFor(?int $id = null) { if ($id === null) { return SwooleCo::getContext(); } return SwooleCo::getContext($id); } /** * {@inheritdoc} */ public static function defer(callable $callable) { SwooleCo::defer($callable); } /** * {@inheritdoc} */ public static function stats(): array { return SwooleCo::stats(); } /** * {@inheritdoc} */ public static function exists(int $id): bool { return SwooleCo::exists($id); } } src/Exception/CoroutineDestroyedException.php000064400000000521151360553570015523 0ustar00setMethod($method); $this->setData($contents); $this->setHeaders($this->encodeHeaders($headers)); $this->execute($path); if ($this->errCode !== 0) { throw new HttpClientException($this->errMsg, $this->errCode); } return new RawResponse( $this->statusCode, $this->decodeHeaders($this->headers ?? []), $this->body, $version ); } /** * @param string[] $headers * @return string[][] */ private function decodeHeaders(array $headers): array { $result = []; foreach ($headers as $name => $header) { // The key of header is lower case. $result[$name][] = $header; } if ($this->set_cookie_headers) { $result['set-cookie'] = $this->set_cookie_headers; } return $result; } /** * Swoole engine not support two dimensional array. * @param string[][] $headers * @return string[] */ private function encodeHeaders(array $headers): array { $result = []; foreach ($headers as $name => $value) { $result[$name] = is_array($value) ? implode(',', $value) : $value; } return $result; } } src/Http/FdGetter.php000064400000000626151360553570010505 0ustar00fd; } } src/Http/RawResponse.php000064400000001544151360553570011251 0ustar00statusCode = $statusCode; $this->headers = $headers; $this->body = $body; $this->version = $version; } } src/Http/Server.php000064400000003305151360553570010244 0ustar00host = $name; $this->port = $port; $this->server = new HttpServer($name, $port, false, true); return $this; } public function handle(callable $callable): static { $this->handler = $callable; return $this; } public function start(): void { $this->server->handle('/', function ($request, $response) { Coroutine::create(function () use ($request, $response) { try { $handler = $this->handler; $handler(Request::loadFromSwooleRequest($request), $response); } catch (Throwable $exception) { $this->logger->critical((string) $exception); } }); }); $this->server->start(); } public function close(): bool { $this->server->shutdown(); return true; } } src/Http/ServerFactory.php000064400000001337151360553570011577 0ustar00logger); return $server->bind($name, $port); } } src/Http/Stream.php000075500000015202151360553570010233 0ustar00size = strlen($this->contents); $this->writable = true; } /** * Reads all data from the stream into a string, from the beginning to end. * This method MUST attempt to seek to the beginning of the stream before * reading data and read the stream until the end is reached. * Warning: This could attempt to load a large amount of data into memory. * This method MUST NOT raise an exception in order to conform with PHP's * string casting operations. * * @see http://php.net/manual/en/language.oop5.magic.php#object.tostring */ public function __toString(): string { try { return $this->getContents(); } catch (\Throwable) { return ''; } } /** * Closes the stream and any underlying resources. */ public function close() { $this->detach(); } /** * Separates any underlying resources from the stream. * After the stream has been detached, the stream is in an unusable state. * * @return null|resource Underlying PHP stream, if any */ public function detach() { $this->contents = ''; $this->size = 0; $this->writable = false; return null; } /** * Get the size of the stream if known. * * @return null|int returns the size in bytes if known, or null if unknown */ public function getSize() { if (! $this->size) { $this->size = strlen($this->getContents()); } return $this->size; } /** * Returns the current position of the file read/write pointer. * * @return int Position of the file pointer * @throws RuntimeException on error */ public function tell() { throw new RuntimeException('Cannot determine the position of a SwooleStream'); } /** * Returns true if the stream is at the end of the stream. * * @return bool */ public function eof() { return $this->getSize() === 0; } /** * Returns whether or not the stream is seekable. * * @return bool */ public function isSeekable() { return false; } /** * Seek to a position in the stream. * * @see http://www.php.net/manual/en/function.fseek.php * @param int $offset Stream offset * @param int $whence Specifies how the cursor position will be calculated * based on the seek offset. Valid values are identical to the built-in * PHP $whence values for `fseek()`. SEEK_SET: Set position equal to * offset bytes SEEK_CUR: Set position to current location plus offset * SEEK_END: Set position to end-of-stream plus offset. * @throws RuntimeException on failure */ public function seek($offset, $whence = SEEK_SET) { throw new RuntimeException('Cannot seek a SwooleStream'); } /** * Seek to the beginning of the stream. * If the stream is not seekable, this method will raise an exception; * otherwise, it will perform a seek(0). * * @see seek() * @see http://www.php.net/manual/en/function.fseek.php * @throws RuntimeException on failure */ public function rewind() { $this->seek(0); } /** * Returns whether or not the stream is writable. * * @return bool */ public function isWritable() { return $this->writable; } /** * Write data to the stream. * * @param string $string the string that is to be written * @return int returns the number of bytes written to the stream * @throws RuntimeException on failure */ public function write($string) { if (! $this->writable) { throw new RuntimeException('Cannot write to a non-writable stream'); } $size = strlen($string); $this->contents .= $string; $this->size += $size; return $size; } /** * Returns whether or not the stream is readable. * * @return bool */ public function isReadable() { return true; } /** * Read data from the stream. * * @param int $length Read up to $length bytes from the object and return * them. Fewer than $length bytes may be returned if underlying stream * call returns fewer bytes. * @return string returns the data read from the stream, or an empty string * if no bytes are available * @throws RuntimeException if an error occurs */ public function read($length) { if ($length >= $this->getSize()) { $result = $this->contents; $this->contents = ''; $this->size = 0; } else { $result = substr($this->contents, 0, $length); $this->contents = substr($this->contents, $length); $this->size = $this->getSize() - $length; } return $result; } /** * Returns the remaining contents in a string. * * @return string * @throws RuntimeException if unable to read or an error occurs while * reading */ public function getContents() { return $this->contents; } /** * Get stream metadata as an associative array or retrieve a specific key. * The keys returned are identical to the keys returned from PHP's * stream_get_meta_data() function. * * @see http://php.net/manual/en/function.stream-get-meta-data.php * @param string $key specific metadata to retrieve * @return null|array|mixed Returns an associative array if no key is * provided. Returns a specific key value if a key is provided and the * value is found, or null if the key is not found. */ public function getMetadata($key = null) { throw new BadMethodCallException('Not implemented'); } } src/Http/V2/Client.php000064400000004736151360553570010514 0ustar00client = new HTTP2Client($host, $port, $ssl); if ($settings) { $this->client->set($settings); } $this->client->connect(); } public function set(array $settings): bool { return $this->client->set($settings); } public function send(RequestInterface $request): int { $res = $this->client->send($this->transformRequest($request)); if ($res === false) { throw new HttpClientException($this->client->errMsg, $this->client->errCode); } return $res; } public function recv(float $timeout = 0): ResponseInterface { return $this->transformResponse($this->client->recv($timeout)); } public function ping(): bool { return $this->client->ping(); } public function close(): bool { return $this->client->close(); } public function isConnected(): bool { return $this->client->connected; } public function write(int $streamId, mixed $data, bool $end = false): bool { return $this->client->write($streamId, $data, $end); } private function transformResponse(SwResponse $request): ResponseInterface { return new Response( $request->streamId, $request->statusCode, $request->headers ?? [], $request->data ); } private function transformRequest(RequestInterface $request): SwRequest { $req = new SwRequest(); $req->method = $request->getMethod(); $req->path = $request->getPath(); $req->headers = $request->getHeaders(); $req->data = $request->getBody(); return $req; } } src/Http/V2/ClientFactory.php000064400000001013151360553570012025 0ustar00path; } public function setPath(string $path): void { $this->path = $path; } public function getMethod(): string { return $this->method; } public function setMethod(string $method): void { $this->method = $method; } public function getBody(): string { return $this->body; } public function setBody(string $body): void { $this->body = $body; } public function getHeaders(): array { return $this->headers; } public function setHeaders(array $headers): void { $this->headers = $headers; } } src/Http/V2/Response.php000064400000001532151360553570011063 0ustar00streamId; } public function getStatusCode(): int { return $this->statusCode; } public function getHeaders(): array { return $this->headers; } public function getBody(): ?string { return $this->body; } } src/ResponseEmitter.php000064400000005472151360553570011216 0ustar00header['Upgrade'] ?? '') === 'websocket') { return; } $this->buildSwooleResponse($connection, $response); $content = $response->getBody(); if ($content instanceof FileInterface) { $connection->sendfile($content->getFilename()); return; } if ($withContent) { $connection->end((string) $content); } else { $connection->end(); } } catch (Throwable $exception) { $this->logger?->critical((string) $exception); } } protected function buildSwooleResponse(Response $swooleResponse, ResponseInterface $response): void { // Headers foreach ($response->getHeaders() as $key => $value) { $swooleResponse->header($key, $value); } if ($response instanceof HyperfResponse) { // Cookies foreach ((array) $response->getCookies() as $domain => $paths) { foreach ($paths ?? [] as $path => $item) { foreach ($item ?? [] as $name => $cookie) { if ($cookie instanceof Cookie) { $value = $cookie->isRaw() ? $cookie->getValue() : rawurlencode($cookie->getValue()); $swooleResponse->rawcookie($cookie->getName(), $value, $cookie->getExpiresTime(), $cookie->getPath(), $cookie->getDomain(), $cookie->isSecure(), $cookie->isHttpOnly(), (string) $cookie->getSameSite()); } } } } // Trailers foreach ($response->getTrailers() ?? [] as $key => $value) { $swooleResponse->trailer($key, $value); } } // Status code $swooleResponse->status($response->getStatusCode(), $response->getReasonPhrase()); } } src/Signal.php000064400000001006151360553570007270 0ustar00getProtocol()) { $socket->setProtocol($protocol); } if ($option->getTimeout() === null) { $res = $socket->connect($option->getHost(), $option->getPort()); } else { $res = $socket->connect($option->getHost(), $option->getPort(), $option->getTimeout()); } if (! $res) { throw new SocketConnectException($socket->errMsg, $socket->errCode); } return $socket; } } src/Socket/SocketOption.php000064400000001542151360553570011731 0ustar00host; } public function getPort(): int { return $this->port; } public function getTimeout(): ?float { return $this->timeout; } public function getProtocol(): array { return $this->protocol; } } src/WaitGroup.php000064400000000500151360553570007772 0ustar00 */ protected $events = []; public function __construct(Response $connection, Request $request) { $this->connection = $connection; $this->connection->upgrade(); } public function on(string $event, callable $callback): void { $this->events[$event] = $callback; } public function start(): void { while (true) { $frame = $this->connection->recv(); if ($frame === false || $frame instanceof CloseFrame || $frame === '') { $callback = $this->events[static::ON_CLOSE]; $callback($this->connection, $this->connection->fd); break; } $callback = $this->events[static::ON_MESSAGE]; $callback($this->connection, $frame); } $this->connection = null; $this->events = []; } }