mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-07-22 04:36:42 +00:00
Peer now implements the peerInfo interface
- changed peer initialisation (now peer can be just a record) - peer.Addr (always listening address) - peer.Hash() memoized hash of the address (pubkey from dialaddr or listenaddr) - LastActive updates with channel - ActiveAddresses convenience filter method returns addresses of all connected peers
This commit is contained in:
parent
72deabeb42
commit
b2bcb32107
1 changed files with 76 additions and 57 deletions
133
p2p/peer.go
133
p2p/peer.go
|
|
@ -49,53 +49,16 @@ func (d peerAddr) RlpData() interface{} {
|
||||||
return []interface{}{d.IP, d.Port, d.Pubkey}
|
return []interface{}{d.IP, d.Port, d.Pubkey}
|
||||||
}
|
}
|
||||||
|
|
||||||
type peerRecord struct {
|
|
||||||
addr *peerAddr
|
|
||||||
hash []byte
|
|
||||||
lastActive time.Time
|
|
||||||
lastActiveC chan time.Time
|
|
||||||
peer *Peer
|
|
||||||
}
|
|
||||||
|
|
||||||
func (self *peerRecord) Addr() *peerAddr {
|
|
||||||
return self.addr
|
|
||||||
}
|
|
||||||
|
|
||||||
func (self *peerRecord) Hash() []byte {
|
|
||||||
if self.hash == nil {
|
|
||||||
self.hash = Hash(self.addr.Pubkey)
|
|
||||||
}
|
|
||||||
return self.hash
|
|
||||||
}
|
|
||||||
|
|
||||||
func (self *peerRecord) LastActive() (lastActive time.Time) {
|
|
||||||
var ok bool
|
|
||||||
select {
|
|
||||||
case lastActive, ok = <-self.lastActiveC:
|
|
||||||
if ok {
|
|
||||||
self.lastActive = lastActive
|
|
||||||
}
|
|
||||||
default:
|
|
||||||
lastActive = self.lastActive
|
|
||||||
}
|
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
||||||
func (self *peerRecord) Connect() error {
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
func (self *peerRecord) Disconnect() error {
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
|
|
||||||
// Peer represents a remote peer.
|
// Peer represents a remote peer.
|
||||||
type Peer struct {
|
type Peer struct {
|
||||||
// Peers have all the log methods.
|
// Peers have all the log methods.
|
||||||
// Use them to display messages related to the peer.
|
// Use them to display messages related to the peer.
|
||||||
*logger.Logger
|
*logger.Logger
|
||||||
|
|
||||||
infolock sync.Mutex
|
lastActive time.Time // updated persisted
|
||||||
|
lastActiveC chan time.Time // for constant querying
|
||||||
|
|
||||||
|
infolock sync.RWMutex
|
||||||
identity ClientIdentity
|
identity ClientIdentity
|
||||||
caps []Cap
|
caps []Cap
|
||||||
listenAddr *peerAddr // what remote peer is listening on
|
listenAddr *peerAddr // what remote peer is listening on
|
||||||
|
|
@ -125,6 +88,7 @@ type Peer struct {
|
||||||
|
|
||||||
// These fields are kept so base protocol can access them.
|
// These fields are kept so base protocol can access them.
|
||||||
// TODO: this should be one or more interfaces
|
// TODO: this should be one or more interfaces
|
||||||
|
hash []byte // hash of pubkey used as address
|
||||||
ourID ClientIdentity // client id of the Server
|
ourID ClientIdentity // client id of the Server
|
||||||
ourListenAddr *peerAddr // listen addr of Server, nil if not listening
|
ourListenAddr *peerAddr // listen addr of Server, nil if not listening
|
||||||
addPeer func(*peerAddr) error // tell server about received peers
|
addPeer func(*peerAddr) error // tell server about received peers
|
||||||
|
|
@ -135,41 +99,96 @@ type Peer struct {
|
||||||
// NewPeer returns a peer for testing purposes.
|
// NewPeer returns a peer for testing purposes.
|
||||||
func NewPeer(id ClientIdentity, caps []Cap) *Peer {
|
func NewPeer(id ClientIdentity, caps []Cap) *Peer {
|
||||||
conn, _ := net.Pipe()
|
conn, _ := net.Pipe()
|
||||||
peer := newPeer(conn, nil, nil)
|
peer := &Peer{}
|
||||||
|
peer.init(conn)
|
||||||
peer.setHandshakeInfo(id, nil, caps)
|
peer.setHandshakeInfo(id, nil, caps)
|
||||||
close(peer.closed)
|
close(peer.closed)
|
||||||
return peer
|
return peer
|
||||||
}
|
}
|
||||||
|
|
||||||
func newServerPeer(server *Server, conn net.Conn, dialAddr *peerAddr) *Peer {
|
func (self *Peer) init(conn net.Conn) {
|
||||||
p := newPeer(conn, server.Protocols, dialAddr)
|
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 (p *Peer) connect(server *Server, conn net.Conn) {
|
||||||
|
p.init(conn)
|
||||||
p.ourID = server.Identity
|
p.ourID = server.Identity
|
||||||
p.addPeer = server.AddPeer
|
p.addPeer = server.AddPeer
|
||||||
p.getPeers = server.GetPeers
|
p.getPeers = server.GetPeers
|
||||||
p.pubkeyHook = server.verifyPeer
|
p.pubkeyHook = server.verifyPeer
|
||||||
p.runBaseProtocol = true
|
p.runBaseProtocol = true
|
||||||
|
p.protocols = server.Protocols
|
||||||
|
|
||||||
// laddr can be updated concurrently by NAT traversal.
|
// laddr can be updated concurrently by NAT traversal.
|
||||||
// newServerPeer must be called with the server lock held.
|
// newServerPeer must be called with the server lock held.
|
||||||
if server.laddr != nil {
|
if server.laddr != nil {
|
||||||
p.ourListenAddr = newPeerAddr(server.laddr, server.Identity.Pubkey())
|
p.ourListenAddr = newPeerAddr(server.laddr, server.Identity.Pubkey())
|
||||||
}
|
}
|
||||||
return p
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func newPeer(conn net.Conn, protocols []Protocol, dialAddr *peerAddr) *Peer {
|
// ActiveAddresses returns addresses of all connected peers.
|
||||||
p := &Peer{
|
// Actually if the peer selector keeps historical info , then once active peers
|
||||||
Logger: logger.NewLogger("P2P " + conn.RemoteAddr().String()),
|
// will be included too.
|
||||||
conn: conn,
|
func ActiveAddresses(peers ...peerInfo) (addrs []*peerAddr) {
|
||||||
dialAddr: dialAddr,
|
for _, peer := range peers {
|
||||||
bufconn: bufio.NewReadWriter(bufio.NewReader(conn), bufio.NewWriter(conn)),
|
if peer != nil {
|
||||||
protocols: protocols,
|
addr := peer.Addr()
|
||||||
running: make(map[string]*proto),
|
// filter out peers that are not listening or
|
||||||
disc: make(chan DiscReason),
|
// have not completed the handshake.
|
||||||
protoErr: make(chan error),
|
// the peer selector can track previously sent peers and exclude them as well.
|
||||||
closed: make(chan struct{}),
|
if addr == nil {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
addrs = append(addrs, addr)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
return p
|
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()
|
||||||
|
if self.dialAddr != nil {
|
||||||
|
pubkey = self.dialAddr.Pubkey
|
||||||
|
} else {
|
||||||
|
if 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
|
// Identity returns the client identity of the remote peer. The
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue