diff --git a/cmd/wnode/main.go b/cmd/wnode/main.go index 0f6c8511fc..80a7955404 100644 --- a/cmd/wnode/main.go +++ b/cmd/wnode/main.go @@ -22,7 +22,6 @@ package main import ( "bufio" "crypto/ecdsa" - crand "crypto/rand" "crypto/sha512" "encoding/binary" "encoding/hex" @@ -81,15 +80,14 @@ var ( // cmd arguments var ( - bootstrapMode = flag.Bool("standalone", false, "boostrap node: don't initiate connection to peers, just wait for incoming connections") - forwarderMode = flag.Bool("forwarder", false, "forwarder mode: only forward messages, neither encrypt nor decrypt messages") + bootstrapMode = flag.Bool("standalone", false, "boostrap node: don't actively connect to peers, wait for incoming connections") + forwarderMode = flag.Bool("forwarder", false, "forwarder mode: only forward messages, neither send nor decrypt messages") mailServerMode = flag.Bool("mailserver", false, "mail server mode: delivers expired messages on demand") requestMail = flag.Bool("mailclient", false, "request expired messages from the bootstrap server") asymmetricMode = flag.Bool("asym", false, "use asymmetric encryption") generateKey = flag.Bool("generatekey", false, "generate and show the private key") fileExMode = flag.Bool("fileexchange", false, "file exchange mode") - fileReader = flag.Bool("filereader", false, "load and decrypt messages saved as files, display as plain text") - testMode = flag.Bool("test", false, "use of predefined parameters for diagnostics (password, etc.)") + testMode = flag.Bool("test", false, "use of predefined parameters for diagnostics") echoMode = flag.Bool("echo", false, "echo mode: prints some arguments for diagnostics") argVerbosity = flag.Int("verbosity", int(log.LvlError), "log verbosity level") @@ -106,7 +104,7 @@ var ( argIDFile = flag.String("idfile", "", "file name with node id (private key)") argEnode = flag.String("boot", "", "bootstrap node you want to connect to (e.g. enode://e454......08d50@52.176.211.200:16428)") argTopic = flag.String("topic", "", "topic in hexadecimal format (e.g. 70a4beef)") - argSaveDir = flag.String("savedir", "", "directory where all incoming messages will be saved as files") + argSaveDir = flag.String("savedir", "", "directory where incoming messages will be saved as files") useLibP2P = flag.Bool("libp2p", false, "Use libp2p as the protocol layer") ) @@ -115,7 +113,6 @@ func main() { processArgs() initialize() run() - shutdown() } func processArgs() { @@ -204,8 +201,6 @@ func initialize() { if len(*argIP) == 0 { argIP = scanLineA("Please enter your IP and port (e.g. 127.0.0.1:30348): ") } - } else if *fileReader { - *bootstrapMode = true } else { if *argPort == 0 { for { @@ -249,6 +244,7 @@ func initialize() { cfg := &whisper.Config{ MaxMessageSize: uint32(*argMaxSize), MinimumAcceptedPOW: *argPoW, + UseDeadlines: *useLibP2P && *useDeadlines, } shh = whisper.New(cfg) diff --git a/whisper/whisperv6/config.go b/whisper/whisperv6/config.go index 61419de007..2416e60394 100644 --- a/whisper/whisperv6/config.go +++ b/whisper/whisperv6/config.go @@ -20,10 +20,12 @@ package whisperv6 type Config struct { MaxMessageSize uint32 `toml:",omitempty"` MinimumAcceptedPOW float64 `toml:",omitempty"` + UseDeadlines bool `toml:",omitempty"` } // DefaultConfig represents (shocker!) the default configuration. var DefaultConfig = Config{ MaxMessageSize: DefaultMaxMessageSize, MinimumAcceptedPOW: DefaultMinimumPoW, + UseDeadlines: true, } diff --git a/whisper/whisperv6/libp2p_glue.go b/whisper/whisperv6/libp2p_glue.go index 450708fb29..304e18919a 100644 --- a/whisper/whisperv6/libp2p_glue.go +++ b/whisper/whisperv6/libp2p_glue.go @@ -41,21 +41,36 @@ import ( // LibP2PStream is a wrapper used to implement the MsgReadWriter // interface for libp2p's streams. type LibP2PStream struct { - lp2pStream inet.Stream - rlpStream *rlp.Stream + UseDeadline bool + lp2pStream inet.Stream + rlpStream *rlp.Stream } -func newLibp2pStream(lp2pStream inet.Stream) p2p.MsgReadWriter { - return &LibP2PStream{lp2pStream: lp2pStream, rlpStream: rlp.NewStream(bufio.NewReader(lp2pStream), 0)} +func newLibp2pStream(server *LibP2PWhisperServer, lp2pStream inet.Stream) p2p.MsgReadWriter { + useDeadline, ok := server.whisper.settings.Load(useDeadlineIdx) + if !ok { + useDeadline = false + } + return &LibP2PStream{ + UseDeadline: useDeadline.(bool), + lp2pStream: lp2pStream, + rlpStream: rlp.NewStream(bufio.NewReader(lp2pStream), 0), + } } // ReadMsg implements the MsgReadWriter interface to read messages // from lilbp2p streams. func (stream *LibP2PStream) ReadMsg() (p2p.Msg, error) { + if stream.UseDeadline { + stream.lp2pStream.SetReadDeadline(time.Now().Add(expirationCycle)) + } msgcode, err := stream.rlpStream.Uint() if err != nil { return p2p.Msg{}, fmt.Errorf("can't read message code: %v", err) } + if stream.UseDeadline { + stream.lp2pStream.SetReadDeadline(time.Now().Add(expirationCycle)) + } _, size, err := stream.rlpStream.Kind() if err != nil { return p2p.Msg{}, fmt.Errorf("can't read message size: %v", err) @@ -78,14 +93,16 @@ func (stream *LibP2PStream) ReadMsg() (p2p.Msg, error) { // WriteMsg implements the MsgReadWriter interface to write messages // to lilbp2p streams. func (stream *LibP2PStream) WriteMsg(msg p2p.Msg) error { - stream.lp2pStream.SetWriteDeadline(time.Now().Add(transmissionCycle)) + if stream.UseDeadline { + stream.lp2pStream.SetWriteDeadline(time.Now().Add(expirationCycle)) + } if err := rlp.Encode(stream.lp2pStream, msg.Code); err != nil { return err } _, err := io.Copy(stream.lp2pStream, msg.Payload) - return err - } + return err +} // LibP2PPeer implements Peer for libp2p type LibP2PPeer struct { @@ -194,16 +211,16 @@ func (server *LibP2PWhisperServer) connectToPeer(p *LibP2PPeer) error { } // Save the stream - lps := newLibp2pStream(s).(*LibP2PStream) + lps := newLibp2pStream(server, s).(*LibP2PStream) p.connectionStream = lps p.ws = p.connectionStream // TODO send my known list of peers - + // Call HandlePeer to perform the handshake go server.whisper.HandlePeer(p, p.connectionStream) - + return err } @@ -211,7 +228,7 @@ func (server *LibP2PWhisperServer) connectToPeer(p *LibP2PPeer) error { func (server *LibP2PWhisperServer) Start() error { server.Host.SetStreamHandler(WhisperProtocolString, func(stream inet.Stream) { log.Info("opening stream from new peer") - + pid := stream.Conn().RemotePeer() var peer Peer for _, p := range server.Peers { @@ -221,7 +238,7 @@ func (server *LibP2PWhisperServer) Start() error { } } - lps := newLibp2pStream(stream).(*LibP2PStream) + lps := newLibp2pStream(server, stream).(*LibP2PStream) // Unknown peer if peer == nil { @@ -241,7 +258,7 @@ func (server *LibP2PWhisperServer) Start() error { if e := server.connectToPeer(p); e != nil { err = e } - } + } return err } diff --git a/whisper/whisperv6/whisper.go b/whisper/whisperv6/whisper.go index a0fdbf4025..92344f4e8b 100644 --- a/whisper/whisperv6/whisper.go +++ b/whisper/whisperv6/whisper.go @@ -54,6 +54,7 @@ const ( minPowToleranceIdx // Minimal PoW tolerated by the whisper node for a limited time bloomFilterIdx // Bloom filter for topics of interest for this node bloomFilterToleranceIdx // Bloom filter tolerated by the whisper node for a limited time + useDeadlineIdx // Indicates whether peer disconnect deadlines should be used (libp2p only) ) // WhisperServer abstracts a server, which could be either DevP2p-based @@ -121,6 +122,7 @@ func New(cfg *Config) *Whisper { whisper.settings.Store(minPowIdx, cfg.MinimumAcceptedPOW) whisper.settings.Store(maxMsgSizeIdx, cfg.MaxMessageSize) whisper.settings.Store(overflowIdx, false) + whisper.settings.Store(useDeadlineIdx, cfg.UseDeadlines) // p2p whisper sub protocol handler whisper.protocol = p2p.Protocol{