Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
58 changes: 25 additions & 33 deletions openc3/python/openc3/streams/tcpip_client_stream.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,8 @@
# if purchased from OpenC3, Inc.

import errno
import os
import select
import socket

from openc3.config.config_parser import ConfigParser
Expand Down Expand Up @@ -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))
45 changes: 45 additions & 0 deletions openc3/python/test/streams/test_tcpip_client_stream.py
Original file line number Diff line number Diff line change
Expand Up @@ -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 *
Expand All @@ -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
Expand Down
Loading