mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-08-18 18:02:24 +00:00
whisper: enable deadlines
Deadlines prevent wnode for working, so they are only enabled when not in wnode mode.
This commit is contained in:
parent
de2686ed10
commit
95d4f6f833
4 changed files with 39 additions and 22 deletions
|
|
@ -22,7 +22,6 @@ package main
|
||||||
import (
|
import (
|
||||||
"bufio"
|
"bufio"
|
||||||
"crypto/ecdsa"
|
"crypto/ecdsa"
|
||||||
crand "crypto/rand"
|
|
||||||
"crypto/sha512"
|
"crypto/sha512"
|
||||||
"encoding/binary"
|
"encoding/binary"
|
||||||
"encoding/hex"
|
"encoding/hex"
|
||||||
|
|
@ -81,15 +80,14 @@ var (
|
||||||
|
|
||||||
// cmd arguments
|
// cmd arguments
|
||||||
var (
|
var (
|
||||||
bootstrapMode = flag.Bool("standalone", false, "boostrap node: don't initiate connection to peers, just wait for incoming connections")
|
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 encrypt nor decrypt messages")
|
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")
|
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")
|
requestMail = flag.Bool("mailclient", false, "request expired messages from the bootstrap server")
|
||||||
asymmetricMode = flag.Bool("asym", false, "use asymmetric encryption")
|
asymmetricMode = flag.Bool("asym", false, "use asymmetric encryption")
|
||||||
generateKey = flag.Bool("generatekey", false, "generate and show the private key")
|
generateKey = flag.Bool("generatekey", false, "generate and show the private key")
|
||||||
fileExMode = flag.Bool("fileexchange", false, "file exchange mode")
|
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")
|
||||||
testMode = flag.Bool("test", false, "use of predefined parameters for diagnostics (password, etc.)")
|
|
||||||
echoMode = flag.Bool("echo", false, "echo mode: prints some arguments for diagnostics")
|
echoMode = flag.Bool("echo", false, "echo mode: prints some arguments for diagnostics")
|
||||||
|
|
||||||
argVerbosity = flag.Int("verbosity", int(log.LvlError), "log verbosity level")
|
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)")
|
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)")
|
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)")
|
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")
|
useLibP2P = flag.Bool("libp2p", false, "Use libp2p as the protocol layer")
|
||||||
)
|
)
|
||||||
|
|
@ -115,7 +113,6 @@ func main() {
|
||||||
processArgs()
|
processArgs()
|
||||||
initialize()
|
initialize()
|
||||||
run()
|
run()
|
||||||
shutdown()
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func processArgs() {
|
func processArgs() {
|
||||||
|
|
@ -204,8 +201,6 @@ func initialize() {
|
||||||
if len(*argIP) == 0 {
|
if len(*argIP) == 0 {
|
||||||
argIP = scanLineA("Please enter your IP and port (e.g. 127.0.0.1:30348): ")
|
argIP = scanLineA("Please enter your IP and port (e.g. 127.0.0.1:30348): ")
|
||||||
}
|
}
|
||||||
} else if *fileReader {
|
|
||||||
*bootstrapMode = true
|
|
||||||
} else {
|
} else {
|
||||||
if *argPort == 0 {
|
if *argPort == 0 {
|
||||||
for {
|
for {
|
||||||
|
|
@ -249,6 +244,7 @@ func initialize() {
|
||||||
cfg := &whisper.Config{
|
cfg := &whisper.Config{
|
||||||
MaxMessageSize: uint32(*argMaxSize),
|
MaxMessageSize: uint32(*argMaxSize),
|
||||||
MinimumAcceptedPOW: *argPoW,
|
MinimumAcceptedPOW: *argPoW,
|
||||||
|
UseDeadlines: *useLibP2P && *useDeadlines,
|
||||||
}
|
}
|
||||||
|
|
||||||
shh = whisper.New(cfg)
|
shh = whisper.New(cfg)
|
||||||
|
|
|
||||||
|
|
@ -20,10 +20,12 @@ package whisperv6
|
||||||
type Config struct {
|
type Config struct {
|
||||||
MaxMessageSize uint32 `toml:",omitempty"`
|
MaxMessageSize uint32 `toml:",omitempty"`
|
||||||
MinimumAcceptedPOW float64 `toml:",omitempty"`
|
MinimumAcceptedPOW float64 `toml:",omitempty"`
|
||||||
|
UseDeadlines bool `toml:",omitempty"`
|
||||||
}
|
}
|
||||||
|
|
||||||
// DefaultConfig represents (shocker!) the default configuration.
|
// DefaultConfig represents (shocker!) the default configuration.
|
||||||
var DefaultConfig = Config{
|
var DefaultConfig = Config{
|
||||||
MaxMessageSize: DefaultMaxMessageSize,
|
MaxMessageSize: DefaultMaxMessageSize,
|
||||||
MinimumAcceptedPOW: DefaultMinimumPoW,
|
MinimumAcceptedPOW: DefaultMinimumPoW,
|
||||||
|
UseDeadlines: true,
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -41,21 +41,36 @@ import (
|
||||||
// LibP2PStream is a wrapper used to implement the MsgReadWriter
|
// LibP2PStream is a wrapper used to implement the MsgReadWriter
|
||||||
// interface for libp2p's streams.
|
// interface for libp2p's streams.
|
||||||
type LibP2PStream struct {
|
type LibP2PStream struct {
|
||||||
lp2pStream inet.Stream
|
UseDeadline bool
|
||||||
rlpStream *rlp.Stream
|
lp2pStream inet.Stream
|
||||||
|
rlpStream *rlp.Stream
|
||||||
}
|
}
|
||||||
|
|
||||||
func newLibp2pStream(lp2pStream inet.Stream) p2p.MsgReadWriter {
|
func newLibp2pStream(server *LibP2PWhisperServer, lp2pStream inet.Stream) p2p.MsgReadWriter {
|
||||||
return &LibP2PStream{lp2pStream: lp2pStream, rlpStream: rlp.NewStream(bufio.NewReader(lp2pStream), 0)}
|
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
|
// ReadMsg implements the MsgReadWriter interface to read messages
|
||||||
// from lilbp2p streams.
|
// from lilbp2p streams.
|
||||||
func (stream *LibP2PStream) ReadMsg() (p2p.Msg, error) {
|
func (stream *LibP2PStream) ReadMsg() (p2p.Msg, error) {
|
||||||
|
if stream.UseDeadline {
|
||||||
|
stream.lp2pStream.SetReadDeadline(time.Now().Add(expirationCycle))
|
||||||
|
}
|
||||||
msgcode, err := stream.rlpStream.Uint()
|
msgcode, err := stream.rlpStream.Uint()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return p2p.Msg{}, fmt.Errorf("can't read message code: %v", err)
|
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()
|
_, size, err := stream.rlpStream.Kind()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return p2p.Msg{}, fmt.Errorf("can't read message size: %v", err)
|
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
|
// WriteMsg implements the MsgReadWriter interface to write messages
|
||||||
// to lilbp2p streams.
|
// to lilbp2p streams.
|
||||||
func (stream *LibP2PStream) WriteMsg(msg p2p.Msg) error {
|
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 {
|
if err := rlp.Encode(stream.lp2pStream, msg.Code); err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
_, err := io.Copy(stream.lp2pStream, msg.Payload)
|
_, err := io.Copy(stream.lp2pStream, msg.Payload)
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
// LibP2PPeer implements Peer for libp2p
|
// LibP2PPeer implements Peer for libp2p
|
||||||
type LibP2PPeer struct {
|
type LibP2PPeer struct {
|
||||||
|
|
@ -194,7 +211,7 @@ func (server *LibP2PWhisperServer) connectToPeer(p *LibP2PPeer) error {
|
||||||
}
|
}
|
||||||
|
|
||||||
// Save the stream
|
// Save the stream
|
||||||
lps := newLibp2pStream(s).(*LibP2PStream)
|
lps := newLibp2pStream(server, s).(*LibP2PStream)
|
||||||
|
|
||||||
p.connectionStream = lps
|
p.connectionStream = lps
|
||||||
p.ws = p.connectionStream
|
p.ws = p.connectionStream
|
||||||
|
|
@ -221,7 +238,7 @@ func (server *LibP2PWhisperServer) Start() error {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
lps := newLibp2pStream(stream).(*LibP2PStream)
|
lps := newLibp2pStream(server, stream).(*LibP2PStream)
|
||||||
|
|
||||||
// Unknown peer
|
// Unknown peer
|
||||||
if peer == nil {
|
if peer == nil {
|
||||||
|
|
@ -241,7 +258,7 @@ func (server *LibP2PWhisperServer) Start() error {
|
||||||
if e := server.connectToPeer(p); e != nil {
|
if e := server.connectToPeer(p); e != nil {
|
||||||
err = e
|
err = e
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -54,6 +54,7 @@ const (
|
||||||
minPowToleranceIdx // Minimal PoW tolerated by the whisper node for a limited time
|
minPowToleranceIdx // Minimal PoW tolerated by the whisper node for a limited time
|
||||||
bloomFilterIdx // Bloom filter for topics of interest for this node
|
bloomFilterIdx // Bloom filter for topics of interest for this node
|
||||||
bloomFilterToleranceIdx // Bloom filter tolerated by the whisper node for a limited time
|
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
|
// 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(minPowIdx, cfg.MinimumAcceptedPOW)
|
||||||
whisper.settings.Store(maxMsgSizeIdx, cfg.MaxMessageSize)
|
whisper.settings.Store(maxMsgSizeIdx, cfg.MaxMessageSize)
|
||||||
whisper.settings.Store(overflowIdx, false)
|
whisper.settings.Store(overflowIdx, false)
|
||||||
|
whisper.settings.Store(useDeadlineIdx, cfg.UseDeadlines)
|
||||||
|
|
||||||
// p2p whisper sub protocol handler
|
// p2p whisper sub protocol handler
|
||||||
whisper.protocol = p2p.Protocol{
|
whisper.protocol = p2p.Protocol{
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue