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
interface
port
eol
record_buffer
reuseaddr
reuseport
s_listening
s_connected
client_addr
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

Read from TCP socket, either up to the next eol in the stream or up to size bytes (ignoring eol even if set!), return the result.