From d79c65afeed91504c191bb69562d7dd0ddb24ae8 Mon Sep 17 00:00:00 2001 From: Jonathan Zernik Date: Wed, 25 Aug 2021 03:54:57 -0700 Subject: [PATCH] Do initial peer sync outside context manager (#1069) * Do initial peer sync outside context manager * Simplify context manager for opening peer connection --- squeaknode/network/connection_manager.py | 4 +-- squeaknode/network/network_manager.py | 3 +- squeaknode/network/peer.py | 42 ++++++++++++++---------- 3 files changed, 28 insertions(+), 21 deletions(-) diff --git a/squeaknode/network/connection_manager.py b/squeaknode/network/connection_manager.py index 771ea70b..3064847d 100644 --- a/squeaknode/network/connection_manager.py +++ b/squeaknode/network/connection_manager.py @@ -87,14 +87,14 @@ class ConnectionManager(object): with self.peers_lock: peer = self.get_peer(address) if peer is not None: - peer.close() + peer.stop() def stop_all_connections(self): """Stop all peer connections. """ with self.peers_lock: for peer in self.peers: - peer.close() + peer.stop() def add_peers_changed_callback(self, name, callback): self.peers_changed_callbacks[name] = callback diff --git a/squeaknode/network/network_manager.py b/squeaknode/network/network_manager.py index 20c3ed4e..8b7de494 100644 --- a/squeaknode/network/network_manager.py +++ b/squeaknode/network/network_manager.py @@ -101,9 +101,10 @@ class NetworkManager(object): ).open_connection(squeak_controller) as peer: self.connection_manager.add_peer(peer) try: + peer.sync(squeak_controller) peer.handle_messages(squeak_controller) except Exception: - logger.exception("Handling messages failed.") + logger.exception("Peer connection failed.") finally: self.connection_manager.remove_peer(peer) logger.debug('Stopped controller for peer address {}.'.format(address)) diff --git a/squeaknode/network/peer.py b/squeaknode/network/peer.py index ca54438e..71e36916 100644 --- a/squeaknode/network/peer.py +++ b/squeaknode/network/peer.py @@ -161,11 +161,19 @@ class Peer(object): logger.info('Received msg {} from {}'.format(msg, self)) return msg - def close(self): - logger.info("closing peer socket: {}".format(self._peer_socket)) + def start(self): + msg_receiver = MessageReceiver( + self._peer_socket, self._recv_msg_queue, self.stopped) + threading.Thread( + target=msg_receiver.recv_msgs, + args=(), + ).start() + + def stop(self): + logger.info("Stopping peer socket: {}".format(self._peer_socket)) try: self._peer_socket.shutdown(socket.SHUT_RDWR) - self._peer_socket.close() + self._peer_socket.stop() except Exception: pass @@ -177,11 +185,11 @@ class Peer(object): self._peer_socket.send(data) except Exception: logger.info('Failed to send msg to {}'.format(self)) - self.close() + self.stop() def handshake(self, squeak_controller): timer = HandshakeTimer( - self.close, + self.stop, str(self), ) timer.start_timer() @@ -239,23 +247,21 @@ class Peer(object): self, squeak_controller) peer_message_handler.handle_msgs() + def sync(self, squeak_controller): + # TODO: getaddrs from peer. + self.update_subscription(squeak_controller) + @contextmanager def open_connection(self, squeak_controller): logger.debug('Setting up peer {} ...'.format(self)) try: - msg_receiver = MessageReceiver( - self._peer_socket, self._recv_msg_queue, self.stopped) - threading.Thread( - target=msg_receiver.recv_msgs, - args=(), - ).start() + self.start() self.handshake(squeak_controller) - self.update_subscription(squeak_controller) self.set_connected() yield self finally: - self.close() - logger.debug('Closed connection to peer {} ...'.format(self)) + self.stop() + logger.debug('Stopped connection to peer {} ...'.format(self)) def __repr__(self): return "Peer(%s)" % (str(self.remote_address)) @@ -328,14 +334,14 @@ class MessageReceiver: class HandshakeTimer: - """Close the peer if handshake is not complete before timeout. + """Stop the peer if handshake is not complete before timeout. """ def __init__(self, - close_fn, + stop_fn, peer_name, ): - self.close_fn = close_fn + self.stop_fn = stop_fn self.peer_name = peer_name self.timer = None @@ -353,4 +359,4 @@ class HandshakeTimer: def stop_peer(self): logger.info("Closing peer from handshake timer.") - self.close_fn() + self.stop_fn()