openrvdas.logger.readers.tcp_reader
No module-level documentation available.
1#!/usr/bin/env python3 2 3import logging 4import mmap 5import socket 6 7from logger.readers.reader import Reader # noqa: E402 8 9# the size of each recv() call in read_to_eol() 10# 11# NOTE: Optimal value should be the system's page size, which is probably 4k. 12# 13READ_BUFFER_SIZE = mmap.PAGESIZE 14 15 16################################################################################ 17class TCPReader(Reader): 18 """ Read TCP packets from network.""" 19 ############################ 20 def __init__(self, interface=None, port=None, eol=None, 21 reuseaddr=True, reuseport=False, 22 encoding='utf-8', encoding_errors='ignore', **kwargs): 23 """ 24 ``` 25 interface IP (or resolvable name) of interface to listen on. None or '' 26 will listen on INADDR_ANY (all interfaces). 27 28 port Port to listen to for packets. REQUIRED 29 30 eol If specified, buffer network reads until the `eol` sequence has 31 been seen, and return the entire record at once. In other words, 32 present the user with a very UDP-ish 1 read gives you 1 whole 33 record feel, even though this is TCP. If not specified, read() 34 calls must specify read sizes and it's up to the user to 35 control the TCP stream. 36 37 reuseaddr Specifies wether we set SO_REUSEADDR on the created socket. This 38 is enabled by default (unlike TCPWriter, UDPWriter, or UDPReader) 39 specifically to avoid leftover TIME_WAIT sockets from 40 interferring with startup. 41 42 reuseport Specifies wether we set SO_REUSEPORT on the created socket. If 43 you don't know you need this, don't enable it. 44 45 encoding - 'utf-8' by default. If empty or None, do not attempt any decoding 46 and return raw bytes. Other possible encodings are listed in online 47 documentation here: 48 https://docs.python.org/3/library/codecs.html#standard-encodings 49 50 encoding_errors - 'ignore' by default. Other error strategies are 'strict', 51 'replace', and 'backslashreplace', described here: 52 https://docs.python.org/3/howto/unicode.html#encodings 53 ``` 54 """ 55 super().__init__(encoding=encoding, encoding_errors=encoding_errors, **kwargs) 56 57 if interface: 58 # resolve once in constructor 59 interface = socket.gethostbyname(interface) 60 else: 61 interface = '' 62 self.interface = interface 63 64 # make sure user passed in `port` 65 # 66 # NOTE: We want the order of the arguments to consistently be (ip, 67 # port, ...) across all the network readers/writers... but we 68 # want `interface` to be optional. All kwargs need to come after 69 # all regular args, so we've assigned a default value of None to 70 # `port`. But don't be confused, it is REQUIRED. 71 # 72 if not port: 73 raise TypeError('must specify `port`') 74 # make sure port gets stored as an int, even if passed in as a string 75 self.port = int(port) 76 77 # prep eol 78 # 79 # NOTE: We're going to be looking for `eol` inside of the byte array 80 # returned by recv(), so it has to be encoded to bytes here 81 # regardless of `encoding`. 82 # 83 # NOTE: We set unescape to True because `eol` might be provided w/ 84 # escaped out characters 85 # 86 if eol: 87 eol = self._encode_str(eol, unescape=True) 88 self.eol = eol 89 90 # Where we'll aggregate incomplete records if an eol char is specified 91 self.record_buffer = b'' 92 93 self.reuseaddr = reuseaddr 94 self.reuseport = reuseport 95 96 # initialize this now, so our socket is ready to accept connections 97 # once constructed 98 self.s_listening = self._open_socket() 99 100 # these get set when we successfully accept() an incoming connection in read() 101 self.s_connected = None 102 self.client_addr = None 103 104 ############################ 105 def __del__(self): 106 if self.s_connected: 107 logging.debug('__del__: closing s_connected') 108 self._close_socket(self.s_connected) 109 110 if self.s_listening: 111 logging.debug('__del__: closing s_listening') 112 self._close_socket(self.s_listening) 113 114 ############################ 115 def _open_socket(self): 116 # create TCP socket 117 s = socket.socket(family=socket.AF_INET, type=socket.SOCK_STREAM, proto=socket.IPPROTO_TCP) 118 119 # set sockopts 120 if self.reuseaddr: 121 s.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, True) 122 if self.reuseport: 123 try: # Raspbian doesn't recognize SO_REUSEPORT 124 s.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEPORT, True) 125 except AttributeError: 126 logging.warning('Unable to set socket REUSEPORT; may be unsupported.') 127 128 # bind to specificed interface 129 s.bind((self.interface, self.port)) 130 131 # start listening 132 s.listen() 133 134 return s 135 136 ############################ 137 def _close_socket(self, s): 138 try: 139 s.shutdown(socket.SHUT_RDWR) 140 s.close() 141 except OSError: 142 logging.debug('Unable to close socket') 143 144 ############################ 145 def _get_connected_socket(self): 146 # make sure we're already listening 147 if not self.s_listening: 148 self.s_listening = self._open_socket() 149 if not self.s_listening: 150 logging.error("failed to open socket") 151 return 152 153 # accept connections 154 # 155 # NOTE: This will block until we have an incoming connection. 156 # 157 s_connected, client_addr = self.s_listening.accept() 158 logging.debug('got connection from %s', client_addr) 159 return s_connected, client_addr 160 161 ############################ 162 def _read_size(self, size): 163 try: 164 # NOTE: This will block until there is *something* to return, but 165 # it won't wait for the full `size` before returning. 166 record = self.s_connected.recv(size) 167 # catch disconnected socket 168 # 169 # NOTE: If remote side disconnects, recv() "successfully" returns 0 170 # bytes. We have to turn that into a failure so we 171 # re-establish comms before trying to recv() again. 172 # 173 if len(record) == 0: 174 raise OSError("socket disconneced") 175 except OSError as e: 176 logging.error('TCPReader recv error: %s', str(e)) 177 # nuke the socket so we reconnect on next read() 178 self._close_socket(self.s_connected) 179 self.s_connected = None 180 return None 181 logging.debug('TCPReader._read_size: received %d bytes', len(record)) 182 return record 183 184 ############################ 185 def _read_to_eol(self): 186 while self.eol not in self.record_buffer: 187 try: 188 record = self.s_connected.recv(READ_BUFFER_SIZE) 189 # catch disconnected socket 190 # 191 # NOTE: If remote side disconnects, recv() "successfully" 192 # returns 0 bytes. We have to turn that into a failure 193 # so we re-establish comms before trying to recv() again. 194 # 195 if len(record) == 0: 196 raise OSError("socket disconneced") 197 except OSError as e: 198 logging.error('TCPReader recv error: %s', str(e)) 199 # nuke the socket so we reconnect on next read() 200 self._close_socket(self.s_connected) 201 self.s_connected = None 202 return None 203 logging.debug('TCPReader._read_to_eol: received %d bytes', len(record)) 204 205 self.record_buffer += record 206 207 # we've got `eol` in our buffer, split out the first record 208 i = self.record_buffer.find(self.eol) 209 # `i` is the index of the BEGINNING of our `eol` sequence 210 record = self.record_buffer[:i] 211 # beginning of NEXT message is at i+len(eol) 212 self.record_buffer = self.record_buffer[i+len(self.eol):] 213 return record 214 215 ############################ 216 def read(self, size=None): 217 """Read from TCP socket, either up to the next `eol` in the stream or up to 218 `size` bytes (ignoring `eol` even if set!), return the result. 219 """ 220 # If socket isn't ready, set it up. If something fails, return w/out reading. 221 if not self.s_connected: 222 self.s_connected, self.client_addr = self._get_connected_socket() 223 if not self.s_connected: 224 logging.error('TCPReader.read: unable to get connected socket') 225 return 226 227 if size: 228 record = self._read_size(size) 229 elif self.eol: 230 record = self._read_to_eol() 231 else: 232 # invalid, need `eol` or `size` 233 logging.error('need either `eol` or `size`') 234 return 235 236 record = self._decode_bytes(record) 237 return record
READ_BUFFER_SIZE =
4096
class
TCPReader(logger.readers.reader.Reader):
18class TCPReader(Reader): 19 """ Read TCP packets from network.""" 20 ############################ 21 def __init__(self, interface=None, port=None, eol=None, 22 reuseaddr=True, reuseport=False, 23 encoding='utf-8', encoding_errors='ignore', **kwargs): 24 """ 25 ``` 26 interface IP (or resolvable name) of interface to listen on. None or '' 27 will listen on INADDR_ANY (all interfaces). 28 29 port Port to listen to for packets. REQUIRED 30 31 eol If specified, buffer network reads until the `eol` sequence has 32 been seen, and return the entire record at once. In other words, 33 present the user with a very UDP-ish 1 read gives you 1 whole 34 record feel, even though this is TCP. If not specified, read() 35 calls must specify read sizes and it's up to the user to 36 control the TCP stream. 37 38 reuseaddr Specifies wether we set SO_REUSEADDR on the created socket. This 39 is enabled by default (unlike TCPWriter, UDPWriter, or UDPReader) 40 specifically to avoid leftover TIME_WAIT sockets from 41 interferring with startup. 42 43 reuseport Specifies wether we set SO_REUSEPORT on the created socket. If 44 you don't know you need this, don't enable it. 45 46 encoding - 'utf-8' by default. If empty or None, do not attempt any decoding 47 and return raw bytes. Other possible encodings are listed in online 48 documentation here: 49 https://docs.python.org/3/library/codecs.html#standard-encodings 50 51 encoding_errors - 'ignore' by default. Other error strategies are 'strict', 52 'replace', and 'backslashreplace', described here: 53 https://docs.python.org/3/howto/unicode.html#encodings 54 ``` 55 """ 56 super().__init__(encoding=encoding, encoding_errors=encoding_errors, **kwargs) 57 58 if interface: 59 # resolve once in constructor 60 interface = socket.gethostbyname(interface) 61 else: 62 interface = '' 63 self.interface = interface 64 65 # make sure user passed in `port` 66 # 67 # NOTE: We want the order of the arguments to consistently be (ip, 68 # port, ...) across all the network readers/writers... but we 69 # want `interface` to be optional. All kwargs need to come after 70 # all regular args, so we've assigned a default value of None to 71 # `port`. But don't be confused, it is REQUIRED. 72 # 73 if not port: 74 raise TypeError('must specify `port`') 75 # make sure port gets stored as an int, even if passed in as a string 76 self.port = int(port) 77 78 # prep eol 79 # 80 # NOTE: We're going to be looking for `eol` inside of the byte array 81 # returned by recv(), so it has to be encoded to bytes here 82 # regardless of `encoding`. 83 # 84 # NOTE: We set unescape to True because `eol` might be provided w/ 85 # escaped out characters 86 # 87 if eol: 88 eol = self._encode_str(eol, unescape=True) 89 self.eol = eol 90 91 # Where we'll aggregate incomplete records if an eol char is specified 92 self.record_buffer = b'' 93 94 self.reuseaddr = reuseaddr 95 self.reuseport = reuseport 96 97 # initialize this now, so our socket is ready to accept connections 98 # once constructed 99 self.s_listening = self._open_socket() 100 101 # these get set when we successfully accept() an incoming connection in read() 102 self.s_connected = None 103 self.client_addr = None 104 105 ############################ 106 def __del__(self): 107 if self.s_connected: 108 logging.debug('__del__: closing s_connected') 109 self._close_socket(self.s_connected) 110 111 if self.s_listening: 112 logging.debug('__del__: closing s_listening') 113 self._close_socket(self.s_listening) 114 115 ############################ 116 def _open_socket(self): 117 # create TCP socket 118 s = socket.socket(family=socket.AF_INET, type=socket.SOCK_STREAM, proto=socket.IPPROTO_TCP) 119 120 # set sockopts 121 if self.reuseaddr: 122 s.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, True) 123 if self.reuseport: 124 try: # Raspbian doesn't recognize SO_REUSEPORT 125 s.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEPORT, True) 126 except AttributeError: 127 logging.warning('Unable to set socket REUSEPORT; may be unsupported.') 128 129 # bind to specificed interface 130 s.bind((self.interface, self.port)) 131 132 # start listening 133 s.listen() 134 135 return s 136 137 ############################ 138 def _close_socket(self, s): 139 try: 140 s.shutdown(socket.SHUT_RDWR) 141 s.close() 142 except OSError: 143 logging.debug('Unable to close socket') 144 145 ############################ 146 def _get_connected_socket(self): 147 # make sure we're already listening 148 if not self.s_listening: 149 self.s_listening = self._open_socket() 150 if not self.s_listening: 151 logging.error("failed to open socket") 152 return 153 154 # accept connections 155 # 156 # NOTE: This will block until we have an incoming connection. 157 # 158 s_connected, client_addr = self.s_listening.accept() 159 logging.debug('got connection from %s', client_addr) 160 return s_connected, client_addr 161 162 ############################ 163 def _read_size(self, size): 164 try: 165 # NOTE: This will block until there is *something* to return, but 166 # it won't wait for the full `size` before returning. 167 record = self.s_connected.recv(size) 168 # catch disconnected socket 169 # 170 # NOTE: If remote side disconnects, recv() "successfully" returns 0 171 # bytes. We have to turn that into a failure so we 172 # re-establish comms before trying to recv() again. 173 # 174 if len(record) == 0: 175 raise OSError("socket disconneced") 176 except OSError as e: 177 logging.error('TCPReader recv error: %s', str(e)) 178 # nuke the socket so we reconnect on next read() 179 self._close_socket(self.s_connected) 180 self.s_connected = None 181 return None 182 logging.debug('TCPReader._read_size: received %d bytes', len(record)) 183 return record 184 185 ############################ 186 def _read_to_eol(self): 187 while self.eol not in self.record_buffer: 188 try: 189 record = self.s_connected.recv(READ_BUFFER_SIZE) 190 # catch disconnected socket 191 # 192 # NOTE: If remote side disconnects, recv() "successfully" 193 # returns 0 bytes. We have to turn that into a failure 194 # so we re-establish comms before trying to recv() again. 195 # 196 if len(record) == 0: 197 raise OSError("socket disconneced") 198 except OSError as e: 199 logging.error('TCPReader recv error: %s', str(e)) 200 # nuke the socket so we reconnect on next read() 201 self._close_socket(self.s_connected) 202 self.s_connected = None 203 return None 204 logging.debug('TCPReader._read_to_eol: received %d bytes', len(record)) 205 206 self.record_buffer += record 207 208 # we've got `eol` in our buffer, split out the first record 209 i = self.record_buffer.find(self.eol) 210 # `i` is the index of the BEGINNING of our `eol` sequence 211 record = self.record_buffer[:i] 212 # beginning of NEXT message is at i+len(eol) 213 self.record_buffer = self.record_buffer[i+len(self.eol):] 214 return record 215 216 ############################ 217 def read(self, size=None): 218 """Read from TCP socket, either up to the next `eol` in the stream or up to 219 `size` bytes (ignoring `eol` even if set!), return the result. 220 """ 221 # If socket isn't ready, set it up. If something fails, return w/out reading. 222 if not self.s_connected: 223 self.s_connected, self.client_addr = self._get_connected_socket() 224 if not self.s_connected: 225 logging.error('TCPReader.read: unable to get connected socket') 226 return 227 228 if size: 229 record = self._read_size(size) 230 elif self.eol: 231 record = self._read_to_eol() 232 else: 233 # invalid, need `eol` or `size` 234 logging.error('need either `eol` or `size`') 235 return 236 237 record = self._decode_bytes(record) 238 return record
Read TCP packets from network.
TCPReader( interface=None, port=None, eol=None, reuseaddr=True, reuseport=False, encoding='utf-8', encoding_errors='ignore', **kwargs)
21 def __init__(self, interface=None, port=None, eol=None, 22 reuseaddr=True, reuseport=False, 23 encoding='utf-8', encoding_errors='ignore', **kwargs): 24 """ 25 ``` 26 interface IP (or resolvable name) of interface to listen on. None or '' 27 will listen on INADDR_ANY (all interfaces). 28 29 port Port to listen to for packets. REQUIRED 30 31 eol If specified, buffer network reads until the `eol` sequence has 32 been seen, and return the entire record at once. In other words, 33 present the user with a very UDP-ish 1 read gives you 1 whole 34 record feel, even though this is TCP. If not specified, read() 35 calls must specify read sizes and it's up to the user to 36 control the TCP stream. 37 38 reuseaddr Specifies wether we set SO_REUSEADDR on the created socket. This 39 is enabled by default (unlike TCPWriter, UDPWriter, or UDPReader) 40 specifically to avoid leftover TIME_WAIT sockets from 41 interferring with startup. 42 43 reuseport Specifies wether we set SO_REUSEPORT on the created socket. If 44 you don't know you need this, don't enable it. 45 46 encoding - 'utf-8' by default. If empty or None, do not attempt any decoding 47 and return raw bytes. Other possible encodings are listed in online 48 documentation here: 49 https://docs.python.org/3/library/codecs.html#standard-encodings 50 51 encoding_errors - 'ignore' by default. Other error strategies are 'strict', 52 'replace', and 'backslashreplace', described here: 53 https://docs.python.org/3/howto/unicode.html#encodings 54 ``` 55 """ 56 super().__init__(encoding=encoding, encoding_errors=encoding_errors, **kwargs) 57 58 if interface: 59 # resolve once in constructor 60 interface = socket.gethostbyname(interface) 61 else: 62 interface = '' 63 self.interface = interface 64 65 # make sure user passed in `port` 66 # 67 # NOTE: We want the order of the arguments to consistently be (ip, 68 # port, ...) across all the network readers/writers... but we 69 # want `interface` to be optional. All kwargs need to come after 70 # all regular args, so we've assigned a default value of None to 71 # `port`. But don't be confused, it is REQUIRED. 72 # 73 if not port: 74 raise TypeError('must specify `port`') 75 # make sure port gets stored as an int, even if passed in as a string 76 self.port = int(port) 77 78 # prep eol 79 # 80 # NOTE: We're going to be looking for `eol` inside of the byte array 81 # returned by recv(), so it has to be encoded to bytes here 82 # regardless of `encoding`. 83 # 84 # NOTE: We set unescape to True because `eol` might be provided w/ 85 # escaped out characters 86 # 87 if eol: 88 eol = self._encode_str(eol, unescape=True) 89 self.eol = eol 90 91 # Where we'll aggregate incomplete records if an eol char is specified 92 self.record_buffer = b'' 93 94 self.reuseaddr = reuseaddr 95 self.reuseport = reuseport 96 97 # initialize this now, so our socket is ready to accept connections 98 # once constructed 99 self.s_listening = self._open_socket() 100 101 # these get set when we successfully accept() an incoming connection in read() 102 self.s_connected = None 103 self.client_addr = None
interface IP (or resolvable name) of interface to listen on. None or ''
will listen on INADDR_ANY (all interfaces).
port Port to listen to for packets. REQUIRED
eol If specified, buffer network reads until the `eol` sequence has
been seen, and return the entire record at once. In other words,
present the user with a very UDP-ish 1 read gives you 1 whole
record feel, even though this is TCP. If not specified, read()
calls must specify read sizes and it's up to the user to
control the TCP stream.
reuseaddr Specifies wether we set SO_REUSEADDR on the created socket. This
is enabled by default (unlike TCPWriter, UDPWriter, or UDPReader)
specifically to avoid leftover TIME_WAIT sockets from
interferring with startup.
reuseport Specifies wether we set SO_REUSEPORT on the created socket. If
you don't know you need this, don't enable it.
encoding - 'utf-8' by default. If empty or None, do not attempt any decoding
and return raw bytes. Other possible encodings are listed in online
documentation here:
https://docs.python.org/3/library/codecs.html#standard-encodings
encoding_errors - 'ignore' by default. Other error strategies are 'strict',
'replace', and 'backslashreplace', described here:
https://docs.python.org/3/howto/unicode.html#encodings
def
read(self, size=None):
217 def read(self, size=None): 218 """Read from TCP socket, either up to the next `eol` in the stream or up to 219 `size` bytes (ignoring `eol` even if set!), return the result. 220 """ 221 # If socket isn't ready, set it up. If something fails, return w/out reading. 222 if not self.s_connected: 223 self.s_connected, self.client_addr = self._get_connected_socket() 224 if not self.s_connected: 225 logging.error('TCPReader.read: unable to get connected socket') 226 return 227 228 if size: 229 record = self._read_size(size) 230 elif self.eol: 231 record = self._read_to_eol() 232 else: 233 # invalid, need `eol` or `size` 234 logging.error('need either `eol` or `size`') 235 return 236 237 record = self._decode_bytes(record) 238 return record