mirror of
https://github.com/yzernik/squeaknode.git
synced 2026-08-17 13:07:32 +02:00
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
This commit is contained in:
parent
31a90a6ceb
commit
9f3391a578
4 changed files with 35 additions and 50 deletions
|
|
@ -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)
|
||||
|
||||
|
|
|
|||
|
|
@ -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.")
|
||||
|
|
|
|||
|
|
@ -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()
|
||||
|
|
|
|||
|
|
@ -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()
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue