Do handshake in peer handler, not in connection class (#1326)

* Do handshake in peer handler, not in connection class

* Do handle connection synchronously

* Start connection context manager and handler in separate thread

* Update frontend build
This commit is contained in:
Jonathan Zernik 2021-09-17 22:03:00 -07:00 committed by GitHub
parent 97ac3c49cc
commit cc950eb173
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
13 changed files with 97 additions and 71 deletions

View file

@ -42,6 +42,7 @@ const portDefaultValue = '0';
export default function ConnectPeerDialog({
open,
handleClose,
handlePeerConnected,
...props
}) {
const classes = useStyles();
@ -80,6 +81,7 @@ export default function ConnectPeerDialog({
const connectPeer = (peerName, host, port) => {
connectSqueakPeerRequest(host, port, (response) => {
// goToPeerPage(history, response.getPeerId());
handlePeerConnected();
},
handleConnectPeerError);
};

View file

@ -317,6 +317,7 @@ export default function Peers() {
<ConnectPeerDialog
open={connectPeerDialogOpen}
handleClose={handleCloseConnectPeerDialog}
handlePeerConnected={getConnectedPeers}
/>
</>
);

View file

@ -1,14 +1,14 @@
{
"files": {
"main.js": "/static/js/main.c43d636b.chunk.js",
"main.js.map": "/static/js/main.c43d636b.chunk.js.map",
"main.js": "/static/js/main.f43898a0.chunk.js",
"main.js.map": "/static/js/main.f43898a0.chunk.js.map",
"runtime-main.js": "/static/js/runtime-main.9f0ba400.js",
"runtime-main.js.map": "/static/js/runtime-main.9f0ba400.js.map",
"static/css/2.ea4ba2f0.chunk.css": "/static/css/2.ea4ba2f0.chunk.css",
"static/js/2.4667e88e.chunk.js": "/static/js/2.4667e88e.chunk.js",
"static/js/2.4667e88e.chunk.js.map": "/static/js/2.4667e88e.chunk.js.map",
"index.html": "/index.html",
"precache-manifest.e1af39cbecc57b2d1cbd97f192625aa0.js": "/precache-manifest.e1af39cbecc57b2d1cbd97f192625aa0.js",
"precache-manifest.930ede73e84f1fad0f4bd9bc17ca9884.js": "/precache-manifest.930ede73e84f1fad0f4bd9bc17ca9884.js",
"service-worker.js": "/service-worker.js",
"static/css/2.ea4ba2f0.chunk.css.map": "/static/css/2.ea4ba2f0.chunk.css.map",
"static/js/2.4667e88e.chunk.js.LICENSE.txt": "/static/js/2.4667e88e.chunk.js.LICENSE.txt",
@ -20,6 +20,6 @@
"static/js/runtime-main.9f0ba400.js",
"static/css/2.ea4ba2f0.chunk.css",
"static/js/2.4667e88e.chunk.js",
"static/js/main.c43d636b.chunk.js"
"static/js/main.f43898a0.chunk.js"
]
}

View file

@ -1 +1 @@
<!doctype html><html lang="en"><head><meta charset="utf-8"/><link rel="shortcut icon" href="/favicon.ico"/><meta name="viewport" content="width=device-width,initial-scale=1,shrink-to-fit=no"/><meta name="theme-color" content="#000000"/><link rel="manifest" href="/manifest.json"/><title>Squeak Node</title><meta name="description" content="Squeak Node is a frontend for accessing a squeak node"><meta name="keywords" content="squeak, bitcoin, lightning"><meta name="author" content="Flatlogic LLC."><link href="/static/css/2.ea4ba2f0.chunk.css" rel="stylesheet"></head><body style="font-family:Roboto,sans-serif"><noscript>You need to enable JavaScript to run this app.</noscript><div id="root"></div><script>!function(e){function r(r){for(var n,f,l=r[0],a=r[1],i=r[2],c=0,s=[];c<l.length;c++)f=l[c],Object.prototype.hasOwnProperty.call(o,f)&&o[f]&&s.push(o[f][0]),o[f]=0;for(n in a)Object.prototype.hasOwnProperty.call(a,n)&&(e[n]=a[n]);for(p&&p(r);s.length;)s.shift()();return u.push.apply(u,i||[]),t()}function t(){for(var e,r=0;r<u.length;r++){for(var t=u[r],n=!0,l=1;l<t.length;l++){var a=t[l];0!==o[a]&&(n=!1)}n&&(u.splice(r--,1),e=f(f.s=t[0]))}return e}var n={},o={1:0},u=[];function f(r){if(n[r])return n[r].exports;var t=n[r]={i:r,l:!1,exports:{}};return e[r].call(t.exports,t,t.exports,f),t.l=!0,t.exports}f.m=e,f.c=n,f.d=function(e,r,t){f.o(e,r)||Object.defineProperty(e,r,{enumerable:!0,get:t})},f.r=function(e){"undefined"!=typeof Symbol&&Symbol.toStringTag&&Object.defineProperty(e,Symbol.toStringTag,{value:"Module"}),Object.defineProperty(e,"__esModule",{value:!0})},f.t=function(e,r){if(1&r&&(e=f(e)),8&r)return e;if(4&r&&"object"==typeof e&&e&&e.__esModule)return e;var t=Object.create(null);if(f.r(t),Object.defineProperty(t,"default",{enumerable:!0,value:e}),2&r&&"string"!=typeof e)for(var n in e)f.d(t,n,function(r){return e[r]}.bind(null,n));return t},f.n=function(e){var r=e&&e.__esModule?function(){return e.default}:function(){return e};return f.d(r,"a",r),r},f.o=function(e,r){return Object.prototype.hasOwnProperty.call(e,r)},f.p="/";var l=this["webpackJsonpsqueak-node-frontend"]=this["webpackJsonpsqueak-node-frontend"]||[],a=l.push.bind(l);l.push=r,l=l.slice();for(var i=0;i<l.length;i++)r(l[i]);var p=a;t()}([])</script><script src="/static/js/2.4667e88e.chunk.js"></script><script src="/static/js/main.c43d636b.chunk.js"></script></body></html>
<!doctype html><html lang="en"><head><meta charset="utf-8"/><link rel="shortcut icon" href="/favicon.ico"/><meta name="viewport" content="width=device-width,initial-scale=1,shrink-to-fit=no"/><meta name="theme-color" content="#000000"/><link rel="manifest" href="/manifest.json"/><title>Squeak Node</title><meta name="description" content="Squeak Node is a frontend for accessing a squeak node"><meta name="keywords" content="squeak, bitcoin, lightning"><meta name="author" content="Flatlogic LLC."><link href="/static/css/2.ea4ba2f0.chunk.css" rel="stylesheet"></head><body style="font-family:Roboto,sans-serif"><noscript>You need to enable JavaScript to run this app.</noscript><div id="root"></div><script>!function(e){function r(r){for(var n,f,l=r[0],a=r[1],i=r[2],c=0,s=[];c<l.length;c++)f=l[c],Object.prototype.hasOwnProperty.call(o,f)&&o[f]&&s.push(o[f][0]),o[f]=0;for(n in a)Object.prototype.hasOwnProperty.call(a,n)&&(e[n]=a[n]);for(p&&p(r);s.length;)s.shift()();return u.push.apply(u,i||[]),t()}function t(){for(var e,r=0;r<u.length;r++){for(var t=u[r],n=!0,l=1;l<t.length;l++){var a=t[l];0!==o[a]&&(n=!1)}n&&(u.splice(r--,1),e=f(f.s=t[0]))}return e}var n={},o={1:0},u=[];function f(r){if(n[r])return n[r].exports;var t=n[r]={i:r,l:!1,exports:{}};return e[r].call(t.exports,t,t.exports,f),t.l=!0,t.exports}f.m=e,f.c=n,f.d=function(e,r,t){f.o(e,r)||Object.defineProperty(e,r,{enumerable:!0,get:t})},f.r=function(e){"undefined"!=typeof Symbol&&Symbol.toStringTag&&Object.defineProperty(e,Symbol.toStringTag,{value:"Module"}),Object.defineProperty(e,"__esModule",{value:!0})},f.t=function(e,r){if(1&r&&(e=f(e)),8&r)return e;if(4&r&&"object"==typeof e&&e&&e.__esModule)return e;var t=Object.create(null);if(f.r(t),Object.defineProperty(t,"default",{enumerable:!0,value:e}),2&r&&"string"!=typeof e)for(var n in e)f.d(t,n,function(r){return e[r]}.bind(null,n));return t},f.n=function(e){var r=e&&e.__esModule?function(){return e.default}:function(){return e};return f.d(r,"a",r),r},f.o=function(e,r){return Object.prototype.hasOwnProperty.call(e,r)},f.p="/";var l=this["webpackJsonpsqueak-node-frontend"]=this["webpackJsonpsqueak-node-frontend"]||[],a=l.push.bind(l);l.push=r,l=l.slice();for(var i=0;i<l.length;i++)r(l[i]);var p=a;t()}([])</script><script src="/static/js/2.4667e88e.chunk.js"></script><script src="/static/js/main.f43898a0.chunk.js"></script></body></html>

View file

@ -1,6 +1,6 @@
self.__precacheManifest = (self.__precacheManifest || []).concat([
{
"revision": "7be392c085038d37b2dad89d648d10be",
"revision": "0076751edb4460adffbb7fc695e68f98",
"url": "/index.html"
},
{
@ -16,8 +16,8 @@ self.__precacheManifest = (self.__precacheManifest || []).concat([
"url": "/static/js/2.4667e88e.chunk.js.LICENSE.txt"
},
{
"revision": "0bc74bc27b6f5f3c9147",
"url": "/static/js/main.c43d636b.chunk.js"
"revision": "017707fdf1f5d031857e",
"url": "/static/js/main.f43898a0.chunk.js"
},
{
"revision": "cc9816de5a8639d377ea",

View file

@ -14,7 +14,7 @@
importScripts("https://storage.googleapis.com/workbox-cdn/releases/4.3.1/workbox-sw.js");
importScripts(
"/precache-manifest.e1af39cbecc57b2d1cbd97f192625aa0.js"
"/precache-manifest.930ede73e84f1fad0f4bd9bc17ca9884.js"
);
self.addEventListener('message', (event) => {

File diff suppressed because one or more lines are too long

File diff suppressed because one or more lines are too long

View file

@ -20,7 +20,6 @@
# OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE
# SOFTWARE.
import logging
import threading
from contextlib import contextmanager
from squeak.messages import msg_getaddr
@ -29,8 +28,6 @@ from squeak.messages import msg_subscribe
from squeaknode.network.peer import Peer
from squeaknode.network.peer_message_handler import PeerMessageHandler
HANDSHAKE_TIMEOUT = 30
logger = logging.getLogger(__name__)
@ -45,7 +42,6 @@ class Connection(object):
@contextmanager
def connect(self, connection_manager):
self.handshake()
logger.debug("Adding peer.")
connection_manager.add_peer(self.peer)
try:
@ -78,55 +74,9 @@ class Connection(object):
getaddr_msg = msg_getaddr()
self.peer.send_msg(getaddr_msg)
def handshake(self):
timer = HandshakeTimer(
self.peer.stop,
str(self),
)
timer.start_timer()
if self.peer.outgoing:
self.peer.send_version()
self.peer.receive_version()
if not self.peer.outgoing:
self.peer.send_version()
self.peer.set_connected()
logger.debug("HANDSHAKE COMPLETE-----------")
timer.stop_timer()
def handle_messages(self):
peer_message_handler = PeerMessageHandler(
self.peer,
self.squeak_controller,
)
peer_message_handler.handle_msgs()
class HandshakeTimer:
"""Stop the peer if handshake is not complete before timeout.
"""
def __init__(self,
stop_fn,
peer_name,
):
self.stop_fn = stop_fn
self.peer_name = peer_name
self.timer = None
def start_timer(self):
self.timer = threading.Timer(
HANDSHAKE_TIMEOUT,
self.stop_peer,
)
self.timer.name = "handshake_timere_thread_{}".format(self.peer_name)
self.timer.start()
def stop_timer(self):
logger.debug("Canceling handshake timer.")
self.timer.cancel()
def stop_peer(self):
logger.info("Closing peer from handshake timer.")
self.stop_fn()

View file

@ -122,7 +122,12 @@ class NetworkManager(object):
)
def subscribe_connected_peers(self, stopped):
yield from self.connection_manager.yield_peers_changed(stopped)
# yield from self.connection_manager.yield_peers_changed(stopped)
for item in self.connection_manager.yield_peers_changed(stopped):
logger.info("subscribe_connected_peers yielding item: {}".format(
item,
))
yield item
def subscribe_connected_peer(self, peer_address: PeerAddress, stopped):
yield from self.connection_manager.yield_single_peer_changed(

View file

@ -21,7 +21,6 @@
# SOFTWARE.
import logging
import socket
import threading
import socks
@ -77,11 +76,11 @@ class PeerClient(object):
outgoing: bool,
):
"""Handle a newly connected peer socket."""
threading.Thread(
target=self.peer_handler.handle_connection,
args=(peer_socket, peer_address, outgoing),
name="peer_client_connection_thread",
).start()
self.peer_handler.handle_connection(
peer_socket,
peer_address,
outgoing,
)
def get_socket(self):
if self.tor_proxy_ip:

View file

@ -21,6 +21,7 @@
# SOFTWARE.
import logging
import socket
import threading
from squeaknode.core.peer_address import PeerAddress
from squeaknode.network.connection import Connection
@ -30,6 +31,9 @@ from squeaknode.network.peer import Peer
logger = logging.getLogger(__name__)
HANDSHAKE_TIMEOUT = 30
class PeerHandler():
"""Handles new peer connection.
"""
@ -62,8 +66,43 @@ class PeerHandler():
self.connection_manager.single_peer_changed_listener,
)
try:
self.do_handshake(peer)
except Exception:
peer.stop()
raise
threading.Thread(
target=self.start_connection,
args=(peer,),
name="handle_peer_connection_thread",
).start()
def do_handshake(self, peer: Peer):
"""Do a handshake with a peer.
"""
timer = HandshakeTimer(
peer.stop,
str(self),
)
timer.start_timer()
if peer.outgoing:
peer.send_version()
# raise Exception("Fooooooo!")
peer.receive_version()
if not peer.outgoing:
peer.send_version()
peer.set_connected()
logger.debug("HANDSHAKE COMPLETE-----------")
timer.stop_timer()
def start_connection(self, peer: Peer):
"""Start a connection
"""
logger.debug(
'Setting up connection for peer address {} ...'.format(address))
'Setting up connection for peer {}'.format(peer))
try:
with Connection(peer, self.squeak_controller).connect(
self.connection_manager
@ -72,4 +111,34 @@ class PeerHandler():
finally:
peer.stop()
logger.debug(
'Stopped connection for peer address {}.'.format(address))
'Stopped connection for peer {}.'.format(peer),
)
class HandshakeTimer:
"""Stop the peer if handshake is not complete before timeout.
"""
def __init__(self,
stop_fn,
peer_name,
):
self.stop_fn = stop_fn
self.peer_name = peer_name
self.timer = None
def start_timer(self):
self.timer = threading.Timer(
HANDSHAKE_TIMEOUT,
self.stop_peer,
)
self.timer.name = "handshake_timere_thread_{}".format(self.peer_name)
self.timer.start()
def stop_timer(self):
logger.debug("Canceling handshake timer.")
self.timer.cancel()
def stop_peer(self):
logger.info("Closing peer from handshake timer.")
self.stop_fn()