From 9f3391a5783583a46cec727aeed9f11df85eaab2 Mon Sep 17 00:00:00 2001 From: Jonathan Zernik Date: Tue, 17 Aug 2021 16:13:32 -0700 Subject: [PATCH] Remove stopped event squeak node (#956) * Remove stopped event from squeak node class * Removed stop parameter from web server * Remove stop event from rpc server * Remove old comments * Remove more old comments --- .../admin/squeak_admin_server_servicer.py | 29 ++++++++-------- squeaknode/admin/webapp/app.py | 16 +++++---- .../node/process_received_payments_worker.py | 7 ++-- squeaknode/node/squeak_node.py | 33 +++---------------- 4 files changed, 35 insertions(+), 50 deletions(-) diff --git a/squeaknode/admin/squeak_admin_server_servicer.py b/squeaknode/admin/squeak_admin_server_servicer.py index e11f11fc..d4d736f3 100644 --- a/squeaknode/admin/squeak_admin_server_servicer.py +++ b/squeaknode/admin/squeak_admin_server_servicer.py @@ -13,11 +13,25 @@ logger = logging.getLogger(__name__) class SqueakAdminServerServicer(squeak_admin_pb2_grpc.SqueakAdminServicer): """Provides methods that implement functionality of squeak admin server.""" - def __init__(self, host, port, handler, stopped): + def __init__(self, host, port, handler): self.host = host self.port = port self.handler = handler - self.stopped = stopped + self.server = None + + def start(self): + self.server = grpc.server(futures.ThreadPoolExecutor(max_workers=10)) + squeak_admin_pb2_grpc.add_SqueakAdminServicer_to_server( + self, self.server) + self.server.add_insecure_port("{}:{}".format(self.host, self.port)) + logger.info("Starting SqueakAdminServerServicer...") + self.server.start() + + def stop(self): + if self.server is None: + return + self.server.stop(None) + logger.info("Stopped SqueakAdminServerServicer.") def LndGetInfo(self, request, context): return self.handler.handle_lnd_get_info(request) @@ -215,17 +229,6 @@ class SqueakAdminServerServicer(squeak_admin_pb2_grpc.SqueakAdminServicer): def UnlikeSqueak(self, request, context): return self.handler.handle_unlike_squeak(request) - def serve(self): - server = grpc.server(futures.ThreadPoolExecutor(max_workers=10)) - squeak_admin_pb2_grpc.add_SqueakAdminServicer_to_server(self, server) - server.add_insecure_port("{}:{}".format(self.host, self.port)) - logger.info("Starting SqueakAdminServerServicer...") - server.start() - # server.wait_for_termination() - self.stopped.wait() - server.stop(None) - logger.info("Stopped SqueakAdminServerServicer.") - def GetLikedSqueakDisplays(self, request, context): return self.handler.handle_get_(request) diff --git a/squeaknode/admin/webapp/app.py b/squeaknode/admin/webapp/app.py index 91c77a26..6d5e83eb 100644 --- a/squeaknode/admin/webapp/app.py +++ b/squeaknode/admin/webapp/app.py @@ -585,7 +585,6 @@ class SqueakAdminWebServer: login_disabled, allow_cors, handler, - stopped, ): self.host = host self.port = port @@ -593,7 +592,7 @@ class SqueakAdminWebServer: self.login_disabled = login_disabled self.allow_cors = allow_cors self.app = create_app(handler, username, password) - self.stopped = stopped + self.server = None def get_app(self): @@ -614,8 +613,8 @@ class SqueakAdminWebServer: return self.app - def serve(self): - server = make_server( + def start(self): + self.server = make_server( self.host, self.port, self.get_app(), @@ -624,9 +623,12 @@ class SqueakAdminWebServer: logger.info("Starting SqueakAdminWebServer...") threading.Thread( - target=server.serve_forever, + target=self.server.serve_forever, ).start() - self.stopped.wait() + + def stop(self): + if self.server is None: + return logger.info("Stopping SqueakAdminWebServer....") - server.shutdown() + self.server.shutdown() logger.info("Stopped SqueakAdminWebServer.") diff --git a/squeaknode/node/process_received_payments_worker.py b/squeaknode/node/process_received_payments_worker.py index 6daadc2c..c922f512 100644 --- a/squeaknode/node/process_received_payments_worker.py +++ b/squeaknode/node/process_received_payments_worker.py @@ -5,9 +5,9 @@ logger = logging.getLogger(__name__) class ProcessReceivedPaymentsWorker: - def __init__(self, payment_processor, stopped: threading.Event): + def __init__(self, payment_processor): self.payment_processor = payment_processor - self.stopped = stopped + self.stopped = threading.Event() def start_running(self): threading.Thread( @@ -16,6 +16,9 @@ class ProcessReceivedPaymentsWorker: name="process_received_payments_thread", ).start() + def stop_running(self): + self.stopped.set() + def process_subscribed_invoices(self): logger.info("Starting ProcessReceivedPaymentsWorker...") self.payment_processor.start_processing() diff --git a/squeaknode/node/squeak_node.py b/squeaknode/node/squeak_node.py index f20d18cc..a3b2fd77 100644 --- a/squeaknode/node/squeak_node.py +++ b/squeaknode/node/squeak_node.py @@ -1,5 +1,4 @@ import logging -import threading from squeak.params import SelectParams @@ -31,7 +30,6 @@ class SqueakNode: def __init__(self, config: SqueaknodeConfig): self.config = config - self.stopped = threading.Event() self._initialize() def _initialize(self): @@ -56,11 +54,11 @@ class SqueakNode: def start_running(self): # start admin rpc server if self.config.admin.rpc_enabled: - start_admin_rpc_server(self.admin_rpc_server) + self.admin_rpc_server.start() # start admin web server if self.config.webadmin.enabled: - start_admin_web_server(self.admin_web_server) + self.admin_web_server.start() # Start peer socket server and peer client self.network_manager.start(self.squeak_controller) @@ -74,10 +72,10 @@ class SqueakNode: self.received_payment_processor_worker.start_running() def stop_running(self): - self.stopped.set() - - # TODO: Use explicit stop to stop all components + self.admin_web_server.stop() + self.admin_rpc_server.stop() self.network_manager.stop() + self.received_payment_processor_worker.stop_running() def start_peer_connection_worker(self): logger.info("Starting peer connection worker...") @@ -187,7 +185,6 @@ class SqueakNode: self.config.admin.rpc_host, self.config.admin.rpc_port, self.admin_handler, - self.stopped, ) def initialize_admin_web_server(self): @@ -200,29 +197,9 @@ class SqueakNode: self.config.webadmin.login_disabled, self.config.webadmin.allow_cors, self.admin_handler, - self.stopped, ) def initialize_received_payment_processor_worker(self): self.received_payment_processor_worker = ProcessReceivedPaymentsWorker( self.payment_processor, - self.stopped, ) - - -def start_admin_rpc_server(rpc_server): - logger.info("Starting admin RPC server...") - thread = threading.Thread( - target=rpc_server.serve, - args=(), - ) - thread.start() - - -def start_admin_web_server(admin_web_server): - logger.info("Starting admin web server...") - thread = threading.Thread( - target=admin_web_server.serve, - args=(), - ) - thread.start()