mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-08-12 15:03:45 +00:00
all: Fix logging
Signed-off-by: Lewis Marshall <lewis@lmars.net>
This commit is contained in:
parent
bcf2622f29
commit
4c6648ef49
25 changed files with 177 additions and 194 deletions
|
|
@ -20,7 +20,7 @@ import (
|
||||||
"fmt"
|
"fmt"
|
||||||
"sync"
|
"sync"
|
||||||
|
|
||||||
"github.com/ethereum/go-ethereum/logger/glog"
|
"github.com/ethereum/go-ethereum/log"
|
||||||
"github.com/ethereum/go-ethereum/p2p"
|
"github.com/ethereum/go-ethereum/p2p"
|
||||||
"github.com/ethereum/go-ethereum/p2p/discover"
|
"github.com/ethereum/go-ethereum/p2p/discover"
|
||||||
)
|
)
|
||||||
|
|
@ -121,7 +121,7 @@ func (self *SimNode) Disconnect(rid []byte) error {
|
||||||
// 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
|
||||||
glog.V(6).Infof("dropped peer %v", id)
|
log.Trace(fmt.Sprintf("dropped peer %v", id))
|
||||||
|
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
@ -158,10 +158,10 @@ func (self *SimNode) Connect(rid []byte) error {
|
||||||
|
|
||||||
func (self *SimNode) RunProtocol(id *NodeId, rw, rrw p2p.MsgReadWriter, peer *Peer) error {
|
func (self *SimNode) RunProtocol(id *NodeId, rw, rrw p2p.MsgReadWriter, peer *Peer) error {
|
||||||
if self.Run == nil {
|
if self.Run == nil {
|
||||||
glog.V(6).Infof("no protocol starting on peer %v (connection with %v)", self.Id, id)
|
log.Trace(fmt.Sprintf("no protocol starting on peer %v (connection with %v)", self.Id, id))
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
glog.V(6).Infof("protocol starting on peer %v (connection with %v)", self.Id, id)
|
log.Trace(fmt.Sprintf("protocol starting on peer %v (connection with %v)", self.Id, id))
|
||||||
p := p2p.NewPeer(id.NodeID, Name(id.Bytes()), []p2p.Cap{})
|
p := p2p.NewPeer(id.NodeID, Name(id.Bytes()), []p2p.Cap{})
|
||||||
go func() {
|
go func() {
|
||||||
self.network.DidConnect(self.Id, id)
|
self.network.DidConnect(self.Id, id)
|
||||||
|
|
@ -169,7 +169,7 @@ func (self *SimNode) RunProtocol(id *NodeId, rw, rrw p2p.MsgReadWriter, peer *Pe
|
||||||
<-peer.Readyc
|
<-peer.Readyc
|
||||||
self.Disconnect(id.Bytes())
|
self.Disconnect(id.Bytes())
|
||||||
peer.Errc <- err
|
peer.Errc <- err
|
||||||
glog.V(6).Infof("protocol quit on peer %v (connection with %v broken: %v)", self.Id, id, err)
|
log.Trace(fmt.Sprintf("protocol quit on peer %v (connection with %v broken: %v)", self.Id, id, err))
|
||||||
self.network.DidDisconnect(self.Id, id)
|
self.network.DidDisconnect(self.Id, id)
|
||||||
}()
|
}()
|
||||||
return nil
|
return nil
|
||||||
|
|
|
||||||
|
|
@ -34,8 +34,7 @@ import (
|
||||||
"reflect"
|
"reflect"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/ethereum/go-ethereum/logger"
|
"github.com/ethereum/go-ethereum/log"
|
||||||
"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/discover"
|
"github.com/ethereum/go-ethereum/p2p/discover"
|
||||||
)
|
)
|
||||||
|
|
@ -223,7 +222,7 @@ func (self *Peer) Register(msg interface{}, handler func(interface{}) error) uin
|
||||||
if !found {
|
if !found {
|
||||||
panic(fmt.Sprintf("message type '%v' unknown ", typ))
|
panic(fmt.Sprintf("message type '%v' unknown ", typ))
|
||||||
}
|
}
|
||||||
glog.V(logger.Detail).Infof("register handle for %v", typ)
|
log.Trace(fmt.Sprintf("register handle for %v", typ))
|
||||||
self.handlers[typ] = append(self.handlers[typ], handler)
|
self.handlers[typ] = append(self.handlers[typ], handler)
|
||||||
return code
|
return code
|
||||||
}
|
}
|
||||||
|
|
@ -243,7 +242,7 @@ func (self *Peer) Run() error {
|
||||||
err := <-self.Errc
|
err := <-self.Errc
|
||||||
d := &Disconnect{err}
|
d := &Disconnect{err}
|
||||||
for _, f := range self.handlers[reflect.TypeOf(d)] {
|
for _, f := range self.handlers[reflect.TypeOf(d)] {
|
||||||
glog.V(6).Infof("disconnect hook for %v", d)
|
log.Trace(fmt.Sprintf("disconnect hook for %v", d))
|
||||||
f(err)
|
f(err)
|
||||||
}
|
}
|
||||||
return err
|
return err
|
||||||
|
|
@ -268,7 +267,7 @@ func (self *Peer) Send(msg interface{}) error {
|
||||||
if !found {
|
if !found {
|
||||||
return errorf(ErrInvalidMsgType, "%v", code)
|
return errorf(ErrInvalidMsgType, "%v", code)
|
||||||
}
|
}
|
||||||
glog.V(logger.Detail).Infof("=> msg #%d TO %v : %v", code, self.ID(), msg)
|
log.Trace(fmt.Sprintf("=> msg #%d TO %v : %v", code, self.ID(), msg))
|
||||||
go func() {
|
go func() {
|
||||||
self.wErrc <- p2p.Send(self.rw, uint64(code), msg)
|
self.wErrc <- p2p.Send(self.rw, uint64(code), msg)
|
||||||
}()
|
}()
|
||||||
|
|
@ -304,7 +303,7 @@ func (self *Peer) handleIncoming() (interface{}, error) {
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
glog.V(logger.Detail).Infof("<= %v", msg)
|
log.Trace(fmt.Sprintf("<= %v", msg))
|
||||||
// make sure that the payload has been fully consumed
|
// make sure that the payload has been fully consumed
|
||||||
defer msg.Discard()
|
defer msg.Discard()
|
||||||
|
|
||||||
|
|
@ -326,7 +325,7 @@ func (self *Peer) handleIncoming() (interface{}, error) {
|
||||||
if err := msg.Decode(val.Interface()); err != nil {
|
if err := msg.Decode(val.Interface()); err != nil {
|
||||||
return nil, errorf(ErrDecode, "<= %v: %v", msg, err)
|
return nil, errorf(ErrDecode, "<= %v: %v", msg, err)
|
||||||
}
|
}
|
||||||
glog.V(logger.Detail).Infof("<= %v FROM %v %v %v", msg, self.ID(), req, typ)
|
log.Trace(fmt.Sprintf("<= %v FROM %v %v %v", msg, self.ID(), req, typ))
|
||||||
|
|
||||||
// call the registered handler callbacks
|
// call the registered handler callbacks
|
||||||
// a registered callback take the decoded message as argument as an interface
|
// a registered callback take the decoded message as argument as an interface
|
||||||
|
|
@ -335,11 +334,11 @@ func (self *Peer) handleIncoming() (interface{}, error) {
|
||||||
// chosen based on the proper type in the first place
|
// chosen based on the proper type in the first place
|
||||||
handlers := self.handlers[typ]
|
handlers := self.handlers[typ]
|
||||||
if len(handlers) == 0 {
|
if len(handlers) == 0 {
|
||||||
glog.V(6).Infof("no handler (msg code %v)", msg.Code)
|
log.Trace(fmt.Sprintf("no handler (msg code %v)", msg.Code))
|
||||||
// return nil, errorf(ErrNoHandler, "(msg code %v)", msg.Code)
|
// return nil, errorf(ErrNoHandler, "(msg code %v)", msg.Code)
|
||||||
} else {
|
} else {
|
||||||
for i, f := range handlers {
|
for i, f := range handlers {
|
||||||
glog.V(6).Infof("handler %v for %v", i, typ)
|
log.Trace(fmt.Sprintf("handler %v for %v", i, typ))
|
||||||
err = f(req.Interface())
|
err = f(req.Interface())
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, errorf(ErrHandler, "(msg code %v): %v", msg.Code, err)
|
return nil, errorf(ErrHandler, "(msg code %v): %v", msg.Code, err)
|
||||||
|
|
|
||||||
|
|
@ -2,19 +2,18 @@ package protocols
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"fmt"
|
"fmt"
|
||||||
|
"os"
|
||||||
"testing"
|
"testing"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/ethereum/go-ethereum/logger"
|
"github.com/ethereum/go-ethereum/log"
|
||||||
"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/adapters"
|
||||||
p2ptest "github.com/ethereum/go-ethereum/p2p/testing"
|
p2ptest "github.com/ethereum/go-ethereum/p2p/testing"
|
||||||
)
|
)
|
||||||
|
|
||||||
func init() {
|
func init() {
|
||||||
glog.SetV(logger.Error)
|
log.Root().SetHandler(log.LvlFilterHandler(log.LvlError, log.StreamHandler(os.Stderr, log.TerminalFormat(false))))
|
||||||
glog.SetToStderr(true)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// handshake message type
|
// handshake message type
|
||||||
|
|
@ -65,13 +64,13 @@ func newProtocol(pp *p2ptest.TestPeerPool) adapters.ProtoCall {
|
||||||
peer.Register(&kill{}, func(msg interface{}) error {
|
peer.Register(&kill{}, func(msg interface{}) error {
|
||||||
id := msg.(*kill).C
|
id := msg.(*kill).C
|
||||||
pp.Get(id).Drop(fmt.Errorf("killed"))
|
pp.Get(id).Drop(fmt.Errorf("killed"))
|
||||||
glog.V(logger.Detail).Infof("id %v killed", id)
|
log.Trace(fmt.Sprintf("id %v killed", id))
|
||||||
return nil
|
return nil
|
||||||
})
|
})
|
||||||
|
|
||||||
// for testing we can trigger self induced disconnect upon receiving drop message
|
// for testing we can trigger self induced disconnect upon receiving drop message
|
||||||
peer.Register(&drop{}, func(msg interface{}) error {
|
peer.Register(&drop{}, func(msg interface{}) error {
|
||||||
glog.V(logger.Detail).Infof("dropped")
|
log.Trace("dropped")
|
||||||
return fmt.Errorf("dropped")
|
return fmt.Errorf("dropped")
|
||||||
})
|
})
|
||||||
|
|
||||||
|
|
@ -107,11 +106,11 @@ func newProtocol(pp *p2ptest.TestPeerPool) adapters.ProtoCall {
|
||||||
return peer.Send(lhs)
|
return peer.Send(lhs)
|
||||||
})
|
})
|
||||||
|
|
||||||
glog.V(logger.Detail).Infof("adding peer %v", peer)
|
log.Trace(fmt.Sprintf("adding peer %v", peer))
|
||||||
pp.Add(peer)
|
pp.Add(peer)
|
||||||
defer pp.Remove(peer)
|
defer pp.Remove(peer)
|
||||||
err = peer.Run()
|
err = peer.Run()
|
||||||
glog.V(logger.Detail).Infof("peer %v protocol quitting: %v", peer, err)
|
log.Trace(fmt.Sprintf("peer %v protocol quitting: %v", peer, err))
|
||||||
|
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
@ -280,7 +279,7 @@ func runMultiplePeers(t *testing.T, peer int, errs ...error) {
|
||||||
// time.Sleep(1)
|
// time.Sleep(1)
|
||||||
for !pp.Has(s.Ids[0]) {
|
for !pp.Has(s.Ids[0]) {
|
||||||
time.Sleep(1)
|
time.Sleep(1)
|
||||||
glog.V(logger.Detail).Infof("missing peer test-0: %v (%v)", pp, s.Ids)
|
log.Trace(fmt.Sprintf("missing peer test-0: %v (%v)", pp, s.Ids))
|
||||||
}
|
}
|
||||||
// if !pp.Has(s.Ids[0]) {
|
// if !pp.Has(s.Ids[0]) {
|
||||||
// t.Fatalf("missing peer test-0: %v (%v)", pp, s.Ids)
|
// t.Fatalf("missing peer test-0: %v (%v)", pp, s.Ids)
|
||||||
|
|
|
||||||
|
|
@ -1,16 +1,17 @@
|
||||||
package main
|
package main
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"os"
|
||||||
"runtime"
|
"runtime"
|
||||||
|
|
||||||
"github.com/ethereum/go-ethereum/logger/glog"
|
"github.com/ethereum/go-ethereum/log"
|
||||||
"github.com/ethereum/go-ethereum/p2p/simulations"
|
"github.com/ethereum/go-ethereum/p2p/simulations"
|
||||||
)
|
)
|
||||||
|
|
||||||
func main() {
|
func main() {
|
||||||
runtime.GOMAXPROCS(runtime.NumCPU())
|
runtime.GOMAXPROCS(runtime.NumCPU())
|
||||||
glog.SetV(6)
|
|
||||||
glog.SetToStderr(true)
|
log.Root().SetHandler(log.LvlFilterHandler(log.LvlTrace, log.StreamHandler(os.Stderr, log.TerminalFormat(false))))
|
||||||
|
|
||||||
c, quitc := simulations.NewSessionController(simulations.DefaultNet)
|
c, quitc := simulations.NewSessionController(simulations.DefaultNet)
|
||||||
simulations.StartRestApiServer("8888", c)
|
simulations.StartRestApiServer("8888", c)
|
||||||
|
|
|
||||||
|
|
@ -9,8 +9,7 @@ import (
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/ethereum/go-ethereum/event"
|
"github.com/ethereum/go-ethereum/event"
|
||||||
"github.com/ethereum/go-ethereum/logger"
|
"github.com/ethereum/go-ethereum/log"
|
||||||
"github.com/ethereum/go-ethereum/logger/glog"
|
|
||||||
"github.com/ethereum/go-ethereum/p2p/adapters"
|
"github.com/ethereum/go-ethereum/p2p/adapters"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
@ -41,7 +40,7 @@ func NewJournal() *Journal {
|
||||||
// used for journalling history of a network
|
// used for journalling history of a network
|
||||||
// the goroutine terminates when the journal is closed
|
// the goroutine terminates when the journal is closed
|
||||||
func (self *Journal) Subscribe(eventer *event.TypeMux, types ...interface{}) {
|
func (self *Journal) Subscribe(eventer *event.TypeMux, types ...interface{}) {
|
||||||
glog.V(logger.Info).Infof("subscribe")
|
log.Info("subscribe")
|
||||||
sub := eventer.Subscribe(types...)
|
sub := eventer.Subscribe(types...)
|
||||||
go func() {
|
go func() {
|
||||||
defer sub.Unsubscribe()
|
defer sub.Unsubscribe()
|
||||||
|
|
@ -151,7 +150,7 @@ func (self *Journal) TimedRead(acc float64, f func(interface{}) bool) (read int)
|
||||||
panic("events not ordered")
|
panic("events not ordered")
|
||||||
}
|
}
|
||||||
interval := ev.Time.Sub(lastEvent)
|
interval := ev.Time.Sub(lastEvent)
|
||||||
glog.V(6).Infof("reset timer to interval %v", interval)
|
log.Trace(fmt.Sprintf("reset timer to interval %v", interval))
|
||||||
timer.Reset(time.Duration(acc) * interval)
|
timer.Reset(time.Duration(acc) * interval)
|
||||||
lastEvent = ev.Time
|
lastEvent = ev.Time
|
||||||
data = ev.Data
|
data = ev.Data
|
||||||
|
|
@ -171,7 +170,7 @@ func (self *Journal) TimedRead(acc float64, f func(interface{}) bool) (read int)
|
||||||
}
|
}
|
||||||
read += n
|
read += n
|
||||||
if n == 0 || !f(data) {
|
if n == 0 || !f(data) {
|
||||||
glog.V(6).Infof("timed read ends (read %v entries)", read)
|
log.Trace(fmt.Sprintf("timed read ends (read %v entries)", read))
|
||||||
break
|
break
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -192,7 +191,7 @@ func (self *Journal) reset(n int) {
|
||||||
if n >= length-1 {
|
if n >= length-1 {
|
||||||
n = length - 1
|
n = length - 1
|
||||||
}
|
}
|
||||||
glog.V(6).Infof("cursor reset from %v to %v/%v (%v)", self.cursor, n, len(self.Events), self.counter)
|
log.Trace(fmt.Sprintf("cursor reset from %v to %v/%v (%v)", self.cursor, n, len(self.Events), self.counter))
|
||||||
self.Events = self.Events[self.cursor:]
|
self.Events = self.Events[self.cursor:]
|
||||||
self.cursor = 0
|
self.cursor = 0
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -7,8 +7,7 @@ import (
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/ethereum/go-ethereum/event"
|
"github.com/ethereum/go-ethereum/event"
|
||||||
"github.com/ethereum/go-ethereum/logger"
|
"github.com/ethereum/go-ethereum/log"
|
||||||
"github.com/ethereum/go-ethereum/logger/glog"
|
|
||||||
"github.com/ethereum/go-ethereum/p2p/adapters"
|
"github.com/ethereum/go-ethereum/p2p/adapters"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
@ -26,7 +25,7 @@ func NewMockersController(eventer *event.TypeMux, defaultMockerConfig *MockerCon
|
||||||
if len(conf.Id) == 0 {
|
if len(conf.Id) == 0 {
|
||||||
conf.Id = fmt.Sprintf("%d", parent.id)
|
conf.Id = fmt.Sprintf("%d", parent.id)
|
||||||
}
|
}
|
||||||
glog.V(logger.Detail).Infof("new mocker controller on %v", conf.Id)
|
log.Trace(fmt.Sprintf("new mocker controller on %v", conf.Id))
|
||||||
if parent != nil {
|
if parent != nil {
|
||||||
parent.SetResource(conf.Id, c)
|
parent.SetResource(conf.Id, c)
|
||||||
parent.id++
|
parent.id++
|
||||||
|
|
@ -127,7 +126,7 @@ func MockEvents(eventer *event.TypeMux, ids []*adapters.NodeId, conf *MockerConf
|
||||||
|
|
||||||
rounds := 0
|
rounds := 0
|
||||||
for _ = range conf.ticker.C {
|
for _ = range conf.ticker.C {
|
||||||
glog.V(logger.Detail).Infof("rates: %v/%v, %v (%v/%v)", switchonRate, dropoutRate, newConnCount, connFailRate, disconnRate)
|
log.Trace(fmt.Sprintf("rates: %v/%v, %v (%v/%v)", switchonRate, dropoutRate, newConnCount, connFailRate, disconnRate))
|
||||||
// here switchon rate will depend
|
// here switchon rate will depend
|
||||||
nodesUp := len(offNodes) / switchonRate
|
nodesUp := len(offNodes) / switchonRate
|
||||||
missing := nodesTarget - len(onNodes)
|
missing := nodesTarget - len(onNodes)
|
||||||
|
|
@ -149,7 +148,7 @@ func MockEvents(eventer *event.TypeMux, ids []*adapters.NodeId, conf *MockerConf
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
connsDown := len(onConns) / disconnRate
|
connsDown := len(onConns) / disconnRate
|
||||||
glog.V(logger.Detail).Infof("Nodes Up: %v, Down: %v [ON: %v/%v]\nConns Up: %v, Down: %v [ON: %v/%v(%v)]", nodesUp, nodesDown, len(onNodes), len(onNodes)+len(offNodes), connsUp, connsDown, len(onConns), len(conns)-len(onConns), len(conns))
|
log.Trace(fmt.Sprintf("Nodes Up: %v, Down: %v [ON: %v/%v]\nConns Up: %v, Down: %v [ON: %v/%v(%v)]", nodesUp, nodesDown, len(onNodes), len(onNodes)+len(offNodes), connsUp, connsDown, len(onConns), len(conns)-len(onConns), len(conns)))
|
||||||
|
|
||||||
for i := 0; len(onNodes) > 0 && i < nodesDown; i++ {
|
for i := 0; len(onNodes) > 0 && i < nodesDown; i++ {
|
||||||
c := rand.Intn(len(onNodes))
|
c := rand.Intn(len(onNodes))
|
||||||
|
|
|
||||||
|
|
@ -32,8 +32,7 @@ import (
|
||||||
"sync"
|
"sync"
|
||||||
|
|
||||||
"github.com/ethereum/go-ethereum/event"
|
"github.com/ethereum/go-ethereum/event"
|
||||||
"github.com/ethereum/go-ethereum/logger"
|
"github.com/ethereum/go-ethereum/log"
|
||||||
"github.com/ethereum/go-ethereum/logger/glog"
|
|
||||||
"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"
|
||||||
)
|
)
|
||||||
|
|
@ -81,7 +80,7 @@ func NewNetworkController(net NetworkControl, nodesController *ResourceControlle
|
||||||
// GET /<networkId>/
|
// GET /<networkId>/
|
||||||
Retrieve: &ResourceHandler{
|
Retrieve: &ResourceHandler{
|
||||||
Handle: func(msg interface{}, parent *ResourceController) (interface{}, error) {
|
Handle: func(msg interface{}, parent *ResourceController) (interface{}, error) {
|
||||||
glog.V(logger.Detail).Infof("msg: %v", msg)
|
log.Trace(fmt.Sprintf("msg: %v", msg))
|
||||||
cyConfig, ok := msg.(*CyConfig)
|
cyConfig, ok := msg.(*CyConfig)
|
||||||
if ok {
|
if ok {
|
||||||
return UpdateCy(cyConfig, journal)
|
return UpdateCy(cyConfig, journal)
|
||||||
|
|
@ -213,7 +212,7 @@ func (self *Network) SetNaf(naf func(*NodeConfig) adapters.NodeAdapter) {
|
||||||
// and launches a goroutine that reads control events from an eventer Subsription channel
|
// and launches a goroutine that reads control events from an eventer Subsription channel
|
||||||
// and executes the events
|
// and executes the events
|
||||||
func (self *Network) Subscribe(eventer *event.TypeMux, types ...interface{}) {
|
func (self *Network) Subscribe(eventer *event.TypeMux, types ...interface{}) {
|
||||||
glog.V(logger.Info).Infof("subscribe")
|
log.Info("subscribe")
|
||||||
sub := eventer.Subscribe(types...)
|
sub := eventer.Subscribe(types...)
|
||||||
go func() {
|
go func() {
|
||||||
defer sub.Unsubscribe()
|
defer sub.Unsubscribe()
|
||||||
|
|
@ -229,38 +228,38 @@ func (self *Network) Subscribe(eventer *event.TypeMux, types ...interface{}) {
|
||||||
}
|
}
|
||||||
|
|
||||||
func (self *Network) execute(in *event.TypeMuxEvent) {
|
func (self *Network) execute(in *event.TypeMuxEvent) {
|
||||||
glog.V(logger.Detail).Infof("execute event %v", in)
|
log.Trace(fmt.Sprintf("execute event %v", in))
|
||||||
ev := in.Data
|
ev := in.Data
|
||||||
if ne, ok := ev.(*NodeEvent); ok {
|
if ne, ok := ev.(*NodeEvent); ok {
|
||||||
if ne.Action == "up" {
|
if ne.Action == "up" {
|
||||||
err := self.NewNode(&NodeConfig{Id: ne.node.Id})
|
err := self.NewNode(&NodeConfig{Id: ne.node.Id})
|
||||||
if err != nil {
|
if err != nil {
|
||||||
glog.V(logger.Detail).Infof("error execute event %v: %v", ne, err)
|
log.Trace(fmt.Sprintf("error execute event %v: %v", ne, err))
|
||||||
}
|
}
|
||||||
err = self.Start(ne.node.Id)
|
err = self.Start(ne.node.Id)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
glog.V(logger.Detail).Infof("error execute event %v: %v", ne, err)
|
log.Trace(fmt.Sprintf("error execute event %v: %v", ne, err))
|
||||||
}
|
}
|
||||||
} else {
|
} else {
|
||||||
err := self.Stop(ne.node.Id)
|
err := self.Stop(ne.node.Id)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
glog.V(logger.Detail).Infof("error execute event %v: %v", ne, err)
|
log.Trace(fmt.Sprintf("error execute event %v: %v", ne, err))
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
} else if ce, ok := ev.(*ConnEvent); ok {
|
} else if ce, ok := ev.(*ConnEvent); ok {
|
||||||
if ce.Action == "up" {
|
if ce.Action == "up" {
|
||||||
err := self.Connect(ce.conn.One, ce.conn.Other)
|
err := self.Connect(ce.conn.One, ce.conn.Other)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
glog.V(logger.Detail).Infof("error execute event %v: %v", ne, err)
|
log.Trace(fmt.Sprintf("error execute event %v: %v", ne, err))
|
||||||
}
|
}
|
||||||
} else {
|
} else {
|
||||||
err := self.Disconnect(ce.conn.One, ce.conn.Other)
|
err := self.Disconnect(ce.conn.One, ce.conn.Other)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
glog.V(logger.Detail).Infof("error execute event %v: %v", ne, err)
|
log.Trace(fmt.Sprintf("error execute event %v: %v", ne, err))
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
} else {
|
} else {
|
||||||
glog.V(logger.Detail).Infof("event: %#v", ev)
|
log.Trace(fmt.Sprintf("event: %#v", ev))
|
||||||
panic("unhandled event")
|
panic("unhandled event")
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -418,7 +417,7 @@ func (self *Network) NewNode(conf *NodeConfig) error {
|
||||||
na: na,
|
na: na,
|
||||||
}
|
}
|
||||||
self.Nodes = append(self.Nodes, node)
|
self.Nodes = append(self.Nodes, node)
|
||||||
glog.V(6).Infof("node %v created", id)
|
log.Trace(fmt.Sprintf("node %v created", id))
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -469,7 +468,7 @@ func (self *Network) Start(id *adapters.NodeId) error {
|
||||||
if node.Up {
|
if node.Up {
|
||||||
return fmt.Errorf("node %v already up", id)
|
return fmt.Errorf("node %v already up", id)
|
||||||
}
|
}
|
||||||
glog.V(6).Infof("starting node %v: %v adapter %v", id, node.Up, node.Adapter())
|
log.Trace(fmt.Sprintf("starting node %v: %v adapter %v", id, node.Up, node.Adapter()))
|
||||||
sa, ok := node.Adapter().(adapters.StartAdapter)
|
sa, ok := node.Adapter().(adapters.StartAdapter)
|
||||||
if ok {
|
if ok {
|
||||||
err := sa.Start()
|
err := sa.Start()
|
||||||
|
|
@ -478,7 +477,7 @@ func (self *Network) Start(id *adapters.NodeId) error {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
node.Up = true
|
node.Up = true
|
||||||
glog.V(logger.Info).Infof("started node %v: %v", id, node.Up)
|
log.Info(fmt.Sprintf("started node %v: %v", id, node.Up))
|
||||||
|
|
||||||
self.events.Post(&NodeEvent{
|
self.events.Post(&NodeEvent{
|
||||||
Action: "up",
|
Action: "up",
|
||||||
|
|
@ -505,7 +504,7 @@ func (self *Network) Stop(id *adapters.NodeId) error {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
node.Up = false
|
node.Up = false
|
||||||
glog.V(logger.Info).Infof("stop node %v: %v", id, node.Up)
|
log.Info(fmt.Sprintf("stop node %v: %v", id, node.Up))
|
||||||
|
|
||||||
self.events.Post(&NodeEvent{
|
self.events.Post(&NodeEvent{
|
||||||
Action: "down",
|
Action: "down",
|
||||||
|
|
|
||||||
|
|
@ -7,8 +7,7 @@ import (
|
||||||
"strings"
|
"strings"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/ethereum/go-ethereum/logger"
|
"github.com/ethereum/go-ethereum/log"
|
||||||
"github.com/ethereum/go-ethereum/logger/glog"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
type Controller interface {
|
type Controller interface {
|
||||||
|
|
@ -25,16 +24,16 @@ func StartRestApiServer(port string, c Controller) {
|
||||||
})
|
})
|
||||||
fd, err := net.Listen("tcp", ":"+port)
|
fd, err := net.Listen("tcp", ":"+port)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
glog.Errorf("Can't listen on :%s: %v", port, err)
|
log.Error(fmt.Sprintf("Can't listen on :%s: %v", port, err))
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
go http.Serve(fd, serveMux)
|
go http.Serve(fd, serveMux)
|
||||||
glog.V(logger.Info).Infof("Swarm Network Controller HTTP server started on localhost:%s", port)
|
log.Info(fmt.Sprintf("Swarm Network Controller HTTP server started on localhost:%s", port))
|
||||||
}
|
}
|
||||||
|
|
||||||
func handle(w http.ResponseWriter, r *http.Request, c Controller) {
|
func handle(w http.ResponseWriter, r *http.Request, c Controller) {
|
||||||
requestURL := r.URL
|
requestURL := r.URL
|
||||||
glog.V(logger.Debug).Infof("HTTP %s request URL: '%s', Host: '%s', Path: '%s', Referer: '%s', Accept: '%s'", r.Method, r.RequestURI, requestURL.Host, requestURL.Path, r.Referer(), r.Header.Get("Accept"))
|
log.Debug(fmt.Sprintf("HTTP %s request URL: '%s', Host: '%s', Path: '%s', Referer: '%s', Accept: '%s'", r.Method, r.RequestURI, requestURL.Host, requestURL.Path, r.Referer(), r.Header.Get("Accept")))
|
||||||
uri := requestURL.Path
|
uri := requestURL.Path
|
||||||
w.Header().Set("Content-Type", "text/json")
|
w.Header().Set("Content-Type", "text/json")
|
||||||
w.Header().Set("Access-Control-Allow-Origin", "*")
|
w.Header().Set("Access-Control-Allow-Origin", "*")
|
||||||
|
|
|
||||||
|
|
@ -10,8 +10,7 @@ import (
|
||||||
"sync"
|
"sync"
|
||||||
|
|
||||||
"github.com/ethereum/go-ethereum/crypto"
|
"github.com/ethereum/go-ethereum/crypto"
|
||||||
"github.com/ethereum/go-ethereum/logger"
|
"github.com/ethereum/go-ethereum/log"
|
||||||
"github.com/ethereum/go-ethereum/logger/glog"
|
|
||||||
"github.com/ethereum/go-ethereum/p2p/adapters"
|
"github.com/ethereum/go-ethereum/p2p/adapters"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
@ -98,7 +97,7 @@ func NewSessionController(nethook func(*NetworkConfig) (NetworkControl, *Resourc
|
||||||
Handle: func(msg interface{}, parent *ResourceController) (interface{}, error) {
|
Handle: func(msg interface{}, parent *ResourceController) (interface{}, error) {
|
||||||
conf := msg.(*NetworkConfig)
|
conf := msg.(*NetworkConfig)
|
||||||
netC, nodesC := nethook(conf)
|
netC, nodesC := nethook(conf)
|
||||||
glog.V(logger.Info).Infof("new network controller on %v", conf.Id)
|
log.Info(fmt.Sprintf("new network controller on %v", conf.Id))
|
||||||
m := NewNetworkController(netC, nodesC)
|
m := NewNetworkController(netC, nodesC)
|
||||||
if parent != nil {
|
if parent != nil {
|
||||||
parent.SetResource(conf.Id, m)
|
parent.SetResource(conf.Id, m)
|
||||||
|
|
@ -110,7 +109,7 @@ func NewSessionController(nethook func(*NetworkConfig) (NetworkControl, *Resourc
|
||||||
|
|
||||||
Destroy: &ResourceHandler{
|
Destroy: &ResourceHandler{
|
||||||
Handle: func(msg interface{}, parent *ResourceController) (interface{}, error) {
|
Handle: func(msg interface{}, parent *ResourceController) (interface{}, error) {
|
||||||
glog.V(logger.Debug).Infof("destroy handler called")
|
log.Debug("destroy handler called")
|
||||||
// this can quit the entire app (shut down the backend server)
|
// this can quit the entire app (shut down the backend server)
|
||||||
quitc <- true
|
quitc <- true
|
||||||
return empty, nil
|
return empty, nil
|
||||||
|
|
|
||||||
|
|
@ -7,11 +7,12 @@ import (
|
||||||
"io"
|
"io"
|
||||||
"io/ioutil"
|
"io/ioutil"
|
||||||
"net/http"
|
"net/http"
|
||||||
|
"os"
|
||||||
"testing"
|
"testing"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/ethereum/go-ethereum/event"
|
"github.com/ethereum/go-ethereum/event"
|
||||||
"github.com/ethereum/go-ethereum/logger/glog"
|
"github.com/ethereum/go-ethereum/log"
|
||||||
"github.com/ethereum/go-ethereum/p2p/adapters"
|
"github.com/ethereum/go-ethereum/p2p/adapters"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
@ -27,8 +28,8 @@ var quitc chan bool
|
||||||
var controller *ResourceController
|
var controller *ResourceController
|
||||||
|
|
||||||
func init() {
|
func init() {
|
||||||
glog.SetV(0)
|
log.Root().SetHandler(log.LvlFilterHandler(log.LvlError, log.StreamHandler(os.Stderr, log.TerminalFormat(false))))
|
||||||
glog.SetToStderr(true)
|
|
||||||
controller, quitc = NewSessionController(DefaultNet)
|
controller, quitc = NewSessionController(DefaultNet)
|
||||||
StartRestApiServer(port, controller)
|
StartRestApiServer(port, controller)
|
||||||
}
|
}
|
||||||
|
|
@ -160,9 +161,9 @@ func TestUpdate(t *testing.T) {
|
||||||
}
|
}
|
||||||
|
|
||||||
func mockNewNodes(eventer *event.TypeMux, ids []*adapters.NodeId) {
|
func mockNewNodes(eventer *event.TypeMux, ids []*adapters.NodeId) {
|
||||||
glog.V(6).Infof("mock starting")
|
log.Trace("mock starting")
|
||||||
for _, id := range ids {
|
for _, id := range ids {
|
||||||
glog.V(6).Infof("mock adding node %v", id)
|
log.Trace(fmt.Sprintf("mock adding node %v", id))
|
||||||
eventer.Post(&NodeEvent{
|
eventer.Post(&NodeEvent{
|
||||||
Action: "up",
|
Action: "up",
|
||||||
Type: "node",
|
Type: "node",
|
||||||
|
|
|
||||||
|
|
@ -1,11 +1,10 @@
|
||||||
package testing
|
package testing
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"fmt"
|
||||||
"sync"
|
"sync"
|
||||||
|
|
||||||
"github.com/ethereum/go-ethereum/logger"
|
"github.com/ethereum/go-ethereum/log"
|
||||||
"github.com/ethereum/go-ethereum/logger/glog"
|
|
||||||
|
|
||||||
"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"
|
||||||
)
|
)
|
||||||
|
|
@ -28,7 +27,7 @@ func NewTestPeerPool() *TestPeerPool {
|
||||||
func (self *TestPeerPool) Add(p TestPeer) {
|
func (self *TestPeerPool) Add(p TestPeer) {
|
||||||
self.lock.Lock()
|
self.lock.Lock()
|
||||||
defer self.lock.Unlock()
|
defer self.lock.Unlock()
|
||||||
glog.V(logger.Detail).Infof("pp add peer %v", p.ID())
|
log.Trace(fmt.Sprintf("pp add peer %v", p.ID()))
|
||||||
self.peers[p.ID()] = p
|
self.peers[p.ID()] = p
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -5,8 +5,7 @@ import (
|
||||||
"sync"
|
"sync"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/ethereum/go-ethereum/logger"
|
"github.com/ethereum/go-ethereum/log"
|
||||||
"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/adapters"
|
||||||
)
|
)
|
||||||
|
|
@ -78,9 +77,9 @@ func (self *ProtocolSession) trigger(trig Trigger) error {
|
||||||
errc := make(chan error)
|
errc := make(chan error)
|
||||||
|
|
||||||
go func() {
|
go func() {
|
||||||
glog.V(logger.Detail).Infof("trigger %v (%v)....", trig.Msg, trig.Code)
|
log.Trace(fmt.Sprintf("trigger %v (%v)....", trig.Msg, trig.Code))
|
||||||
errc <- p2p.Send(peer, trig.Code, trig.Msg)
|
errc <- p2p.Send(peer, trig.Code, trig.Msg)
|
||||||
glog.V(logger.Detail).Infof("triggered %v (%v)", trig.Msg, trig.Code)
|
log.Trace(fmt.Sprintf("triggered %v (%v)", trig.Msg, trig.Code))
|
||||||
}()
|
}()
|
||||||
|
|
||||||
t := trig.Timeout
|
t := trig.Timeout
|
||||||
|
|
@ -111,7 +110,7 @@ func (self *ProtocolSession) expect(exp Expect) error {
|
||||||
|
|
||||||
errc := make(chan error)
|
errc := make(chan error)
|
||||||
go func() {
|
go func() {
|
||||||
glog.V(logger.Detail).Infof("waiting for msg, %v", exp.Msg)
|
log.Trace(fmt.Sprintf("waiting for msg, %v", exp.Msg))
|
||||||
errc <- p2p.ExpectMsg(peer, exp.Code, exp.Msg)
|
errc <- p2p.ExpectMsg(peer, exp.Code, exp.Msg)
|
||||||
}()
|
}()
|
||||||
|
|
||||||
|
|
@ -122,7 +121,7 @@ func (self *ProtocolSession) expect(exp Expect) error {
|
||||||
alarm := time.NewTimer(t)
|
alarm := time.NewTimer(t)
|
||||||
select {
|
select {
|
||||||
case err := <-errc:
|
case err := <-errc:
|
||||||
glog.V(logger.Detail).Infof("expected msg arrives with error %v", err)
|
log.Trace(fmt.Sprintf("expected msg arrives with error %v", err))
|
||||||
return err
|
return err
|
||||||
case <-alarm.C:
|
case <-alarm.C:
|
||||||
return fmt.Errorf("timout expecting %v sent to peer %v", exp.Msg, exp.Peer)
|
return fmt.Errorf("timout expecting %v sent to peer %v", exp.Msg, exp.Peer)
|
||||||
|
|
@ -158,7 +157,7 @@ func (self *ProtocolSession) TestExchanges(exchanges ...Exchange) error {
|
||||||
defer wg.Done()
|
defer wg.Done()
|
||||||
err := self.expect(exp)
|
err := self.expect(exp)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
glog.V(logger.Detail).Infof("expect msg fails %v", err)
|
log.Trace(fmt.Sprintf("expect msg fails %v", err))
|
||||||
errc <- err
|
errc <- err
|
||||||
}
|
}
|
||||||
}(ex)
|
}(ex)
|
||||||
|
|
@ -178,7 +177,7 @@ func (self *ProtocolSession) TestExchanges(exchanges ...Exchange) error {
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return fmt.Errorf("exchange failed with: %v", err)
|
return fmt.Errorf("exchange failed with: %v", err)
|
||||||
} else {
|
} else {
|
||||||
glog.V(logger.Detail).Infof("exchange %v: '%v' run successfully", i, e.Label)
|
log.Trace(fmt.Sprintf("exchange %v: '%v' run successfully", i, e.Label))
|
||||||
}
|
}
|
||||||
case <-alarm.C:
|
case <-alarm.C:
|
||||||
return fmt.Errorf("exchange %v: '%v' timed out", i, e.Label)
|
return fmt.Errorf("exchange %v: '%v' timed out", i, e.Label)
|
||||||
|
|
|
||||||
|
|
@ -1,10 +1,10 @@
|
||||||
package testing
|
package testing
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"fmt"
|
||||||
"testing"
|
"testing"
|
||||||
|
|
||||||
"github.com/ethereum/go-ethereum/logger"
|
"github.com/ethereum/go-ethereum/log"
|
||||||
"github.com/ethereum/go-ethereum/logger/glog"
|
|
||||||
"github.com/ethereum/go-ethereum/p2p/adapters"
|
"github.com/ethereum/go-ethereum/p2p/adapters"
|
||||||
"github.com/ethereum/go-ethereum/p2p/simulations"
|
"github.com/ethereum/go-ethereum/p2p/simulations"
|
||||||
)
|
)
|
||||||
|
|
@ -21,7 +21,7 @@ func NewProtocolTester(t *testing.T, id *adapters.NodeId, n int, run adapters.Pr
|
||||||
naf := func(conf *simulations.NodeConfig) adapters.NodeAdapter {
|
naf := func(conf *simulations.NodeConfig) adapters.NodeAdapter {
|
||||||
na := adapters.NewSimNode(conf.Id, net)
|
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)
|
log.Trace(fmt.Sprintf("adapter run function set to protocol for node %v (=%v)", conf.Id, id))
|
||||||
na.Run = run
|
na.Run = run
|
||||||
}
|
}
|
||||||
return na
|
return na
|
||||||
|
|
@ -56,11 +56,11 @@ func (self *ProtocolTester) Start(id *adapters.NodeId) error {
|
||||||
}
|
}
|
||||||
node := self.network.GetNode(id)
|
node := self.network.GetNode(id)
|
||||||
if node == nil {
|
if node == nil {
|
||||||
glog.V(logger.Detail).Infof("node for peer %v not found", id)
|
log.Trace(fmt.Sprintf("node for peer %v not found", id))
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
if node.Adapter() == nil {
|
if node.Adapter() == nil {
|
||||||
glog.V(logger.Detail).Infof("node adapter for peer %v not found", id)
|
log.Trace(fmt.Sprintf("node adapter for peer %v not found", id))
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
return nil
|
return nil
|
||||||
|
|
@ -68,15 +68,15 @@ func (self *ProtocolTester) Start(id *adapters.NodeId) error {
|
||||||
|
|
||||||
func (self *ProtocolTester) Connect(ids ...*adapters.NodeId) {
|
func (self *ProtocolTester) Connect(ids ...*adapters.NodeId) {
|
||||||
for _, id := range ids {
|
for _, id := range ids {
|
||||||
glog.V(logger.Detail).Infof("start node %v", id)
|
log.Trace(fmt.Sprintf("start node %v", id))
|
||||||
err := self.Start(id)
|
err := self.Start(id)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
glog.V(logger.Detail).Infof("error starting peer %v: %v", id, err)
|
log.Trace(fmt.Sprintf("error starting peer %v: %v", id, err))
|
||||||
}
|
}
|
||||||
glog.V(logger.Detail).Infof("connect to %v", id)
|
log.Trace(fmt.Sprintf("connect to %v", id))
|
||||||
err = self.na.Connect(id.Bytes())
|
err = self.na.Connect(id.Bytes())
|
||||||
if err != nil {
|
if err != nil {
|
||||||
glog.V(logger.Detail).Infof("error connecting to peer %v: %v", id, err)
|
log.Trace(fmt.Sprintf("error connecting to peer %v: %v", id, err))
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -19,13 +19,13 @@ import (
|
||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
"math/rand"
|
"math/rand"
|
||||||
|
"os"
|
||||||
"runtime"
|
"runtime"
|
||||||
"sync"
|
"sync"
|
||||||
"testing"
|
"testing"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/ethereum/go-ethereum/logger"
|
"github.com/ethereum/go-ethereum/log"
|
||||||
"github.com/ethereum/go-ethereum/logger/glog"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
const (
|
const (
|
||||||
|
|
@ -34,8 +34,7 @@ const (
|
||||||
)
|
)
|
||||||
|
|
||||||
func init() {
|
func init() {
|
||||||
glog.SetV(0)
|
log.Root().SetHandler(log.LvlFilterHandler(log.LvlError, log.StreamHandler(os.Stderr, log.TerminalFormat(false))))
|
||||||
glog.SetToStderr(true)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
type testBVAddr struct {
|
type testBVAddr struct {
|
||||||
|
|
@ -303,7 +302,7 @@ func TestPotMergeOne(t *testing.T) {
|
||||||
pot1.Add(NewTestAddr("00", 0))
|
pot1.Add(NewTestAddr("00", 0))
|
||||||
pot2 := NewPot(nil, 0)
|
pot2 := NewPot(nil, 0)
|
||||||
pot2.Add(NewTestAddr("01", 0))
|
pot2.Add(NewTestAddr("01", 0))
|
||||||
glog.V(logger.Debug).Infof("\n%v\n%v", pot2, pot1)
|
log.Debug(fmt.Sprintf("\n%v\n%v", pot2, pot1))
|
||||||
pot1.Merge(pot2)
|
pot1.Merge(pot2)
|
||||||
count := 0
|
count := 0
|
||||||
pot1.Each(func(val PotVal, i int) bool {
|
pot1.Each(func(val PotVal, i int) bool {
|
||||||
|
|
@ -326,7 +325,7 @@ func TestPotMergeCommon(t *testing.T) {
|
||||||
max1 := rand.Intn(mergeTestChoose) + 1
|
max1 := rand.Intn(mergeTestChoose) + 1
|
||||||
n0 := NewPot(nil, 0)
|
n0 := NewPot(nil, 0)
|
||||||
n1 := NewPot(nil, 0)
|
n1 := NewPot(nil, 0)
|
||||||
glog.V(3).Infof("round %v: %v - %v", i, max0, max1)
|
log.Trace(fmt.Sprintf("round %v: %v - %v", i, max0, max1))
|
||||||
m := make(map[string]bool)
|
m := make(map[string]bool)
|
||||||
for j := 0; j < max0; {
|
for j := 0; j < max0; {
|
||||||
r := rand.Intn(max0)
|
r := rand.Intn(max0)
|
||||||
|
|
@ -357,9 +356,9 @@ func TestPotMergeCommon(t *testing.T) {
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
expSize := len(m)
|
expSize := len(m)
|
||||||
glog.V(4).Infof("%v-0: pin: %v, size: %v", i, n0.Pin(), max0)
|
log.Trace(fmt.Sprintf("%v-0: pin: %v, size: %v", i, n0.Pin(), max0))
|
||||||
glog.V(4).Infof("%v-1: pin: %v, size: %v", i, n1.Pin(), max1)
|
log.Trace(fmt.Sprintf("%v-1: pin: %v, size: %v", i, n1.Pin(), max1))
|
||||||
glog.V(4).Infof("%v: merged tree size: %v, newly added: %v", i, expSize, expAdded)
|
log.Trace(fmt.Sprintf("%v: merged tree size: %v, newly added: %v", i, expSize, expAdded))
|
||||||
n, common := Union(n0, n1)
|
n, common := Union(n0, n1)
|
||||||
added := n1.Size() - common
|
added := n1.Size() - common
|
||||||
size := n.Size()
|
size := n.Size()
|
||||||
|
|
@ -388,7 +387,7 @@ func TestPotMergeScale(t *testing.T) {
|
||||||
max1 := rand.Intn(maxEachNeighbour) + 1
|
max1 := rand.Intn(maxEachNeighbour) + 1
|
||||||
n0 := NewPot(nil, 0)
|
n0 := NewPot(nil, 0)
|
||||||
n1 := NewPot(nil, 0)
|
n1 := NewPot(nil, 0)
|
||||||
glog.V(3).Infof("round %v: %v - %v", i, max0, max1)
|
log.Trace(fmt.Sprintf("round %v: %v - %v", i, max0, max1))
|
||||||
m := make(map[string]bool)
|
m := make(map[string]bool)
|
||||||
for j := 0; j < max0; {
|
for j := 0; j < max0; {
|
||||||
v := randomTestBVAddr(keylen, j)
|
v := randomTestBVAddr(keylen, j)
|
||||||
|
|
@ -418,9 +417,9 @@ func TestPotMergeScale(t *testing.T) {
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
expSize := len(m)
|
expSize := len(m)
|
||||||
glog.V(4).Infof("%v-0: pin: %v, size: %v", i, n0.Pin(), max0)
|
log.Trace(fmt.Sprintf("%v-0: pin: %v, size: %v", i, n0.Pin(), max0))
|
||||||
glog.V(4).Infof("%v-1: pin: %v, size: %v", i, n1.Pin(), max1)
|
log.Trace(fmt.Sprintf("%v-1: pin: %v, size: %v", i, n1.Pin(), max1))
|
||||||
glog.V(4).Infof("%v: merged tree size: %v, newly added: %v", i, expSize, expAdded)
|
log.Trace(fmt.Sprintf("%v: merged tree size: %v, newly added: %v", i, expSize, expAdded))
|
||||||
n, common := Union(n0, n1)
|
n, common := Union(n0, n1)
|
||||||
added := n1.Size() - common
|
added := n1.Size() - common
|
||||||
size := n.Size()
|
size := n.Size()
|
||||||
|
|
@ -476,7 +475,7 @@ func TestPotEachNeighbourSync(t *testing.T) {
|
||||||
}
|
}
|
||||||
count := rand.Intn(size/2) + size/2
|
count := rand.Intn(size/2) + size/2
|
||||||
val := randomTestAddr(keylen, max+1)
|
val := randomTestAddr(keylen, max+1)
|
||||||
glog.V(3).Infof("%v: pin: %v, size: %v, val: %v, count: %v", i, n.Pin(), size, val, count)
|
log.Trace(fmt.Sprintf("%v: pin: %v, size: %v, val: %v, count: %v", i, n.Pin(), size, val, count))
|
||||||
err := testPotEachNeighbour(n, val, count, checkPo(val), checkOrder(val), checkValues(m, val))
|
err := testPotEachNeighbour(n, val, count, checkPo(val), checkOrder(val), checkValues(m, val))
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
|
|
@ -525,16 +524,16 @@ func TestPotEachNeighbourAsync(t *testing.T) {
|
||||||
mu := sync.Mutex{}
|
mu := sync.Mutex{}
|
||||||
m := make(map[string]bool)
|
m := make(map[string]bool)
|
||||||
maxPos := rand.Intn(keylen)
|
maxPos := rand.Intn(keylen)
|
||||||
glog.V(3).Infof("%v: pin: %v, size: %v, val: %v, count: %v, maxPos: %v", i, n.Pin(), size, val, count, maxPos)
|
log.Trace(fmt.Sprintf("%v: pin: %v, size: %v, val: %v, count: %v, maxPos: %v", i, n.Pin(), size, val, count, maxPos))
|
||||||
msize := 0
|
msize := 0
|
||||||
remember := func(v PotVal, po int) error {
|
remember := func(v PotVal, po int) error {
|
||||||
// mu.Lock()
|
// mu.Lock()
|
||||||
// defer mu.Unlock()
|
// defer mu.Unlock()
|
||||||
if po > maxPos {
|
if po > maxPos {
|
||||||
// glog.V(4).Infof("NOT ADD %v", v)
|
// log.Trace(fmt.Sprintf("NOT ADD %v", v))
|
||||||
return errNoCount
|
return errNoCount
|
||||||
}
|
}
|
||||||
// glog.V(4).Infof("ADD %v, %v", v, msize)
|
// log.Trace(fmt.Sprintf("ADD %v, %v", v, msize))
|
||||||
m[v.String()] = true
|
m[v.String()] = true
|
||||||
msize++
|
msize++
|
||||||
return nil
|
return nil
|
||||||
|
|
@ -544,14 +543,14 @@ func TestPotEachNeighbourAsync(t *testing.T) {
|
||||||
}
|
}
|
||||||
err := testPotEachNeighbour(n, val, count, remember)
|
err := testPotEachNeighbour(n, val, count, remember)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
glog.V(6).Info(err)
|
log.Error(err.Error())
|
||||||
}
|
}
|
||||||
d := 0
|
d := 0
|
||||||
forget := func(v PotVal, po int) {
|
forget := func(v PotVal, po int) {
|
||||||
mu.Lock()
|
mu.Lock()
|
||||||
defer mu.Unlock()
|
defer mu.Unlock()
|
||||||
d++
|
d++
|
||||||
// glog.V(4).Infof("DEL %v", v)
|
// log.Trace(fmt.Sprintf("DEL %v", v))
|
||||||
delete(m, v.String())
|
delete(m, v.String())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -3,9 +3,7 @@ package network
|
||||||
import (
|
import (
|
||||||
"fmt"
|
"fmt"
|
||||||
|
|
||||||
"github.com/ethereum/go-ethereum/logger"
|
"github.com/ethereum/go-ethereum/log"
|
||||||
"github.com/ethereum/go-ethereum/logger/glog"
|
|
||||||
// "github.com/ethereum/go-ethereum/p2p/discover"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
// discovery bzz overlay extension doing peer relaying
|
// discovery bzz overlay extension doing peer relaying
|
||||||
|
|
@ -46,11 +44,11 @@ func NewDiscovery(p Peer, o Overlay) *discPeer {
|
||||||
// NotifyPeer notifies the receiver remote end of a peer p or PO po.
|
// NotifyPeer notifies the receiver remote end of a peer p or PO po.
|
||||||
// callback for overlay driver
|
// callback for overlay driver
|
||||||
func (self *discPeer) NotifyPeer(p Peer, po uint8) error {
|
func (self *discPeer) NotifyPeer(p Peer, po uint8) error {
|
||||||
glog.V(logger.Warn).Infof("peers %v", self.peers)
|
log.Warn(fmt.Sprintf("peers %v", self.peers))
|
||||||
if po < self.proxLimit || self.seen(p) {
|
if po < self.proxLimit || self.seen(p) {
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
glog.V(logger.Warn).Infof("notification about %x", p.OverlayAddr())
|
log.Warn(fmt.Sprintf("notification about %x", p.OverlayAddr()))
|
||||||
|
|
||||||
resp := &peersMsg{
|
resp := &peersMsg{
|
||||||
Peers: []*peerAddr{&peerAddr{OAddr: p.OverlayAddr(), UAddr: p.UnderlayAddr()}}, // perhaps the PeerAddr interface is unnecessary generalization
|
Peers: []*peerAddr{&peerAddr{OAddr: p.OverlayAddr(), UAddr: p.UnderlayAddr()}}, // perhaps the PeerAddr interface is unnecessary generalization
|
||||||
|
|
@ -125,7 +123,7 @@ func (self *discPeer) handleSubPeersMsg(msg interface{}) error {
|
||||||
peers = append(peers, &peerAddr{p.OverlayAddr(), p.UnderlayAddr()})
|
peers = append(peers, &peerAddr{p.OverlayAddr(), p.UnderlayAddr()})
|
||||||
return true
|
return true
|
||||||
})
|
})
|
||||||
glog.V(logger.Warn).Infof("found initial %v peers not farther than %v", len(peers), self.proxLimit)
|
log.Warn(fmt.Sprintf("found initial %v peers not farther than %v", len(peers), self.proxLimit))
|
||||||
if len(peers) > 0 {
|
if len(peers) > 0 {
|
||||||
self.Send(&peersMsg{Peers: peers})
|
self.Send(&peersMsg{Peers: peers})
|
||||||
}
|
}
|
||||||
|
|
@ -147,10 +145,10 @@ func (self *discPeer) handlePeersMsg(msg interface{}) error {
|
||||||
}
|
}
|
||||||
|
|
||||||
if len(nas) == 0 {
|
if len(nas) == 0 {
|
||||||
glog.V(logger.Debug).Infof("whoops, no peers in incoming peersMsg from %v", self)
|
log.Debug(fmt.Sprintf("whoops, no peers in incoming peersMsg from %v", self))
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
glog.V(logger.Debug).Infof("got peer addresses from %x, %v (%v)", self.OverlayAddr(), nas, len(nas))
|
log.Debug(fmt.Sprintf("got peer addresses from %x, %v (%v)", self.OverlayAddr(), nas, len(nas)))
|
||||||
return self.overlay.Register(nas...)
|
return self.overlay.Register(nas...)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -172,10 +170,10 @@ func (self *discPeer) handleGetPeersMsg(msg interface{}) error {
|
||||||
return len(peers) < int(req.Max)
|
return len(peers) < int(req.Max)
|
||||||
})
|
})
|
||||||
if len(peers) == 0 {
|
if len(peers) == 0 {
|
||||||
glog.V(logger.Debug).Infof("no peers found for %v", self)
|
log.Debug(fmt.Sprintf("no peers found for %v", self))
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
glog.V(logger.Debug).Infof("%v peers sent to %v", len(peers), self)
|
log.Debug(fmt.Sprintf("%v peers sent to %v", len(peers), self))
|
||||||
resp := &peersMsg{
|
resp := &peersMsg{
|
||||||
Peers: peers,
|
Peers: peers,
|
||||||
}
|
}
|
||||||
|
|
@ -190,7 +188,7 @@ func RequestOrder(k Overlay, order, broadcastSize, maxPeers uint8) {
|
||||||
var i uint8
|
var i uint8
|
||||||
var err error
|
var err error
|
||||||
k.EachLivePeer(nil, 255, func(n Peer, po int) bool {
|
k.EachLivePeer(nil, 255, func(n Peer, po int) bool {
|
||||||
glog.V(logger.Detail).Infof("%T sent to %v", req, n.ID())
|
log.Trace(fmt.Sprintf("%T sent to %v", req, n.ID()))
|
||||||
err = n.Send(req)
|
err = n.Send(req)
|
||||||
if err == nil {
|
if err == nil {
|
||||||
i++
|
i++
|
||||||
|
|
@ -200,7 +198,7 @@ func RequestOrder(k Overlay, order, broadcastSize, maxPeers uint8) {
|
||||||
}
|
}
|
||||||
return true
|
return true
|
||||||
})
|
})
|
||||||
glog.V(logger.Info).Infof("requesting bees of PO%03d from %v/%v (each max %v)", order, i, broadcastSize, maxPeers)
|
log.Info(fmt.Sprintf("requesting bees of PO%03d from %v/%v (each max %v)", order, i, broadcastSize, maxPeers))
|
||||||
}
|
}
|
||||||
|
|
||||||
func (self *discPeer) seen(p PeerAddr) bool {
|
func (self *discPeer) seen(p PeerAddr) bool {
|
||||||
|
|
|
||||||
|
|
@ -1,10 +1,10 @@
|
||||||
package network
|
package network
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"fmt"
|
||||||
"testing"
|
"testing"
|
||||||
|
|
||||||
"github.com/ethereum/go-ethereum/logger"
|
"github.com/ethereum/go-ethereum/log"
|
||||||
"github.com/ethereum/go-ethereum/logger/glog"
|
|
||||||
p2ptest "github.com/ethereum/go-ethereum/p2p/testing"
|
p2ptest "github.com/ethereum/go-ethereum/p2p/testing"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
@ -21,7 +21,7 @@ func TestDiscovery(t *testing.T) {
|
||||||
services := func(p Peer) error {
|
services := func(p Peer) error {
|
||||||
dp := NewDiscovery(p, to)
|
dp := NewDiscovery(p, to)
|
||||||
to.On(dp)
|
to.On(dp)
|
||||||
glog.V(logger.Detail).Infof("kademlia on %v", p)
|
log.Trace(fmt.Sprintf("kademlia on %v", p))
|
||||||
p.DisconnectHook(func(err error) {
|
p.DisconnectHook(func(err error) {
|
||||||
to.Off(p)
|
to.Off(p)
|
||||||
})
|
})
|
||||||
|
|
|
||||||
|
|
@ -17,11 +17,11 @@
|
||||||
package network
|
package network
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"fmt"
|
||||||
"sync"
|
"sync"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/ethereum/go-ethereum/logger"
|
"github.com/ethereum/go-ethereum/log"
|
||||||
"github.com/ethereum/go-ethereum/logger/glog"
|
|
||||||
"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"
|
||||||
)
|
)
|
||||||
|
|
@ -97,7 +97,7 @@ func (self *Hive) Start(connectPeer func(string) error, af func() <-chan time.Ti
|
||||||
self.toggle = make(chan bool)
|
self.toggle = make(chan bool)
|
||||||
self.more = make(chan bool, 1)
|
self.more = make(chan bool, 1)
|
||||||
self.quit = make(chan bool)
|
self.quit = make(chan bool)
|
||||||
glog.V(logger.Debug).Infof("hive started")
|
log.Debug("hive started")
|
||||||
// this loop is doing bootstrapping and maintains a healthy table
|
// this loop is doing bootstrapping and maintains a healthy table
|
||||||
go self.keepAlive(af)
|
go self.keepAlive(af)
|
||||||
go func() {
|
go func() {
|
||||||
|
|
@ -108,17 +108,17 @@ func (self *Hive) Start(connectPeer func(string) error, af func() <-chan time.Ti
|
||||||
// to attempt to write to more (remove Peer when shutting down)
|
// to attempt to write to more (remove Peer when shutting down)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
glog.V(logger.Detail).Infof("hive delegate to overlay driver: suggest addr to connect to")
|
log.Trace("hive delegate to overlay driver: suggest addr to connect to")
|
||||||
addr, order, want := self.SuggestPeer()
|
addr, order, want := self.SuggestPeer()
|
||||||
|
|
||||||
if addr != nil {
|
if addr != nil {
|
||||||
glog.V(logger.Info).Infof("========> connect to bee %v", addr)
|
log.Info(fmt.Sprintf("========> connect to bee %v", addr))
|
||||||
err := connectPeer(NodeId(addr).NodeID.String())
|
err := connectPeer(NodeId(addr).NodeID.String())
|
||||||
if err != nil {
|
if err != nil {
|
||||||
glog.V(logger.Detail).Infof("===X====> connect to bee %v failed: %v", addr, err)
|
log.Error(fmt.Sprintf("===X====> connect to bee %v failed: %v", addr, err))
|
||||||
}
|
}
|
||||||
} else {
|
} else {
|
||||||
glog.V(logger.Detail).Infof("cannot suggest peers")
|
log.Trace("cannot suggest peers")
|
||||||
}
|
}
|
||||||
|
|
||||||
want = want && self.Discovery
|
want = want && self.Discovery
|
||||||
|
|
@ -128,11 +128,11 @@ func (self *Hive) Start(connectPeer func(string) error, af func() <-chan time.Ti
|
||||||
|
|
||||||
select {
|
select {
|
||||||
case self.toggle <- want:
|
case self.toggle <- want:
|
||||||
glog.V(logger.Detail).Infof("keep hive alive: %v", want)
|
log.Trace(fmt.Sprintf("keep hive alive: %v", want))
|
||||||
case <-self.quit:
|
case <-self.quit:
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
glog.V(logger.Info).Infof("%v", self)
|
log.Info(fmt.Sprintf("%v", self))
|
||||||
}
|
}
|
||||||
}()
|
}()
|
||||||
return nil
|
return nil
|
||||||
|
|
@ -148,12 +148,12 @@ func (self *Hive) ticker() <-chan time.Time {
|
||||||
// wake state is toggled by writing to self.toggle
|
// wake state is toggled by writing to self.toggle
|
||||||
// it restarts if the table becomes non-full again due to disconnections
|
// it restarts if the table becomes non-full again due to disconnections
|
||||||
func (self *Hive) keepAlive(af func() <-chan time.Time) {
|
func (self *Hive) keepAlive(af func() <-chan time.Time) {
|
||||||
glog.V(logger.Detail).Infof("keep alive loop started")
|
log.Trace("keep alive loop started")
|
||||||
alarm := af()
|
alarm := af()
|
||||||
for {
|
for {
|
||||||
select {
|
select {
|
||||||
case <-alarm:
|
case <-alarm:
|
||||||
glog.V(logger.Detail).Infof("wake up: make hive alive")
|
log.Trace("wake up: make hive alive")
|
||||||
self.wake()
|
self.wake()
|
||||||
case need := <-self.toggle:
|
case need := <-self.toggle:
|
||||||
if alarm == nil && need {
|
if alarm == nil && need {
|
||||||
|
|
@ -174,17 +174,17 @@ func (self *Hive) keepAlive(af func() <-chan time.Time) {
|
||||||
func (self *Hive) Add(p Peer) error {
|
func (self *Hive) Add(p Peer) error {
|
||||||
defer self.wake()
|
defer self.wake()
|
||||||
dp := NewDiscovery(p, self.Overlay)
|
dp := NewDiscovery(p, self.Overlay)
|
||||||
glog.V(logger.Debug).Infof("to add new bee %v", p)
|
log.Debug(fmt.Sprintf("to add new bee %v", p))
|
||||||
self.On(dp)
|
self.On(dp)
|
||||||
self.String()
|
self.String()
|
||||||
glog.V(logger.Warn).Infof("%v", self)
|
log.Debug(fmt.Sprintf("%v", self))
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// Remove called after peer is disconnected
|
// Remove called after peer is disconnected
|
||||||
func (self *Hive) Remove(p Peer) {
|
func (self *Hive) Remove(p Peer) {
|
||||||
defer self.wake()
|
defer self.wake()
|
||||||
glog.V(logger.Debug).Infof("remove bee %v", p)
|
log.Debug(fmt.Sprintf("remove bee %v", p))
|
||||||
self.Off(p)
|
self.Off(p)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -212,10 +212,10 @@ func (self *Hive) Stop() {
|
||||||
func (self *Hive) wake() {
|
func (self *Hive) wake() {
|
||||||
select {
|
select {
|
||||||
case self.more <- true:
|
case self.more <- true:
|
||||||
glog.V(logger.Detail).Infof("hive woken up")
|
log.Trace("hive woken up")
|
||||||
case <-self.quit:
|
case <-self.quit:
|
||||||
default:
|
default:
|
||||||
glog.V(logger.Detail).Infof("hive already awake")
|
log.Trace("hive already awake")
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -22,8 +22,7 @@ import (
|
||||||
"strings"
|
"strings"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/ethereum/go-ethereum/logger"
|
"github.com/ethereum/go-ethereum/log"
|
||||||
"github.com/ethereum/go-ethereum/logger/glog"
|
|
||||||
"github.com/ethereum/go-ethereum/pot"
|
"github.com/ethereum/go-ethereum/pot"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
@ -118,7 +117,7 @@ func NewKademlia(addr []byte, params *KadParams) *Kademlia {
|
||||||
func (self *Kademlia) Prune(c <-chan time.Time) {
|
func (self *Kademlia) Prune(c <-chan time.Time) {
|
||||||
go func() {
|
go func() {
|
||||||
for _ = range c {
|
for _ = range c {
|
||||||
glog.V(logger.Debug).Infof("pruning...")
|
log.Debug("pruning...")
|
||||||
total := 0
|
total := 0
|
||||||
self.peers.EachBin(self.addr, 0, func(po, size int, f func(func(pot.PotVal, int) bool) bool) bool {
|
self.peers.EachBin(self.addr, 0, func(po, size int, f func(func(pot.PotVal, int) bool) bool) bool {
|
||||||
extra := size - self.MinBinSize
|
extra := size - self.MinBinSize
|
||||||
|
|
@ -136,7 +135,7 @@ func (self *Kademlia) Prune(c <-chan time.Time) {
|
||||||
}
|
}
|
||||||
return true
|
return true
|
||||||
})
|
})
|
||||||
glog.V(logger.Debug).Infof("pruned %v peers", total)
|
log.Debug(fmt.Sprintf("pruned %v peers", total))
|
||||||
}
|
}
|
||||||
}()
|
}()
|
||||||
}
|
}
|
||||||
|
|
@ -166,25 +165,25 @@ func (self *Kademlia) callable(val pot.PotVal) *KadPeer {
|
||||||
kp := val.(*KadPeer)
|
kp := val.(*KadPeer)
|
||||||
// not callable if peer is live or exceeded maxRetries
|
// not callable if peer is live or exceeded maxRetries
|
||||||
if kp.Peer != nil || kp.retries > self.MaxRetries {
|
if kp.Peer != nil || kp.retries > self.MaxRetries {
|
||||||
glog.V(logger.Detail).Infof("peer %v (%T) not callable", kp, kp.Peer)
|
log.Trace(fmt.Sprintf("peer %v (%T) not callable", kp, kp.Peer))
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
// calculate the allowed number of retries based on time lapsed since last seen
|
// calculate the allowed number of retries based on time lapsed since last seen
|
||||||
timeAgo := time.Since(kp.seenAt)
|
timeAgo := time.Since(kp.seenAt)
|
||||||
var retries int
|
var retries int
|
||||||
for delta := int(timeAgo) / self.RetryInterval; delta > 0; delta /= self.RetryExponent {
|
for delta := int(timeAgo) / self.RetryInterval; delta > 0; delta /= self.RetryExponent {
|
||||||
glog.V(logger.Detail).Infof("delta: %v", delta)
|
log.Trace(fmt.Sprintf("delta: %v", delta))
|
||||||
retries++
|
retries++
|
||||||
}
|
}
|
||||||
// this is never called concurrently, so safe to increment
|
// this is never called concurrently, so safe to increment
|
||||||
// peer can be retried again
|
// peer can be retried again
|
||||||
|
|
||||||
if retries < kp.retries {
|
if retries < kp.retries {
|
||||||
glog.V(logger.Detail).Infof("log time needed before retry %v, wait only warrants %v", kp.retries, retries)
|
log.Trace(fmt.Sprintf("log time needed before retry %v, wait only warrants %v", kp.retries, retries))
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
kp.retries++
|
kp.retries++
|
||||||
glog.V(logger.Detail).Infof("peer %v is callable", kp)
|
log.Trace(fmt.Sprintf("peer %v is callable", kp))
|
||||||
|
|
||||||
return kp
|
return kp
|
||||||
}
|
}
|
||||||
|
|
@ -204,7 +203,7 @@ func (self *Kademlia) Register(nas ...PeerAddr) error {
|
||||||
np := pot.NewPot(nil, 0)
|
np := pot.NewPot(nil, 0)
|
||||||
for _, na := range nas {
|
for _, na := range nas {
|
||||||
if bytes.Equal(na.OverlayAddr(), self.addr.OverlayAddr()) {
|
if bytes.Equal(na.OverlayAddr(), self.addr.OverlayAddr()) {
|
||||||
glog.V(logger.Warn).Infof("[%06s] add peers: %x is self.. skipped ", label, self.addr.OverlayAddr())
|
log.Warn(fmt.Sprintf("[%06s] add peers: %x is self.. skipped ", label, self.addr.OverlayAddr()))
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
p := NewKadPeer(na)
|
p := NewKadPeer(na)
|
||||||
|
|
@ -217,7 +216,7 @@ func (self *Kademlia) Register(nas ...PeerAddr) error {
|
||||||
self.peers.Each(func(val pot.PotVal, i int) bool {
|
self.peers.Each(func(val pot.PotVal, i int) bool {
|
||||||
_, found := m[val.String()]
|
_, found := m[val.String()]
|
||||||
// TODO: remove this check
|
// TODO: remove this check
|
||||||
// glog.V(logger.Debug).Infof("-> %v %v", val, i)
|
// log.Debug(fmt.Sprintf("-> %v %v", val, i))
|
||||||
if found {
|
if found {
|
||||||
panic("duplicate found")
|
panic("duplicate found")
|
||||||
}
|
}
|
||||||
|
|
@ -249,19 +248,19 @@ func (self *Kademlia) On(p Peer) {
|
||||||
|
|
||||||
vp, ok := kp.Peer.(KadDiscovery)
|
vp, ok := kp.Peer.(KadDiscovery)
|
||||||
if !ok {
|
if !ok {
|
||||||
glog.V(logger.Detail).Infof("not discovery peer %T", kp)
|
log.Trace(fmt.Sprintf("not discovery peer %T", kp))
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
go vp.NotifyProx(uint8(prox))
|
go vp.NotifyProx(uint8(prox))
|
||||||
f := func(val pot.PotVal, po int) {
|
f := func(val pot.PotVal, po int) {
|
||||||
dp := val.(*KadPeer).Peer.(KadDiscovery)
|
dp := val.(*KadPeer).Peer.(KadDiscovery)
|
||||||
glog.V(logger.Debug).Infof("peer %v notified of %v (%v)", dp, kp, po)
|
log.Debug(fmt.Sprintf("peer %v notified of %v (%v)", dp, kp, po))
|
||||||
dp.NotifyPeer(kp.Peer, uint8(po))
|
dp.NotifyPeer(kp.Peer, uint8(po))
|
||||||
if uint8(prox) != self.lastProxLimit {
|
if uint8(prox) != self.lastProxLimit {
|
||||||
self.lastProxLimit = uint8(prox)
|
self.lastProxLimit = uint8(prox)
|
||||||
dp.NotifyProx(uint8(prox))
|
dp.NotifyProx(uint8(prox))
|
||||||
}
|
}
|
||||||
glog.V(logger.Debug).Infof("peer notified")
|
log.Debug("peer notified")
|
||||||
}
|
}
|
||||||
self.conns.EachNeighbourAsync(kp, 1024, 255, f, false)
|
self.conns.EachNeighbourAsync(kp, 1024, 255, f, false)
|
||||||
}
|
}
|
||||||
|
|
@ -349,7 +348,7 @@ func (self *Kademlia) SuggestPeer() (p PeerAddr, o int, want bool) {
|
||||||
proxLimit := self.proxLimit()
|
proxLimit := self.proxLimit()
|
||||||
// if there is a callable neighbour within the current proxBin, connect
|
// if there is a callable neighbour within the current proxBin, connect
|
||||||
// this makes sure nearest neighbour set is fully connected
|
// this makes sure nearest neighbour set is fully connected
|
||||||
glog.V(logger.Detail).Infof("candidate prox peer checking above PO %v", proxLimit)
|
log.Trace(fmt.Sprintf("candidate prox peer checking above PO %v", proxLimit))
|
||||||
var ppo int
|
var ppo int
|
||||||
self.peers.EachNeighbour(self.addr, func(val pot.PotVal, po int) bool {
|
self.peers.EachNeighbour(self.addr, func(val pot.PotVal, po int) bool {
|
||||||
r := self.callable(val)
|
r := self.callable(val)
|
||||||
|
|
@ -361,15 +360,15 @@ func (self *Kademlia) SuggestPeer() (p PeerAddr, o int, want bool) {
|
||||||
return false
|
return false
|
||||||
})
|
})
|
||||||
if p != nil {
|
if p != nil {
|
||||||
glog.V(logger.Detail).Infof("candidate prox peer found: %v (%v), %v", p, ppo, p)
|
log.Trace(fmt.Sprintf("candidate prox peer found: %v (%v), %v", p, ppo, p))
|
||||||
return p, 0, false
|
return p, 0, false
|
||||||
}
|
}
|
||||||
glog.V(logger.Detail).Infof("no candidate prox peers to connect to (ProxLimit: %v, minProxSize: %v)", proxLimit, self.MinProxBinSize)
|
log.Trace(fmt.Sprintf("no candidate prox peers to connect to (ProxLimit: %v, minProxSize: %v)", proxLimit, self.MinProxBinSize))
|
||||||
|
|
||||||
var bpo []int
|
var bpo []int
|
||||||
prev := -1
|
prev := -1
|
||||||
self.conns.EachBin(self.addr, 0, func(po, size int, f func(func(val pot.PotVal, i int) bool) bool) bool {
|
self.conns.EachBin(self.addr, 0, func(po, size int, f func(func(val pot.PotVal, i int) bool) bool) bool {
|
||||||
glog.V(logger.Detail).Infof("check PO%02d: ", po)
|
log.Trace(fmt.Sprintf("check PO%02d: ", po))
|
||||||
prev++
|
prev++
|
||||||
if po > prev {
|
if po > prev {
|
||||||
size = 0
|
size = 0
|
||||||
|
|
@ -395,7 +394,7 @@ func (self *Kademlia) SuggestPeer() (p PeerAddr, o int, want bool) {
|
||||||
// for each bin we find callable candidate peers
|
// for each bin we find callable candidate peers
|
||||||
f(func(val pot.PotVal, i int) bool {
|
f(func(val pot.PotVal, i int) bool {
|
||||||
r := self.callable(val)
|
r := self.callable(val)
|
||||||
glog.V(logger.Detail).Infof("check PO%02d: ", po)
|
log.Trace(fmt.Sprintf("check PO%02d: ", po))
|
||||||
if r == nil {
|
if r == nil {
|
||||||
return i < proxLimit
|
return i < proxLimit
|
||||||
}
|
}
|
||||||
|
|
@ -472,7 +471,7 @@ func (self *Kademlia) String() string {
|
||||||
rowlen++
|
rowlen++
|
||||||
return rowlen < 4
|
return rowlen < 4
|
||||||
})
|
})
|
||||||
// glog.V(logger.Debug).Infof("po: %v, peerrows length: %v, maxProxDisplay: %v", po, len(peersrows), self.MaxProxDisplay)
|
// log.Debug(fmt.Sprintf("po: %v, peerrows length: %v, maxProxDisplay: %v", po, len(peersrows), self.MaxProxDisplay))
|
||||||
// if po < self.MaxProxDisplay {
|
// if po < self.MaxProxDisplay {
|
||||||
// peersrows[po] = strings.Join(row, " ")
|
// peersrows[po] = strings.Join(row, " ")
|
||||||
// }
|
// }
|
||||||
|
|
|
||||||
|
|
@ -21,8 +21,7 @@ import (
|
||||||
"testing"
|
"testing"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
// "github.com/ethereum/go-ethereum/logger"
|
"github.com/ethereum/go-ethereum/log"
|
||||||
"github.com/ethereum/go-ethereum/logger/glog"
|
|
||||||
"github.com/ethereum/go-ethereum/pot"
|
"github.com/ethereum/go-ethereum/pot"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
@ -111,7 +110,7 @@ func (k *testKademlia) newTestKadPeer(s string) Peer {
|
||||||
}
|
}
|
||||||
|
|
||||||
func overlayStr(a PeerAddr) string {
|
func overlayStr(a PeerAddr) string {
|
||||||
glog.V(6).Infof("PeerAddr: %v (%T)", a, a)
|
log.Error(fmt.Sprintf("PeerAddr: %v (%T)", a, a))
|
||||||
// if a == (*KadPeer)(nil) || a == (*testDiscPeer)(nil) || a == (*bzzPeer)(nil) || a == nil {
|
// if a == (*KadPeer)(nil) || a == (*testDiscPeer)(nil) || a == (*bzzPeer)(nil) || a == nil {
|
||||||
// return "<nil>"
|
// return "<nil>"
|
||||||
// }
|
// }
|
||||||
|
|
@ -122,7 +121,7 @@ func overlayStr(a PeerAddr) string {
|
||||||
// } else {
|
// } else {
|
||||||
// p = a.(*testDiscPeer).Peer
|
// p = a.(*testDiscPeer).Peer
|
||||||
// }
|
// }
|
||||||
// glog.V(6).Infof("PeerAddr: %v (%T)", p, p)
|
// log.Error(fmt.Sprintf("PeerAddr: %v (%T)", p, p))
|
||||||
// if p == (Peer)(nil) || p == (*testDiscPeer)(nil) || p == (*bzzPeer)(nil) {
|
// if p == (Peer)(nil) || p == (*testDiscPeer)(nil) || p == (*bzzPeer)(nil) {
|
||||||
// return "<nil>"
|
// return "<nil>"
|
||||||
// }
|
// }
|
||||||
|
|
|
||||||
|
|
@ -21,7 +21,7 @@ import (
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/ethereum/go-ethereum/crypto"
|
"github.com/ethereum/go-ethereum/crypto"
|
||||||
"github.com/ethereum/go-ethereum/logger/glog"
|
"github.com/ethereum/go-ethereum/log"
|
||||||
"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/adapters"
|
||||||
"github.com/ethereum/go-ethereum/p2p/discover"
|
"github.com/ethereum/go-ethereum/p2p/discover"
|
||||||
|
|
@ -87,7 +87,7 @@ func Bzz(oAddr, uAddr []byte, ct *protocols.CodeMap, services func(Peer) error,
|
||||||
// sets remote peer address
|
// sets remote peer address
|
||||||
err := bee.bzzHandshake()
|
err := bee.bzzHandshake()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
glog.V(6).Infof("handshake error in peer %v: %v", bee.ID(), err)
|
log.Error(fmt.Sprintf("handshake error in peer %v: %v", bee.ID(), err))
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -95,7 +95,7 @@ func Bzz(oAddr, uAddr []byte, ct *protocols.CodeMap, services func(Peer) error,
|
||||||
if services != nil {
|
if services != nil {
|
||||||
err = services(bee)
|
err = services(bee)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
glog.V(6).Infof("protocol service error for peer %v: %v", bee.ID(), err)
|
log.Error(fmt.Sprintf("protocol service error for peer %v: %v", bee.ID(), err))
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -183,7 +183,7 @@ func (self *bzzPeer) bzzHandshake() error {
|
||||||
|
|
||||||
hs, err := self.Handshake(lhs)
|
hs, err := self.Handshake(lhs)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
glog.V(6).Infof("handshake failed: %v", err)
|
log.Error(fmt.Sprintf("handshake failed: %v", err))
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -191,7 +191,7 @@ func (self *bzzPeer) bzzHandshake() error {
|
||||||
self.peerAddr = rhs.Addr
|
self.peerAddr = rhs.Addr
|
||||||
err = checkBzzHandshake(rhs)
|
err = checkBzzHandshake(rhs)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
glog.V(6).Infof("handshake between %v and %v failed: %v", self.localAddr, self.peerAddr, err)
|
log.Error(fmt.Sprintf("handshake between %v and %v failed: %v", self.localAddr, self.peerAddr, err))
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -4,8 +4,7 @@ import (
|
||||||
"fmt"
|
"fmt"
|
||||||
"testing"
|
"testing"
|
||||||
|
|
||||||
"github.com/ethereum/go-ethereum/logger"
|
"github.com/ethereum/go-ethereum/log"
|
||||||
"github.com/ethereum/go-ethereum/logger/glog"
|
|
||||||
"github.com/ethereum/go-ethereum/p2p/adapters"
|
"github.com/ethereum/go-ethereum/p2p/adapters"
|
||||||
"github.com/ethereum/go-ethereum/p2p/protocols"
|
"github.com/ethereum/go-ethereum/p2p/protocols"
|
||||||
p2ptest "github.com/ethereum/go-ethereum/p2p/testing"
|
p2ptest "github.com/ethereum/go-ethereum/p2p/testing"
|
||||||
|
|
@ -166,7 +165,7 @@ func TestBzzPeerPoolAdd(t *testing.T) {
|
||||||
defer s.Stop()
|
defer s.Stop()
|
||||||
|
|
||||||
id := s.Ids[0]
|
id := s.Ids[0]
|
||||||
glog.V(logger.Detail).Infof("handshake with %v", id)
|
log.Trace(fmt.Sprintf("handshake with %v", id))
|
||||||
s.runHandshakes()
|
s.runHandshakes()
|
||||||
|
|
||||||
if !pp.Has(id) {
|
if !pp.Has(id) {
|
||||||
|
|
|
||||||
|
|
@ -1,30 +1,29 @@
|
||||||
package network
|
package network
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"fmt"
|
|
||||||
"bytes"
|
"bytes"
|
||||||
|
"fmt"
|
||||||
|
|
||||||
"github.com/ethereum/go-ethereum/logger"
|
"github.com/ethereum/go-ethereum/log"
|
||||||
"github.com/ethereum/go-ethereum/logger/glog"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
type Pss struct {
|
type Pss struct {
|
||||||
Overlay
|
Overlay
|
||||||
LocalAddr []byte
|
LocalAddr []byte
|
||||||
C chan []byte
|
C chan []byte
|
||||||
}
|
}
|
||||||
|
|
||||||
func NewPss(k Overlay, addr []byte) *Pss {
|
func NewPss(k Overlay, addr []byte) *Pss {
|
||||||
return &Pss{
|
return &Pss{
|
||||||
Overlay: k,
|
Overlay: k,
|
||||||
LocalAddr: addr,
|
LocalAddr: addr,
|
||||||
C: make(chan []byte),
|
C: make(chan []byte),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
type PssMsg struct {
|
type PssMsg struct {
|
||||||
To []byte
|
To []byte
|
||||||
Data []byte
|
Data []byte
|
||||||
}
|
}
|
||||||
|
|
||||||
func (pm *PssMsg) String() string {
|
func (pm *PssMsg) String() string {
|
||||||
|
|
@ -35,7 +34,7 @@ func (ps *Pss) HandlePssMsg(msg interface{}) error {
|
||||||
pssmsg := msg.(*PssMsg)
|
pssmsg := msg.(*PssMsg)
|
||||||
to := pssmsg.To
|
to := pssmsg.To
|
||||||
if bytes.Equal(to, ps.LocalAddr) {
|
if bytes.Equal(to, ps.LocalAddr) {
|
||||||
glog.V(logger.Detail).Infof("Pss to us, yay!", to)
|
log.Trace(fmt.Sprintf("Pss to us, yay! %v", to))
|
||||||
ps.C <- pssmsg.Data
|
ps.C <- pssmsg.Data
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
@ -47,6 +46,6 @@ func (ps *Pss) HandlePssMsg(msg interface{}) error {
|
||||||
}
|
}
|
||||||
return false
|
return false
|
||||||
})
|
})
|
||||||
|
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -5,8 +5,6 @@ import (
|
||||||
"testing"
|
"testing"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
//"github.com/ethereum/go-ethereum/logger"
|
|
||||||
//"github.com/ethereum/go-ethereum/logger/glog"
|
|
||||||
"github.com/ethereum/go-ethereum/p2p/adapters"
|
"github.com/ethereum/go-ethereum/p2p/adapters"
|
||||||
"github.com/ethereum/go-ethereum/p2p/protocols"
|
"github.com/ethereum/go-ethereum/p2p/protocols"
|
||||||
"github.com/ethereum/go-ethereum/p2p/simulations"
|
"github.com/ethereum/go-ethereum/p2p/simulations"
|
||||||
|
|
|
||||||
|
|
@ -6,13 +6,14 @@
|
||||||
package main
|
package main
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"fmt"
|
||||||
"math/rand"
|
"math/rand"
|
||||||
|
"os"
|
||||||
"reflect"
|
"reflect"
|
||||||
"runtime"
|
"runtime"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/ethereum/go-ethereum/logger"
|
"github.com/ethereum/go-ethereum/log"
|
||||||
"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/adapters"
|
||||||
"github.com/ethereum/go-ethereum/p2p/simulations"
|
"github.com/ethereum/go-ethereum/p2p/simulations"
|
||||||
|
|
@ -79,7 +80,7 @@ func (self *Network) NewSimNode(conf *simulations.NodeConfig) adapters.NodeAdapt
|
||||||
services := func(p network.Peer) error {
|
services := func(p network.Peer) error {
|
||||||
dp := network.NewDiscovery(p, to)
|
dp := network.NewDiscovery(p, to)
|
||||||
pp.Add(dp)
|
pp.Add(dp)
|
||||||
glog.V(logger.Detail).Infof("kademlia on %v", dp)
|
log.Trace(fmt.Sprintf("kademlia on %v", dp))
|
||||||
p.DisconnectHook(func(err error) {
|
p.DisconnectHook(func(err error) {
|
||||||
pp.Remove(dp)
|
pp.Remove(dp)
|
||||||
})
|
})
|
||||||
|
|
@ -138,7 +139,7 @@ func nethook(conf *simulations.NetworkConfig) (simulations.NetworkControl, *simu
|
||||||
time.Sleep(time.Duration(n) * time.Millisecond)
|
time.Sleep(time.Duration(n) * time.Millisecond)
|
||||||
net.NewNode(&simulations.NodeConfig{Id: id})
|
net.NewNode(&simulations.NodeConfig{Id: id})
|
||||||
net.Start(id)
|
net.Start(id)
|
||||||
glog.V(logger.Debug).Infof("node %v starting up", id)
|
log.Debug(fmt.Sprintf("node %v starting up", id))
|
||||||
// time.Sleep(1000 * time.Millisecond)
|
// time.Sleep(1000 * time.Millisecond)
|
||||||
// net.Stop(id)
|
// net.Stop(id)
|
||||||
}
|
}
|
||||||
|
|
@ -150,12 +151,12 @@ func nethook(conf *simulations.NetworkConfig) (simulations.NetworkControl, *simu
|
||||||
// // n := rand.Intn(5000)
|
// // n := rand.Intn(5000)
|
||||||
// // n := 3000
|
// // n := 3000
|
||||||
// time.Sleep(time.Duration(n) * time.Millisecond)
|
// time.Sleep(time.Duration(n) * time.Millisecond)
|
||||||
// glog.V(logger.Debug).Infof("node %v shutting down", id)
|
// log.Debug(fmt.Sprintf("node %v shutting down", id))
|
||||||
// net.Stop(id)
|
// net.Stop(id)
|
||||||
// // n = rand.Intn(5000)
|
// // n = rand.Intn(5000)
|
||||||
// n = 2000
|
// n = 2000
|
||||||
// time.Sleep(time.Duration(n) * time.Millisecond)
|
// time.Sleep(time.Duration(n) * time.Millisecond)
|
||||||
// glog.V(logger.Debug).Infof("node %v starting up", id)
|
// log.Debug(fmt.Sprintf("node %v starting up", id))
|
||||||
// net.Start(id)
|
// net.Start(id)
|
||||||
// n = 5000
|
// n = 5000
|
||||||
// }
|
// }
|
||||||
|
|
@ -178,8 +179,8 @@ func nethook(conf *simulations.NetworkConfig) (simulations.NetworkControl, *simu
|
||||||
// var server
|
// var server
|
||||||
func main() {
|
func main() {
|
||||||
runtime.GOMAXPROCS(runtime.NumCPU())
|
runtime.GOMAXPROCS(runtime.NumCPU())
|
||||||
glog.SetV(logger.Detail)
|
|
||||||
glog.SetToStderr(true)
|
log.Root().SetHandler(log.LvlFilterHandler(log.LvlTrace, log.StreamHandler(os.Stderr, log.TerminalFormat(false))))
|
||||||
|
|
||||||
c, quitc := simulations.NewSessionController(nethook)
|
c, quitc := simulations.NewSessionController(nethook)
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -5,9 +5,7 @@ import (
|
||||||
"strings"
|
"strings"
|
||||||
"sync"
|
"sync"
|
||||||
|
|
||||||
// "github.com/ethereum/go-ethereum/p2p/adapters"
|
"github.com/ethereum/go-ethereum/log"
|
||||||
// "github.com/ethereum/go-ethereum/p2p/discover"
|
|
||||||
"github.com/ethereum/go-ethereum/logger/glog"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
const orders = 8
|
const orders = 8
|
||||||
|
|
@ -39,7 +37,7 @@ func (self *testOverlay) register(nas ...PeerAddr) error {
|
||||||
}
|
}
|
||||||
self.posMap[string(addr)] = tna
|
self.posMap[string(addr)] = tna
|
||||||
o := order(addr)
|
o := order(addr)
|
||||||
glog.V(6).Infof("PO: %v, orders: %v", o, orders)
|
log.Trace(fmt.Sprintf("PO: %v, orders: %v", o, orders))
|
||||||
self.pos[o] = append(self.pos[o], tna)
|
self.pos[o] = append(self.pos[o], tna)
|
||||||
}
|
}
|
||||||
return nil
|
return nil
|
||||||
|
|
@ -60,7 +58,7 @@ func (self *testOverlay) On(n Peer) {
|
||||||
} else if na.Peer != nil {
|
} else if na.Peer != nil {
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
glog.V(6).Infof("Online: %v", fmt.Sprintf("%x", addr[:4]))
|
log.Trace(fmt.Sprintf("Online: %x", addr[:4]))
|
||||||
na.Peer = n
|
na.Peer = n
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
@ -133,7 +131,7 @@ func (self *testOverlay) SuggestPeer() (PeerAddr, int, bool) {
|
||||||
if len(ons) < 2 {
|
if len(ons) < 2 {
|
||||||
offs := self.off(po)
|
offs := self.off(po)
|
||||||
if len(offs) > 0 {
|
if len(offs) > 0 {
|
||||||
glog.V(6).Infof("node %v is off", offs[0])
|
log.Trace(fmt.Sprintf("node %v is off", offs[0]))
|
||||||
return offs[0], i, true
|
return offs[0], i, true
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue