From 8c84082e245c75c3badaaea14793f62b1fc4deeb Mon Sep 17 00:00:00 2001 From: Emily Ragan Date: Wed, 30 Sep 2026 15:30:44 -0600 Subject: [PATCH] fix(streams): wait for tcpip client connect instead of retrying Retrying a non-blocking connect() treated EALREADY as connected. On macOS a refused connect still pending on retry returns EALREADY, so the client reported success. Wait for writability with connect_timeout (previously ignored) and read SO_ERROR for the real outcome. Co-Authored-By: Claude Opus 5.5 --- .../openc3/streams/tcpip_client_stream.py | 58 ++++++++----------- .../test/streams/test_tcpip_client_stream.py | 45 ++++++++++++++ 2 files changed, 70 insertions(+), 33 deletions(-) diff --git a/openc3/python/openc3/streams/tcpip_client_stream.py b/openc3/python/openc3/streams/tcpip_client_stream.py index 92a3ce972c..3bbfcb8365 100644 --- a/openc3/python/openc3/streams/tcpip_client_stream.py +++ b/openc3/python/openc3/streams/tcpip_client_stream.py @@ -10,6 +10,8 @@ # if purchased from OpenC3, Inc. import errno +import os +import select import socket from openc3.config.config_parser import ConfigParser @@ -87,36 +89,26 @@ def connect(self): raise super().connect() - def _connect(self, socket, hostname, port): - while True: - try: - socket.connect((hostname, port)) - except BlockingIOError: - # select.select([], [socket], [], self.connect_timeout) - # This is not an error condition - continue - except OSError as error: - if error.errno == errno.EINPROGRESS: - continue - if error.errno == errno.EISCONN or error.errno == errno.EALREADY: - break - else: - raise error - - # except: - # try: - # _, sockets, _ = IO.select(None, [socket], None, self.connect_timeout) # wait 3-way handshake completion - # except IOError, Errno='ENOTSOCK': - # raise "Connect canceled" - # if sockets and !sockets.empty?: - # try: - # socket.connect_nonblock(addr) # check connection failure - # except IOError, Errno='ENOTSOCK': - # raise "Connect canceled" - # except Errno='EINPROGRESS': - # retry - # except Errno='EISCONN', Errno='EALREADY': - # else: - # raise "Connect timeout" - # except IOError, Errno='ENOTSOCK': - # raise "Connect canceled" + def _connect(self, sock, hostname, port): + # The sockets are non-blocking so connect_ex returns immediately. Wait for + # the socket to become writable (3-way handshake done or failed) and then + # read SO_ERROR to learn the outcome. Retrying connect() instead is not + # portable: on macOS a retry while the handshake is still pending returns + # EALREADY, which would wrongly be treated as connected. + try: + result = sock.connect_ex((hostname, port)) + if result in (errno.EINPROGRESS, errno.EWOULDBLOCK, errno.EALREADY): + # Windows reports a failed connect through the exceptional set + _, writeable, exceptional = select.select([], [sock], [sock], self.connect_timeout) + if not writeable and not exceptional: + raise TimeoutError(f"Connect timeout to {hostname}:{port}") + result = sock.getsockopt(socket.SOL_SOCKET, socket.SO_ERROR) + except (ValueError, OSError) as error: + # Python sets fileno() to -1 once a socket is closed, so a disconnect + # from another thread surfaces as ValueError, EBADF or ENOTSOCK + if isinstance(error, ValueError) or error.errno in (errno.EBADF, errno.ENOTSOCK): + raise RuntimeError("Connect canceled") from error + raise + if result not in (0, errno.EISCONN): + # OSError maps the errno to its subclass, e.g. ConnectionRefusedError + raise OSError(result, os.strerror(result)) diff --git a/openc3/python/test/streams/test_tcpip_client_stream.py b/openc3/python/test/streams/test_tcpip_client_stream.py index 3f24f4ee42..c62a5eb855 100644 --- a/openc3/python/test/streams/test_tcpip_client_stream.py +++ b/openc3/python/test/streams/test_tcpip_client_stream.py @@ -9,6 +9,8 @@ # This file may also be used under the terms of a commercial license # if purchased from OpenC3, Inc. +import errno +import socket import socketserver import unittest from unittest.mock import * @@ -28,6 +30,49 @@ def test_complains_if_the_host_is_bad(self): ): TcpipClientStream("asdf", 8888, 8888, 10.0, None) + def unused_port(self): + with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as s: + s.bind(("127.0.0.1", 0)) + return s.getsockname()[1] + + def test_raises_connection_refused_when_nothing_is_listening(self): + port = self.unused_port() + ss = TcpipClientStream("localhost", port, port, 10.0, None) + with self.assertRaises(ConnectionRefusedError): + ss.connect() + self.assertFalse(ss.connected()) + + def test_connects_to_a_listening_server(self): + with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as server: + server.bind(("127.0.0.1", 0)) + server.listen(1) + port = server.getsockname()[1] + ss = TcpipClientStream("localhost", port, port, 10.0, None) + ss.connect() + self.assertTrue(ss.connected()) + self.assertEqual(ss.read_socket, ss.write_socket) + conn, _ = server.accept() + conn.close() + ss.disconnect() + + @patch("openc3.streams.tcpip_client_stream.select.select", return_value=([], [], [])) + def test_raises_a_timeout_when_the_handshake_does_not_complete(self, _select): + ss = TcpipClientStream("localhost", 8888, 8888, 0.1, None, 0.1) + with ( + patch.object(socket.socket, "connect_ex", return_value=errno.EINPROGRESS), + self.assertRaisesRegex(TimeoutError, "Connect timeout"), + ): + ss.connect() + + @patch("openc3.streams.tcpip_client_stream.select.select", side_effect=ValueError("file descriptor cannot be -1")) + def test_raises_canceled_when_the_socket_is_closed_during_connect(self, _select): + ss = TcpipClientStream("localhost", 8888, 8888, 0.1, None) + with ( + patch.object(socket.socket, "connect_ex", return_value=errno.EINPROGRESS), + self.assertRaisesRegex(RuntimeError, "Connect canceled"), + ): + ss.connect() + # TODO: Fails with Traceback (most recent call last): # File "/home/runner/work/cosmos/cosmos/openc3/python/test/streams/test_tcpip_client_stream.py", line 42, in test_uses_the_same_socket_if_read_port_equals_write_port