les: simplified handling of trusted nodes in the server pool

This commit is contained in:
Zsolt Felfoldi 2019-01-09 14:19:21 +01:00
parent 1932f58429
commit 8b6af61abe

View file

@ -126,12 +126,12 @@ type serverPool struct {
discNodes chan *enode.Node discNodes chan *enode.Node
discLookups chan bool discLookups chan bool
trustedNodes []string trustedNodes map[enode.ID]*enode.Node
entries map[enode.ID]*poolEntry entries map[enode.ID]*poolEntry
timeout, enableRetry chan *poolEntry timeout, enableRetry chan *poolEntry
adjustStats chan poolStatAdjust adjustStats chan poolStatAdjust
knownQueue, newQueue, trustedQueue poolEntryQueue knownQueue, newQueue poolEntryQueue
knownSelect, newSelect *weightedRandomSelect knownSelect, newSelect *weightedRandomSelect
knownSelected, newSelected int knownSelected, newSelected int
fastDiscover bool fastDiscover bool
@ -156,10 +156,9 @@ func newServerPool(db ethdb.Database, quit chan struct{}, wg *sync.WaitGroup, tr
knownSelect: newWeightedRandomSelect(), knownSelect: newWeightedRandomSelect(),
newSelect: newWeightedRandomSelect(), newSelect: newWeightedRandomSelect(),
fastDiscover: true, fastDiscover: true,
trustedNodes: trustedNodes, trustedNodes: parseTrustedNodes(trustedNodes),
} }
pool.trustedQueue = newPoolEntryQueue(maxKnownEntries, pool.removeEntry)
pool.knownQueue = newPoolEntryQueue(maxKnownEntries, pool.removeEntry) pool.knownQueue = newPoolEntryQueue(maxKnownEntries, pool.removeEntry)
pool.newQueue = newPoolEntryQueue(maxNewEntries, pool.removeEntry) pool.newQueue = newPoolEntryQueue(maxNewEntries, pool.removeEntry)
return pool return pool
@ -205,7 +204,7 @@ func (pool *serverPool) discoverNodes() {
// Note that whenever a connection has been accepted and a pool entry has been returned, // Note that whenever a connection has been accepted and a pool entry has been returned,
// disconnect should also always be called. // disconnect should also always be called.
func (pool *serverPool) connect(p *peer, node *enode.Node) *poolEntry { func (pool *serverPool) connect(p *peer, node *enode.Node) *poolEntry {
log.Debug("Connect new entry", "enode", p.id) log.Info("Connect new entry", "enode", p.id)
req := &connReq{p: p, node: node, result: make(chan *poolEntry, 1)} req := &connReq{p: p, node: node, result: make(chan *poolEntry, 1)}
select { select {
case pool.connCh <- req: case pool.connCh <- req:
@ -217,7 +216,7 @@ func (pool *serverPool) connect(p *peer, node *enode.Node) *poolEntry {
// registered should be called after a successful handshake // registered should be called after a successful handshake
func (pool *serverPool) registered(entry *poolEntry) { func (pool *serverPool) registered(entry *poolEntry) {
log.Debug("Registered new entry", "enode", entry.node.ID()) log.Info("Registered new entry", "enode", entry.node.ID())
req := &registerReq{entry: entry, done: make(chan struct{})} req := &registerReq{entry: entry, done: make(chan struct{})}
select { select {
case pool.registerCh <- req: case pool.registerCh <- req:
@ -237,7 +236,7 @@ func (pool *serverPool) disconnect(entry *poolEntry) {
stopped = true stopped = true
default: default:
} }
log.Debug("Disconnected old entry", "enode", entry.node.ID()) log.Info("Disconnected old entry", "enode", entry.node.ID())
req := &disconnReq{entry: entry, stopped: stopped, done: make(chan struct{})} req := &disconnReq{entry: entry, stopped: stopped, done: make(chan struct{})}
// Block until disconnection request is served. // Block until disconnection request is served.
@ -305,6 +304,9 @@ func (pool *serverPool) eventLoop() {
} }
} }
entry.state = psNotConnected entry.state = psNotConnected
if pool.trustedNodes[entry.node.ID()] == nil {
pool.server.RemovePeer(entry.node)
}
if entry.knownSelected { if entry.knownSelected {
pool.knownSelected-- pool.knownSelected--
@ -341,8 +343,10 @@ func (pool *serverPool) eventLoop() {
} }
case node := <-pool.discNodes: case node := <-pool.discNodes:
if pool.trustedNodes[node.ID()] == nil {
entry := pool.findOrNewNode(node) entry := pool.findOrNewNode(node)
pool.updateCheckDial(entry) pool.updateCheckDial(entry)
}
case conv := <-pool.discLookups: case conv := <-pool.discLookups:
if conv { if conv {
@ -359,6 +363,10 @@ func (pool *serverPool) eventLoop() {
} }
case req := <-pool.connCh: case req := <-pool.connCh:
if pool.trustedNodes[req.p.ID()] != nil {
// ignore trusted nodes
req.result <- nil
} else {
// Handle peer connection requests. // Handle peer connection requests.
entry := pool.entries[req.p.ID()] entry := pool.entries[req.p.ID()]
if entry == nil { if entry == nil {
@ -382,6 +390,7 @@ func (pool *serverPool) eventLoop() {
entry.addrSelect = *newWeightedRandomSelect() entry.addrSelect = *newWeightedRandomSelect()
entry.addrSelect.update(addr) entry.addrSelect.update(addr)
req.result <- entry req.result <- entry
}
case req := <-pool.registerCh: case req := <-pool.registerCh:
// Handle peer registration requests. // Handle peer registration requests.
@ -426,7 +435,7 @@ func (pool *serverPool) findOrNewNode(node *enode.Node) *poolEntry {
now := mclock.Now() now := mclock.Now()
entry := pool.entries[node.ID()] entry := pool.entries[node.ID()]
if entry == nil { if entry == nil {
log.Debug("Discovered new entry", "id", node.ID()) log.Info("Discovered new entry", "id", node.ID())
entry = &poolEntry{ entry = &poolEntry{
node: node, node: node,
addr: make(map[string]*poolEntryAddress), addr: make(map[string]*poolEntryAddress),
@ -449,7 +458,7 @@ func (pool *serverPool) findOrNewNode(node *enode.Node) *poolEntry {
} }
addr.lastSeen = now addr.lastSeen = now
entry.addrSelect.update(addr) entry.addrSelect.update(addr)
if !entry.known || !entry.trusted { if !entry.known {
pool.newQueue.setLatest(entry) pool.newQueue.setLatest(entry)
} }
return entry return entry
@ -464,45 +473,48 @@ func (pool *serverPool) loadNodes() {
var list []*poolEntry var list []*poolEntry
err = rlp.DecodeBytes(enc, &list) err = rlp.DecodeBytes(enc, &list)
if err != nil { if err != nil {
log.Debug("Failed to decode node list", "err", err) log.Info("Failed to decode node list", "err", err)
return return
} }
for _, e := range list { for _, e := range list {
log.Debug("Loaded server stats", "id", e.node.ID(), "fails", e.lastConnected.fails, log.Info("Loaded server stats", "id", e.node.ID(), "fails", e.lastConnected.fails,
"conn", fmt.Sprintf("%v/%v", e.connectStats.avg, e.connectStats.weight), "conn", fmt.Sprintf("%v/%v", e.connectStats.avg, e.connectStats.weight),
"delay", fmt.Sprintf("%v/%v", time.Duration(e.delayStats.avg), e.delayStats.weight), "delay", fmt.Sprintf("%v/%v", time.Duration(e.delayStats.avg), e.delayStats.weight),
"response", fmt.Sprintf("%v/%v", time.Duration(e.responseStats.avg), e.responseStats.weight), "response", fmt.Sprintf("%v/%v", time.Duration(e.responseStats.avg), e.responseStats.weight),
"timeout", fmt.Sprintf("%v/%v", e.timeoutStats.avg, e.timeoutStats.weight)) "timeout", fmt.Sprintf("%v/%v", e.timeoutStats.avg, e.timeoutStats.weight))
pool.entries[e.node.ID()] = e pool.entries[e.node.ID()] = e
if pool.trustedNodes[e.node.ID()] == nil {
pool.knownQueue.setLatest(e) pool.knownQueue.setLatest(e)
pool.knownSelect.update((*knownEntry)(e)) pool.knownSelect.update((*knownEntry)(e))
} }
} }
}
// connectToTrustedNodes adds trusted server nodes as static trusted peers.
//
// Note: trusted nodes are not handled by the server pool logic, they are not
// added to either the known or new selection pools. They are connected/reconnected
// by p2p.Server whenever possible.
func (pool *serverPool) connectToTrustedNodes() { func (pool *serverPool) connectToTrustedNodes() {
//connect to trusted nodes //connect to trusted nodes
if len(pool.trustedNodes) > 0 {
for _, trusted := range pool.parseTrustedServers() {
e := pool.findOrNewNode(trusted)
e.trusted = true
e.dialed = &poolEntryAddress{ip: trusted.IP(), port: uint16(trusted.TCP())}
pool.entries[e.node.ID()] = e
pool.trustedQueue.setLatest(e)
}
}
}
// parseTrustedServers returns valid and parsed by discovery enodes.
func (pool *serverPool) parseTrustedServers() []*enode.Node {
nodes := make([]*enode.Node, 0, len(pool.trustedNodes))
for _, node := range pool.trustedNodes { for _, node := range pool.trustedNodes {
pool.server.AddTrustedPeer(node)
pool.server.AddPeer(node)
log.Info("Added trusted node", "id", node.ID().String())
}
}
// parseTrustedNodes returns valid and parsed enodes
func parseTrustedNodes(trustedNodes []string) map[enode.ID]*enode.Node {
nodes := make(map[enode.ID]*enode.Node)
for _, node := range trustedNodes {
node, err := enode.ParseV4(node) node, err := enode.ParseV4(node)
if err != nil { if err != nil {
log.Warn("Trusted node URL invalid", "enode", node, "err", err) log.Warn("Trusted node URL invalid", "enode", node, "err", err)
continue continue
} }
nodes = append(nodes, node) nodes[node.ID()] = node
} }
return nodes return nodes
} }
@ -562,10 +574,6 @@ func (pool *serverPool) updateCheckDial(entry *poolEntry) {
// checkDial checks if new dials can/should be made. It tries to select servers both // checkDial checks if new dials can/should be made. It tries to select servers both
// based on good statistics and recent discovery. // based on good statistics and recent discovery.
func (pool *serverPool) checkDial() { func (pool *serverPool) checkDial() {
for _, e := range pool.trustedQueue.queue {
pool.dial(e, false)
}
fillWithKnownSelects := !pool.fastDiscover fillWithKnownSelects := !pool.fastDiscover
for pool.knownSelected < targetKnownSelect { for pool.knownSelected < targetKnownSelect {
entry := pool.knownSelect.choose() entry := pool.knownSelect.choose()
@ -602,8 +610,6 @@ func (pool *serverPool) dial(entry *poolEntry, knownSelected bool) {
return return
} }
entry.state = psDialed entry.state = psDialed
if !entry.trusted {
entry.knownSelected = knownSelected entry.knownSelected = knownSelected
if knownSelected { if knownSelected {
pool.knownSelected++ pool.knownSelected++
@ -611,15 +617,8 @@ func (pool *serverPool) dial(entry *poolEntry, knownSelected bool) {
pool.newSelected++ pool.newSelected++
} }
addr := entry.addrSelect.choose().(*poolEntryAddress) addr := entry.addrSelect.choose().(*poolEntryAddress)
log.Info("Dialing new peer", "lesaddr", entry.node.ID().String()+"@"+addr.strKey(), "set", len(entry.addr), "known", knownSelected)
entry.dialed = addr entry.dialed = addr
}
state := "known"
if entry.trusted {
state = "trusted"
}
log.Debug("Dialing new peer", "lesaddr", entry.node.ID().String()+"@"+entry.dialed.strKey(), "set", len(entry.addr), state, knownSelected)
go func() { go func() {
pool.server.AddPeer(entry.node) pool.server.AddPeer(entry.node)
select { select {
@ -639,8 +638,11 @@ func (pool *serverPool) checkDialTimeout(entry *poolEntry) {
if entry.state != psDialed { if entry.state != psDialed {
return return
} }
log.Debug("Dial timeout", "lesaddr", entry.node.ID().String()+"@"+entry.dialed.strKey()) log.Info("Dial timeout", "lesaddr", entry.node.ID().String()+"@"+entry.dialed.strKey())
entry.state = psNotConnected entry.state = psNotConnected
if pool.trustedNodes[entry.node.ID()] == nil {
pool.server.RemovePeer(entry.node)
}
if entry.knownSelected { if entry.knownSelected {
pool.knownSelected-- pool.knownSelected--
} else { } else {
@ -669,7 +671,6 @@ type poolEntry struct {
lastDiscovered mclock.AbsTime lastDiscovered mclock.AbsTime
known, knownSelected bool known, knownSelected bool
trusted bool
connectStats, delayStats poolStats connectStats, delayStats poolStats
responseStats, timeoutStats poolStats responseStats, timeoutStats poolStats
state int state int