Fix peer server thread shutdown (#903)

This commit is contained in:
Jonathan Zernik 2021-08-10 18:48:00 -07:00 committed by GitHub
parent 7d266cdf5e
commit 36d89bd25e
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
8 changed files with 53 additions and 10 deletions

View file

@ -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());
}
);
}

View file

@ -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

View file

@ -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):

View file

@ -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.")

View file

@ -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):

View file

@ -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(

View file

@ -22,6 +22,7 @@ class SqueakOfferExpiryWorker:
self.start_running,
)
timer.daemon = True
timer.name = "squeak_offer_expiry_thread"
timer.start()
self.remove_expired_offers()

View file

@ -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)