mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-07-27 15:16:43 +00:00
p2p/adapters: Remove Messenger abstraction
Signed-off-by: Lewis Marshall <lewis@lmars.net>
This commit is contained in:
parent
9407ac7566
commit
5dcae2e26e
14 changed files with 63 additions and 190 deletions
|
|
@ -7,7 +7,7 @@ import (
|
||||||
// "time"
|
// "time"
|
||||||
)
|
)
|
||||||
|
|
||||||
func NewRemoteNode(id *NodeId, n Network, m Messenger) *RemoteNode {
|
func NewRemoteNode(id *NodeId, n Network) *RemoteNode {
|
||||||
return &RemoteNode{
|
return &RemoteNode{
|
||||||
ID: id,
|
ID: id,
|
||||||
Network: n,
|
Network: n,
|
||||||
|
|
|
||||||
|
|
@ -25,9 +25,9 @@ import (
|
||||||
"github.com/ethereum/go-ethereum/p2p/discover"
|
"github.com/ethereum/go-ethereum/p2p/discover"
|
||||||
)
|
)
|
||||||
|
|
||||||
func newPeer(m Messenger) *Peer {
|
func newPeer(rw *p2p.MsgPipeRW) *Peer {
|
||||||
return &Peer{
|
return &Peer{
|
||||||
Messenger: m,
|
MsgPipeRW: rw,
|
||||||
Errc: make(chan error, 1),
|
Errc: make(chan error, 1),
|
||||||
Connc: make(chan bool),
|
Connc: make(chan bool),
|
||||||
Readyc: make(chan bool),
|
Readyc: make(chan bool),
|
||||||
|
|
@ -35,7 +35,7 @@ func newPeer(m Messenger) *Peer {
|
||||||
}
|
}
|
||||||
|
|
||||||
type Peer struct {
|
type Peer struct {
|
||||||
Messenger
|
*p2p.MsgPipeRW
|
||||||
Connc chan bool
|
Connc chan bool
|
||||||
Readyc chan bool
|
Readyc chan bool
|
||||||
Errc chan error
|
Errc chan error
|
||||||
|
|
@ -50,25 +50,19 @@ type Network interface {
|
||||||
|
|
||||||
// SimNode is the network adapter that
|
// SimNode is the network adapter that
|
||||||
type SimNode struct {
|
type SimNode struct {
|
||||||
lock sync.RWMutex
|
lock sync.RWMutex
|
||||||
Id *NodeId
|
Id *NodeId
|
||||||
network Network
|
network Network
|
||||||
messenger func(p2p.MsgReadWriter) Messenger
|
peerMap map[discover.NodeID]int
|
||||||
peerMap map[discover.NodeID]int
|
peers []*Peer
|
||||||
peers []*Peer
|
Run ProtoCall
|
||||||
Run ProtoCall
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func (self *SimNode) Messenger(rw p2p.MsgReadWriter) Messenger {
|
func NewSimNode(id *NodeId, n Network) *SimNode {
|
||||||
return self.messenger(rw)
|
|
||||||
}
|
|
||||||
|
|
||||||
func NewSimNode(id *NodeId, n Network, m func(p2p.MsgReadWriter) Messenger) *SimNode {
|
|
||||||
return &SimNode{
|
return &SimNode{
|
||||||
Id: id,
|
Id: id,
|
||||||
network: n,
|
network: n,
|
||||||
messenger: m,
|
peerMap: make(map[discover.NodeID]int),
|
||||||
peerMap: make(map[discover.NodeID]int),
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -94,18 +88,12 @@ func (self *SimNode) getPeer(id *NodeId) *Peer {
|
||||||
return self.peers[i]
|
return self.peers[i]
|
||||||
}
|
}
|
||||||
|
|
||||||
func (self *SimNode) SetPeer(id *NodeId, rw p2p.MsgReadWriter) {
|
func (self *SimNode) setPeer(id *NodeId, rw *p2p.MsgPipeRW) *Peer {
|
||||||
self.lock.Lock()
|
|
||||||
defer self.lock.Unlock()
|
|
||||||
self.setPeer(id, self.Messenger(rw))
|
|
||||||
}
|
|
||||||
|
|
||||||
func (self *SimNode) setPeer(id *NodeId, m Messenger) *Peer {
|
|
||||||
i, found := self.peerMap[id.NodeID]
|
i, found := self.peerMap[id.NodeID]
|
||||||
if !found {
|
if !found {
|
||||||
i = len(self.peers)
|
i = len(self.peers)
|
||||||
self.peerMap[id.NodeID] = i
|
self.peerMap[id.NodeID] = i
|
||||||
p := newPeer(m)
|
p := newPeer(rw)
|
||||||
self.peers = append(self.peers, p)
|
self.peers = append(self.peers, p)
|
||||||
return p
|
return p
|
||||||
}
|
}
|
||||||
|
|
@ -114,7 +102,7 @@ func (self *SimNode) setPeer(id *NodeId, m Messenger) *Peer {
|
||||||
// }
|
// }
|
||||||
// legit reconnect reset disconnection error,
|
// legit reconnect reset disconnection error,
|
||||||
p := self.peers[i]
|
p := self.peers[i]
|
||||||
p.Messenger = m
|
p.MsgPipeRW = rw
|
||||||
p.Connc = make(chan bool)
|
p.Connc = make(chan bool)
|
||||||
p.Readyc = make(chan bool)
|
p.Readyc = make(chan bool)
|
||||||
return p
|
return p
|
||||||
|
|
@ -125,11 +113,11 @@ func (self *SimNode) Disconnect(rid []byte) error {
|
||||||
defer self.lock.Unlock()
|
defer self.lock.Unlock()
|
||||||
id := NewNodeId(rid)
|
id := NewNodeId(rid)
|
||||||
peer := self.getPeer(id)
|
peer := self.getPeer(id)
|
||||||
if peer == nil || peer.Messenger == nil {
|
if peer == nil || peer.MsgPipeRW == nil {
|
||||||
return fmt.Errorf("already disconnected")
|
return fmt.Errorf("already disconnected")
|
||||||
}
|
}
|
||||||
peer.Messenger.Close()
|
peer.MsgPipeRW.Close()
|
||||||
peer.Messenger = nil
|
peer.MsgPipeRW = nil
|
||||||
// na := self.network.GetNodeAdapter(id)
|
// na := self.network.GetNodeAdapter(id)
|
||||||
// peer = na.(*SimNode).GetPeer(self.Id)
|
// peer = na.(*SimNode).GetPeer(self.Id)
|
||||||
// peer.RW = nil
|
// peer.RW = nil
|
||||||
|
|
@ -149,10 +137,10 @@ func (self *SimNode) Connect(rid []byte) error {
|
||||||
rw, rrw := p2p.MsgPipe()
|
rw, rrw := p2p.MsgPipe()
|
||||||
// // run protocol on remote node with self as peer
|
// // run protocol on remote node with self as peer
|
||||||
peer := self.getPeer(id)
|
peer := self.getPeer(id)
|
||||||
if peer != nil && peer.Messenger != nil {
|
if peer != nil && peer.MsgPipeRW != nil {
|
||||||
return fmt.Errorf("already connected %v to peer %v", self.Id, id)
|
return fmt.Errorf("already connected %v to peer %v", self.Id, id)
|
||||||
}
|
}
|
||||||
peer = self.setPeer(id, self.Messenger(rrw))
|
peer = self.setPeer(id, rrw)
|
||||||
close(peer.Connc)
|
close(peer.Connc)
|
||||||
defer close(peer.Readyc)
|
defer close(peer.Readyc)
|
||||||
err := na.(ProtocolRunner).RunProtocol(self.Id, rrw, rw, peer)
|
err := na.(ProtocolRunner).RunProtocol(self.Id, rrw, rw, peer)
|
||||||
|
|
|
||||||
|
|
@ -1,60 +0,0 @@
|
||||||
// Copyright 2016 The go-ethereum Authors
|
|
||||||
// This file is part of the go-ethereum library.
|
|
||||||
//
|
|
||||||
// The go-ethereum library is free software: you can redistribute it and/or modify
|
|
||||||
// it under the terms of the GNU Lesser General Public License as published by
|
|
||||||
// the Free Software Foundation, either version 3 of the License, or
|
|
||||||
// (at your option) any later version.
|
|
||||||
//
|
|
||||||
// The go-ethereum library is distributed in the hope that it will be useful,
|
|
||||||
// but WITHOUT ANY WARRANTY; without even the implied warranty of
|
|
||||||
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
|
||||||
// GNU Lesser General Public License for more details.
|
|
||||||
//
|
|
||||||
// You should have received a copy of the GNU Lesser General Public License
|
|
||||||
// along with the go-ethereum library. If not, see <http://www.gnu.org/licenses/>.
|
|
||||||
|
|
||||||
package adapters
|
|
||||||
|
|
||||||
import (
|
|
||||||
"github.com/ethereum/go-ethereum/p2p"
|
|
||||||
)
|
|
||||||
|
|
||||||
//network adapter's messenger interace
|
|
||||||
// NewPipe() (p2p.MsgReadWriter, p2p.MsgReadWriter)
|
|
||||||
// ClosePipe(rw p2p.MsgReadWriter)
|
|
||||||
|
|
||||||
// protocol Messenger interface
|
|
||||||
// SendMsg(p2p.MsgWriter, uint64, interface{}) error
|
|
||||||
// ReadMsg(p2p.MsgReader) (p2p.Msg, error)
|
|
||||||
|
|
||||||
// peer session test
|
|
||||||
// ExpectMsg(p2p.MsgReader, uint64, interface{}) error
|
|
||||||
// SendMsg(p2p.MsgWriter, uint64, interface{}) error
|
|
||||||
type SimPipe struct {
|
|
||||||
rw p2p.MsgReadWriter
|
|
||||||
}
|
|
||||||
|
|
||||||
func (self *SimPipe) SendMsg(code uint64, msg interface{}) error {
|
|
||||||
return p2p.Send(self.rw, code, msg)
|
|
||||||
}
|
|
||||||
|
|
||||||
func (self *SimPipe) ReadMsg() (p2p.Msg, error) {
|
|
||||||
return self.rw.ReadMsg()
|
|
||||||
}
|
|
||||||
|
|
||||||
func (self *SimPipe) TriggerMsg(code uint64, msg interface{}) error {
|
|
||||||
return p2p.Send(self.rw, code, msg)
|
|
||||||
}
|
|
||||||
|
|
||||||
func (self *SimPipe) ExpectMsg(code uint64, msg interface{}) error {
|
|
||||||
return p2p.ExpectMsg(self.rw, code, msg)
|
|
||||||
}
|
|
||||||
|
|
||||||
func (self *SimPipe) Close() {
|
|
||||||
self.rw.(*p2p.MsgPipeRW).Close()
|
|
||||||
}
|
|
||||||
|
|
||||||
func NewSimPipe(rw p2p.MsgReadWriter) Messenger {
|
|
||||||
return Messenger(&SimPipe{rw})
|
|
||||||
}
|
|
||||||
|
|
@ -31,28 +31,18 @@ type RLPx struct {
|
||||||
id *NodeId
|
id *NodeId
|
||||||
net *p2p.Server
|
net *p2p.Server
|
||||||
addr []byte
|
addr []byte
|
||||||
m func(p2p.MsgReadWriter) Messenger
|
|
||||||
r Reporter
|
r Reporter
|
||||||
}
|
}
|
||||||
|
|
||||||
type RLPxMessenger struct {
|
func NewRLPx(addr []byte, srv *p2p.Server) *RLPx {
|
||||||
rw p2p.MsgReadWriter
|
|
||||||
}
|
|
||||||
|
|
||||||
func NewRLPxMessenger(rw p2p.MsgReadWriter) Messenger {
|
|
||||||
return Messenger(&RLPxMessenger{rw: rw})
|
|
||||||
}
|
|
||||||
|
|
||||||
func NewRLPx(addr []byte, srv *p2p.Server, m func(p2p.MsgReadWriter) Messenger) *RLPx {
|
|
||||||
return &RLPx{
|
return &RLPx{
|
||||||
net: srv,
|
net: srv,
|
||||||
addr: addr,
|
addr: addr,
|
||||||
m: m,
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func NewReportingRLPx(addr []byte, srv *p2p.Server, m func(p2p.MsgReadWriter) Messenger, r Reporter) *RLPx {
|
func NewReportingRLPx(addr []byte, srv *p2p.Server, r Reporter) *RLPx {
|
||||||
rlpx := NewRLPx(addr, srv, m)
|
rlpx := NewRLPx(addr, srv)
|
||||||
rlpx.r = r
|
rlpx.r = r
|
||||||
srv.PeerConnHook = func(p *p2p.Peer) {
|
srv.PeerConnHook = func(p *p2p.Peer) {
|
||||||
r.DidConnect(rlpx.id, &NodeId{p.ID()})
|
r.DidConnect(rlpx.id, &NodeId{p.ID()})
|
||||||
|
|
@ -63,18 +53,6 @@ func NewReportingRLPx(addr []byte, srv *p2p.Server, m func(p2p.MsgReadWriter) Me
|
||||||
return rlpx
|
return rlpx
|
||||||
}
|
}
|
||||||
|
|
||||||
func (self *RLPxMessenger) SendMsg(code uint64, msg interface{}) error {
|
|
||||||
return p2p.Send(self.rw, code, msg)
|
|
||||||
}
|
|
||||||
|
|
||||||
func (self *RLPxMessenger) ReadMsg() (p2p.Msg, error) {
|
|
||||||
return self.rw.ReadMsg()
|
|
||||||
}
|
|
||||||
|
|
||||||
func (self *RLPxMessenger) Close() {
|
|
||||||
|
|
||||||
}
|
|
||||||
|
|
||||||
func (self *RLPx) LocalAddr() []byte {
|
func (self *RLPx) LocalAddr() []byte {
|
||||||
return self.addr
|
return self.addr
|
||||||
}
|
}
|
||||||
|
|
@ -90,10 +68,6 @@ func (self *RLPx) Connect(enode []byte) error {
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (self *RLPx) Messenger(rw p2p.MsgReadWriter) Messenger {
|
|
||||||
return self.m(rw)
|
|
||||||
}
|
|
||||||
|
|
||||||
//func (self *RLPx) Disconnect(p *p2p.Peer, rw p2p.MsgReadWriter) error {
|
//func (self *RLPx) Disconnect(p *p2p.Peer, rw p2p.MsgReadWriter) error {
|
||||||
func (self *RLPx) Disconnect(b []byte) error {
|
func (self *RLPx) Disconnect(b []byte) error {
|
||||||
//p.Disconnect(p2p.DiscSubprotocolError)
|
//p.Disconnect(p2p.DiscSubprotocolError)
|
||||||
|
|
|
||||||
|
|
@ -63,20 +63,12 @@ func (self *NodeId) Label() string {
|
||||||
return self.String()[:lablen]
|
return self.String()[:lablen]
|
||||||
}
|
}
|
||||||
|
|
||||||
type Messenger interface {
|
|
||||||
SendMsg(uint64, interface{}) error
|
|
||||||
ReadMsg() (p2p.Msg, error)
|
|
||||||
Close()
|
|
||||||
}
|
|
||||||
|
|
||||||
type NodeAdapter interface {
|
type NodeAdapter interface {
|
||||||
Connect([]byte) error
|
Connect([]byte) error
|
||||||
Disconnect([]byte) error
|
Disconnect([]byte) error
|
||||||
// Disconnect(*p2p.Peer, p2p.MsgReadWriter)
|
// Disconnect(*p2p.Peer, p2p.MsgReadWriter)
|
||||||
LocalAddr() []byte
|
LocalAddr() []byte
|
||||||
ParseAddr([]byte, string) ([]byte, error)
|
ParseAddr([]byte, string) ([]byte, error)
|
||||||
// Messenger() Messenger <<<... old version
|
|
||||||
Messenger(p2p.MsgReadWriter) Messenger
|
|
||||||
}
|
}
|
||||||
|
|
||||||
type ProtocolRunner interface {
|
type ProtocolRunner interface {
|
||||||
|
|
|
||||||
|
|
@ -37,7 +37,6 @@ import (
|
||||||
"github.com/ethereum/go-ethereum/logger"
|
"github.com/ethereum/go-ethereum/logger"
|
||||||
"github.com/ethereum/go-ethereum/logger/glog"
|
"github.com/ethereum/go-ethereum/logger/glog"
|
||||||
"github.com/ethereum/go-ethereum/p2p"
|
"github.com/ethereum/go-ethereum/p2p"
|
||||||
"github.com/ethereum/go-ethereum/p2p/adapters"
|
|
||||||
"github.com/ethereum/go-ethereum/p2p/discover"
|
"github.com/ethereum/go-ethereum/p2p/discover"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
@ -160,17 +159,13 @@ func (self *CodeMap) Register(msgs ...interface{}) {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func NewProtocol(protocolname string, protocolversion uint, run func(*Peer) error, na adapters.NodeAdapter, ct *CodeMap, peerInfo func(id discover.NodeID) interface{}, nodeInfo func() interface{}) *p2p.Protocol {
|
func NewProtocol(protocolname string, protocolversion uint, run func(*Peer) error, ct *CodeMap, peerInfo func(id discover.NodeID) interface{}, nodeInfo func() interface{}) *p2p.Protocol {
|
||||||
|
|
||||||
// PeerInfo is an optional helper method to retrieve protocol specific metadata
|
// PeerInfo is an optional helper method to retrieve protocol specific metadata
|
||||||
// about a certain peer in the network. If an info retrieval function is set,
|
// about a certain peer in the network. If an info retrieval function is set,
|
||||||
// but returns nil, it is assumed that the protocol handshake is still running.
|
// but returns nil, it is assumed that the protocol handshake is still running.
|
||||||
r := func(p *p2p.Peer, rw p2p.MsgReadWriter) error {
|
r := func(p *p2p.Peer, rw p2p.MsgReadWriter) error {
|
||||||
|
return run(NewPeer(p, ct, rw))
|
||||||
m := na.Messenger(rw)
|
|
||||||
|
|
||||||
peer := NewPeer(p, ct, m)
|
|
||||||
return run(peer)
|
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -192,7 +187,6 @@ type Disconnect struct {
|
||||||
// a remote peer
|
// a remote peer
|
||||||
type Peer struct {
|
type Peer struct {
|
||||||
ct *CodeMap // CodeMap for the protocol
|
ct *CodeMap // CodeMap for the protocol
|
||||||
m adapters.Messenger // defines senf and receive
|
|
||||||
*p2p.Peer // the p2p.Peer object representing the remote
|
*p2p.Peer // the p2p.Peer object representing the remote
|
||||||
rw p2p.MsgReadWriter // p2p.MsgReadWriter to send messages to and read messages from
|
rw p2p.MsgReadWriter // p2p.MsgReadWriter to send messages to and read messages from
|
||||||
handlers map[reflect.Type][]func(interface{}) error // message type -> message handler callback(s) map
|
handlers map[reflect.Type][]func(interface{}) error // message type -> message handler callback(s) map
|
||||||
|
|
@ -204,11 +198,11 @@ type Peer struct {
|
||||||
// this constructor is called by the p2p.Protocol#Run function
|
// this constructor is called by the p2p.Protocol#Run function
|
||||||
// the first two arguments are comming the arguments passed to p2p.Protocol.Run function
|
// the first two arguments are comming the arguments passed to p2p.Protocol.Run function
|
||||||
// the third argument is the CodeMap describing the protocol messages and options
|
// the third argument is the CodeMap describing the protocol messages and options
|
||||||
func NewPeer(p *p2p.Peer, ct *CodeMap, m adapters.Messenger) *Peer {
|
func NewPeer(p *p2p.Peer, ct *CodeMap, rw p2p.MsgReadWriter) *Peer {
|
||||||
return &Peer{
|
return &Peer{
|
||||||
ct: ct,
|
ct: ct,
|
||||||
m: m,
|
|
||||||
Peer: p,
|
Peer: p,
|
||||||
|
rw: rw,
|
||||||
Errc: make(chan error),
|
Errc: make(chan error),
|
||||||
wErrc: make(chan error),
|
wErrc: make(chan error),
|
||||||
handlers: make(map[reflect.Type][]func(interface{}) error),
|
handlers: make(map[reflect.Type][]func(interface{}) error),
|
||||||
|
|
@ -276,7 +270,7 @@ func (self *Peer) Send(msg interface{}) error {
|
||||||
}
|
}
|
||||||
glog.V(logger.Detail).Infof("=> msg #%d TO %v : %v", code, self.ID(), msg)
|
glog.V(logger.Detail).Infof("=> msg #%d TO %v : %v", code, self.ID(), msg)
|
||||||
go func() {
|
go func() {
|
||||||
self.wErrc <- self.m.SendMsg(uint64(code), msg)
|
self.wErrc <- p2p.Send(self.rw, uint64(code), msg)
|
||||||
}()
|
}()
|
||||||
var err error
|
var err error
|
||||||
select {
|
select {
|
||||||
|
|
@ -306,7 +300,7 @@ func (self *Peer) DisconnectHook(f func(error)) {
|
||||||
// checks message size, out-of-range message codes, handles decoding with reflection,
|
// checks message size, out-of-range message codes, handles decoding with reflection,
|
||||||
// call handlers as callback onside
|
// call handlers as callback onside
|
||||||
func (self *Peer) handleIncoming() (interface{}, error) {
|
func (self *Peer) handleIncoming() (interface{}, error) {
|
||||||
msg, err := self.m.ReadMsg()
|
msg, err := self.rw.ReadMsg()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -2,7 +2,6 @@ package protocols
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"fmt"
|
"fmt"
|
||||||
"sync"
|
|
||||||
"testing"
|
"testing"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
|
|
@ -57,14 +56,11 @@ const networkId = "420"
|
||||||
// newProtocol sets up a protocol
|
// newProtocol sets up a protocol
|
||||||
// the run function here demonstrates a typical protocol using peerPool, handshake
|
// the run function here demonstrates a typical protocol using peerPool, handshake
|
||||||
// and messages registered to handlers
|
// and messages registered to handlers
|
||||||
func newProtocol(pp *p2ptest.TestPeerPool, wg *sync.WaitGroup) func(adapters.NodeAdapter) adapters.ProtoCall {
|
func newProtocol(pp *p2ptest.TestPeerPool) func(adapters.NodeAdapter) adapters.ProtoCall {
|
||||||
ct := NewCodeMap("test", 42, 1024, &protoHandshake{}, &hs0{}, &kill{}, &drop{})
|
ct := NewCodeMap("test", 42, 1024, &protoHandshake{}, &hs0{}, &kill{}, &drop{})
|
||||||
return func(na adapters.NodeAdapter) adapters.ProtoCall {
|
return func(na adapters.NodeAdapter) adapters.ProtoCall {
|
||||||
return func(p *p2p.Peer, rw p2p.MsgReadWriter) error {
|
return func(p *p2p.Peer, rw p2p.MsgReadWriter) error {
|
||||||
if wg != nil {
|
peer := NewPeer(p, ct, rw)
|
||||||
wg.Add(1)
|
|
||||||
}
|
|
||||||
peer := NewPeer(p, ct, na.Messenger(rw))
|
|
||||||
|
|
||||||
// demonstrates use of peerPool, killing another peer connection as a response to a message
|
// demonstrates use of peerPool, killing another peer connection as a response to a message
|
||||||
peer.Register(&kill{}, func(msg interface{}) error {
|
peer.Register(&kill{}, func(msg interface{}) error {
|
||||||
|
|
@ -116,9 +112,6 @@ func newProtocol(pp *p2ptest.TestPeerPool, wg *sync.WaitGroup) func(adapters.Nod
|
||||||
pp.Add(peer)
|
pp.Add(peer)
|
||||||
defer pp.Remove(peer)
|
defer pp.Remove(peer)
|
||||||
err = peer.Run()
|
err = peer.Run()
|
||||||
if wg != nil {
|
|
||||||
wg.Done()
|
|
||||||
}
|
|
||||||
glog.V(logger.Detail).Infof("peer %v protocol quitting: %v", peer, err)
|
glog.V(logger.Detail).Infof("peer %v protocol quitting: %v", peer, err)
|
||||||
|
|
||||||
return err
|
return err
|
||||||
|
|
@ -126,9 +119,9 @@ func newProtocol(pp *p2ptest.TestPeerPool, wg *sync.WaitGroup) func(adapters.Nod
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func protocolTester(t *testing.T, pp *p2ptest.TestPeerPool, wg *sync.WaitGroup) *p2ptest.ProtocolTester {
|
func protocolTester(t *testing.T, pp *p2ptest.TestPeerPool) *p2ptest.ProtocolTester {
|
||||||
id := adapters.RandomNodeId()
|
id := adapters.RandomNodeId()
|
||||||
return p2ptest.NewProtocolTester(t, id, 2, newProtocol(pp, wg))
|
return p2ptest.NewProtocolTester(t, id, 2, newProtocol(pp))
|
||||||
}
|
}
|
||||||
|
|
||||||
func protoHandshakeExchange(id *adapters.NodeId, proto *protoHandshake) []p2ptest.Exchange {
|
func protoHandshakeExchange(id *adapters.NodeId, proto *protoHandshake) []p2ptest.Exchange {
|
||||||
|
|
@ -157,7 +150,7 @@ func protoHandshakeExchange(id *adapters.NodeId, proto *protoHandshake) []p2ptes
|
||||||
|
|
||||||
func runProtoHandshake(t *testing.T, proto *protoHandshake, errs ...error) {
|
func runProtoHandshake(t *testing.T, proto *protoHandshake, errs ...error) {
|
||||||
pp := p2ptest.NewTestPeerPool()
|
pp := p2ptest.NewTestPeerPool()
|
||||||
s := protocolTester(t, pp, nil)
|
s := protocolTester(t, pp)
|
||||||
// TODO: make this more than one handshake
|
// TODO: make this more than one handshake
|
||||||
id := s.Ids[0]
|
id := s.Ids[0]
|
||||||
s.TestExchanges(protoHandshakeExchange(id, proto)...)
|
s.TestExchanges(protoHandshakeExchange(id, proto)...)
|
||||||
|
|
@ -206,7 +199,7 @@ func moduleHandshakeExchange(id *adapters.NodeId, resp uint) []p2ptest.Exchange
|
||||||
|
|
||||||
func runModuleHandshake(t *testing.T, resp uint, errs ...error) {
|
func runModuleHandshake(t *testing.T, resp uint, errs ...error) {
|
||||||
pp := p2ptest.NewTestPeerPool()
|
pp := p2ptest.NewTestPeerPool()
|
||||||
s := protocolTester(t, pp, nil)
|
s := protocolTester(t, pp)
|
||||||
id := s.Ids[0]
|
id := s.Ids[0]
|
||||||
s.TestExchanges(protoHandshakeExchange(id, &protoHandshake{42, "420"})...)
|
s.TestExchanges(protoHandshakeExchange(id, &protoHandshake{42, "420"})...)
|
||||||
s.TestExchanges(moduleHandshakeExchange(id, resp)...)
|
s.TestExchanges(moduleHandshakeExchange(id, resp)...)
|
||||||
|
|
@ -279,9 +272,8 @@ func testMultiPeerSetup(a, b *adapters.NodeId) []p2ptest.Exchange {
|
||||||
}
|
}
|
||||||
|
|
||||||
func runMultiplePeers(t *testing.T, peer int, errs ...error) {
|
func runMultiplePeers(t *testing.T, peer int, errs ...error) {
|
||||||
wg := &sync.WaitGroup{}
|
|
||||||
pp := p2ptest.NewTestPeerPool()
|
pp := p2ptest.NewTestPeerPool()
|
||||||
s := protocolTester(t, pp, wg)
|
s := protocolTester(t, pp)
|
||||||
|
|
||||||
s.TestExchanges(testMultiPeerSetup(s.Ids[0], s.Ids[1])...)
|
s.TestExchanges(testMultiPeerSetup(s.Ids[0], s.Ids[1])...)
|
||||||
// after some exchanges of messages, we can test state changes
|
// after some exchanges of messages, we can test state changes
|
||||||
|
|
@ -320,7 +312,6 @@ func runMultiplePeers(t *testing.T, peer int, errs ...error) {
|
||||||
},
|
},
|
||||||
},
|
},
|
||||||
})
|
})
|
||||||
wg.Wait()
|
|
||||||
// check the actual discconnect errors on the individual peers
|
// check the actual discconnect errors on the individual peers
|
||||||
var disconnects []*p2ptest.Disconnect
|
var disconnects []*p2ptest.Disconnect
|
||||||
for i, err := range errs {
|
for i, err := range errs {
|
||||||
|
|
|
||||||
|
|
@ -34,7 +34,6 @@ import (
|
||||||
"github.com/ethereum/go-ethereum/event"
|
"github.com/ethereum/go-ethereum/event"
|
||||||
"github.com/ethereum/go-ethereum/logger"
|
"github.com/ethereum/go-ethereum/logger"
|
||||||
"github.com/ethereum/go-ethereum/logger/glog"
|
"github.com/ethereum/go-ethereum/logger/glog"
|
||||||
"github.com/ethereum/go-ethereum/p2p"
|
|
||||||
"github.com/ethereum/go-ethereum/p2p/adapters"
|
"github.com/ethereum/go-ethereum/p2p/adapters"
|
||||||
"github.com/ethereum/go-ethereum/p2p/discover"
|
"github.com/ethereum/go-ethereum/p2p/discover"
|
||||||
)
|
)
|
||||||
|
|
@ -181,15 +180,14 @@ func NewDebugController(journal *Journal) Controller {
|
||||||
// messaging is implemented in the particular NodeAdapter interface
|
// messaging is implemented in the particular NodeAdapter interface
|
||||||
type Network struct {
|
type Network struct {
|
||||||
// input trigger events and other events
|
// input trigger events and other events
|
||||||
events *event.TypeMux // generated events a journal can subsribe to
|
events *event.TypeMux // generated events a journal can subsribe to
|
||||||
lock sync.RWMutex
|
lock sync.RWMutex
|
||||||
nodeMap map[discover.NodeID]int
|
nodeMap map[discover.NodeID]int
|
||||||
connMap map[string]int
|
connMap map[string]int
|
||||||
Nodes []*Node `json:"nodes"`
|
Nodes []*Node `json:"nodes"`
|
||||||
Conns []*Conn `json:"conns"`
|
Conns []*Conn `json:"conns"`
|
||||||
messenger func(p2p.MsgReadWriter) adapters.Messenger
|
quitc chan bool
|
||||||
quitc chan bool
|
conf *NetworkConfig
|
||||||
conf *NetworkConfig
|
|
||||||
//
|
//
|
||||||
// adapters.Messenger
|
// adapters.Messenger
|
||||||
// node adapter function that creates the node model for
|
// node adapter function that creates the node model for
|
||||||
|
|
@ -199,12 +197,11 @@ type Network struct {
|
||||||
|
|
||||||
func NewNetwork(conf *NetworkConfig) *Network {
|
func NewNetwork(conf *NetworkConfig) *Network {
|
||||||
return &Network{
|
return &Network{
|
||||||
conf: conf,
|
conf: conf,
|
||||||
events: &event.TypeMux{},
|
events: &event.TypeMux{},
|
||||||
nodeMap: make(map[discover.NodeID]int),
|
nodeMap: make(map[discover.NodeID]int),
|
||||||
connMap: make(map[string]int),
|
connMap: make(map[string]int),
|
||||||
messenger: adapters.NewSimPipe,
|
quitc: make(chan bool),
|
||||||
quitc: make(chan bool),
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -430,7 +427,7 @@ func (self *Network) Config() *NetworkConfig {
|
||||||
}
|
}
|
||||||
func (self *Network) NewSimNode(conf *NodeConfig) adapters.NodeAdapter {
|
func (self *Network) NewSimNode(conf *NodeConfig) adapters.NodeAdapter {
|
||||||
id := conf.Id
|
id := conf.Id
|
||||||
na := adapters.NewSimNode(id, self, self.messenger)
|
na := adapters.NewSimNode(id, self)
|
||||||
return na
|
return na
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -7,6 +7,7 @@ import (
|
||||||
|
|
||||||
"github.com/ethereum/go-ethereum/logger"
|
"github.com/ethereum/go-ethereum/logger"
|
||||||
"github.com/ethereum/go-ethereum/logger/glog"
|
"github.com/ethereum/go-ethereum/logger/glog"
|
||||||
|
"github.com/ethereum/go-ethereum/p2p"
|
||||||
"github.com/ethereum/go-ethereum/p2p/adapters"
|
"github.com/ethereum/go-ethereum/p2p/adapters"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
@ -71,15 +72,14 @@ func (self *ProtocolSession) trigger(trig Trigger) error {
|
||||||
if peer == nil {
|
if peer == nil {
|
||||||
panic(fmt.Sprintf("trigger: peer %v does not exist (1- %v)", trig.Peer, len(self.Ids)))
|
panic(fmt.Sprintf("trigger: peer %v does not exist (1- %v)", trig.Peer, len(self.Ids)))
|
||||||
}
|
}
|
||||||
m := peer.Messenger
|
if peer.MsgPipeRW == nil {
|
||||||
if m == nil {
|
|
||||||
return fmt.Errorf("trigger: peer %v unreachable", trig.Peer)
|
return fmt.Errorf("trigger: peer %v unreachable", trig.Peer)
|
||||||
}
|
}
|
||||||
errc := make(chan error)
|
errc := make(chan error)
|
||||||
|
|
||||||
go func() {
|
go func() {
|
||||||
glog.V(logger.Detail).Infof("trigger %v (%v)....", trig.Msg, trig.Code)
|
glog.V(logger.Detail).Infof("trigger %v (%v)....", trig.Msg, trig.Code)
|
||||||
errc <- m.(TestMessenger).TriggerMsg(trig.Code, trig.Msg)
|
errc <- p2p.Send(peer, trig.Code, trig.Msg)
|
||||||
glog.V(logger.Detail).Infof("triggered %v (%v)", trig.Msg, trig.Code)
|
glog.V(logger.Detail).Infof("triggered %v (%v)", trig.Msg, trig.Code)
|
||||||
}()
|
}()
|
||||||
|
|
||||||
|
|
@ -105,15 +105,14 @@ func (self *ProtocolSession) expect(exp Expect) error {
|
||||||
if peer == nil {
|
if peer == nil {
|
||||||
panic(fmt.Sprintf("expect: peer %v does not exist (1- %v)", exp.Peer, len(self.Ids)))
|
panic(fmt.Sprintf("expect: peer %v does not exist (1- %v)", exp.Peer, len(self.Ids)))
|
||||||
}
|
}
|
||||||
m := peer.Messenger
|
if peer.MsgPipeRW == nil {
|
||||||
if m == nil {
|
|
||||||
return fmt.Errorf("trigger: peer %v unreachable", exp.Peer)
|
return fmt.Errorf("trigger: peer %v unreachable", exp.Peer)
|
||||||
}
|
}
|
||||||
|
|
||||||
errc := make(chan error)
|
errc := make(chan error)
|
||||||
go func() {
|
go func() {
|
||||||
glog.V(logger.Detail).Infof("waiting for msg, %v", exp.Msg)
|
glog.V(logger.Detail).Infof("waiting for msg, %v", exp.Msg)
|
||||||
errc <- m.(TestMessenger).ExpectMsg(exp.Code, exp.Msg)
|
errc <- p2p.ExpectMsg(peer, exp.Code, exp.Msg)
|
||||||
}()
|
}()
|
||||||
|
|
||||||
t := exp.Timeout
|
t := exp.Timeout
|
||||||
|
|
@ -210,7 +209,7 @@ func (self *ProtocolSession) TestDisconnected(disconnects ...*Disconnect) error
|
||||||
func (self *ProtocolSession) Stop() {
|
func (self *ProtocolSession) Stop() {
|
||||||
for _, id := range self.Ids {
|
for _, id := range self.Ids {
|
||||||
p := self.GetPeer(id)
|
p := self.GetPeer(id)
|
||||||
if p != nil && p.Messenger != nil {
|
if p != nil && p.MsgPipeRW != nil {
|
||||||
p.Close()
|
p.Close()
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -17,10 +17,9 @@ type ProtocolTester struct {
|
||||||
|
|
||||||
func NewProtocolTester(t *testing.T, id *adapters.NodeId, n int, run func(id adapters.NodeAdapter) adapters.ProtoCall) *ProtocolTester {
|
func NewProtocolTester(t *testing.T, id *adapters.NodeId, n int, run func(id adapters.NodeAdapter) adapters.ProtoCall) *ProtocolTester {
|
||||||
|
|
||||||
simPipe := adapters.NewSimPipe
|
|
||||||
net := simulations.NewNetwork(&simulations.NetworkConfig{})
|
net := simulations.NewNetwork(&simulations.NetworkConfig{})
|
||||||
naf := func(conf *simulations.NodeConfig) adapters.NodeAdapter {
|
naf := func(conf *simulations.NodeConfig) adapters.NodeAdapter {
|
||||||
na := adapters.NewSimNode(conf.Id, net, simPipe)
|
na := adapters.NewSimNode(conf.Id, net)
|
||||||
if conf.Id.NodeID == id.NodeID {
|
if conf.Id.NodeID == id.NodeID {
|
||||||
glog.V(logger.Detail).Infof("adapter run function set to protocol for node %v (=%v)", conf.Id, id)
|
glog.V(logger.Detail).Infof("adapter run function set to protocol for node %v (=%v)", conf.Id, id)
|
||||||
na.Run = run(na)
|
na.Run = run(na)
|
||||||
|
|
|
||||||
|
|
@ -108,7 +108,7 @@ func Bzz(localAddr []byte, na adapters.NodeAdapter, ct *protocols.CodeMap, servi
|
||||||
return bee.Run()
|
return bee.Run()
|
||||||
}
|
}
|
||||||
|
|
||||||
return protocols.NewProtocol(ProtocolName, Version, run, na, ct, peerInfo, nodeInfo)
|
return protocols.NewProtocol(ProtocolName, Version, run, ct, peerInfo, nodeInfo)
|
||||||
}
|
}
|
||||||
|
|
||||||
/*
|
/*
|
||||||
|
|
|
||||||
|
|
@ -170,7 +170,6 @@ func newPssBaseTester(t *testing.T, addr *peerAddr, n int) *pssTester {
|
||||||
ct.Register(&getPeersMsg{})
|
ct.Register(&getPeersMsg{})
|
||||||
ct.Register(&subPeersMsg{}) // why is this public?
|
ct.Register(&subPeersMsg{}) // why is this public?
|
||||||
|
|
||||||
simPipe := adapters.NewSimPipe
|
|
||||||
kp := NewKadParams()
|
kp := NewKadParams()
|
||||||
kp.MinProxBinSize = 3
|
kp.MinProxBinSize = 3
|
||||||
to := NewKademlia(addr.OverlayAddr(), kp)
|
to := NewKademlia(addr.OverlayAddr(), kp)
|
||||||
|
|
@ -178,7 +177,7 @@ func newPssBaseTester(t *testing.T, addr *peerAddr, n int) *pssTester {
|
||||||
ps := NewPss(to, addr.OverlayAddr())
|
ps := NewPss(to, addr.OverlayAddr())
|
||||||
net := simulations.NewNetwork(&simulations.NetworkConfig{})
|
net := simulations.NewNetwork(&simulations.NetworkConfig{})
|
||||||
naf := func(conf *simulations.NodeConfig) adapters.NodeAdapter {
|
naf := func(conf *simulations.NodeConfig) adapters.NodeAdapter {
|
||||||
na := adapters.NewSimNode(conf.Id, net, simPipe)
|
na := adapters.NewSimNode(conf.Id, net)
|
||||||
return na
|
return na
|
||||||
}
|
}
|
||||||
net.SetNaf(naf)
|
net.SetNaf(naf)
|
||||||
|
|
|
||||||
|
|
@ -109,7 +109,7 @@ func newNode(id *adapters.NodeId, net *simulations.Network, trigger chan *adapte
|
||||||
kademlia := newKademlia(addr.OverlayAddr())
|
kademlia := newKademlia(addr.OverlayAddr())
|
||||||
hive := newHive(kademlia)
|
hive := newHive(kademlia)
|
||||||
codeMap := network.BzzCodeMap(network.DiscoveryMsgs...)
|
codeMap := network.BzzCodeMap(network.DiscoveryMsgs...)
|
||||||
nodeAdapter := adapters.NewSimNode(id, net, adapters.NewSimPipe)
|
nodeAdapter := adapters.NewSimNode(id, net)
|
||||||
node := &node{
|
node := &node{
|
||||||
Hive: hive,
|
Hive: hive,
|
||||||
NodeAdapter: nodeAdapter,
|
NodeAdapter: nodeAdapter,
|
||||||
|
|
|
||||||
|
|
@ -57,7 +57,7 @@ func (self *SimNode) RunProtocol(id *adapters.NodeId, rw, rrw p2p.MsgReadWriter,
|
||||||
// NewSimNode creates adapters for nodes in the simulation.
|
// NewSimNode creates adapters for nodes in the simulation.
|
||||||
func (self *Network) NewSimNode(conf *simulations.NodeConfig) adapters.NodeAdapter {
|
func (self *Network) NewSimNode(conf *simulations.NodeConfig) adapters.NodeAdapter {
|
||||||
id := conf.Id
|
id := conf.Id
|
||||||
na := adapters.NewSimNode(id, self.Network, adapters.NewSimPipe)
|
na := adapters.NewSimNode(id, self.Network)
|
||||||
addr := network.NewPeerAddrFromNodeId(id)
|
addr := network.NewPeerAddrFromNodeId(id)
|
||||||
kp := network.NewKadParams()
|
kp := network.NewKadParams()
|
||||||
|
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue