<?php
* This file is part of workerman.
*
* Licensed under The MIT License
* For full copyright and license information, please see the MIT-LICENSE.txt
* Redistributions of files must retain the above copyright notice.
*
* @author walkor<walkor@workerman.net>
* @copyright walkor<walkor@workerman.net>
* @link http://www.workerman.net/
* @license http://www.opensource.org/licenses/mit-license.php MIT License
*/
declare(strict_types=1);
namespace Workerman\Connection;
use JsonSerializable;
use RuntimeException;
use stdClass;
use Throwable;
use Workerman\Events\EventInterface;
use Workerman\Protocols\Http;
use Workerman\Protocols\Http\Request;
use Workerman\Timer;
use Workerman\Worker;
use function ceil;
use function count;
use function fclose;
use function feof;
use function fread;
use function function_exists;
use function fwrite;
use function is_object;
use function is_resource;
use function key;
use function method_exists;
use function posix_getpid;
use function restore_error_handler;
use function set_error_handler;
use function stream_set_blocking;
use function stream_set_read_buffer;
use function stream_socket_shutdown;
use function stream_socket_enable_crypto;
use function stream_socket_get_name;
use function strlen;
use function strrchr;
use function strrpos;
use function substr;
use function var_export;
use const PHP_INT_MAX;
use const STREAM_CRYPTO_METHOD_SSLv23_CLIENT;
use const STREAM_CRYPTO_METHOD_SSLv23_SERVER;
use const STREAM_CRYPTO_METHOD_SSLv2_CLIENT;
use const STREAM_CRYPTO_METHOD_SSLv2_SERVER;
use const STREAM_SHUT_WR;
* TcpConnection.
* @property string $websocketType
* @property string|null $websocketClientProtocol
* @property string|null $websocketOrigin
*/
class TcpConnection extends ConnectionInterface implements JsonSerializable
{
* Read buffer size.
*
* @var int
*/
public const READ_BUFFER_SIZE = 87380;
* Status initial.
*
* @var int
*/
public const STATUS_INITIAL = 0;
* Status connecting.
*
* @var int
*/
public const STATUS_CONNECTING = 1;
* Status connection established.
*
* @var int
*/
public const STATUS_ESTABLISHED = 2;
* Status ending (graceful close: write -> FIN -> linger/drain -> close).
*
* @var int
*/
public const STATUS_ENDING = 4;
* Status closing.
*
* @var int
*/
public const STATUS_CLOSING = 8;
* Status closed.
*
* @var int
*/
public const STATUS_CLOSED = 16;
* Maximum string length for cache
*
* @var int
*/
public const MAX_CACHE_STRING_LENGTH = 2048;
* Maximum cache size.
*
* @var int
*/
public const MAX_CACHE_SIZE = 512;
* Tcp keepalive interval.
*/
public const TCP_KEEPALIVE_INTERVAL = 55;
* Emitted when socket connection is successfully established.
*
* @var ?callable
*/
public $onConnect = null;
* Emitted before websocket handshake (Only called when protocol is ws).
*
* @var ?callable
*/
public $onWebSocketConnect = null;
* Emitted after websocket handshake (Only called when protocol is ws).
*
* @var ?callable
*/
public $onWebSocketConnected = null;
* Emitted when websocket connection is closed (Only called when protocol is ws).
*
* @var ?callable
*/
public $onWebSocketClose = null;
* Emitted when data is received.
*
* @var ?callable
*/
public $onMessage = null;
* Emitted when the other end of the socket sends a FIN packet.
*
* @var ?callable
*/
public $onClose = null;
* Emitted when an error occurs with connection.
*
* @var ?callable
*/
public $onError = null;
* Emitted when the send buffer becomes full.
*
* @var ?callable
*/
public $onBufferFull = null;
* Emitted when send buffer becomes empty.
*
* @var ?callable
*/
public $onBufferDrain = null;
* Transport (tcp/udp/unix/ssl).
*
* @var string
*/
public string $transport = 'tcp';
* Which worker belong to.
*
* @var ?Worker
*/
public ?Worker $worker = null;
* Bytes read.
*
* @var int
*/
public int $bytesRead = 0;
* Bytes written.
*
* @var int
*/
public int $bytesWritten = 0;
* Connection->id.
*
* @var int
*/
public int $id = 0;
* A copy of $worker->id which used to clean up the connection in worker->connections
*
* @var int
*/
protected int $realId = 0;
* Sets the maximum send buffer size for the current connection.
* OnBufferFull callback will be emitted When send buffer is full.
*
* @var int
*/
public int $maxSendBufferSize = 1048576;
* Context.
*
* @var ?stdClass
*/
public ?stdClass $context = null;
* Internal use only. Do not access or modify from application code.
*
* @internal Framework internal API
* @deprecated Do not set this property, use $response->header() or $response->widthHeaders() instead
* @var array
*/
public array $headers = [];
* Is safe.
*
* @var bool
*/
protected bool $isSafe = true;
* Default send buffer size.
*
* @var int
*/
public static int $defaultMaxSendBufferSize = 1048576;
* Sets the maximum acceptable packet size for the current connection.
*
* @var int
*/
public int $maxPackageSize = 1048576;
* Default maximum acceptable packet size.
*
* @var int
*/
public static int $defaultMaxPackageSize = 10485760;
* Default linger timeout for graceful end (seconds).
*
* @var float
*/
public static float $defaultLingerTimeout = 1.0;
* Linger timeout for graceful end (seconds).
*
* @var float
*/
public float $lingerTimeout = 1.0;
* Id recorder.
*
* @var int
*/
protected static int $idRecorder = 1;
* Socket
*
* @var resource
*/
protected $socket = null;
* Send buffer.
*
* @var string
*/
protected string $sendBuffer = '';
* Receive buffer.
*
* @var string
*/
protected string $recvBuffer = '';
* Current package length.
*
* @var int
*/
protected int $currentPackageLength = 0;
* Connection status.
*
* @var int
*/
protected int $status = self::STATUS_ESTABLISHED;
* Linger timer id for end().
*
* @var int
*/
protected int $endLingerTimerId = 0;
* Whether write side has been shutdown (FIN sent) during end().
*
* @var bool
*/
protected bool $endWriteShutdown = false;
* Remote address.
*
* @var string
*/
protected string $remoteAddress = '';
* Is paused.
*
* @var bool
*/
protected bool $isPaused = false;
* SSL handshake completed or not.
*
* @var bool
*/
protected bool|int $sslHandshakeCompleted = false;
* All connection instances.
*
* @var array
*/
public static array $connections = [];
* Status to string.
*
* @var array
*/
public const STATUS_TO_STRING = [
self::STATUS_INITIAL => 'INITIAL',
self::STATUS_CONNECTING => 'CONNECTING',
self::STATUS_ESTABLISHED => 'ESTABLISHED',
self::STATUS_CLOSING => 'CLOSING',
self::STATUS_ENDING => 'ENDING',
self::STATUS_CLOSED => 'CLOSED',
];
* Construct.
*
* @param EventInterface $eventLoop
* @param resource $socket
* @param string $remoteAddress
*/
public function __construct(EventInterface $eventLoop, $socket, string $remoteAddress = '')
{
++self::$statistics['connection_count'];
$this->id = $this->realId = self::$idRecorder++;
if (self::$idRecorder === PHP_INT_MAX) {
self::$idRecorder = 0;
}
$this->socket = $socket;
stream_set_blocking($this->socket, false);
stream_set_read_buffer($this->socket, 0);
$this->eventLoop = $eventLoop;
$this->eventLoop->onReadable($this->socket, $this->baseRead(...));
$this->maxSendBufferSize = self::$defaultMaxSendBufferSize;
$this->maxPackageSize = self::$defaultMaxPackageSize;
$this->lingerTimeout = self::$defaultLingerTimeout;
$this->remoteAddress = $remoteAddress;
static::$connections[$this->id] = $this;
$this->context = new stdClass();
}
* Get status.
*
* @param bool $rawOutput
*
* @return int|string
*/
public function getStatus(bool $rawOutput = true): int|string
{
if ($rawOutput) {
return $this->status;
}
return self::STATUS_TO_STRING[$this->status];
}
* Sends data on the connection.
*
* @param mixed $sendBuffer
* @param bool $raw
* @return bool|null
*/
public function send(mixed $sendBuffer, bool $raw = false): bool|null
{
if ($this->status === self::STATUS_ENDING || $this->status === self::STATUS_CLOSING || $this->status === self::STATUS_CLOSED) {
return false;
}
if (false === $raw && $this->protocol !== null) {
try {
$sendBuffer = $this->protocol::encode($sendBuffer, $this);
} catch(Throwable $e) {
$this->error($e);
}
if ($sendBuffer === '') {
return null;
}
}
if ($this->status !== self::STATUS_ESTABLISHED ||
($this->transport === 'ssl' && $this->sslHandshakeCompleted !== true)
) {
if ($this->sendBuffer && $this->bufferIsFull()) {
++self::$statistics['send_fail'];
return false;
}
$this->sendBuffer .= $sendBuffer;
$this->checkBufferWillFull();
return null;
}
if ($this->sendBuffer === '') {
$len = 0;
try {
$len = @fwrite($this->socket, $sendBuffer);
} catch (Throwable $e) {
Worker::log($e);
}
if ($len === strlen($sendBuffer)) {
$this->bytesWritten += $len;
return true;
}
if ($len > 0) {
$this->sendBuffer = substr($sendBuffer, $len);
$this->bytesWritten += $len;
} else {
if (!is_resource($this->socket) || feof($this->socket)) {
++self::$statistics['send_fail'];
if ($this->onError) {
try {
($this->onError)($this, static::SEND_FAIL, 'client closed');
} catch (Throwable $e) {
$this->error($e);
}
}
$this->destroy();
return false;
}
$this->sendBuffer = $sendBuffer;
}
$this->eventLoop->onWritable($this->socket, $this->baseWrite(...));
$this->checkBufferWillFull();
return null;
}
if ($this->bufferIsFull()) {
++self::$statistics['send_fail'];
return false;
}
$this->sendBuffer .= $sendBuffer;
$this->checkBufferWillFull();
return null;
}
* Get remote IP.
*
* @return string
*/
public function getRemoteIp(): string
{
$pos = strrpos($this->remoteAddress, ':');
if ($pos) {
return substr($this->remoteAddress, 0, $pos);
}
return '';
}
* Get remote port.
*
* @return int
*/
public function getRemotePort(): int
{
if ($this->remoteAddress) {
return (int)substr(strrchr($this->remoteAddress, ':'), 1);
}
return 0;
}
* Get remote address.
*
* @return string
*/
public function getRemoteAddress(): string
{
return $this->remoteAddress;
}
* Get local IP.
*
* @return string
*/
public function getLocalIp(): string
{
$address = $this->getLocalAddress();
$pos = strrpos($address, ':');
if (!$pos) {
return '';
}
return substr($address, 0, $pos);
}
* Get local port.
*
* @return int
*/
public function getLocalPort(): int
{
$address = $this->getLocalAddress();
$pos = strrpos($address, ':');
if (!$pos) {
return 0;
}
return (int)substr(strrchr($address, ':'), 1);
}
* Get local address.
*
* @return string
*/
public function getLocalAddress(): string
{
if (!is_resource($this->socket)) {
return '';
}
return (string)@stream_socket_get_name($this->socket, false);
}
* Get send buffer queue size.
*
* @return integer
*/
public function getSendBufferQueueSize(): int
{
return strlen($this->sendBuffer);
}
* Get receive buffer queue size.
*
* @return integer
*/
public function getRecvBufferQueueSize(): int
{
return strlen($this->recvBuffer);
}
* Pauses the reading of data. That is onMessage will not be emitted. Useful to throttle back an upload.
*
* @return void
*/
public function pauseRecv(): void
{
if($this->eventLoop !== null){
$this->eventLoop->offReadable($this->socket);
}
$this->isPaused = true;
}
* Resumes reading after a call to pauseRecv.
*
* @return void
*/
public function resumeRecv(): void
{
if ($this->isPaused === true) {
$this->eventLoop->onReadable($this->socket, $this->baseRead(...));
$this->isPaused = false;
$this->baseRead($this->socket, false);
}
}
* Base read handler.
*
* @param resource $socket
* @param bool $checkEof
* @return void
*/
public function baseRead($socket, bool $checkEof = true): void
{
static $requests = [];
if ($this->transport === 'ssl' && $this->sslHandshakeCompleted !== true) {
if ($this->doSslHandshake($socket)) {
$this->sslHandshakeCompleted = true;
if ($this->sendBuffer) {
$this->eventLoop->onWritable($socket, $this->baseWrite(...));
}
} else {
return;
}
}
$buffer = '';
try {
$buffer = @fread($socket, self::READ_BUFFER_SIZE);
} catch (Throwable) {
}
if ($buffer === '' || $buffer === false) {
if ($checkEof && (!is_resource($socket) || feof($socket) || $buffer === false)) {
$this->destroy();
return;
}
} else {
$this->bytesRead += strlen($buffer);
if ($this->status === self::STATUS_ENDING) {
return;
}
if ($this->recvBuffer === '') {
if (!isset($buffer[static::MAX_CACHE_STRING_LENGTH]) && isset($requests[$buffer])) {
++self::$statistics['total_request'];
if ($this->protocol === Http::class) {
$request = $requests[$buffer];
$request->connection = $this;
try {
($this->onMessage)($this, $request);
} catch (Throwable $e) {
$this->error($e);
}
$request = clone $request;
$request->destroy();
$requests[$buffer] = $request;
return;
}
$request = $requests[$buffer];
try {
($this->onMessage)($this, $request);
} catch (Throwable $e) {
$this->error($e);
}
return;
}
$this->recvBuffer = $buffer;
} else {
$this->recvBuffer .= $buffer;
}
}
if ($this->protocol !== null) {
while ($this->recvBuffer !== '' && !$this->isPaused) {
if ($this->currentPackageLength) {
$recvBufferLength = strlen($this->recvBuffer);
if ($this->currentPackageLength > $recvBufferLength) {
break;
}
} else {
try {
$this->currentPackageLength = $this->protocol::input($this->recvBuffer, $this);
} catch (Throwable $e) {
$this->currentPackageLength = -1;
Worker::safeEcho((string)$e);
}
if ($this->currentPackageLength === 0) {
break;
} elseif ($this->currentPackageLength > 0 && $this->currentPackageLength <= $this->maxPackageSize) {
$recvBufferLength = strlen($this->recvBuffer);
if ($this->currentPackageLength > $recvBufferLength) {
break;
}
}
else {
Worker::safeEcho((string)(new RuntimeException("Protocol $this->protocol Error package. package_length=" . var_export($this->currentPackageLength, true))));
$this->destroy();
return;
}
}
++self::$statistics['total_request'];
if ($recvBufferLength === $this->currentPackageLength) {
$oneRequestBuffer = $this->recvBuffer;
$this->recvBuffer = '';
} else {
$oneRequestBuffer = substr($this->recvBuffer, 0, $this->currentPackageLength);
$this->recvBuffer = substr($this->recvBuffer, $this->currentPackageLength);
}
$this->currentPackageLength = 0;
try {
if (!isset($oneRequestBuffer[static::MAX_CACHE_STRING_LENGTH]) && isset($requests[$oneRequestBuffer])) {
$request = $requests[$oneRequestBuffer];
if ($request instanceof Request) {
$request->connection = $this;
($this->onMessage)($this, $request);
$request = clone $request;
$request->destroy();
$requests[$oneRequestBuffer] = $request;
} else {
($this->onMessage)($this, $request);
}
continue;
}
$request = $this->protocol::decode($oneRequestBuffer, $this);
if ((!is_object($request) || $request instanceof Request) && !isset($oneRequestBuffer[static::MAX_CACHE_STRING_LENGTH])) {
($this->onMessage)($this, $request);
if ($request instanceof Request) {
$request = clone $request;
$request->destroy();
}
$requests[$oneRequestBuffer] = $request;
if (count($requests) > static::MAX_CACHE_SIZE) {
unset($requests[key($requests)]);
}
continue;
}
($this->onMessage)($this, $request);
} catch (Throwable $e) {
$this->error($e);
}
}
return;
}
if ($this->recvBuffer === '' || $this->isPaused) {
return;
}
++self::$statistics['total_request'];
try {
($this->onMessage)($this, $this->recvBuffer);
} catch (Throwable $e) {
$this->error($e);
}
$this->recvBuffer = '';
}
* Base write handler.
*
* @return void
*/
public function baseWrite(): void
{
$len = 0;
try {
if ($this->transport === 'ssl') {
$len = @fwrite($this->socket, $this->sendBuffer, 8192);
} else {
$len = @fwrite($this->socket, $this->sendBuffer);
}
} catch (Throwable) {
}
if ($len === strlen($this->sendBuffer)) {
$this->bytesWritten += $len;
$this->eventLoop->offWritable($this->socket);
$this->sendBuffer = '';
if ($this->onBufferDrain) {
try {
($this->onBufferDrain)($this);
} catch (Throwable $e) {
$this->error($e);
}
}
if ($this->status === self::STATUS_ENDING) {
$this->endMaybeShutdownWrite();
}
if ($this->status === self::STATUS_CLOSING) {
if (!empty($this->context->streamSending)) {
return;
}
$this->destroy();
}
return;
}
if ($len > 0) {
$this->bytesWritten += $len;
$this->sendBuffer = substr($this->sendBuffer, $len);
} else {
++self::$statistics['send_fail'];
$this->destroy();
}
}
* SSL handshake.
*
* @param resource $socket
* @return bool|int
*/
public function doSslHandshake($socket): bool|int
{
if (!is_resource($socket) || feof($socket)) {
$this->destroy();
return false;
}
$async = $this instanceof AsyncTcpConnection;
* We disabled ssl3 because https://blog.qualys.com/ssllabs/2014/10/15/ssl-3-is-dead-killed-by-the-poodle-attack.
* You can enable ssl3 by the codes below.
*/
$type = STREAM_CRYPTO_METHOD_SSLv2_CLIENT | STREAM_CRYPTO_METHOD_SSLv23_CLIENT | STREAM_CRYPTO_METHOD_SSLv3_CLIENT;
}else{
$type = STREAM_CRYPTO_METHOD_SSLv2_SERVER | STREAM_CRYPTO_METHOD_SSLv23_SERVER | STREAM_CRYPTO_METHOD_SSLv3_SERVER;
}*/
if ($async) {
$type = STREAM_CRYPTO_METHOD_SSLv2_CLIENT | STREAM_CRYPTO_METHOD_SSLv23_CLIENT;
} else {
$type = STREAM_CRYPTO_METHOD_SSLv2_SERVER | STREAM_CRYPTO_METHOD_SSLv23_SERVER;
}
set_error_handler(static function (int $code, string $msg): bool {
if (!Worker::$daemonize) {
Worker::safeEcho(sprintf("SSL handshake error: %s\n", $msg));
}
return true;
});
$ret = stream_socket_enable_crypto($socket, true, $type);
restore_error_handler();
if (false === $ret) {
$this->destroy();
return false;
}
if (0 === $ret) {
return 0;
}
return true;
}
* This method pulls all the data out of a readable stream, and writes it to the supplied destination.
*
* @param self $dest
* @return void
*/
public function pipe(self $dest): void
{
$this->onMessage = fn ($source, $data) => $dest->send($data);
$this->onClose = fn () => $dest->close();
$dest->onBufferFull = fn () => $this->pauseRecv();
$dest->onBufferDrain = fn() => $this->resumeRecv();
}
* Remove $length of data from receive buffer.
*
* @param int $length
* @return void
*/
public function consumeRecvBuffer(int $length): void
{
$this->recvBuffer = substr($this->recvBuffer, $length);
}
* Close connection.
*
* @param mixed $data
* @param bool $raw
* @return void
*/
public function close(mixed $data = null, bool $raw = false): void
{
if ($this->status === self::STATUS_INITIAL || $this->status === self::STATUS_CONNECTING) {
$this->destroy();
return;
}
if ($this->status === self::STATUS_CLOSING || $this->status === self::STATUS_CLOSED) {
return;
}
if ($data !== null) {
$this->send($data, $raw);
}
$this->status = self::STATUS_CLOSING;
if ($this->sendBuffer === '') {
$this->destroy();
} else {
$this->pauseRecv();
}
}
* Graceful end connection.
* It tries to: send response -> wait sendBuffer empty -> shutdown write(FIN) -> linger/drain reads -> close().
*
* @param mixed $data
* @param bool $raw
* @return void
*/
public function end(mixed $data = null, bool $raw = false): void
{
if ($this->status === self::STATUS_INITIAL || $this->status === self::STATUS_CONNECTING) {
$this->destroy();
return;
}
if ($this->status === self::STATUS_ENDING || $this->status === self::STATUS_CLOSING || $this->status === self::STATUS_CLOSED) {
return;
}
if ($data !== null) {
$this->send($data, $raw);
}
$this->status = self::STATUS_ENDING;
$this->onMessage = static function (self $connection, mixed $data = null): void {};
$this->recvBuffer = '';
$this->currentPackageLength = 0;
if ($this->sendBuffer === '') {
$this->endMaybeShutdownWrite();
return;
}
}
* If in ENDING and sendBuffer is empty, shutdown write side and start linger timer.
*
* @return void
*/
protected function endMaybeShutdownWrite(): void
{
if ($this->status !== self::STATUS_ENDING || $this->endWriteShutdown || $this->sendBuffer !== '') {
return;
}
if (is_resource($this->socket)) {
try {
@stream_socket_shutdown($this->socket, STREAM_SHUT_WR);
} catch (Throwable) {
}
}
$this->endWriteShutdown = true;
$timeout = $this->lingerTimeout;
if ($timeout <= 0) {
$this->close();
return;
}
$this->endLingerTimerId = Timer::delay($timeout, function (): void {
$this->endLingerTimerId = 0;
if ($this->status === self::STATUS_CLOSED) {
return;
}
$this->close();
});
}
* Is ipv4.
*
* return bool.
*/
public function isIpV4(): bool
{
if ($this->transport === 'unix') {
return false;
}
return !str_contains($this->getRemoteIp(), ':');
}
* Is ipv6.
*
* return bool.
*/
public function isIpV6(): bool
{
if ($this->transport === 'unix') {
return false;
}
return str_contains($this->getRemoteIp(), ':');
}
* Get the real socket.
*
* @return resource
*/
public function getSocket()
{
return $this->socket;
}
* Check whether send buffer will be full.
*
* @return void
*/
protected function checkBufferWillFull(): void
{
if ($this->onBufferFull && $this->maxSendBufferSize <= strlen($this->sendBuffer)) {
try {
($this->onBufferFull)($this);
} catch (Throwable $e) {
$this->error($e);
}
}
}
* Whether send buffer is full.
*
* @return bool
*/
protected function bufferIsFull(): bool
{
if ($this->maxSendBufferSize <= strlen($this->sendBuffer)) {
if ($this->onError) {
try {
($this->onError)($this, static::SEND_FAIL, 'send buffer full and drop package');
} catch (Throwable $e) {
$this->error($e);
}
}
return true;
}
return false;
}
* Whether send buffer is Empty.
*
* @return bool
*/
public function bufferIsEmpty(): bool
{
return empty($this->sendBuffer);
}
* Destroy connection.
*
* @return void
*/
public function destroy(): void
{
if ($this->status === self::STATUS_CLOSED) {
return;
}
if($this->eventLoop !== null){
$this->eventLoop->offReadable($this->socket);
$this->eventLoop->offWritable($this->socket);
if (DIRECTORY_SEPARATOR === '\\' && method_exists($this->eventLoop, 'offExcept')) {
$this->eventLoop->offExcept($this->socket);
}
}
try {
@fclose($this->socket);
} catch (Throwable) {
}
$this->status = self::STATUS_CLOSED;
if ($this->onClose) {
try {
($this->onClose)($this);
} catch (Throwable $e) {
$this->error($e);
}
}
if ($this->protocol && method_exists($this->protocol, 'onClose')) {
try {
$this->protocol::onClose($this);
} catch (Throwable $e) {
$this->error($e);
}
}
$this->sendBuffer = $this->recvBuffer = '';
$this->currentPackageLength = 0;
$this->isPaused = $this->sslHandshakeCompleted = false;
$this->endWriteShutdown = false;
if ($this->status === self::STATUS_CLOSED) {
$this->onMessage = $this->onClose = $this->onError = $this->onBufferFull = $this->onBufferDrain = $this->eventLoop = $this->errorHandler = null;
if ($this->worker) {
unset($this->worker->connections[$this->realId]);
}
$this->worker = null;
unset(static::$connections[$this->realId]);
}
}
* Get the json_encode information.
*
* @return array
*/
public function jsonSerialize(): array
{
return [
'id' => $this->id,
'status' => $this->getStatus(),
'transport' => $this->transport,
'getRemoteIp' => $this->getRemoteIp(),
'remotePort' => $this->getRemotePort(),
'getRemoteAddress' => $this->getRemoteAddress(),
'getLocalIp' => $this->getLocalIp(),
'getLocalPort' => $this->getLocalPort(),
'getLocalAddress' => $this->getLocalAddress(),
'isIpV4' => $this->isIpV4(),
'isIpV6' => $this->isIpV6(),
];
}
* __unserialize.
*
* @param array $data
* @return void
*/
public function __unserialize(array $data): void
{
$this->isSafe = false;
}
* Destruct.
*
* @return void
*/
public function __destruct()
{
static $mod;
if (!$this->isSafe) {
return;
}
self::$statistics['connection_count']--;
if (Worker::getGracefulStop()) {
$mod ??= ceil((self::$statistics['connection_count'] + 1) / 3);
if (0 === self::$statistics['connection_count'] % $mod) {
$pid = function_exists('posix_getpid') ? posix_getpid() : 0;
Worker::log('worker[' . $pid . '] remains ' . self::$statistics['connection_count'] . ' connection(s)');
}
if (0 === self::$statistics['connection_count']) {
Worker::stopAll();
}
}
}
}