From 34354ae88f97a7f62ad2cfb31fda0932666da5e3 Mon Sep 17 00:00:00 2001 From: Mathijs Smit Date: Fri, 5 Jun 2026 11:57:52 +0200 Subject: [PATCH] Add opt-in ReconnectingTransport for recreated Unix sockets. Invalidate stale fds on connection failure, make UnixSocketTransport reconnectable, and wrap it for automatic recovery in long-lived loops. --- README.md | 29 ++- src/Transport/ReconnectingTransport.php | 98 ++++++++ src/Transport/SocketTransport.php | 28 ++- src/Transport/UnixSocketTransport.php | 27 ++- tests/Integration/FileSocketViciServer.php | 228 +++++++++++++++++ tests/Integration/MockViciServer.php | 28 ++- .../ReconnectingTransportIntegrationTest.php | 80 ++++++ .../Transport/ReconnectingTransportTest.php | 229 ++++++++++++++++++ tests/Unit/Transport/StreamTransportTest.php | 17 ++ .../Transport/UnixSocketTransportTest.php | 150 ++++++++++++ 10 files changed, 901 insertions(+), 13 deletions(-) create mode 100644 src/Transport/ReconnectingTransport.php create mode 100644 tests/Integration/FileSocketViciServer.php create mode 100644 tests/Integration/ReconnectingTransportIntegrationTest.php create mode 100644 tests/Unit/Transport/ReconnectingTransportTest.php create mode 100644 tests/Unit/Transport/UnixSocketTransportTest.php diff --git a/README.md b/README.md index b4a698c..4870c3c 100644 --- a/README.md +++ b/README.md @@ -5,7 +5,7 @@ A pure-PHP client implementation of strongSwan's [VICI protocol](https://github.com/strongswan/strongswan/blob/master/src/libcharon/plugins/vici/README.md). Use it from PHP to monitor, configure, and control the IKE daemon `charon`. - Covers every command and event documented in the VICI README. -- Pluggable transport: Unix domain socket (default) or TCP, plus a generic `StreamTransport` for injection and testing. +- Pluggable transport: Unix domain socket (default) or TCP, plus a generic `StreamTransport` for injection and testing. An opt-in `ReconnectingTransport` wrapper recovers from charon restarts on Unix sockets. - Blocking `Session` for commands, plus an `EventListener` for long-running event subscriptions. - Streaming list commands (`list-sas`, `list-conns`, ...) expose event streams as PHP generators. - Fully typed, PHPStan level 8 clean, zero runtime dependencies. @@ -68,6 +68,33 @@ $stream = stream_socket_client('unix:///run/strongswan/charon.vici'); $session = new Session(new StreamTransport($stream, readTimeout: 10.0)); ``` +### Long-lived connections (Unix socket reconnect) + +For daemons or `while (true)` loops where charon may restart and recreate the Unix socket file, use `ReconnectingTransport`. It reconnects automatically when a single `send()`, `receive()`, or `hasData()` call fails with `ConnectionException`: + +```php +use Bk203\Vici\Session; +use Bk203\Vici\Transport\ReconnectingTransport; + +$session = new Session(new ReconnectingTransport( + path: '/var/run/charon.vici', + readTimeout: 30.0, +)); + +while (true) { + sleep(60); + $info = $session->version(); +} +``` + +`ReconnectingTransport` is opt-in; `new Session()` alone still uses a plain `UnixSocketTransport` with no automatic recovery. + +**v1 limitations** + +- Retry covers one transport I/O call. Multi-packet commands (`streamedRequest()`, `EventListener::listen()`) can still fail mid-operation; catch `ConnectionException` and restart the command or listener loop. +- Daemon-side `EVENT_REGISTER` state is not replayed after reconnect. Re-register events or wait for a future Session-level restore helper. +- `TimeoutException` is not retried (slow charon is not treated as a dead socket). + ## Common workflows ### Load a connection diff --git a/src/Transport/ReconnectingTransport.php b/src/Transport/ReconnectingTransport.php new file mode 100644 index 0000000..72e562c --- /dev/null +++ b/src/Transport/ReconnectingTransport.php @@ -0,0 +1,98 @@ +inner = new UnixSocketTransport( + $this->path, + $this->connectTimeout, + $this->readTimeout, + ); + } + + public function send(string $bytes): void + { + $this->withReconnect(function (UnixSocketTransport $inner) use ($bytes): void { + $inner->send($bytes); + }); + } + + public function receive(?float $timeout = null): string + { + return $this->withReconnect(static fn (UnixSocketTransport $inner): string => $inner->receive($timeout)); + } + + public function hasData(float $timeout = 0.0): bool + { + return $this->withReconnect(static fn (UnixSocketTransport $inner): bool => $inner->hasData($timeout)); + } + + public function isConnected(): bool + { + return $this->inner->isConnected(); + } + + public function close(): void + { + $this->inner->close(); + } + + public function getStream() + { + return $this->inner->getStream(); + } + + /** + * @template T + * + * @param callable(UnixSocketTransport): T $operation + * + * @return T + */ + private function withReconnect(callable $operation): mixed + { + try { + return $operation($this->inner); + } catch (ConnectionException $first) { + $last = $first; + + for ($attempt = 0; $attempt < $this->maxReconnectAttempts; $attempt++) { + usleep($this->reconnectDelayMs * 1000); + + try { + $this->inner->reconnect(); + } catch (ConnectionException $e) { + $last = $e; + continue; + } + + try { + return $operation($this->inner); + } catch (ConnectionException $e) { + $last = $e; + } + } + + throw $last; + } + } +} diff --git a/src/Transport/SocketTransport.php b/src/Transport/SocketTransport.php index 4e0cd8f..dc81e11 100644 --- a/src/Transport/SocketTransport.php +++ b/src/Transport/SocketTransport.php @@ -72,7 +72,7 @@ final public function hasData(float $timeout = 0.0): bool $ready = @stream_select($read, $write, $except, $sec, $usec); if ($ready === false) { - throw new ConnectionException('stream_select() failed on VICI transport.'); + $this->connectionFailed('stream_select() failed on VICI transport.'); } return $ready > 0; @@ -96,6 +96,20 @@ final public function getStream() return \is_resource($this->stream) ? $this->stream : null; } + protected function invalidate(): void + { + $this->close(); + } + + /** + * @return never + */ + protected function connectionFailed(string $message): void + { + $this->invalidate(); + throw new ConnectionException($message); + } + private function writeAll(string $data): void { $stream = $this->requireStream(); @@ -110,9 +124,9 @@ private function writeAll(string $data): void throw new TimeoutException('Timed out writing to VICI socket.'); } if (feof($stream)) { - throw new ConnectionException('VICI socket closed during write.'); + $this->connectionFailed('VICI socket closed during write.'); } - throw new ConnectionException('Failed to write to VICI socket.'); + $this->connectionFailed('Failed to write to VICI socket.'); } $written += $chunk; } @@ -137,7 +151,7 @@ private function readAll(int $length, ?float $timeout): string $usec = (int) round(($remaining - $sec) * 1_000_000); $ready = @stream_select($read, $write, $except, $sec, $usec); if ($ready === false) { - throw new ConnectionException('stream_select() failed on VICI transport.'); + $this->connectionFailed('stream_select() failed on VICI transport.'); } if ($ready === 0) { throw new TimeoutException('Timed out reading from VICI socket.'); @@ -148,11 +162,11 @@ private function readAll(int $length, ?float $timeout): string \assert($need > 0); $chunk = @fread($stream, $need); if ($chunk === false) { - throw new ConnectionException('Failed to read from VICI socket.'); + $this->connectionFailed('Failed to read from VICI socket.'); } if ($chunk === '') { if (feof($stream)) { - throw new ConnectionException('VICI socket closed during read.'); + $this->connectionFailed('VICI socket closed during read.'); } $meta = stream_get_meta_data($stream); if ($meta['timed_out']) { @@ -173,7 +187,7 @@ private function readAll(int $length, ?float $timeout): string private function requireStream() { if (!\is_resource($this->stream)) { - throw new ConnectionException('VICI transport is not connected.'); + $this->connectionFailed('VICI transport is not connected.'); } return $this->stream; } diff --git a/src/Transport/UnixSocketTransport.php b/src/Transport/UnixSocketTransport.php index 04d55ae..70cfec0 100644 --- a/src/Transport/UnixSocketTransport.php +++ b/src/Transport/UnixSocketTransport.php @@ -22,15 +22,27 @@ public function __construct( $this->connect(); } - private function connect(): void + public function reconnect(): void { + $this->close(); + $this->connect(); + } + + protected function connect(): void + { + $deadline = microtime(true) + $this->connectTimeout; + while (!file_exists($this->path) && microtime(true) < $deadline) { + usleep(100_000); + } + + $remaining = max(0.0, $deadline - microtime(true)); $errno = 0; $errstr = ''; $stream = @stream_socket_client( 'unix://' . $this->path, $errno, $errstr, - $this->connectTimeout, + $remaining, \STREAM_CLIENT_CONNECT, ); @@ -43,13 +55,20 @@ private function connect(): void )); } + $this->applyStreamOptions($stream); + $this->stream = $stream; + } + + /** + * @param resource $stream + */ + private function applyStreamOptions($stream): void + { stream_set_blocking($stream, true); if ($this->readTimeout !== null) { $sec = (int) floor($this->readTimeout); $usec = (int) round(($this->readTimeout - $sec) * 1_000_000); stream_set_timeout($stream, $sec, $usec); } - - $this->stream = $stream; } } diff --git a/tests/Integration/FileSocketViciServer.php b/tests/Integration/FileSocketViciServer.php new file mode 100644 index 0000000..ab31e9b --- /dev/null +++ b/tests/Integration/FileSocketViciServer.php @@ -0,0 +1,228 @@ +path = $path ?? sys_get_temp_dir() . '/vici-file-test-' . uniqid('', true) . '.sock'; + $this->codec = new PacketCodec(); + $this->encoder = new MessageEncoder(); + $this->decoder = new MessageDecoder(); + $this->startListening(); + } + + public function getPath(): string + { + return $this->path; + } + + public function acceptClient(float $timeout = 2.0): void + { + if (!\is_resource($this->listenSocket)) { + Assert::fail('FileSocketViciServer is not listening.'); + } + + $listenSocket = $this->listenSocket; + $read = [$listenSocket]; + $write = null; + $except = null; + $ready = stream_select($read, $write, $except, (int) floor($timeout), 0); + if ($ready === false || $ready === 0) { + Assert::fail('Timed out waiting for VICI client connection.'); + } + + $client = stream_socket_accept($listenSocket, $timeout); + if ($client === false) { + Assert::fail('Failed to accept VICI client connection.'); + } + + $this->serverStream = $client; + } + + public function simulateRestart(): void + { + if (\is_resource($this->serverStream)) { + @fclose($this->serverStream); + } + $this->serverStream = null; + + if (\is_resource($this->listenSocket)) { + @fclose($this->listenSocket); + } + $this->listenSocket = null; + + if (file_exists($this->path)) { + @unlink($this->path); + } + + $this->startListening(); + } + + public function close(): void + { + if (\is_resource($this->serverStream)) { + @fclose($this->serverStream); + } + if (\is_resource($this->listenSocket)) { + @fclose($this->listenSocket); + } + if (file_exists($this->path)) { + @unlink($this->path); + } + } + + /** + * @return array + */ + public function expectCommand(string $command, float $timeout = 1.0): array + { + $packet = $this->readPacket($timeout); + Assert::assertSame(PacketType::CMD_REQUEST, $packet->type, \sprintf( + 'Expected CMD_REQUEST "%s"; got %s.', + $command, + $packet->type->name, + )); + Assert::assertSame($command, $packet->name); + + return $packet->payload === '' ? [] : $this->decoder->decode($packet->payload); + } + + /** + * @param array $message + */ + public function sendCmdResponse(array $message = []): void + { + $this->writePacket(Packet::cmdResponse($this->encoder->encode($message))); + } + + private function readPacket(float $timeout = 1.0): Packet + { + if (!\is_resource($this->serverStream)) { + Assert::fail('FileSocketViciServer has no accepted client.'); + } + + $serverStream = $this->serverStream; + $deadline = microtime(true) + $timeout; + $read = [$serverStream]; + $write = null; + $except = null; + $ready = stream_select($read, $write, $except, (int) floor($timeout), 0); + if ($ready === false || $ready === 0) { + Assert::fail('Timed out waiting for VICI packet from client.'); + } + + $header = $this->readExactly(4, $deadline); + /** @var array{1: int} $unpacked */ + $unpacked = unpack('N', $header); + $length = $unpacked[1]; + $body = $length === 0 ? '' : $this->readExactly($length, $deadline); + + return $this->codec->decode($body); + } + + private function writePacket(Packet $packet): void + { + if (!\is_resource($this->serverStream)) { + throw new \RuntimeException('FileSocketViciServer has no accepted client.'); + } + + $serverStream = $this->serverStream; + $bytes = $this->codec->encode($packet); + $frame = pack('N', \strlen($bytes)) . $bytes; + $written = 0; + $total = \strlen($frame); + while ($written < $total) { + $chunk = fwrite($serverStream, substr($frame, $written)); + if ($chunk === false || $chunk === 0) { + throw new \RuntimeException('FileSocketViciServer failed to write packet.'); + } + $written += $chunk; + } + fflush($serverStream); + } + + private function readExactly(int $length, float $deadline): string + { + if (!\is_resource($this->serverStream)) { + Assert::fail('FileSocketViciServer has no accepted client.'); + } + + $serverStream = $this->serverStream; + $buf = ''; + while (\strlen($buf) < $length) { + $remaining = $deadline - microtime(true); + if ($remaining <= 0) { + Assert::fail(\sprintf( + 'Timed out reading %d bytes from client (got %d).', + $length, + \strlen($buf), + )); + } + $read = [$serverStream]; + $write = null; + $except = null; + $sec = (int) floor($remaining); + $usec = (int) round(($remaining - $sec) * 1_000_000); + $ready = stream_select($read, $write, $except, $sec, $usec); + if ($ready === false || $ready === 0) { + Assert::fail('Timed out reading from client.'); + } + $need = $length - \strlen($buf); + \assert($need > 0); + $chunk = fread($serverStream, $need); + if ($chunk === false || $chunk === '') { + if (feof($serverStream)) { + throw new \RuntimeException('Client closed connection unexpectedly.'); + } + continue; + } + $buf .= $chunk; + } + + return $buf; + } + + private function startListening(): void + { + if (file_exists($this->path)) { + @unlink($this->path); + } + + $listenSocket = @stream_socket_server('unix://' . $this->path); + if ($listenSocket === false) { + throw new \RuntimeException('Failed to create FileSocketViciServer listener.'); + } + + $this->listenSocket = $listenSocket; + } +} diff --git a/tests/Integration/MockViciServer.php b/tests/Integration/MockViciServer.php index ddac89c..123b707 100644 --- a/tests/Integration/MockViciServer.php +++ b/tests/Integration/MockViciServer.php @@ -28,7 +28,7 @@ final class MockViciServer /** @var resource */ private $serverStream; - private readonly StreamTransport $clientTransport; + private StreamTransport $clientTransport; private readonly PacketCodec $codec; private readonly MessageEncoder $encoder; private readonly MessageDecoder $decoder; @@ -65,6 +65,32 @@ public function close(): void $this->clientTransport->close(); } + /** + * Close the current socket pair and replace it with a fresh one, as if + * charon had restarted. Callers using {@see StreamTransport} must adopt + * {@see getClientTransport()} again; {@see ReconnectingTransport} clients + * recover via automatic reconnect on the next I/O. + */ + public function simulateRestart(): void + { + if (\is_resource($this->serverStream)) { + @fclose($this->serverStream); + } + $this->clientTransport->close(); + + $pair = stream_socket_pair( + \STREAM_PF_UNIX, + \STREAM_SOCK_STREAM, + \STREAM_IPPROTO_IP, + ); + if ($pair === false) { + throw new \RuntimeException('Failed to create socket pair for MockViciServer restart.'); + } + [$clientSide, $serverSide] = $pair; + $this->serverStream = $serverSide; + $this->clientTransport = new StreamTransport($clientSide, readTimeout: 2.0); + } + // ------------------------------------------------------------------ // Inbound (read what the client sent). // ------------------------------------------------------------------ diff --git a/tests/Integration/ReconnectingTransportIntegrationTest.php b/tests/Integration/ReconnectingTransportIntegrationTest.php new file mode 100644 index 0000000..38d9c41 --- /dev/null +++ b/tests/Integration/ReconnectingTransportIntegrationTest.php @@ -0,0 +1,80 @@ +server = new FileSocketViciServer(); + + $this->session = new Session(new ReconnectingTransport( + path: $this->server->getPath(), + connectTimeout: 2.0, + readTimeout: 2.0, + maxReconnectAttempts: 3, + reconnectDelayMs: 50, + )); + $this->server->acceptClient(); + } + + protected function tearDown(): void + { + $this->session->close(); + $this->server->close(); + } + + public function testVersionSurvivesSimulatedServerRestart(): void + { + $this->server->sendCmdResponse([ + 'daemon' => 'charon', + 'version' => '5.9.13', + 'sysname' => 'Linux', + 'release' => '6.1.0', + 'machine' => 'x86_64', + ]); + + $version = $this->session->version(); + self::assertSame('charon', $version['daemon']); + self::assertSame([], $this->server->expectCommand('version')); + + $this->server->simulateRestart(); + + if (!\function_exists('pcntl_fork')) { + self::markTestSkipped('pcntl extension required to accept during reconnect.'); + } + + $server = $this->server; + $pid = pcntl_fork(); + if ($pid === -1) { + self::markTestSkipped('pcntl_fork() failed.'); + } + if ($pid === 0) { + $server->acceptClient(5.0); + $server->expectCommand('version'); + $server->sendCmdResponse([ + 'daemon' => 'charon', + 'version' => '5.9.14', + 'sysname' => 'Linux', + 'release' => '6.1.0', + 'machine' => 'x86_64', + ]); + exit(0); + } + + $version = $this->session->version(); + pcntl_waitpid($pid, $status); + + self::assertSame('5.9.14', $version['version']); + self::assertSame(0, pcntl_wexitstatus($status)); + } +} diff --git a/tests/Unit/Transport/ReconnectingTransportTest.php b/tests/Unit/Transport/ReconnectingTransportTest.php new file mode 100644 index 0000000..2551ec2 --- /dev/null +++ b/tests/Unit/Transport/ReconnectingTransportTest.php @@ -0,0 +1,229 @@ +socketPath = sys_get_temp_dir() . '/vici-reconnect-test-' . uniqid('', true) . '.sock'; + $this->startListener(); + } + + protected function tearDown(): void + { + if (\is_resource($this->serverStream)) { + @fclose($this->serverStream); + } + if (\is_resource($this->listenSocket)) { + @fclose($this->listenSocket); + } + if (file_exists($this->socketPath)) { + @unlink($this->socketPath); + } + } + + public function testReceiveRetriesAfterServerRestart(): void + { + $transport = new ReconnectingTransport( + path: $this->socketPath, + connectTimeout: 2.0, + readTimeout: 2.0, + maxReconnectAttempts: 3, + reconnectDelayMs: 50, + ); + $this->acceptClient(); + + $serverStream = $this->requireServerStream(); + fwrite($serverStream, pack('N', 5) . 'first'); + fflush($serverStream); + self::assertSame('first', $transport->receive()); + + $this->simulateServerRestart(); + $pid = $this->forkAcceptAndWrite(pack('N', 6) . 'second'); + self::assertSame('second', $transport->receive()); + $this->waitFork($pid); + + $transport->close(); + } + + public function testSendRetriesAfterServerRestart(): void + { + $transport = new ReconnectingTransport( + path: $this->socketPath, + connectTimeout: 2.0, + readTimeout: 2.0, + maxReconnectAttempts: 3, + reconnectDelayMs: 50, + ); + $this->acceptClient(); + + $transport->send('payload'); + $serverStream = $this->requireServerStream(); + $header = fread($serverStream, 4); + self::assertIsString($header); + /** @var array{1: int} $unpacked */ + $unpacked = unpack('N', $header); + self::assertSame(7, $unpacked[1]); + self::assertSame('payload', fread($serverStream, 7)); + + $this->simulateServerRestart(); + $pid = $this->forkAcceptAndReadExpectedFrame(5, 'again'); + + $transport->send('again'); + $this->waitFork($pid); + + $transport->close(); + } + + private function simulateServerRestart(): void + { + if (\is_resource($this->serverStream)) { + @fclose($this->serverStream); + } + $this->serverStream = null; + if (\is_resource($this->listenSocket)) { + @fclose($this->listenSocket); + } + $this->listenSocket = null; + if (file_exists($this->socketPath)) { + @unlink($this->socketPath); + } + $this->startListener(); + } + + private function startListener(): void + { + if (file_exists($this->socketPath)) { + @unlink($this->socketPath); + } + + $listenSocket = @stream_socket_server('unix://' . $this->socketPath); + if ($listenSocket === false) { + self::fail('Failed to create unix socket listener.'); + } + + $this->listenSocket = $listenSocket; + } + + private function acceptClient(float $timeout = 2.0): void + { + $listenSocket = $this->requireListenSocket(); + $read = [$listenSocket]; + $write = null; + $except = null; + $ready = stream_select($read, $write, $except, (int) floor($timeout), 0); + if ($ready === false || $ready === 0) { + self::fail('Timed out waiting for client connection.'); + } + + $client = stream_socket_accept($listenSocket, $timeout); + if ($client === false) { + self::fail('Failed to accept client connection.'); + } + + $this->serverStream = $client; + } + + private function forkAcceptAndWrite(string $bytes): int + { + if (!\function_exists('pcntl_fork')) { + self::markTestSkipped('pcntl extension required to accept during reconnect.'); + } + + $listenSocket = $this->requireListenSocket(); + $pid = pcntl_fork(); + if ($pid === -1) { + self::markTestSkipped('pcntl_fork() failed.'); + } + if ($pid === 0) { + $client = stream_socket_accept($listenSocket, 5.0); + if ($client === false) { + exit(1); + } + fwrite($client, $bytes); + fflush($client); + exit(0); + } + + return $pid; + } + + private function forkAcceptAndReadExpectedFrame(int $length, string $payload): int + { + if (!\function_exists('pcntl_fork')) { + self::markTestSkipped('pcntl extension required to accept during reconnect.'); + } + + $listenSocket = $this->requireListenSocket(); + $pid = pcntl_fork(); + if ($pid === -1) { + self::markTestSkipped('pcntl_fork() failed.'); + } + if ($pid === 0) { + $client = stream_socket_accept($listenSocket, 5.0); + if ($client === false) { + exit(1); + } + $header = fread($client, 4); + $body = $length > 0 ? fread($client, $length) : ''; + if (!\is_string($header) || !\is_string($body)) { + exit(1); + } + /** @var array{1: int} $unpacked */ + $unpacked = unpack('N', $header); + if ($unpacked[1] !== $length || $body !== $payload) { + exit(1); + } + exit(0); + } + + return $pid; + } + + private function waitFork(int $pid): void + { + $status = 0; + pcntl_waitpid($pid, $status); + self::assertSame(0, pcntl_wexitstatus($status)); + } + + /** + * @return resource + */ + private function requireListenSocket() + { + if (!\is_resource($this->listenSocket)) { + self::fail('Unix socket listener is not open.'); + } + + return $this->listenSocket; + } + + /** + * @return resource + */ + private function requireServerStream() + { + if (!\is_resource($this->serverStream)) { + self::fail('No accepted client connection.'); + } + + return $this->serverStream; + } +} diff --git a/tests/Unit/Transport/StreamTransportTest.php b/tests/Unit/Transport/StreamTransportTest.php index 3f3cad8..3dd2e1b 100644 --- a/tests/Unit/Transport/StreamTransportTest.php +++ b/tests/Unit/Transport/StreamTransportTest.php @@ -132,10 +132,27 @@ public function testReceiveDetectsRemoteClose(): void try { $transport->receive(0.5); } finally { + self::assertFalse($transport->isConnected()); $transport->close(); } } + public function testConnectionFailureInvalidatesStream(): void + { + [$transport, $peer] = $this->makePair(); + self::assertTrue($transport->isConnected()); + + fclose($peer); + + try { + $transport->receive(0.5); + self::fail('Expected ConnectionException.'); + } catch (\Bk203\Vici\Exception\ConnectionException) { + self::assertFalse($transport->isConnected()); + self::assertNull($transport->getStream()); + } + } + public function testIsConnectedReflectsState(): void { [$transport, $peer] = $this->makePair(); diff --git a/tests/Unit/Transport/UnixSocketTransportTest.php b/tests/Unit/Transport/UnixSocketTransportTest.php new file mode 100644 index 0000000..54c48ec --- /dev/null +++ b/tests/Unit/Transport/UnixSocketTransportTest.php @@ -0,0 +1,150 @@ +socketPath = sys_get_temp_dir() . '/vici-test-' . uniqid('', true) . '.sock'; + $this->startListener(); + } + + protected function tearDown(): void + { + if (\is_resource($this->serverStream)) { + @fclose($this->serverStream); + } + if (\is_resource($this->listenSocket)) { + @fclose($this->listenSocket); + } + if (file_exists($this->socketPath)) { + @unlink($this->socketPath); + } + } + + public function testReconnectOpensNewConnectionAfterPeerCloses(): void + { + $transport = new UnixSocketTransport($this->socketPath, connectTimeout: 2.0, readTimeout: 2.0); + $this->acceptClient(); + + $serverStream = $this->requireServerStream(); + fwrite($serverStream, pack('N', 5) . 'hello'); + fflush($serverStream); + self::assertSame('hello', $transport->receive()); + + $this->closeListener(); + + $this->startListener(); + + $transport->reconnect(); + $this->acceptClient(); + self::assertTrue($transport->isConnected()); + + $serverStream = $this->requireServerStream(); + fwrite($serverStream, pack('N', 5) . 'world'); + fflush($serverStream); + self::assertSame('world', $transport->receive()); + + $transport->close(); + } + + public function testReconnectThrowsWhenSocketFileMissing(): void + { + $transport = new UnixSocketTransport($this->socketPath, connectTimeout: 0.2, readTimeout: 2.0); + $this->acceptClient(); + + $this->closeListener(); + + $this->expectException(ConnectionException::class); + $transport->reconnect(); + } + + private function startListener(): void + { + if (file_exists($this->socketPath)) { + @unlink($this->socketPath); + } + + $listenSocket = @stream_socket_server('unix://' . $this->socketPath); + if ($listenSocket === false) { + self::fail('Failed to create unix socket listener.'); + } + + $this->listenSocket = $listenSocket; + } + + private function acceptClient(float $timeout = 2.0): void + { + $listenSocket = $this->requireListenSocket(); + $read = [$listenSocket]; + $write = null; + $except = null; + $ready = stream_select($read, $write, $except, (int) floor($timeout), 0); + if ($ready === false || $ready === 0) { + self::fail('Timed out waiting for client connection.'); + } + + $client = stream_socket_accept($listenSocket, $timeout); + if ($client === false) { + self::fail('Failed to accept client connection.'); + } + + $this->serverStream = $client; + } + + private function closeListener(): void + { + if (\is_resource($this->serverStream)) { + @fclose($this->serverStream); + } + $this->serverStream = null; + if (\is_resource($this->listenSocket)) { + @fclose($this->listenSocket); + } + $this->listenSocket = null; + if (file_exists($this->socketPath)) { + @unlink($this->socketPath); + } + } + + /** + * @return resource + */ + private function requireListenSocket() + { + if (!\is_resource($this->listenSocket)) { + self::fail('Unix socket listener is not open.'); + } + + return $this->listenSocket; + } + + /** + * @return resource + */ + private function requireServerStream() + { + if (!\is_resource($this->serverStream)) { + self::fail('No accepted client connection.'); + } + + return $this->serverStream; + } +}