.gitattributes000064400000000105151360554260007442 0ustar00/.github export-ignore /examples export-ignore /tests export-ignore .gitignore000064400000000054151360554260006542 0ustar00/vendor/ composer.lock *.cache *.log .idea/*.php-cs-fixer.php000064400000005300151360554260007645 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, 'single_line_empty_body' => false, ]) ->setFinder( PhpCsFixer\Finder::create() ->exclude('vendor') ->in(__DIR__) ) ->setUsingCache(false); .phpstorm.meta.php000064400000000164151360554260010144 0ustar00" version="1.0" license="MIT" app.name="Hyperf" ARG timezone ARG PHP_VERSION ENV TIMEZONE=${timezone:-"Asia/Shanghai"} ENV COMPOSER_ROOT_VERSION="v2.0.0" # update RUN set -ex \ # show php version and extensions && php -v \ && php -m \ && php --ri swoole \ # ---------- some config ---------- && cd /etc/php* \ # - 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 README.md000064400000000234151360554260006031 0ustar00# Swoole Engine ![Swoole Engine Test](https://github.com/hyperf/engine/workflows/Swoole%20Engine%20Test/badge.svg) ``` composer require hyperf/engine ``` composer.json000064400000003163151360554260007300 0ustar00{ "name": "hyperf/engine", "type": "library", "license": "MIT", "keywords": [ "php", "hyperf", "engine", "swoole" ], "description": "Coroutine engine provided by swoole.", "autoload": { "psr-4": { "Hyperf\\Engine\\": "src/" }, "files": [ "src/Functions.php" ] }, "autoload-dev": { "psr-4": { "HyperfTest\\": "tests" } }, "require": { "php": ">=8.0", "hyperf/engine-contract": "~1.9.0" }, "require-dev": { "friendsofphp/php-cs-fixer": "^3.0", "hyperf/guzzle": "^3.0", "hyperf/http-message": " 3.0.*", "mockery/mockery": "^1.5", "phpstan/phpstan": "^1.0", "phpunit/phpunit": "^9.4", "swoole/ide-helper": "dev-master" }, "suggest": { "ext-sockets": "*", "ext-swoole": ">=5.0", "psr/http-message": "Required to use WebSocket Frame.", "hyperf/http-message": "Required to use ResponseEmitter." }, "conflict": { "ext-swoole": "<5.0" }, "minimum-stability": "dev", "prefer-stable": true, "config": { "optimize-autoloader": true, "sort-packages": true }, "extra": { "branch-alias": { "dev-master": "2.10-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.xml000064400000000755151360554260006773 0ustar00 ./tests/ src/Channel.php000064400000003322151360554260007423 0ustar00capacity; } 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.php000064400000001111151360554260010765 0ustar00 [ SocketFactoryInterface::class => SocketFactory::class, ], ]; } } src/Constant.php000064400000001042151360554260007641 0ustar00callable = $callable; } public static function create(callable $callable, ...$data): static { $coroutine = new static($callable); $coroutine->execute(...$data); return $coroutine; } public function execute(...$data): static { $this->id = SwooleCo::create($this->callable, ...$data); return $this; } public function getId(): int { if (is_null($this->id)) { throw new RuntimeException('Coroutine was not be executed.'); } return $this->id; } public static function id(): int { return SwooleCo::getCid(); } public static function pid(?int $id = null): int { 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): void { SwooleCo::set($config); } public static function getContextFor(?int $id = null): ?ArrayObject { if ($id === null) { return SwooleCo::getContext(); } return SwooleCo::getContext($id); } public static function defer(callable $callable): void { SwooleCo::defer($callable); } /** * Yield the current coroutine. * @param mixed $data only Support Swow * @return bool */ public static function yield(mixed $data = null): mixed { return SwooleCo::yield(); } /** * Resume the coroutine by coroutine Id. * @param mixed $data only Support Swow * @return bool */ public static function resumeById(int $id, mixed ...$data): mixed { return SwooleCo::resume($id); } /** * Get the coroutine stats. */ public static function stats(): array { return SwooleCo::stats(); } public static function exists(int $id = null): bool { return SwooleCo::exists($id); } } src/DefaultOption.php000064400000000717151360554260010635 0ustar00getFin()) { $flags |= SWOOLE_WEBSOCKET_FLAG_FIN; } if ($frame->getRSV1()) { $flags |= SWOOLE_WEBSOCKET_FLAG_RSV1; } if ($frame->getRSV2()) { $flags |= SWOOLE_WEBSOCKET_FLAG_RSV2; } if ($frame->getRSV3()) { $flags |= SWOOLE_WEBSOCKET_FLAG_RSV3; } if ($frame->getMask()) { $flags |= SWOOLE_WEBSOCKET_FLAG_MASK; } return $flags; } src/Http/Client.php000064400000004050151360554260010207 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/EventStream.php000064400000002153151360554260011230 0ustar00connection->getSocket(); $socket->header('Content-Type', 'text/event-stream; charset=utf-8'); $socket->header('Transfer-Encoding', 'chunked'); $socket->header('Cache-Control', 'no-cache'); foreach ($response?->getHeaders() ?? [] as $name => $values) { $socket->header($name, implode(', ', $values)); } } public function write(string $data): self { $this->connection->write($data); return $this; } public function end(): void { $this->connection->end(); } } src/Http/FdGetter.php000064400000000627151360554260010503 0ustar00fd; } } src/Http/RawResponse.php000064400000001613151360554260011243 0ustar00statusCode; } public function getHeaders(): array { return $this->headers; } public function getBody(): string { return $this->body; } public function getVersion(): string { return $this->version; } } src/Http/Server.php000064400000003313151360554260010240 0ustar00host = $name; $this->port = $port; $this->server = new HttpServer($name, $port, reuse_port: 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.php000064400000001340151360554260011566 0ustar00logger); return $server->bind($name, $port); } } src/Http/Stream.php000075500000015131151360554260010231 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(): void { $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(): ?int { 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(): int { throw new RuntimeException('Cannot determine the position of a SwooleStream'); } /** * Returns true if the stream is at the end of the stream. */ public function eof(): bool { return $this->getSize() === 0; } /** * Returns whether or not the stream is seekable. */ public function isSeekable(): bool { 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): void { 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). * * @throws RuntimeException on failure * @see http://www.php.net/manual/en/function.fseek.php * @see seek() */ public function rewind(): void { $this->seek(0); } /** * Returns whether or not the stream is writable. */ public function isWritable(): bool { 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): int { 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. */ public function isReadable(): bool { 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): string { 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. * * @throws RuntimeException if unable to read or an error occurs while * reading */ public function getContents(): string { 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.php000064400000005361151360554260010504 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 { $response = $this->client->recv($timeout); if ($response === false) { throw new HttpClientException($this->client->errMsg, $this->client->errCode); } return $this->transformResponse($response); } public function write(int $streamId, mixed $data, bool $end = false): bool { return $this->client->write($streamId, $data, $end); } public function ping(): bool { return $this->client->ping(); } public function close(): bool { return $this->client->close(); } public function isConnected(): bool { return $this->client->connected; } 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(); $req->pipeline = $request->isPipeline(); $req->usePipelineRead = $request->isPipeline(); return $req; } } src/Http/V2/ClientFactory.php000064400000001014151360554260012023 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; } public function isPipeline(): bool { return $this->pipeline; } public function setPipeline(bool $pipeline): void { $this->pipeline = $pipeline; } } src/Http/V2/Response.php000064400000001533151360554260011061 0ustar00streamId; } public function getStatusCode(): int { return $this->statusCode; } public function getHeaders(): array { return $this->headers; } public function getBody(): ?string { return $this->body; } } src/Http/WritableConnection.php000064400000001407151360554260012565 0ustar00response->write($data); } /** * @return Response */ public function getSocket(): mixed { return $this->response; } public function end(): bool { return $this->response->end(); } } src/ResponseEmitter.php000064400000006075151360554260011213 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()); } protected function isMethodsExists(object $object, array $methods): bool { foreach ($methods as $method) { if (! method_exists($object, $method)) { return false; } } return true; } } src/SafeSocket.php000064400000007060151360554260010105 0ustar00channel = new Channel($capacity); } /** * @throws SocketTimeoutException when send data timeout * @throws SocketClosedException when the client is closed */ public function sendAll(string $data, float $timeout = 0): int|false { $this->loop(); $res = $this->channel->push([$data, $timeout], $timeout); if ($res === false) { if ($this->channel->isClosing()) { $this->throw && throw new SocketClosedException('The channel is closed.'); } if ($this->channel->isTimeout()) { $this->throw && throw new SocketTimeoutException('The channel is full.'); } return false; } return strlen($data); } /** * @throws SocketTimeoutException when send data timeout * @throws SocketClosedException when the client is closed */ public function recvAll(int $length = 65536, float $timeout = 0): string|false { $res = $this->socket->recvAll($length, $timeout); if (! $res) { if ($this->socket->errCode === SOCKET_ETIMEDOUT) { $this->throw && throw new SocketTimeoutException('Recv timeout'); } $this->throw && throw new SocketClosedException('The socket is closed.'); } return $res; } /** * @throws SocketTimeoutException when send data timeout * @throws SocketClosedException when the client is closed */ public function recvPacket(float $timeout = 0): string|false { $res = $this->socket->recvPacket($timeout); if (! $res) { if ($this->socket->errCode === SOCKET_ETIMEDOUT) { $this->throw && throw new SocketTimeoutException('Recv timeout'); } $this->throw && throw new SocketClosedException('The socket is closed.'); } return $res; } public function close(): bool { $this->channel->close(); return $this->socket->close(); } public function setLogger(?LoggerInterface $logger): static { $this->logger = $logger; return $this; } protected function loop(): void { if ($this->loop) { return; } $this->loop = true; Coroutine::create(function () { try { while (true) { $data = $this->channel->pop(-1); if ($this->channel->isClosing()) { return; } [$data, $timeout] = $data; $this->socket->sendAll($data, $timeout); } } catch (Throwable $exception) { $this->logger?->critical((string) $exception); } }); } } src/Signal.php000064400000001007151360554260007266 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.php000064400000001543151360554260011727 0ustar00host; } public function getPort(): int { return $this->port; } public function getTimeout(): ?float { return $this->timeout; } public function getProtocol(): array { return $this->protocol; } } src/WebSocket/Frame.php000064400000011613151360554260010775 0ustar00setPayloadData($payloadData); } } public function __toString() { return $this->toString(); } public function getOpcode(): int { return $this->opcode; } public function setOpcode(int $opcode): static { $this->opcode = $opcode; return $this; } public function withOpcode(int $opcode): static { return (clone $this)->setOpcode($opcode); } public function getFin(): bool { return $this->fin; } public function setFin(bool $fin): static { $this->fin = $fin; return $this; } public function withFin(bool $fin): static { return (clone $this)->setFin($fin); } public function getRSV1(): bool { return $this->rsv1; } public function setRSV1(bool $rsv1): static { $this->rsv1 = $rsv1; return $this; } public function withRSV1(bool $rsv1): static { return (clone $this)->setRSV1($rsv1); } public function getRSV2(): bool { return $this->rsv2; } public function setRSV2(bool $rsv2): static { $this->rsv2 = $rsv2; return $this; } public function withRSV2(bool $rsv2): static { return (clone $this)->setRSV2($rsv2); } public function getRSV3(): bool { return $this->rsv3; } public function setRSV3(bool $rsv3): static { $this->rsv3 = $rsv3; return $this; } public function withRSV3(bool $rsv3): static { return (clone $this)->setRSV3($rsv3); } public function getPayloadLength(): int { return $this->payloadData?->getSize() ?? 0; } public function setPayloadLength(int $payloadLength): static { $this->payloadLength = $payloadLength; return $this; } public function withPayloadLength(int $payloadLength): static { return (clone $this)->setPayloadLength($payloadLength); } public function getMask(): bool { return ! empty($this->maskingKey); } public function getMaskingKey(): string { return $this->maskingKey; } public function setMaskingKey(string $maskingKey): static { $this->maskingKey = $maskingKey; return $this; } public function withMaskingKey(string $maskingKey): static { return (clone $this)->setMaskingKey($maskingKey); } public function getPayloadData(): StreamInterface { return $this->payloadData; } public function setPayloadData(mixed $payloadData): static { $this->payloadData = new Stream((string) $payloadData); return $this; } public function withPayloadData(mixed $payloadData): static { return (clone $this)->setPayloadData($payloadData); } public function toString(bool $withoutPayloadData = false): string { return SwooleFrame::pack( (string) $this->getPayloadData(), $this->getOpcode(), swoole_get_flags_from_frame($this) ); } public static function from(mixed $frame): static { if (! $frame instanceof SwooleFrame) { throw new InvalidArgumentException('The frame is invalid.'); } return new Frame( (bool) ($frame->flags & SWOOLE_WEBSOCKET_FLAG_FIN), (bool) ($frame->flags & SWOOLE_WEBSOCKET_FLAG_RSV1), (bool) ($frame->flags & SWOOLE_WEBSOCKET_FLAG_RSV2), (bool) ($frame->flags & SWOOLE_WEBSOCKET_FLAG_RSV3), $frame->opcode, strlen($frame->data), $frame->flags & SWOOLE_WEBSOCKET_FLAG_MASK ? '258E' : '', $frame->data ); } } src/WebSocket/Opcode.php000064400000000727151360554260011160 0ustar00getPayloadData(); $flags = swoole_get_flags_from_frame($frame); if ($this->connection instanceof SwooleResponse) { $this->connection->push($data, $frame->getOpcode(), $flags); return true; } if ($this->connection instanceof Server) { $this->connection->push($this->fd, $data, $frame->getOpcode(), $flags); return true; } throw new InvalidArgumentException('The websocket connection is invalid.'); } public function init(mixed $frame): static { switch (true) { case is_int($frame): $this->fd = $frame; break; case $frame instanceof Request || $frame instanceof SwooleFrame: $this->fd = $frame->fd; break; } return $this; } public function getFd(): int { return $this->fd; } public function close(): bool { if ($this->connection instanceof SwooleResponse) { return $this->connection->close(); } if ($this->connection instanceof Server) { return $this->connection->disconnect($this->fd); } return false; } } src/WebSocket/WebSocket.php000064400000003523151360554260011632 0ustar00 */ protected array $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) { /** @var false|string|SwFrame $frame */ $frame = $this->connection->recv(); if ($frame === false || $frame instanceof CloseFrame || $frame === '') { if ($callback = $this->events[static::ON_CLOSE] ?? null) { $callback($this->connection, $this->connection->fd); } break; } switch ($frame->opcode) { case Opcode::PING: $this->connection->push('', Opcode::PONG); break; case Opcode::PONG: break; default: if ($callback = $this->events[static::ON_MESSAGE] ?? null) { $callback($this->connection, $frame); } } } $this->connection = null; $this->events = []; } }