This commit is contained in:
Guillaume Ballet 2018-01-26 11:43:22 +00:00 committed by GitHub
commit 738e8cfc11
5 changed files with 316 additions and 267 deletions

View file

@ -21,6 +21,7 @@ package main
import ( import (
"bufio" "bufio"
"context"
"crypto/ecdsa" "crypto/ecdsa"
"crypto/sha512" "crypto/sha512"
"encoding/binary" "encoding/binary"
@ -39,24 +40,38 @@ import (
"github.com/ethereum/go-ethereum/console" "github.com/ethereum/go-ethereum/console"
"github.com/ethereum/go-ethereum/crypto" "github.com/ethereum/go-ethereum/crypto"
"github.com/ethereum/go-ethereum/log" "github.com/ethereum/go-ethereum/log"
"github.com/ethereum/go-ethereum/p2p"
"github.com/ethereum/go-ethereum/p2p/discover" "github.com/ethereum/go-ethereum/p2p/discover"
"github.com/ethereum/go-ethereum/p2p/nat"
"github.com/ethereum/go-ethereum/whisper/mailserver" "github.com/ethereum/go-ethereum/whisper/mailserver"
whisper "github.com/ethereum/go-ethereum/whisper/whisperv5" whisper "github.com/ethereum/go-ethereum/whisper/whisperv5"
"github.com/ethereum/go-ethereum/whisper/whisperv6"
"golang.org/x/crypto/pbkdf2" "golang.org/x/crypto/pbkdf2"
libp2pCrypto "github.com/libp2p/go-libp2p-crypto"
libp2pHost "github.com/libp2p/go-libp2p-host"
peer "github.com/libp2p/go-libp2p-peer"
ma "github.com/multiformats/go-multiaddr"
// inet "gx/ipfs/QmU4vCDZTPLDqSDKguWbHCiUe46mZUtmM2g2suBZ9NE8ko/go-libp2p-net"
// swarm "gx/ipfs/QmUhvp4VoQ9cKDVLqAxciEKdm8ymBx2Syx4C1Tv6SmSTPa/go-libp2p-swarm"
// peer "gx/ipfs/QmWNY7dV54ZDYmTA1ykVdwNCqC11mpU4zSUp6XDpLTH9eG/go-libp2p-peer"
// ma "gx/ipfs/QmW8s4zTsUoX1Q6CeYxVKPyqSKbF7H1YDUyTostBtZ8DaG/go-multiaddr"
inet "github.com/libp2p/go-libp2p-net"
) )
const quitCommand = "~Q" const quitCommand = "~Q"
// singletons // singletons
var ( var (
server *p2p.Server network *inet.Network
host libp2pHost.Host
shh *whisper.Whisper shh *whisper.Whisper
done chan struct{} done chan struct{}
mailServer mailserver.WMailServer mailServer mailserver.WMailServer
input = bufio.NewReader(os.Stdin) input = bufio.NewReader(os.Stdin)
bootnode peer.ID
) )
// encryption // encryption
@ -64,7 +79,7 @@ var (
symKey []byte symKey []byte
pub *ecdsa.PublicKey pub *ecdsa.PublicKey
asymKey *ecdsa.PrivateKey asymKey *ecdsa.PrivateKey
nodeid *ecdsa.PrivateKey nodeid libp2pCrypto.PrivKey
topic whisper.TopicType topic whisper.TopicType
asymKeyID string asymKeyID string
filterID string filterID string
@ -91,7 +106,6 @@ var (
argPoW = flag.Float64("pow", whisper.DefaultMinimumPoW, "PoW for normal messages in float format (e.g. 2.7)") argPoW = flag.Float64("pow", whisper.DefaultMinimumPoW, "PoW for normal messages in float format (e.g. 2.7)")
argServerPoW = flag.Float64("mspow", whisper.DefaultMinimumPoW, "PoW requirement for Mail Server request") argServerPoW = flag.Float64("mspow", whisper.DefaultMinimumPoW, "PoW requirement for Mail Server request")
argIP = flag.String("ip", "", "IP address and port of this node (e.g. 127.0.0.1:30303)")
argPub = flag.String("pub", "", "public key for asymmetric encryption") argPub = flag.String("pub", "", "public key for asymmetric encryption")
argDBPath = flag.String("dbpath", "", "path to the server's DB directory") argDBPath = flag.String("dbpath", "", "path to the server's DB directory")
argIDFile = flag.String("idfile", "", "file name with node id (private key)") argIDFile = flag.String("idfile", "", "file name with node id (private key)")
@ -111,16 +125,14 @@ func processArgs() {
if len(*argIDFile) > 0 { if len(*argIDFile) > 0 {
var err error var err error
nodeid, err = crypto.LoadECDSA(*argIDFile) keyData, err := ioutil.ReadFile(*argIDFile)
if err != nil { if err != nil {
utils.Fatalf("Failed to load file [%s]: %s.", *argIDFile, err) utils.Fatalf("Failed to load file [%s]: %s.", *argIDFile, err)
} }
}
const enodePrefix = "enode://" nodeid, err = libp2pCrypto.UnmarshalPrivateKey(keyData)
if len(*argEnode) > 0 { if err != nil {
if (*argEnode)[:len(enodePrefix)] != enodePrefix { utils.Fatalf("Failed to load file [%s]: %s.", *argIDFile, err)
*argEnode = enodePrefix + *argEnode
} }
} }
@ -157,7 +169,7 @@ func echo() {
fmt.Printf("workTime = %d \n", *argWorkTime) fmt.Printf("workTime = %d \n", *argWorkTime)
fmt.Printf("pow = %f \n", *argPoW) fmt.Printf("pow = %f \n", *argPoW)
fmt.Printf("mspow = %f \n", *argServerPoW) fmt.Printf("mspow = %f \n", *argServerPoW)
fmt.Printf("ip = %s \n", *argIP) fmt.Printf("ip = %v \n", host.Addrs()[0].Protocols())
fmt.Printf("pub = %s \n", common.ToHex(crypto.FromECDSAPub(pub))) fmt.Printf("pub = %s \n", common.ToHex(crypto.FromECDSAPub(pub)))
fmt.Printf("idfile = %s \n", *argIDFile) fmt.Printf("idfile = %s \n", *argIDFile)
fmt.Printf("dbpath = %s \n", *argDBPath) fmt.Printf("dbpath = %s \n", *argDBPath)
@ -168,7 +180,6 @@ func initialize() {
log.Root().SetHandler(log.LvlFilterHandler(log.Lvl(*argVerbosity), log.StreamHandler(os.Stderr, log.TerminalFormat(false)))) log.Root().SetHandler(log.LvlFilterHandler(log.Lvl(*argVerbosity), log.StreamHandler(os.Stderr, log.TerminalFormat(false))))
done = make(chan struct{}) done = make(chan struct{})
var peers []*discover.Node
var err error var err error
if *generateKey { if *generateKey {
@ -186,16 +197,23 @@ func initialize() {
msPassword = "wwww" msPassword = "wwww"
} }
if *bootstrapMode { if !*bootstrapMode {
if len(*argIP) == 0 {
argIP = scanLineA("Please enter your IP and port (e.g. 127.0.0.1:30348): ")
}
} else {
if len(*argEnode) == 0 { if len(*argEnode) == 0 {
argEnode = scanLineA("Please enter the peer's enode: ") argEnode = scanLineA("Please enter the peer's address: ")
} }
peer := discover.MustParseNode(*argEnode) peeraddr, err := ma.NewMultiaddr(*argEnode)
peers = append(peers, peer) if err != nil {
utils.Fatalf("Invalid address: %s", err)
}
pid, err := peeraddr.ValueForProtocol(ma.P_IPFS)
if err != nil {
utils.Fatalf("Could not find peer id: %s", err)
}
peerid, err := peer.IDB58Decode(pid)
if err != nil {
utils.Fatalf("Could not decode peer id: %s", err)
}
bootnode = peerid
} }
cfg := &whisper.Config{ cfg := &whisper.Config{
@ -242,57 +260,31 @@ func initialize() {
utils.Fatalf("Failed to retrieve a new key pair: %s", err) utils.Fatalf("Failed to retrieve a new key pair: %s", err)
} }
if nodeid == nil { if nodeid != nil {
tmpID, err := shh.NewKeyPair() nodeid, _, err = libp2pCrypto.GenerateKeyPair(libp2pCrypto.Ed25519, 384)
if err != nil {
utils.Fatalf("Failed to generate a new key pair: %s", err)
}
nodeid, err = shh.GetPrivateKey(tmpID)
if err != nil { if err != nil {
utils.Fatalf("Failed to retrieve a new key pair: %s", err) utils.Fatalf("Failed to retrieve a new key pair: %s", err)
} }
} }
maxPeers := 80
if *bootstrapMode {
maxPeers = 800
}
server = &p2p.Server{
Config: p2p.Config{
PrivateKey: nodeid,
MaxPeers: maxPeers,
Name: common.MakeName("wnode", "5.0"),
Protocols: shh.Protocols(),
ListenAddr: *argIP,
NAT: nat.Any(),
BootstrapNodes: peers,
StaticNodes: peers,
TrustedNodes: peers,
},
}
} }
func startServer() { func startServer() {
err := server.Start()
if err != nil {
utils.Fatalf("Failed to start Whisper peer: %s.", err)
}
fmt.Printf("my public key: %s \n", common.ToHex(crypto.FromECDSAPub(&asymKey.PublicKey))) fmt.Printf("my public key: %s \n", common.ToHex(crypto.FromECDSAPub(&asymKey.PublicKey)))
fmt.Println(server.NodeInfo().Enode) fmt.Println(host.ID())
if *bootstrapMode { if *bootstrapMode {
configureNode() fmt.Print("Bootstrap ")
fmt.Println("Bootstrap Whisper node started")
} else { } else {
fmt.Println("Whisper node started") // first see if we can establish connection
// first see if we can establish connection, then ask for user input _, err := host.NewStream(context.Background(), bootnode, whisperv6.WhisperProtocolString)
waitForConnection(true) if err != nil {
configureNode() utils.Fatalf("Error connecting to bootstrap node: %s", err)
} }
}
fmt.Println("Whisper node started")
configureNode()
if !*forwarderMode { if !*forwarderMode {
fmt.Printf("Please type the message. To quit type: '%s'\n", quitCommand) fmt.Printf("Please type the message. To quit type: '%s'\n", quitCommand)
} }
@ -382,27 +374,10 @@ func generateTopic(password []byte) {
} }
} }
func waitForConnection(timeout bool) {
var cnt int
var connected bool
for !connected {
time.Sleep(time.Millisecond * 50)
connected = server.PeerCount() > 0
if timeout {
cnt++
if cnt > 1000 {
utils.Fatalf("Timeout expired, failed to connect")
}
}
}
fmt.Println("Connected to peer.")
}
func run() { func run() {
defer mailServer.Close() defer mailServer.Close()
startServer() startServer()
defer server.Stop() defer host.Close()
shh.Start(nil) shh.Start(nil)
defer shh.Stop() defer shh.Stop()
@ -601,7 +576,7 @@ func requestExpiredMessagesLoop() {
if err != nil { if err != nil {
utils.Fatalf("Failed to save symmetric key for mail request: %s", err) utils.Fatalf("Failed to save symmetric key for mail request: %s", err)
} }
peerID = extractIdFromEnode(*argEnode) peerID = extractIDFromEnode(*argEnode)
shh.AllowP2PMessagesFromPeer(peerID) shh.AllowP2PMessagesFromPeer(peerID)
for { for {
@ -627,11 +602,19 @@ func requestExpiredMessagesLoop() {
data = data[:8] data = data[:8]
} }
nodeiddata, err := nodeid.Bytes()
if err != nil {
utils.Fatalf("Error converting libp2p key to geth key: %s", err)
}
nodeidkey, err := crypto.ToECDSA(nodeiddata)
if err != nil {
utils.Fatalf("Error converting libp2p key to geth key: %s", err)
}
var params whisper.MessageParams var params whisper.MessageParams
params.PoW = *argServerPoW params.PoW = *argServerPoW
params.Payload = data params.Payload = data
params.KeySym = key params.KeySym = key
params.Src = nodeid params.Src = nodeidkey
params.WorkTime = 5 params.WorkTime = 5
msg, err := whisper.NewSentMessage(&params) msg, err := whisper.NewSentMessage(&params)
@ -652,7 +635,7 @@ func requestExpiredMessagesLoop() {
} }
} }
func extractIdFromEnode(s string) []byte { func extractIDFromEnode(s string) []byte {
n, err := discover.ParseNode(s) n, err := discover.ParseNode(s)
if err != nil { if err != nil {
utils.Fatalf("Failed to parse enode: %s", err) utils.Fatalf("Failed to parse enode: %s", err)

View file

@ -21,6 +21,7 @@ import (
"crypto/ecdsa" "crypto/ecdsa"
"errors" "errors"
"fmt" "fmt"
"strings"
"sync" "sync"
"time" "time"
@ -28,8 +29,9 @@ import (
"github.com/ethereum/go-ethereum/common/hexutil" "github.com/ethereum/go-ethereum/common/hexutil"
"github.com/ethereum/go-ethereum/crypto" "github.com/ethereum/go-ethereum/crypto"
"github.com/ethereum/go-ethereum/log" "github.com/ethereum/go-ethereum/log"
"github.com/ethereum/go-ethereum/p2p/discover"
"github.com/ethereum/go-ethereum/rpc" "github.com/ethereum/go-ethereum/rpc"
peer "github.com/libp2p/go-libp2p-peer"
) )
const ( const (
@ -126,14 +128,32 @@ func (api *PublicWhisperAPI) SetBloomFilter(ctx context.Context, bloom hexutil.B
return true, api.w.SetBloomFilter(bloom) return true, api.w.SetBloomFilter(bloom)
} }
// This is a helper for turning an address into a peer.ID. It
// is temporary as there seems to be an official way to to it
// with `Multiaddres`.
func peerIDFromAddress(addr string) (peer.ID, error) {
comps := strings.Split(addr, "/")
var pid peer.ID
var err error
if comps[len(comps)-1] == "" && len(comps) > 1 {
pid, err = peer.IDFromString(comps[len(comps)-2])
} else {
pid, err = peer.IDFromString(comps[len(comps)-1])
}
if err != nil {
log.Error(fmt.Sprintf("Error getting id from enode: %s", err))
}
return pid, err
}
// MarkTrustedPeer marks a peer trusted, which will allow it to send historic (expired) messages. // MarkTrustedPeer marks a peer trusted, which will allow it to send historic (expired) messages.
// Note: This function is not adding new nodes, the node needs to exists as a peer. // Note: This function is not adding new nodes, the node needs to exists as a peer.
func (api *PublicWhisperAPI) MarkTrustedPeer(ctx context.Context, enode string) (bool, error) { func (api *PublicWhisperAPI) MarkTrustedPeer(ctx context.Context, url string) (bool, error) {
n, err := discover.ParseNode(enode) pid, err := peerIDFromAddress(url)
if err != nil {
return false, err return err == nil, api.w.AllowP2PMessagesFromPeer(pid)
}
return true, api.w.AllowP2PMessagesFromPeer(n.ID[:])
} }
// NewKeyPair generates a new public and private key pair for message decryption and encryption. // NewKeyPair generates a new public and private key pair for message decryption and encryption.
@ -303,11 +323,11 @@ func (api *PublicWhisperAPI) Post(ctx context.Context, req NewMessage) (bool, er
// send to specific node (skip PoW check) // send to specific node (skip PoW check)
if len(req.TargetPeer) > 0 { if len(req.TargetPeer) > 0 {
n, err := discover.ParseNode(req.TargetPeer) pid, err := peerIDFromAddress(req.TargetPeer)
if err != nil { if err != nil {
return false, fmt.Errorf("failed to parse target peer: %s", err) return false, fmt.Errorf("failed to parse target peer: %s", err)
} }
return true, api.w.SendP2PMessage(n.ID[:], env) return true, api.w.SendP2PMessage(pid, env)
} }
// ensure that the message PoW meets the node's minimum accepted PoW // ensure that the message PoW meets the node's minimum accepted PoW

View file

@ -72,6 +72,9 @@ const (
DefaultTTL = 50 // seconds DefaultTTL = 50 // seconds
DefaultSyncAllowance = 10 // seconds DefaultSyncAllowance = 10 // seconds
WhisperPort = 534848
WhisperProtocolString = "/whisper/6.0"
) )
type unknownVersionError uint64 type unknownVersionError uint64

View file

@ -26,12 +26,14 @@ import (
"github.com/ethereum/go-ethereum/p2p" "github.com/ethereum/go-ethereum/p2p"
"github.com/ethereum/go-ethereum/rlp" "github.com/ethereum/go-ethereum/rlp"
set "gopkg.in/fatih/set.v0" set "gopkg.in/fatih/set.v0"
libp2pPeer "github.com/libp2p/go-libp2p-peer"
) )
// peer represents a whisper protocol peer connection. // peer represents a whisper protocol peer connection.
type Peer struct { type Peer struct {
host *Whisper host *Whisper
peer *p2p.Peer pid libp2pPeer.ID
ws p2p.MsgReadWriter ws p2p.MsgReadWriter
trusted bool trusted bool
@ -44,10 +46,10 @@ type Peer struct {
} }
// newPeer creates a new whisper peer object, but does not run the handshake itself. // newPeer creates a new whisper peer object, but does not run the handshake itself.
func newPeer(host *Whisper, remote *p2p.Peer, rw p2p.MsgReadWriter) *Peer { func newPeer(host *Whisper, remote *libp2pPeer.ID, rw p2p.MsgReadWriter) *Peer {
return &Peer{ return &Peer{
host: host, host: host,
peer: remote, pid: *remote,
ws: rw, ws: rw,
trusted: false, trusted: false,
powRequirement: 0.0, powRequirement: 0.0,
@ -60,13 +62,13 @@ func newPeer(host *Whisper, remote *p2p.Peer, rw p2p.MsgReadWriter) *Peer {
// into the network. // into the network.
func (p *Peer) start() { func (p *Peer) start() {
go p.update() go p.update()
log.Trace("start", "peer", p.ID()) log.Trace("start", "peer", p.pid)
} }
// stop terminates the peer updater, stopping message forwarding to it. // stop terminates the peer updater, stopping message forwarding to it.
func (p *Peer) stop() { func (p *Peer) stop() {
close(p.quit) close(p.quit)
log.Trace("stop", "peer", p.ID()) log.Trace("stop", "peer", p.pid)
} }
// handshake sends the protocol initiation status message to the remote peer and // handshake sends the protocol initiation status message to the remote peer and
@ -87,19 +89,19 @@ func (p *Peer) handshake() error {
return err return err
} }
if packet.Code != statusCode { if packet.Code != statusCode {
return fmt.Errorf("peer [%x] sent packet %x before status packet", p.ID(), packet.Code) return fmt.Errorf("peer [%x] sent packet %x before status packet", p.pid, packet.Code)
} }
s := rlp.NewStream(packet.Payload, uint64(packet.Size)) s := rlp.NewStream(packet.Payload, uint64(packet.Size))
_, err = s.List() _, err = s.List()
if err != nil { if err != nil {
return fmt.Errorf("peer [%x] sent bad status message: %v", p.ID(), err) return fmt.Errorf("peer [%x] sent bad status message: %v", p.pid, err)
} }
peerVersion, err := s.Uint() peerVersion, err := s.Uint()
if err != nil { if err != nil {
return fmt.Errorf("peer [%x] sent bad status message (unable to decode version): %v", p.ID(), err) return fmt.Errorf("peer [%x] sent bad status message (unable to decode version): %v", p.pid, err)
} }
if peerVersion != ProtocolVersion { if peerVersion != ProtocolVersion {
return fmt.Errorf("peer [%x]: protocol version mismatch %d != %d", p.ID(), peerVersion, ProtocolVersion) return fmt.Errorf("peer [%x]: protocol version mismatch %d != %d", p.pid, peerVersion, ProtocolVersion)
} }
// only version is mandatory, subsequent parameters are optional // only version is mandatory, subsequent parameters are optional
@ -107,7 +109,7 @@ func (p *Peer) handshake() error {
if err == nil { if err == nil {
pow := math.Float64frombits(powRaw) pow := math.Float64frombits(powRaw)
if math.IsInf(pow, 0) || math.IsNaN(pow) || pow < 0.0 { if math.IsInf(pow, 0) || math.IsNaN(pow) || pow < 0.0 {
return fmt.Errorf("peer [%x] sent bad status message: invalid pow", p.ID()) return fmt.Errorf("peer [%x] sent bad status message: invalid pow", p.pid)
} }
p.powRequirement = pow p.powRequirement = pow
@ -116,7 +118,7 @@ func (p *Peer) handshake() error {
if err == nil { if err == nil {
sz := len(bloom) sz := len(bloom)
if sz != bloomFilterSize && sz != 0 { if sz != bloomFilterSize && sz != 0 {
return fmt.Errorf("peer [%x] sent bad status message: wrong bloom filter size %d", p.ID(), sz) return fmt.Errorf("peer [%x] sent bad status message: wrong bloom filter size %d", p.pid, sz)
} }
if isFullNode(bloom) { if isFullNode(bloom) {
p.bloomFilter = nil p.bloomFilter = nil
@ -127,7 +129,7 @@ func (p *Peer) handshake() error {
} }
if err := <-errc; err != nil { if err := <-errc; err != nil {
return fmt.Errorf("peer [%x] failed to send status packet: %v", p.ID(), err) return fmt.Errorf("peer [%x] failed to send status packet: %v", p.pid, err)
} }
return nil return nil
} }
@ -147,7 +149,7 @@ func (p *Peer) update() {
case <-transmit.C: case <-transmit.C:
if err := p.broadcast(); err != nil { if err := p.broadcast(); err != nil {
log.Trace("broadcast failed", "reason", err, "peer", p.ID()) log.Trace("broadcast failed", "reason", err, "peer", p.pid)
return return
} }
@ -210,11 +212,6 @@ func (p *Peer) broadcast() error {
return nil return nil
} }
func (p *Peer) ID() []byte {
id := p.peer.ID()
return id[:]
}
func (p *Peer) notifyAboutPowRequirementChange(pow float64) error { func (p *Peer) notifyAboutPowRequirementChange(pow float64) error {
i := math.Float64bits(pow) i := math.Float64bits(pow)
return p2p.Send(p.ws, powRequirementCode, i) return p2p.Send(p.ws, powRequirementCode, i)

View file

@ -17,11 +17,12 @@
package whisperv6 package whisperv6
import ( import (
"bytes" "context"
"crypto/ecdsa" "crypto/ecdsa"
crand "crypto/rand" crand "crypto/rand"
"crypto/sha256" "crypto/sha256"
"fmt" "fmt"
"io/ioutil"
"math" "math"
"runtime" "runtime"
"sync" "sync"
@ -33,6 +34,12 @@ import (
"github.com/ethereum/go-ethereum/p2p" "github.com/ethereum/go-ethereum/p2p"
"github.com/ethereum/go-ethereum/rlp" "github.com/ethereum/go-ethereum/rlp"
"github.com/ethereum/go-ethereum/rpc" "github.com/ethereum/go-ethereum/rpc"
libp2p "github.com/libp2p/go-libp2p"
libp2pCrypto "github.com/libp2p/go-libp2p-crypto"
host "github.com/libp2p/go-libp2p-host"
inet "github.com/libp2p/go-libp2p-net"
peer "github.com/libp2p/go-libp2p-peer"
protocol "github.com/libp2p/go-libp2p-protocol"
"github.com/syndtr/goleveldb/leveldb/errors" "github.com/syndtr/goleveldb/leveldb/errors"
"golang.org/x/crypto/pbkdf2" "golang.org/x/crypto/pbkdf2"
"golang.org/x/sync/syncmap" "golang.org/x/sync/syncmap"
@ -59,7 +66,8 @@ const (
// Whisper represents a dark communication interface through the Ethereum // Whisper represents a dark communication interface through the Ethereum
// network, using its very own P2P communication layer. // network, using its very own P2P communication layer.
type Whisper struct { type Whisper struct {
protocol p2p.Protocol // Protocol description and parameters protocol protocol.ID // Protocol description and parameters
host host.Host
filters *Filters // Message filters installed with Subscribe function filters *Filters // Message filters installed with Subscribe function
privateKeys map[string]*ecdsa.PrivateKey // Private key storage privateKeys map[string]*ecdsa.PrivateKey // Private key storage
@ -93,6 +101,28 @@ func New(cfg *Config) *Whisper {
cfg = &DefaultConfig cfg = &DefaultConfig
} }
priv, _, err := libp2pCrypto.GenerateKeyPair(libp2pCrypto.Ed25519, 384)
// TODO check if this doesn't reveal my private key, which would
// be very bad. Alternatively, use the crypto stuff from libp2p,
// hopefully this can be converted back and forth with the golang
// crypto primitives.
// pid, _ := peer.IDFromString(fmt.Sprintf("%v", nodeid.D))
// n, _ := swarm.NewNetwork(context.Background(), []ma.Multiaddr{}, pid, nil, nil)
// fmt.Println(n)
opts := []libp2p.Option{
libp2p.ListenAddrStrings(fmt.Sprintf("/ip4/0.0.0.0/tcp/%d", WhisperPort)),
libp2p.Identity(priv),
}
h, err := libp2p.New(context.Background(), opts...)
if err != nil {
// TODO return error too
log.Error("Error setting up the libp2p network: %s", err)
return nil
}
log.Info("Host address is at %s", h.Addrs()[0].String())
whisper := &Whisper{ whisper := &Whisper{
privateKeys: make(map[string]*ecdsa.PrivateKey), privateKeys: make(map[string]*ecdsa.PrivateKey),
symKeys: make(map[string][]byte), symKeys: make(map[string][]byte),
@ -103,29 +133,45 @@ func New(cfg *Config) *Whisper {
p2pMsgQueue: make(chan *Envelope, messageQueueLimit), p2pMsgQueue: make(chan *Envelope, messageQueueLimit),
quit: make(chan struct{}), quit: make(chan struct{}),
syncAllowance: DefaultSyncAllowance, syncAllowance: DefaultSyncAllowance,
host: h,
protocol: WhisperProtocolString,
} }
h.SetStreamHandler(whisper.protocol, func(stream inet.Stream) {
defer stream.Close()
data, err := ioutil.ReadAll(stream)
if err != nil {
log.Error("Error reading stream data %s", err)
return
}
if len(data) < 2 {
log.Error("Invalid data received: it has to be at least two bytes")
return
}
messageType := data[0]
var envelope Envelope
err = rlp.DecodeBytes(data[1:], envelope)
if err != nil {
log.Error(fmt.Sprintf("Error decoding payload: %s", err))
return
}
p, err := whisper.getPeer(stream.Conn().RemotePeer())
if err != nil {
log.Error(fmt.Sprintf("Could not find peer: %s", err))
return
}
whisper.ProcessIncomingMessage(messageType, data[1:], p)
})
whisper.filters = NewFilters(whisper) whisper.filters = NewFilters(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)
// p2p whisper sub protocol handler
whisper.protocol = p2p.Protocol{
Name: ProtocolName,
Version: uint(ProtocolVersion),
Length: NumberOfMessageCodes,
Run: whisper.HandlePeer,
NodeInfo: func() interface{} {
return map[string]interface{}{
"version": ProtocolVersionStr,
"maxMessageSize": whisper.MaxMessageSize(),
"minimumPoW": whisper.MinPow(),
}
},
}
return whisper return whisper
} }
@ -208,16 +254,6 @@ func (w *Whisper) RegisterServer(server MailServer) {
w.mailServer = server w.mailServer = server
} }
// Protocols returns the whisper sub-protocols ran by this particular client.
func (w *Whisper) Protocols() []p2p.Protocol {
return []p2p.Protocol{w.protocol}
}
// Version returns the whisper sub-protocols version number.
func (w *Whisper) Version() uint {
return w.protocol.Version
}
// SetMaxMessageSize sets the maximal message size allowed by this node // SetMaxMessageSize sets the maximal message size allowed by this node
func (w *Whisper) SetMaxMessageSize(size uint32) error { func (w *Whisper) SetMaxMessageSize(size uint32) error {
if size > MaxMessageSize { if size > MaxMessageSize {
@ -282,7 +318,7 @@ func (w *Whisper) notifyPeersAboutPowRequirementChange(pow float64) {
err = p.notifyAboutPowRequirementChange(pow) err = p.notifyAboutPowRequirementChange(pow)
} }
if err != nil { if err != nil {
log.Warn("failed to notify peer about new pow requirement", "peer", p.ID(), "error", err) log.Warn("failed to notify peer about new pow requirement", "peer", p.pid, "error", err)
} }
} }
} }
@ -296,7 +332,7 @@ func (w *Whisper) notifyPeersAboutBloomFilterChange(bloom []byte) {
err = p.notifyAboutBloomFilterChange(bloom) err = p.notifyAboutBloomFilterChange(bloom)
} }
if err != nil { if err != nil {
log.Warn("failed to notify peer about new bloom filter", "peer", p.ID(), "error", err) log.Warn("failed to notify peer about new bloom filter", "peer", p.pid, "error", err)
} }
} }
} }
@ -314,12 +350,11 @@ func (w *Whisper) getPeers() []*Peer {
} }
// getPeer retrieves peer by ID // getPeer retrieves peer by ID
func (w *Whisper) getPeer(peerID []byte) (*Peer, error) { func (w *Whisper) getPeer(peerID peer.ID) (*Peer, error) {
w.peerMu.Lock() w.peerMu.Lock()
defer w.peerMu.Unlock() defer w.peerMu.Unlock()
for p := range w.peers { for p := range w.peers {
id := p.peer.ID() if peerID == p.pid {
if bytes.Equal(peerID, id[:]) {
return p, nil return p, nil
} }
} }
@ -328,7 +363,7 @@ func (w *Whisper) getPeer(peerID []byte) (*Peer, error) {
// AllowP2PMessagesFromPeer marks specific peer trusted, // AllowP2PMessagesFromPeer marks specific peer trusted,
// which will allow it to send historic (expired) messages. // which will allow it to send historic (expired) messages.
func (w *Whisper) AllowP2PMessagesFromPeer(peerID []byte) error { func (w *Whisper) AllowP2PMessagesFromPeer(peerID peer.ID) error {
p, err := w.getPeer(peerID) p, err := w.getPeer(peerID)
if err != nil { if err != nil {
return err return err
@ -342,7 +377,7 @@ func (w *Whisper) AllowP2PMessagesFromPeer(peerID []byte) error {
// request and respond with a number of peer-to-peer messages (possibly expired), // request and respond with a number of peer-to-peer messages (possibly expired),
// which are not supposed to be forwarded any further. // which are not supposed to be forwarded any further.
// The whisper protocol is agnostic of the format and contents of envelope. // The whisper protocol is agnostic of the format and contents of envelope.
func (w *Whisper) RequestHistoricMessages(peerID []byte, envelope *Envelope) error { func (w *Whisper) RequestHistoricMessages(peerID peer.ID, envelope *Envelope) error {
p, err := w.getPeer(peerID) p, err := w.getPeer(peerID)
if err != nil { if err != nil {
return err return err
@ -352,7 +387,7 @@ func (w *Whisper) RequestHistoricMessages(peerID []byte, envelope *Envelope) err
} }
// SendP2PMessage sends a peer-to-peer message to a specific peer. // SendP2PMessage sends a peer-to-peer message to a specific peer.
func (w *Whisper) SendP2PMessage(peerID []byte, envelope *Envelope) error { func (w *Whisper) SendP2PMessage(peerID peer.ID, envelope *Envelope) error {
p, err := w.getPeer(peerID) p, err := w.getPeer(peerID)
if err != nil { if err != nil {
return err return err
@ -362,7 +397,19 @@ func (w *Whisper) SendP2PMessage(peerID []byte, envelope *Envelope) error {
// SendP2PDirect sends a peer-to-peer message to a specific peer. // SendP2PDirect sends a peer-to-peer message to a specific peer.
func (w *Whisper) SendP2PDirect(peer *Peer, envelope *Envelope) error { func (w *Whisper) SendP2PDirect(peer *Peer, envelope *Envelope) error {
return p2p.Send(peer.ws, p2pMessageCode, envelope) stream, err := w.host.NewStream(context.Background(), peer.pid, w.protocol)
if err != nil {
return err
}
defer stream.Close()
b, err := rlp.EncodeToBytes(envelope)
if err != nil {
return err
}
stream.Write(append([]byte{p2pMessageCode}, b...))
return nil
} }
// NewKeyPair generates a new cryptographic identity for the client, and injects // NewKeyPair generates a new cryptographic identity for the client, and injects
@ -621,18 +668,18 @@ func (w *Whisper) Stop() error {
// HandlePeer is called by the underlying P2P layer when the whisper sub-protocol // HandlePeer is called by the underlying P2P layer when the whisper sub-protocol
// connection is negotiated. // connection is negotiated.
func (wh *Whisper) HandlePeer(peer *p2p.Peer, rw p2p.MsgReadWriter) error { func (w *Whisper) HandlePeer(peer *peer.ID, rw p2p.MsgReadWriter) error {
// Create the new peer and start tracking it // Create the new peer and start tracking it
whisperPeer := newPeer(wh, peer, rw) whisperPeer := newPeer(w, peer, rw)
wh.peerMu.Lock() w.peerMu.Lock()
wh.peers[whisperPeer] = struct{}{} w.peers[whisperPeer] = struct{}{}
wh.peerMu.Unlock() w.peerMu.Unlock()
defer func() { defer func() {
wh.peerMu.Lock() w.peerMu.Lock()
delete(wh.peers, whisperPeer) delete(w.peers, whisperPeer)
wh.peerMu.Unlock() w.peerMu.Unlock()
}() }()
// Run the peer handshake and state updates // Run the peer handshake and state updates
@ -642,41 +689,31 @@ func (wh *Whisper) HandlePeer(peer *p2p.Peer, rw p2p.MsgReadWriter) error {
whisperPeer.start() whisperPeer.start()
defer whisperPeer.stop() defer whisperPeer.stop()
return wh.runMessageLoop(whisperPeer, rw) return w.runMessageLoop(whisperPeer, rw)
} }
// runMessageLoop reads and processes inbound messages directly to merge into client-global state. // ProcessIncomingMessage gets the payload of a p2p message and decodes it
func (wh *Whisper) runMessageLoop(p *Peer, rw p2p.MsgReadWriter) error { // depending on its type.
for { func (w *Whisper) ProcessIncomingMessage(code byte, data []byte, p *Peer) error {
// fetch the next packet switch code {
packet, err := rw.ReadMsg()
if err != nil {
log.Warn("message loop", "peer", p.peer.ID(), "err", err)
return err
}
if packet.Size > wh.MaxMessageSize() {
log.Warn("oversized message received", "peer", p.peer.ID())
return errors.New("oversized message received")
}
switch packet.Code {
case statusCode: case statusCode:
// this should not happen, but no need to panic; just ignore this message. // this should not happen, but no need to panic; just ignore this message.
log.Warn("unxepected status message received", "peer", p.peer.ID()) log.Warn("unxepected status message received", "peer", p.pid)
case messagesCode: case messagesCode:
// decode the contained envelopes // decode the contained envelopes
var envelopes []*Envelope var envelopes []*Envelope
if err := packet.Decode(&envelopes); err != nil { err := rlp.DecodeBytes(data, envelopes)
log.Warn("failed to decode envelopes, peer will be disconnected", "peer", p.peer.ID(), "err", err) if err != nil {
log.Warn("failed to decode envelopes, peer will be disconnected", "peer", p.pid, "err", err)
return errors.New("invalid envelopes") return errors.New("invalid envelopes")
} }
trouble := false trouble := false
for _, env := range envelopes { for _, env := range envelopes {
cached, err := wh.add(env) cached, err := w.add(env)
if err != nil { if err != nil {
trouble = true trouble = true
log.Error("bad envelope received, peer will be disconnected", "peer", p.peer.ID(), "err", err) log.Error("bad envelope received, peer will be disconnected", "peer", p.pid, "err", err)
} }
if cached { if cached {
p.mark(env) p.mark(env)
@ -687,34 +724,20 @@ func (wh *Whisper) runMessageLoop(p *Peer, rw p2p.MsgReadWriter) error {
return errors.New("invalid envelope") return errors.New("invalid envelope")
} }
case powRequirementCode: case powRequirementCode:
s := rlp.NewStream(packet.Payload, uint64(packet.Size)) var i uint64
i, err := s.Uint() err := rlp.DecodeBytes(data, i)
if err != nil { if err != nil {
log.Warn("failed to decode powRequirementCode message, peer will be disconnected", "peer", p.peer.ID(), "err", err) log.Warn("failed to decode powRequirementCode message, peer will be disconnected", "peer", p.pid, "err", err)
return errors.New("invalid powRequirementCode message") return errors.New("invalid powRequirementCode message")
} }
f := math.Float64frombits(i) f := math.Float64frombits(i)
if math.IsInf(f, 0) || math.IsNaN(f) || f < 0.0 { if math.IsInf(f, 0) || math.IsNaN(f) || f < 0.0 {
log.Warn("invalid value in powRequirementCode message, peer will be disconnected", "peer", p.peer.ID(), "err", err) log.Warn("invalid value in powRequirementCode message, peer will be disconnected", "peer", p.pid, "err", err)
return errors.New("invalid value in powRequirementCode message") return errors.New("invalid value in powRequirementCode message")
} }
p.powRequirement = f p.powRequirement = f
case bloomFilterExCode: case bloomFilterExCode:
var bloom []byte // to be implemented
err := packet.Decode(&bloom)
if err == nil && len(bloom) != bloomFilterSize {
err = fmt.Errorf("wrong bloom filter size %d", len(bloom))
}
if err != nil {
log.Warn("failed to decode bloom filter exchange message, peer will be disconnected", "peer", p.peer.ID(), "err", err)
return errors.New("invalid bloom filter exchange message")
}
if isFullNode(bloom) {
p.bloomFilter = nil
} else {
p.bloomFilter = bloom
}
case p2pMessageCode: case p2pMessageCode:
// peer-to-peer message, sent directly to peer bypassing PoW checks, etc. // peer-to-peer message, sent directly to peer bypassing PoW checks, etc.
// this message is not supposed to be forwarded to other peers, and // this message is not supposed to be forwarded to other peers, and
@ -722,104 +745,127 @@ func (wh *Whisper) runMessageLoop(p *Peer, rw p2p.MsgReadWriter) error {
// these messages are only accepted from the trusted peer. // these messages are only accepted from the trusted peer.
if p.trusted { if p.trusted {
var envelope Envelope var envelope Envelope
if err := packet.Decode(&envelope); err != nil { err := rlp.DecodeBytes(data, envelope)
log.Warn("failed to decode direct message, peer will be disconnected", "peer", p.peer.ID(), "err", err) if err != nil {
log.Warn("failed to decode direct message, peer will be disconnected", "peer", p.pid, "err", err)
return errors.New("invalid direct message") return errors.New("invalid direct message")
} }
wh.postEvent(&envelope, true) w.postEvent(&envelope, true)
} }
case p2pRequestCode: case p2pRequestCode:
// Must be processed if mail server is implemented. Otherwise ignore. // Must be processed if mail server is implemented. Otherwise ignore.
if wh.mailServer != nil { if w.mailServer != nil {
var request Envelope var request Envelope
if err := packet.Decode(&request); err != nil { err := rlp.DecodeBytes(data, request)
log.Warn("failed to decode p2p request message, peer will be disconnected", "peer", p.peer.ID(), "err", err) if err != nil {
log.Warn("failed to decode p2p request message, peer will be disconnected", "peer", p.pid, "err", err)
return errors.New("invalid p2p request") return errors.New("invalid p2p request")
} }
wh.mailServer.DeliverMail(p, &request) w.mailServer.DeliverMail(p, &request)
} }
default: default:
// New message types might be implemented in the future versions of Whisper. // New message types might be implemented in the future versions of Whisper.
// For forward compatibility, just ignore. // For forward compatibility, just ignore.
} }
packet.Discard() return nil
}
// runMessageLoop reads and processes inbound messages directly to merge into client-global state.
func (w *Whisper) runMessageLoop(p *Peer, rw p2p.MsgReadWriter) error {
for {
// fetch the next packet
packet, err := rw.ReadMsg()
if err != nil {
log.Warn("message loop", "peer", p.pid, "err", err)
return err
}
if packet.Size > w.MaxMessageSize() {
log.Warn("oversized message received", "peer", p.pid)
return errors.New("oversized message received")
}
data, err := ioutil.ReadAll(packet.Payload)
err = w.ProcessIncomingMessage(byte(packet.Code), data, p)
if err != nil {
return err
}
} }
} }
// add inserts a new envelope into the message pool to be distributed within the // add inserts a new envelope into the message pool to be distributed within the
// whisper network. It also inserts the envelope into the expiration pool at the // whisper network. It also inserts the envelope into the expiration pool at the
// appropriate time-stamp. In case of error, connection should be dropped. // appropriate time-stamp. In case of error, connection should be dropped.
func (wh *Whisper) add(envelope *Envelope) (bool, error) { func (w *Whisper) add(envelope *Envelope) (bool, error) {
now := uint32(time.Now().Unix()) now := uint32(time.Now().Unix())
sent := envelope.Expiry - envelope.TTL sent := envelope.Expiry - envelope.TTL
if sent > now { if sent > now {
if sent-DefaultSyncAllowance > now { if sent-DefaultSyncAllowance > now {
return false, fmt.Errorf("envelope created in the future [%x]", envelope.Hash()) return false, fmt.Errorf("envelope created in the future [%x]", envelope.Hash())
} else { }
// recalculate PoW, adjusted for the time difference, plus one second for latency // recalculate PoW, adjusted for the time difference, plus one second for latency
envelope.calculatePoW(sent - now + 1) envelope.calculatePoW(sent - now + 1)
} }
}
if envelope.Expiry < now { if envelope.Expiry < now {
if envelope.Expiry+DefaultSyncAllowance*2 < now { if envelope.Expiry+DefaultSyncAllowance*2 < now {
return false, fmt.Errorf("very old message") return false, fmt.Errorf("very old message")
} else { }
log.Debug("expired envelope dropped", "hash", envelope.Hash().Hex()) log.Debug("expired envelope dropped", "hash", envelope.Hash().Hex())
return false, nil // drop envelope without error return false, nil // drop envelope without error
} }
}
if uint32(envelope.size()) > wh.MaxMessageSize() { if uint32(envelope.size()) > w.MaxMessageSize() {
return false, fmt.Errorf("huge messages are not allowed [%x]", envelope.Hash()) return false, fmt.Errorf("huge messages are not allowed [%x]", envelope.Hash())
} }
if envelope.PoW() < wh.MinPow() { if envelope.PoW() < w.MinPow() {
// maybe the value was recently changed, and the peers did not adjust yet. // maybe the value was recently changed, and the peers did not adjust yet.
// in this case the previous value is retrieved by MinPowTolerance() // in this case the previous value is retrieved by MinPowTolerance()
// for a short period of peer synchronization. // for a short period of peer synchronization.
if envelope.PoW() < wh.MinPowTolerance() { if envelope.PoW() < w.MinPowTolerance() {
return false, fmt.Errorf("envelope with low PoW received: PoW=%f, hash=[%v]", envelope.PoW(), envelope.Hash().Hex()) return false, fmt.Errorf("envelope with low PoW received: PoW=%f, hash=[%v]", envelope.PoW(), envelope.Hash().Hex())
} }
} }
if !bloomFilterMatch(wh.BloomFilter(), envelope.Bloom()) { if !bloomFilterMatch(w.BloomFilter(), envelope.Bloom()) {
// maybe the value was recently changed, and the peers did not adjust yet. // maybe the value was recently changed, and the peers did not adjust yet.
// in this case the previous value is retrieved by BloomFilterTolerance() // in this case the previous value is retrieved by BloomFilterTolerance()
// for a short period of peer synchronization. // for a short period of peer synchronization.
if !bloomFilterMatch(wh.BloomFilterTolerance(), envelope.Bloom()) { if !bloomFilterMatch(w.BloomFilterTolerance(), envelope.Bloom()) {
return false, fmt.Errorf("envelope does not match bloom filter, hash=[%v], bloom: \n%x \n%x \n%x", return false, fmt.Errorf("envelope does not match bloom filter, hash=[%v], bloom: \n%x \n%x \n%x",
envelope.Hash().Hex(), wh.BloomFilter(), envelope.Bloom(), envelope.Topic) envelope.Hash().Hex(), w.BloomFilter(), envelope.Bloom(), envelope.Topic)
} }
} }
hash := envelope.Hash() hash := envelope.Hash()
wh.poolMu.Lock() w.poolMu.Lock()
_, alreadyCached := wh.envelopes[hash] _, alreadyCached := w.envelopes[hash]
if !alreadyCached { if !alreadyCached {
wh.envelopes[hash] = envelope w.envelopes[hash] = envelope
if wh.expirations[envelope.Expiry] == nil { if w.expirations[envelope.Expiry] == nil {
wh.expirations[envelope.Expiry] = set.NewNonTS() w.expirations[envelope.Expiry] = set.NewNonTS()
} }
if !wh.expirations[envelope.Expiry].Has(hash) { if !w.expirations[envelope.Expiry].Has(hash) {
wh.expirations[envelope.Expiry].Add(hash) w.expirations[envelope.Expiry].Add(hash)
} }
} }
wh.poolMu.Unlock() w.poolMu.Unlock()
if alreadyCached { if alreadyCached {
log.Trace("whisper envelope already cached", "hash", envelope.Hash().Hex()) log.Trace("whisper envelope already cached", "hash", envelope.Hash().Hex())
} else { } else {
log.Trace("cached whisper envelope", "hash", envelope.Hash().Hex()) log.Trace("cached whisper envelope", "hash", envelope.Hash().Hex())
wh.statsMu.Lock() w.statsMu.Lock()
wh.stats.memoryUsed += envelope.size() w.stats.memoryUsed += envelope.size()
wh.statsMu.Unlock() w.statsMu.Unlock()
wh.postEvent(envelope, false) // notify the local node about the new message w.postEvent(envelope, false) // notify the local node about the new message
if wh.mailServer != nil { if w.mailServer != nil {
wh.mailServer.Archive(envelope) w.mailServer.Archive(envelope)
} }
} }
return true, nil return true, nil