From 36d89bd25e2fc963750dffa17bcbab649c46fbab Mon Sep 17 00:00:00 2001 From: Jonathan Zernik Date: Tue, 10 Aug 2021 18:48:00 -0700 Subject: [PATCH] Fix peer server thread shutdown (#903) --- frontend/src/squeakclient/requests.js | 14 +++++++++ squeaknode/admin/webapp/app.py | 8 +++++ squeaknode/network/peer_client.py | 1 + squeaknode/network/peer_server.py | 29 ++++++++++++------- .../node/process_received_payments_worker.py | 1 + squeaknode/node/squeak_node.py | 3 ++ squeaknode/node/squeak_offer_expiry_worker.py | 1 + squeaknode/node/squeak_peer_sync_worker.py | 6 ++++ 8 files changed, 53 insertions(+), 10 deletions(-) diff --git a/frontend/src/squeakclient/requests.js b/frontend/src/squeakclient/requests.js index cfc8be43..21f0a1b6 100644 --- a/frontend/src/squeakclient/requests.js +++ b/frontend/src/squeakclient/requests.js @@ -108,6 +108,8 @@ import { UnlikeSqueakReply, GetLikedSqueakDisplaysRequest, GetLikedSqueakDisplaysReply, + GetConnectedPeersRequest, + GetConnectedPeersReply, } from "../proto/squeak_admin_pb" console.log('The value of REACT_APP_SERVER_PORT is:', process.env.REACT_APP_SERVER_PORT); @@ -804,3 +806,15 @@ export function getLikedSqueakDisplaysRequest(handleResponse) { } ); } + +export function getConnectedPeersRequest(handleResponse) { + var request = new GetConnectedPeersRequest(); + makeRequest( + 'getconnectedpeers', + request, + GetConnectedPeersReply.deserializeBinary, + (response) => { + handleResponse(response.getConnectedPeersList()); + } + ); +} diff --git a/squeaknode/admin/webapp/app.py b/squeaknode/admin/webapp/app.py index f7e31fa8..1e8cc729 100644 --- a/squeaknode/admin/webapp/app.py +++ b/squeaknode/admin/webapp/app.py @@ -531,6 +531,14 @@ def create_app(handler, username, password): handler.handle_get_liked_squeak_display_entries, ) + @app.route("/getconnectedpeers", methods=["POST"]) + @login_required + def getconnectedpeers(): + return handle_request( + squeak_admin_pb2.GetConnectedPeersRequest(), + handler.handle_get_connected_peers, + ) + return app diff --git a/squeaknode/network/peer_client.py b/squeaknode/network/peer_client.py index eaecb865..dcbfc6bf 100644 --- a/squeaknode/network/peer_client.py +++ b/squeaknode/network/peer_client.py @@ -52,6 +52,7 @@ class PeerClient(object): threading.Thread( target=self.make_connection, args=(ip, port), + name="peer_client_connection_thread", ).start() def disconnect_address(self, address): diff --git a/squeaknode/network/peer_server.py b/squeaknode/network/peer_server.py index 3352febc..8319e0eb 100644 --- a/squeaknode/network/peer_server.py +++ b/squeaknode/network/peer_server.py @@ -21,24 +21,33 @@ class PeerServer(object): self.ip = socket.gethostbyname('localhost') self.port = port or squeak.params.params.DEFAULT_PORT self.connection_manager = connection_manager + self.listen_socket = socket.socket() def start(self, peer_handler): self.peer_handler = peer_handler # Start Listen thread - threading.Thread(target=self.accept_connections).start() + threading.Thread( + target=self.accept_connections, + name="peer_server_listen_thread", + ).start() def stop(self): # TODO: stop accepting connections thread. # TODO: stop every peer in connection manager. - pass + # pass + logger.info("Stopping peer server listener thread...") + self.listen_socket.shutdown(socket.SHUT_RDWR) + self.listen_socket.close() def accept_connections(self): - listen_socket = socket.socket() - listen_socket.bind(('', self.port)) - listen_socket.listen() - while True: - peer_socket, address = listen_socket.accept() - peer_socket.setblocking(True) - self.peer_handler.handle_connection( - peer_socket, address, outgoing=False) + try: + self.listen_socket.bind(('', self.port)) + self.listen_socket.listen() + while True: + peer_socket, address = self.listen_socket.accept() + peer_socket.setblocking(True) + self.peer_handler.handle_connection( + peer_socket, address, outgoing=False) + except Exception: + logger.info("Stopped accepting incoming connections.") diff --git a/squeaknode/node/process_received_payments_worker.py b/squeaknode/node/process_received_payments_worker.py index c9b76e2a..6daadc2c 100644 --- a/squeaknode/node/process_received_payments_worker.py +++ b/squeaknode/node/process_received_payments_worker.py @@ -13,6 +13,7 @@ class ProcessReceivedPaymentsWorker: threading.Thread( target=self.process_subscribed_invoices, # daemon=True, + name="process_received_payments_thread", ).start() def process_subscribed_invoices(self): diff --git a/squeaknode/node/squeak_node.py b/squeaknode/node/squeak_node.py index 378c9c7c..ab1bf95b 100644 --- a/squeaknode/node/squeak_node.py +++ b/squeaknode/node/squeak_node.py @@ -143,6 +143,9 @@ class SqueakNode: def stop_running(self): self.stopped.set() + # TODO: Use explicit stop to stop all components + self.peer_server.stop() + def load_lightning_client(config) -> LNDLightningClient: return LNDLightningClient( diff --git a/squeaknode/node/squeak_offer_expiry_worker.py b/squeaknode/node/squeak_offer_expiry_worker.py index 9907ee28..a92aa2a9 100644 --- a/squeaknode/node/squeak_offer_expiry_worker.py +++ b/squeaknode/node/squeak_offer_expiry_worker.py @@ -22,6 +22,7 @@ class SqueakOfferExpiryWorker: self.start_running, ) timer.daemon = True + timer.name = "squeak_offer_expiry_thread" timer.start() self.remove_expired_offers() diff --git a/squeaknode/node/squeak_peer_sync_worker.py b/squeaknode/node/squeak_peer_sync_worker.py index 2aa4a874..05478e68 100644 --- a/squeaknode/node/squeak_peer_sync_worker.py +++ b/squeaknode/node/squeak_peer_sync_worker.py @@ -17,6 +17,7 @@ class SqueakPeerSyncWorker: def sync_timeline(self): logger.info("Syncing timeline with peers...") + self.print_running_threads() self.squeak_controller.sync_timeline() self.squeak_controller.share_squeaks() @@ -27,5 +28,10 @@ class SqueakPeerSyncWorker: self.start_running, ) timer.daemon = True + timer.name = "squeak_peer_sync_thread" timer.start() self.sync_timeline() + + def print_running_threads(self): + for thread in threading.enumerate(): + print(thread.name)