This commit is contained in:
Viktor Trón 2015-01-20 16:27:03 +00:00
commit 7bd37f74f2
13 changed files with 662 additions and 233 deletions

View file

@ -199,7 +199,7 @@ func (ui *UiLib) Connect(button qml.Object) {
}
func (ui *UiLib) ConnectToPeer(addr string) {
if err := ui.eth.SuggestPeer(addr); err != nil {
if err := ui.eth.SuggestPeer(addr, []byte{}); err != nil {
guilogger.Infoln(err)
}
}

View file

@ -2,7 +2,6 @@ package eth
import (
"fmt"
"net"
"sync"
"github.com/ethereum/go-ethereum/core"
@ -21,6 +20,8 @@ const (
seedNodeAddress = "poc-8.ethdev.com:30303"
)
var seednodeId []byte = nil
type Config struct {
Name string
Version string
@ -244,7 +245,7 @@ func (s *Ethereum) Start(seed bool) error {
// TODO: read peers here
if seed {
logger.Infof("Connect to seed node %v", seedNodeAddress)
if err := s.SuggestPeer(seedNodeAddress); err != nil {
if err := s.SuggestPeer(seedNodeAddress, seednodeId); err != nil {
return err
}
}
@ -253,14 +254,8 @@ func (s *Ethereum) Start(seed bool) error {
return nil
}
func (self *Ethereum) SuggestPeer(addr string) error {
netaddr, err := net.ResolveTCPAddr("tcp", addr)
if err != nil {
logger.Errorf("couldn't resolve %s:", addr, err)
return err
}
self.net.SuggestPeer(netaddr.IP, netaddr.Port, nil)
func (self *Ethereum) SuggestPeer(addr string, pubkey []byte) error {
self.net.SuggestPeer(addr, pubkey)
return nil
}

View file

@ -1,23 +0,0 @@
package eth
import (
"encoding/json"
"github.com/ethereum/go-ethereum/ethutil"
)
func WritePeers(path string, addresses []string) {
if len(addresses) > 0 {
data, _ := json.MarshalIndent(addresses, "", " ")
ethutil.WriteFile(path, data)
}
}
func ReadPeers(path string) (ips []string, err error) {
var data string
data, err = ethutil.ReadAllFile(path)
if err != nil {
json.Unmarshal([]byte(data), &ips)
}
return
}

View file

@ -202,7 +202,7 @@ func (self *JSRE) addPeer(call otto.FunctionCall) otto.Value {
if err != nil {
return otto.FalseValue()
}
self.ethereum.SuggestPeer(host)
self.ethereum.SuggestPeer(host, nil)
return otto.TrueValue()
}

287
p2p/cademlia.go Normal file
View file

@ -0,0 +1,287 @@
package p2p
import (
"fmt"
"math"
"sync"
"time"
ethlogger "github.com/ethereum/go-ethereum/logger"
)
var cadlogger = ethlogger.NewLogger("CAD")
const (
hashBytes = 20
rowLength = 10
maxProx = 20
)
var maxAge = 180 * time.Nanosecond
var purgeInterval = 300 * time.Second
type Cademlia struct {
Hash []byte
HashBytes int
RowLength int
MaxProx int
MaxProxBinSize int
MaxAge time.Duration
PurgeInterval time.Duration
proxLimit int
proxSize int
rows []*row
lock sync.RWMutex
quitC chan bool
}
// public constructor with compulsory arguments
// hash is a byte slice of length equal to self.HashBytes
func NewCademlia(hash []byte) *Cademlia {
return &Cademlia{
Hash: hash, // compulsory fields without default
}
}
// Start brings up a pool of peers potentially from an offline persisted source
// and sets default values for optional parameters
func (self *Cademlia) Start() error {
self.lock.Lock()
defer self.lock.Unlock()
if self.quitC != nil {
return nil
}
// these + self.Hash can and should be checked against the
// saved file/db
if self.HashBytes == 0 {
self.HashBytes = hashBytes
}
if self.MaxProx == 0 {
self.MaxProx = maxProx
}
if self.RowLength == 0 {
self.RowLength = rowLength
}
// runtime parameters
if self.MaxProxBinSize == 0 {
self.MaxProxBinSize = self.RowLength
}
if self.MaxAge == time.Duration(0) {
self.MaxAge = maxAge
}
if self.PurgeInterval == time.Duration(0) {
self.PurgeInterval = purgeInterval
}
self.rows = make([]*row, self.MaxProx)
for i, _ := range self.rows {
self.rows[i] = &row{} // will initialise row{int(0),[]*entry(nil),sync.Mutex}
}
self.quitC = make(chan bool)
go self.purgeLoop()
return nil
}
// Stop saves the routing table into a persistant form
func (self *Cademlia) Stop() {
self.lock.Lock()
defer self.lock.Unlock()
if self.quitC == nil {
return
}
close(self.quitC)
self.quitC = nil
}
// AddPeer is the entry point where new peers are suggested for addition to the peer pool
// peers conform to the peerrInfo interface
// AddPeer(peer) returns an error if it deems the peer unworthy
func (self *Cademlia) AddPeer(peer peerInfo) (err error) {
self.lock.Lock()
defer self.lock.Unlock()
index := self.ProximityBin(peer.Hash())
row := self.rows[index]
added := row.insert(&entry{peer: peer})
if added {
if index >= self.proxLimit {
go self.adjustProx(index, 1)
}
cadlogger.Infof("accept peer %x...", peer.Hash()[:8])
} else {
err = fmt.Errorf("no worse peer found")
cadlogger.Infof("reject peer %x..: %v", peer.Hash()[:8], err)
}
return
}
// adjust Prox (proxLimit and proxSize after an insertion of add entries into row r)
func (self *Cademlia) adjustProx(r int, add int) {
self.lock.Lock()
defer self.lock.Unlock()
if r >= self.proxLimit &&
self.proxSize+add > self.MaxProxBinSize &&
self.rows[r].len() > 0 {
self.proxLimit++
} else {
self.proxSize += add
}
}
// updates Prox (proxLimit and proxSize after purging entries)
func (self *Cademlia) updateProx() {
self.lock.Lock()
defer self.lock.Unlock()
var sum int
for i := self.MaxProx - 1; i >= 0; i-- {
l := self.rows[i].len()
sum += l
if sum <= self.MaxProxBinSize || l == 0 {
self.proxSize = sum
}
}
}
// GetPeers(target) returns the list of peers belonging to the same proximity bin as the target. The most proximate bin will be the union of the bins between proxLimit and MaxProx. proxLimit is dynamically adjusted so that 1) there is no empty rows in bin < proxLimit and 2) the sum of all items are the maximum possible but lower than MaxProxBinSize
func (self *Cademlia) GetPeers(target []byte) (peers []peerInfo) {
self.lock.RLock()
defer self.lock.RUnlock()
index := self.ProximityBin(target)
var entries []*entry
if index >= self.proxLimit {
for i := self.proxLimit; i < self.MaxProx; i++ {
entries = append(entries, self.rows[i].row...)
}
} else {
entries = self.rows[index].row
}
for _, entry := range entries {
peers = append(peers, entry.peer)
}
return
}
// entry wrapper type for peer object adding potentially persisted metadata for offline permanent record
type entry struct {
peer peerInfo
// metadata
}
// in situ mutable row
type row struct {
length int
row []*entry
lock sync.RWMutex
}
func (self *row) len() int {
self.lock.RLock()
defer self.lock.RUnlock()
return self.length
}
// insert adds a peer to a row either by appending to existing items if row length does not exceed RowLength, or by replacing the worst entry in the row
func (self *row) insert(entry *entry) (added bool) {
self.lock.Lock()
defer self.lock.Unlock()
if len(self.row) >= self.length { // >= allows us to add peers beyond the Rowlength limitation
worst := self.worst()
self.row[worst].peer.Disconnect(DiscSubprotocolError)
self.row[worst] = entry
} else {
self.row = append(self.row, entry)
added = true
self.length++
}
return
}
// worst expunges the single worst entry in a row, where worst entry is with a peer that has not been active the longests
func (self *row) worst() (index int) {
var oldest time.Time
for i, entry := range self.row {
if (oldest == time.Time{}) || entry.peer.LastActive().Before(oldest) {
oldest = entry.peer.LastActive()
index = i
}
}
return
}
// expunges entries from a row that were last active more that MaxAge ago
// calls Disconnect on entry.peer
func (self *row) purge(recently time.Time) {
self.lock.Lock()
var newRow []*entry
for _, entry := range self.row {
if !entry.peer.LastActive().Before(recently) {
newRow = append(newRow, entry)
} else {
entry.peer.Disconnect(DiscSubprotocolError)
}
}
self.row = newRow
self.length = len(newRow)
self.lock.Unlock()
}
func Hash(key []byte) []byte {
return key
}
func (self *Cademlia) purgeLoop() {
ticker := time.Tick(self.PurgeInterval)
for {
select {
case <-self.quitC:
return
case <-ticker:
self.lock.Lock()
for _, r := range self.rows {
r.purge(time.Now().Add(-self.MaxAge))
}
self.updateProx()
self.lock.Unlock()
}
}
}
/*
Taking the proximity value relative to a fix point x classifies the points in the space (n byte long byte sequences) into bins the items in which are each at most half as distant from x as items in the previous bin. Given a sample of uniformly distrbuted items (a hash function over arbitrary sequence) the proximity scale maps onto series of subsets with cardinalities on a negative exponential scale.
It also has the property that any two item belonging to the same bin are at most half as distant from each other as they are from x.
If we think of random sample of items in the bins as connections in a network of interconnected nodes than relative proximity can serve as the basis for local decisions for graph traversal where the task is to find a route between two points. Since in every step of forwarding, the finite distance halves, there is a guaranteed constant maximum limit on the number of hops needed to reach one node from the other.
*/
func (self *Cademlia) ProximityBin(other []byte) (ret int) {
return int(math.Min(float64(self.MaxProx), float64(self.Proximity(self.Hash, other))))
}
/*
The distance metric MSB(x, y) of two equal length bytesequences x an y is the value of the
binary integer cast of the xor-ed bytesequence most significant bit first.
Proximity(x, y) counts the common zeros in the front of this distance measure.
*/
func (self *Cademlia) Proximity(one, other []byte) (ret int) {
xor := Xor(one, other)
for i := 0; i < self.HashBytes; i++ {
for j := 0; j < 8; j++ {
if (xor[i]>>uint8(7-j))&0x1 != 0 {
return i*8 + j
}
}
}
return self.HashBytes*8 - 1
}
func Xor(one, other []byte) (xor []byte) {
for i := 0; i < len(one); i++ {
xor[i] = one[i] ^ other[i]
}
return
}

View file

@ -45,8 +45,8 @@ func (d peerAddr) String() string {
return fmt.Sprintf("%v:%d", d.IP, d.Port)
}
func (d *peerAddr) RlpData() interface{} {
return []interface{}{string(d.IP), d.Port, d.Pubkey}
func (d peerAddr) RlpData() interface{} {
return []interface{}{d.IP, d.Port, d.Pubkey}
}
// Peer represents a remote peer.
@ -55,7 +55,10 @@ type Peer struct {
// Use them to display messages related to the peer.
*logger.Logger
infolock sync.Mutex
lastActive time.Time // updated persisted
lastActiveC chan time.Time // for constant querying
infolock sync.RWMutex
identity ClientIdentity
caps []Cap
listenAddr *peerAddr // what remote peer is listening on
@ -85,51 +88,92 @@ type Peer struct {
// These fields are kept so base protocol can access them.
// TODO: this should be one or more interfaces
hash []byte // hash of pubkey used as address
ourID ClientIdentity // client id of the Server
ourListenAddr *peerAddr // listen addr of Server, nil if not listening
newPeerAddr chan<- *peerAddr // tell server about received peers
otherPeers func() []*Peer // should return the list of all peers
pubkeyHook func(*peerAddr) error // called at end of handshake to validate pubkey
addPeer func(*peerAddr) error // tell server about received peers
getPeers func(...[]byte) []*peerAddr // should return the list of all peers
verifyPeerHook func(*Peer) error // called at end of handshake to validate peer
}
// NewPeer returns a peer for testing purposes.
func NewPeer(id ClientIdentity, caps []Cap) *Peer {
conn, _ := net.Pipe()
peer := newPeer(conn, nil, nil)
peer := &Peer{}
peer.init(conn)
peer.setHandshakeInfo(id, nil, caps)
close(peer.closed)
return peer
}
func newServerPeer(server *Server, conn net.Conn, dialAddr *peerAddr) *Peer {
p := newPeer(conn, server.Protocols, dialAddr)
p.ourID = server.Identity
p.newPeerAddr = server.peerConnect
p.otherPeers = server.Peers
p.pubkeyHook = server.verifyPeer
p.runBaseProtocol = true
// laddr can be updated concurrently by NAT traversal.
// newServerPeer must be called with the server lock held.
if server.laddr != nil {
p.ourListenAddr = newPeerAddr(server.laddr, server.Identity.Pubkey())
}
return p
func (self *Peer) init(conn net.Conn) {
self.conn = conn
self.Logger = logger.NewLogger("P2P " + conn.RemoteAddr().String())
self.bufconn = bufio.NewReadWriter(bufio.NewReader(conn), bufio.NewWriter(conn))
self.running = make(map[string]*proto)
self.disc = make(chan DiscReason)
self.protoErr = make(chan error)
self.closed = make(chan struct{})
}
func newPeer(conn net.Conn, protocols []Protocol, dialAddr *peerAddr) *Peer {
p := &Peer{
Logger: logger.NewLogger("P2P " + conn.RemoteAddr().String()),
conn: conn,
dialAddr: dialAddr,
bufconn: bufio.NewReadWriter(bufio.NewReader(conn), bufio.NewWriter(conn)),
protocols: protocols,
running: make(map[string]*proto),
disc: make(chan DiscReason),
protoErr: make(chan error),
closed: make(chan struct{}),
// ActiveAddresses returns addresses of all connected peers.
// Actually if the peer selector keeps historical info , then once active peers
// will be included too.
func ActiveAddresses(peers ...peerInfo) (addrs []*peerAddr) {
for _, peer := range peers {
if peer != nil {
addr := peer.Addr()
// filter out peers that are not listening or
// have not completed the handshake.
// the peer selector can track previously sent peers and exclude them as well.
if addr == nil {
continue
}
return p
addrs = append(addrs, addr)
}
}
return addrs
}
// implements the peerInfo interface
func (self *Peer) Addr() *peerAddr {
self.infolock.RLock()
defer self.infolock.RUnlock()
return self.listenAddr
}
func (self *Peer) Hash() []byte {
if self.hash == nil {
self.hash = Hash(self.Pubkey())
}
return self.hash
}
func (self *Peer) Pubkey() (pubkey []byte) {
self.infolock.Lock()
defer self.infolock.Unlock()
switch {
case self.identity != nil:
pubkey = self.identity.Pubkey()
case self.dialAddr != nil:
pubkey = self.dialAddr.Pubkey
case self.listenAddr != nil:
pubkey = self.listenAddr.Pubkey
}
return
}
func (self *Peer) LastActive() (lastActive time.Time) {
var ok bool
select {
case lastActive, ok = <-self.lastActiveC:
if ok {
self.lastActive = lastActive
}
default:
lastActive = self.lastActive
}
return
}
// Identity returns the client identity of the remote peer. The

View file

@ -14,7 +14,11 @@ const (
errP2PVersionMismatch
errPubkeyMissing
errPubkeyInvalid
errPubkeyForbidden
errPubkeyMismatch
errBlacklistedPeer
errSelfConnection
errConnectedPeer
errRejectedPeer
errProtocolBreach
errPingTimeout
errInvalidNetworkId
@ -31,7 +35,11 @@ var errorToString = map[int]string{
errP2PVersionMismatch: "P2P Version Mismatch",
errPubkeyMissing: "Public key missing",
errPubkeyInvalid: "Public key invalid",
errPubkeyForbidden: "Public key forbidden",
errPubkeyMismatch: "Public key mismatch",
errBlacklistedPeer: "Blacklisted peer",
errSelfConnection: "Self connection",
errConnectedPeer: "Connected peer",
errRejectedPeer: "Rejected peer",
errProtocolBreach: "Protocol Breach",
errPingTimeout: "Ping timeout",
errInvalidNetworkId: "Invalid network id",
@ -117,9 +125,9 @@ func discReasonForError(err error) DiscReason {
switch peerError.Code {
case errP2PVersionMismatch:
return DiscIncompatibleVersion
case errPubkeyMissing, errPubkeyInvalid:
case errPubkeyMissing, errPubkeyMismatch, errPubkeyInvalid:
return DiscInvalidIdentity
case errPubkeyForbidden:
case errBlacklistedPeer, errSelfConnection, errConnectedPeer, errRejectedPeer:
return DiscUselessPeer
case errInvalidMsgCode, errMagicTokenMismatch, errProtocolBreach:
return DiscProtocolError

79
p2p/peer_selector.go Normal file
View file

@ -0,0 +1,79 @@
package p2p
import (
"encoding/json"
"path"
"time"
"github.com/ethereum/go-ethereum/ethutil"
)
type peerInfo interface {
Addr() *peerAddr
Hash() []byte
LastActive() time.Time
Disconnect(DiscReason)
}
type peerSelector interface {
AddPeer(peer peerInfo) error
GetPeers(target ...[]byte) []peerInfo
Start() error
Stop() error
}
type BaseSelector struct {
DirPath string
getPeers func() []peerInfo
peers []peerInfo
}
func (self *BaseSelector) AddPeer(peer peerInfo) error {
return nil
}
func (self *BaseSelector) GetPeers(target ...[]byte) []peerInfo {
return self.getPeers()
}
func (self *BaseSelector) Start() error {
if len(self.DirPath) > 0 {
path := path.Join(self.DirPath, "peers.json")
peers, err := ReadPeers(path)
if err != nil {
return err
}
self.peers = peers
}
return nil
}
func (self *BaseSelector) Stop() error {
if len(self.DirPath) > 0 {
path := path.Join(self.DirPath, "peers.json")
if err := WritePeers(path, self.peers); err != nil {
return err
}
}
return nil
}
func WritePeers(path string, addresses []peerInfo) error {
if len(addresses) > 0 {
data, err := json.MarshalIndent(addresses, "", " ")
if err == nil {
ethutil.WriteFile(path, data)
}
return err
}
return nil
}
func ReadPeers(path string) (peers []peerInfo, err error) {
var data string
data, err = ethutil.ReadAllFile(path)
if err == nil {
json.Unmarshal([]byte(data), &peers)
}
return
}

View file

@ -30,9 +30,11 @@ var discard = Protocol{
func testPeer(protos []Protocol) (net.Conn, *Peer, <-chan error) {
conn1, conn2 := net.Pipe()
peer := newPeer(conn1, protos, nil)
peer := &Peer{}
peer.init(conn1)
peer.protocols = protos
peer.ourID = &peerId{}
peer.pubkeyHook = func(*peerAddr) error { return nil }
peer.verifyPeerHook = func(*Peer) error { return nil }
errc := make(chan error, 1)
go func() {
_, err := peer.loop()

View file

@ -114,6 +114,7 @@ func (bp *baseProtocol) loop(quit <-chan error) error {
ping := time.NewTimer(pingTimeout)
activity := bp.peer.activity.Subscribe(time.Time{})
lastActive := time.Time{}
lastActiveC := bp.peer.lastActiveC
defer ping.Stop()
defer activity.Unsubscribe()
@ -125,6 +126,7 @@ func (bp *baseProtocol) loop(quit <-chan error) error {
select {
case err = <-quit:
return err
case lastActiveC <- lastActive:
case <-getPeersTick.C:
err = EncodeMsg(bp.rw, getPeersMsg)
case event := <-activity.Chan():
@ -169,15 +171,35 @@ func (bp *baseProtocol) handle(rw MsgReadWriter) error {
case pongMsg:
case getPeersMsg:
peers := bp.peerList()
var target [][]byte
if err := msg.Decode(&target); err != nil {
return newPeerError(errInvalidMsg, "%v", err)
}
peers := bp.peer.getPeers(target...)
if len(target) == 0 {
// then add ourselves to the list
ourAddr := bp.peer.ourListenAddr
if ourAddr != nil && !ourAddr.IP.IsLoopback() && !ourAddr.IP.IsUnspecified() {
peers = append(peers, ourAddr)
}
}
addrs := make([]interface{}, 0, len(peers))
// encode and filter out requesting peer
for _, addr := range peers {
if addr != bp.peer.Addr() {
addrs = append(addrs, addr)
}
}
// this is dangerous. the spec says that we should _delay_
// sending the response if no new information is available.
// this means that would need to send a response later when
// new peers become available.
//
// TODO: add event mechanism to notify baseProtocol for new peers
if len(peers) > 0 {
return EncodeMsg(bp.rw, peersMsg, peers...)
if len(addrs) > 0 {
return EncodeMsg(bp.rw, peersMsg, addrs...)
}
case peersMsg:
@ -187,7 +209,7 @@ func (bp *baseProtocol) handle(rw MsgReadWriter) error {
}
for _, addr := range peers {
bp.peer.Debugf("received peer suggestion: %v", addr)
bp.peer.newPeerAddr <- addr
bp.peer.addPeer(addr)
}
default:
@ -227,13 +249,9 @@ func (bp *baseProtocol) readHandshake() error {
// verify that the peer we wanted to connect to
// actually holds the target public key.
if da.Pubkey != nil && !bytes.Equal(da.Pubkey, hs.NodeID) {
return newPeerError(errPubkeyForbidden, "dial address pubkey mismatch")
return newPeerError(errPubkeyMismatch, "dial address pubkey mismatch: %x vs %x", da.Pubkey, hs.NodeID)
}
}
pa := newPeerAddr(bp.peer.conn.RemoteAddr(), hs.NodeID)
if err := bp.peer.pubkeyHook(pa); err != nil {
return newPeerError(errPubkeyForbidden, "%v", err)
}
// TODO: remove Caps with empty name
var addr *peerAddr
if hs.ListenPort != 0 {
@ -241,6 +259,9 @@ func (bp *baseProtocol) readHandshake() error {
addr.Port = hs.ListenPort
}
bp.peer.setHandshakeInfo(&hs, addr, hs.Caps)
if err := bp.peer.verifyPeerHook(bp.peer); err != nil {
return err
}
bp.peer.startSubprotocols(hs.Caps)
return nil
}
@ -264,25 +285,3 @@ func (bp *baseProtocol) handshakeMsg() Msg {
bp.peer.ourID.Pubkey()[1:],
)
}
func (bp *baseProtocol) peerList() []interface{} {
peers := bp.peer.otherPeers()
ds := make([]interface{}, 0, len(peers))
for _, p := range peers {
p.infolock.Lock()
addr := p.listenAddr
p.infolock.Unlock()
// filter out this peer and peers that are not listening or
// have not completed the handshake.
// TODO: track previously sent peers and exclude them as well.
if p == bp.peer || addr == nil {
continue
}
ds = append(ds, addr)
}
ourAddr := bp.peer.ourListenAddr
if ourAddr != nil && !ourAddr.IP.IsLoopback() && !ourAddr.IP.IsUnspecified() {
ds = append(ds, ourAddr)
}
return ds
}

View file

@ -29,10 +29,10 @@ func (self *peerId) Pubkey() (pubkey []byte) {
func newTestPeer() (peer *Peer) {
peer = NewPeer(&peerId{}, []Cap{})
peer.pubkeyHook = func(*peerAddr) error { return nil }
peer.verifyPeerHook = func(*Peer) error { return nil }
peer.ourID = &peerId{}
peer.listenAddr = &peerAddr{}
peer.otherPeers = func() []*Peer { return nil }
peer.getPeers = func(...[]byte) []*peerAddr { return nil }
return
}
@ -79,10 +79,10 @@ func TestBaseProtocolPeers(t *testing.T) {
// run first peer (in background)
peer1 := newTestPeer()
peer1.ourListenAddr = listenAddr
peer1.otherPeers = func() []*Peer {
pl := make([]*Peer, len(peerList))
peer1.getPeers = func(...[]byte) []*peerAddr {
pl := make([]*peerAddr, len(peerList))
for i, addr := range peerList {
pl[i] = &Peer{listenAddr: addr}
pl[i] = addr
}
return pl
}
@ -94,7 +94,10 @@ func TestBaseProtocolPeers(t *testing.T) {
// run second peer
peer2 := newTestPeer()
peer2.newPeerAddr = addrChan // feed peer suggestions into matcher
peer2.addPeer = func(addr *peerAddr) error {
addrChan <- addr // feed peer suggestions into matcher
return nil
}
if err := runBaseProtocol(peer2, rw2); err != ErrPipeClosed {
t.Errorf("peer2 terminated with unexpected error: %v", err)
}
@ -107,7 +110,7 @@ func TestBaseProtocolPeers(t *testing.T) {
func TestBaseProtocolDisconnect(t *testing.T) {
peer := NewPeer(&peerId{}, nil)
peer.ourID = &peerId{}
peer.pubkeyHook = func(*peerAddr) error { return nil }
peer.verifyPeerHook = func(*Peer) error { return nil }
rw1, rw2 := MsgPipe()
done := make(chan struct{})

View file

@ -60,11 +60,16 @@ type Server struct {
Dialer *net.Dialer
// If NoDial is true, the server will not dial any peers.
// this maybe used in test environments where we want to prevent a node from
// connecting (and synchronising) with other nodes
NoDial bool
// peer selector
PeerSelector peerSelector
// Hook for testing. This is useful because we can inhibit
// the whole protocol stack.
newPeerFunc peerFunc
// newPeerFunc peerFunc
lock sync.RWMutex
running bool
@ -74,9 +79,10 @@ type Server struct {
peerSlots chan int
peerCount int
connectFunc func(*Peer, net.Conn)
quit chan struct{}
wg sync.WaitGroup
peerConnect chan *peerAddr
peerDisconnect chan *Peer
}
@ -90,9 +96,7 @@ type NAT interface {
String() string
}
type peerFunc func(srv *Server, c net.Conn, dialAddr *peerAddr) *Peer
// Peers returns all connected peers.
// Peers returns all currently connected peers.
func (srv *Server) Peers() (peers []*Peer) {
srv.lock.RLock()
defer srv.lock.RUnlock()
@ -104,6 +108,21 @@ func (srv *Server) Peers() (peers []*Peer) {
return
}
// GetPeers returns addresses near target if given or supported by the client
// or falls back to all actively connected peers
func (srv *Server) GetPeers(target ...[]byte) []*peerAddr {
if len(target) == 1 { // delegate to selector
return ActiveAddresses(srv.PeerSelector.GetPeers(target[0])...)
} else {
// in fact it is not clear why the selector would not want to reply to this case as well
var peers []peerInfo
for _, peer := range srv.Peers() {
peers = append(peers, peer)
}
return ActiveAddresses(peers...)
}
}
// PeerCount returns the number of connected peers.
func (srv *Server) PeerCount() int {
srv.lock.RLock()
@ -111,13 +130,75 @@ func (srv *Server) PeerCount() int {
return srv.peerCount
}
// SuggestPeer injects an address into the outbound address pool.
func (srv *Server) SuggestPeer(ip net.IP, port int, nodeID []byte) {
addr := &peerAddr{ip, uint64(port), nodeID}
// SuggestPeer is a convenient method that does dns resolution
// and passes on the request to AddPeer
func (srv *Server) SuggestPeer(addr string, pubkey []byte) error {
netaddr, err := net.ResolveTCPAddr("tcp", addr)
if err != nil {
srvlog.Errorf("couldn't resolve %s:", addr, err)
return err
}
peerAddr := &peerAddr{netaddr.IP, uint64(netaddr.Port), pubkey}
return srv.AddPeer(peerAddr)
}
// AddPeer takes a peerAddr address as argument.
// If not found among connected peers turns to the peerSelector
// to decide if it is a worthwhile connection
func (srv *Server) AddPeer(addr *peerAddr) (err error) {
if srv.NoDial {
return fmt.Errorf("no dial out")
}
// need to look up nodeID first
peer := &Peer{
dialAddr: addr,
lastActiveC: make(chan time.Time),
lastActive: time.Now().Add(-24 * time.Hour),
}
if err = srv.PeerSelector.AddPeer(peer); err == nil {
srvlog.Infof("peer %v accepted by peer selection", addr)
err = srv.dialPeer(peer)
} else {
srvlog.Infof("peer %v rejected by peer selection", addr)
}
return
}
func (srv *Server) dialPeer(peer *Peer) (err error) {
timeout := time.After(5 * time.Second)
select {
case srv.peerConnect <- addr:
default: // don't block
srvlog.Warnf("peer suggestion %v ignored", addr)
case <-timeout:
err = fmt.Errorf("Too many connections. No slot available")
case slot := <-srv.peerSlots: // there is a slot available
srvlog.Infof("Dialing %v (slot %d)\n", peer.dialAddr, slot)
conn, dialErr := srv.Dialer.Dial(peer.dialAddr.Network(), peer.dialAddr.String())
if dialErr != nil {
err = fmt.Errorf("Dial error: %v", dialErr)
srvlog.Errorln(err)
srv.peerSlots <- slot
return
}
srvlog.Infof("Connected to %v (slot %d)\n", peer.dialAddr, slot)
peer.slot = slot
srv.connectFunc(peer, conn)
go srv.addPeer(peer)
}
return
}
func (srv *Server) connect(p *Peer, conn net.Conn) {
p.init(conn)
p.ourID = srv.Identity
p.addPeer = srv.AddPeer
p.getPeers = srv.GetPeers
p.verifyPeerHook = srv.verifyPeer
p.runBaseProtocol = true
p.protocols = srv.Protocols
// laddr can be updated concurrently by NAT traversal.
if srv.laddr != nil {
p.ourListenAddr = newPeerAddr(srv.laddr, srv.Identity.Pubkey())
}
}
@ -162,11 +243,10 @@ func (srv *Server) Start() (err error) {
srv.quit = make(chan struct{})
srv.peers = make([]*Peer, srv.MaxPeers)
srv.peerSlots = make(chan int, srv.MaxPeers)
srv.peerConnect = make(chan *peerAddr, outboundAddressPoolSize)
srv.peerDisconnect = make(chan *Peer)
if srv.newPeerFunc == nil {
srv.newPeerFunc = newServerPeer
}
// if srv.newPeerFunc == nil {
// srv.newPeerFunc = newServerPeer
// }
if srv.Blacklist == nil {
srv.Blacklist = NewBlacklist()
}
@ -179,12 +259,13 @@ func (srv *Server) Start() (err error) {
return err
}
}
if !srv.NoDial {
srv.wg.Add(1)
go srv.dialLoop()
if srv.PeerSelector == nil {
srv.PeerSelector = &BaseSelector{}
}
if srv.NoDial && srv.ListenAddr == "" {
srvlog.Warnln("I will be kind-of useless, neither dialing nor listening.")
if srv.connectFunc == nil {
srv.connectFunc = srv.connect
}
// make all slots available
@ -266,8 +347,10 @@ func (srv *Server) listenLoop() {
srv.peerSlots <- slot
return
}
srvlog.Debugf("Accepted conn %v (slot %d)\n", conn.RemoteAddr(), slot)
srv.addPeer(conn, nil, slot)
srvlog.Debugf("Accepted conn %v (slot %d) - peer selector check after handshake", conn.RemoteAddr(), slot)
peer := &Peer{slot: slot}
srv.connectFunc(peer, conn)
srv.addPeer(peer)
case <-srv.quit:
return
}
@ -313,77 +396,22 @@ func (srv *Server) removePortMapping(port int) {
srv.NAT.DeletePortMapping("tcp", port, port)
}
func (srv *Server) dialLoop() {
defer srv.wg.Done()
var (
suggest chan *peerAddr
slot *int
slots = srv.peerSlots
)
for {
select {
case i := <-slots:
// we need a peer in slot i, slot reserved
slot = &i
// now we can watch for candidate peers in the next loop
suggest = srv.peerConnect
// do not consume more until candidate peer is found
slots = nil
case desc := <-suggest:
// candidate peer found, will dial out asyncronously
// if connection fails slot will be released
srvlog.Infof("dial %v (%v)", desc, *slot)
go srv.dialPeer(desc, *slot)
// we can watch if more peers needed in the next loop
slots = srv.peerSlots
// until then we dont care about candidate peers
suggest = nil
case <-srv.quit:
// give back the currently reserved slot
if slot != nil {
srv.peerSlots <- *slot
}
return
}
}
}
// connect to peer via dial out
func (srv *Server) dialPeer(desc *peerAddr, slot int) {
srvlog.Debugf("Dialing %v (slot %d)\n", desc, slot)
conn, err := srv.Dialer.Dial(desc.Network(), desc.String())
if err != nil {
srvlog.Errorf("Dial error: %v", err)
srv.peerSlots <- slot
return
}
go srv.addPeer(conn, desc, slot)
}
// creates the new peer object and inserts it into its slot
func (srv *Server) addPeer(conn net.Conn, desc *peerAddr, slot int) *Peer {
func (srv *Server) addPeer(peer *Peer) {
srv.lock.Lock()
defer srv.lock.Unlock()
if !srv.running {
conn.Close()
srv.peerSlots <- slot // release slot
return nil
}
peer := srv.newPeerFunc(srv, conn, desc)
peer.slot = slot
srv.peers[slot] = peer
srvlog.Debugf("Add peer %v (slot %v)\n", peer, peer.slot)
srv.peers[peer.slot] = peer
srv.peerCount++
go func() { peer.loop(); srv.peerDisconnect <- peer }()
return peer
return
}
// removes peer: sending disconnect msg, stop peer, remove rom list/table, release slot
func (srv *Server) removePeer(peer *Peer) {
srv.lock.Lock()
defer srv.lock.Unlock()
srvlog.Debugf("Removing %v (slot %v)\n", peer, peer.slot)
srvlog.Debugf("Remove peer %v (slot %v)\n", peer, peer.slot)
if srv.peers[peer.slot] != peer {
srvlog.Warnln("Invalid peer to remove:", peer)
return
@ -395,23 +423,29 @@ func (srv *Server) removePeer(peer *Peer) {
srv.peerSlots <- peer.slot
}
func (srv *Server) verifyPeer(addr *peerAddr) error {
if srv.Blacklist.Exists(addr.Pubkey) {
return errors.New("blacklisted")
func (srv *Server) verifyPeer(peer *Peer) error {
pubkey := peer.Pubkey()
if srv.Blacklist.Exists(pubkey) {
return newPeerError(errBlacklistedPeer, "")
}
if bytes.Equal(srv.Identity.Pubkey()[1:], addr.Pubkey) {
return newPeerError(errPubkeyForbidden, "not allowed to connect to srv")
if bytes.Equal(srv.Identity.Pubkey()[1:], pubkey) {
return newPeerError(errSelfConnection, "")
}
srv.lock.RLock()
defer srv.lock.RUnlock()
for _, peer := range srv.peers {
if peer != nil {
id := peer.Identity()
if id != nil && bytes.Equal(id.Pubkey(), addr.Pubkey) {
return errors.New("already connected")
if id != nil && bytes.Equal(id.Pubkey(), pubkey) {
return newPeerError(errConnectedPeer, "")
}
}
}
srv.lock.RUnlock()
if peer.dialAddr == nil {
if err := srv.PeerSelector.AddPeer(peer); err != nil {
return newPeerError(errRejectedPeer, "%v", err)
}
}
return nil
}

View file

@ -9,12 +9,12 @@ import (
"time"
)
func startTestServer(t *testing.T, pf peerFunc) *Server {
func startTestServer(t *testing.T, cb func(*Peer, net.Conn)) *Server {
server := &Server{
Identity: &peerId{},
MaxPeers: 10,
ListenAddr: "127.0.0.1:0",
newPeerFunc: pf,
connectFunc: cb,
}
if err := server.Start(); err != nil {
t.Fatalf("Could not start server: %v", err)
@ -27,16 +27,9 @@ func TestServerListen(t *testing.T) {
// start the test server
connected := make(chan *Peer)
srv := startTestServer(t, func(srv *Server, conn net.Conn, dialAddr *peerAddr) *Peer {
if conn == nil {
t.Error("peer func called with nil conn")
}
if dialAddr != nil {
t.Error("peer func called with non-nil dialAddr")
}
peer := newPeer(conn, nil, dialAddr)
srv := startTestServer(t, func(peer *Peer, conn net.Conn) {
peer.init(conn)
connected <- peer
return peer
})
defer close(connected)
defer srv.Stop()
@ -50,10 +43,18 @@ func TestServerListen(t *testing.T) {
select {
case peer := <-connected:
if peer.conn == nil {
t.Error("peer setup with nil conn")
} else {
if peer.conn.LocalAddr().String() != conn.RemoteAddr().String() {
t.Errorf("peer started with wrong conn: got %v, want %v",
peer.conn.LocalAddr(), conn.RemoteAddr())
}
}
if peer.dialAddr != nil {
t.Error("peer setup with non-nil dialAddr")
}
case <-time.After(1 * time.Second):
t.Error("server did not accept within one second")
}
@ -72,7 +73,7 @@ func TestServerDial(t *testing.T) {
go func() {
conn, err := listener.Accept()
if err != nil {
t.Error("acccept error:", err)
t.Error("accept error:", err)
}
conn.Close()
accepted <- conn
@ -80,29 +81,29 @@ func TestServerDial(t *testing.T) {
// start the test server
connected := make(chan *Peer)
srv := startTestServer(t, func(srv *Server, conn net.Conn, dialAddr *peerAddr) *Peer {
if conn == nil {
t.Error("peer func called with nil conn")
}
peer := newPeer(conn, nil, dialAddr)
connected <- peer
return peer
srv := startTestServer(t, func(peer *Peer, conn net.Conn) {
peer.init(conn)
go func() { connected <- peer }()
})
defer close(connected)
defer srv.Stop()
// tell the server to connect.
connAddr := newPeerAddr(listener.Addr(), nil)
srv.peerConnect <- connAddr
srv.AddPeer(connAddr)
select {
case conn := <-accepted:
select {
case peer := <-connected:
if conn == nil {
t.Error("peer func called with nil conn")
} else {
if peer.conn.RemoteAddr().String() != conn.LocalAddr().String() {
t.Errorf("peer started with wrong conn: got %v, want %v",
peer.conn.RemoteAddr(), conn.LocalAddr())
}
}
if peer.dialAddr != connAddr {
t.Errorf("peer started with wrong dialAddr: got %v, want %v",
peer.dialAddr, connAddr)
@ -119,11 +120,11 @@ func TestServerDial(t *testing.T) {
func TestServerBroadcast(t *testing.T) {
defer testlog(t).detach()
var connected sync.WaitGroup
srv := startTestServer(t, func(srv *Server, c net.Conn, dialAddr *peerAddr) *Peer {
peer := newPeer(c, []Protocol{discard}, dialAddr)
srv := startTestServer(t, func(peer *Peer, conn net.Conn) {
peer.init(conn)
peer.protocols = []Protocol{discard}
peer.startSubprotocols([]Cap{discard.cap()})
connected.Done()
return peer
})
defer srv.Stop()