From 8b6af61abe9ccb9b4c8dc1988724f0faea4329f3 Mon Sep 17 00:00:00 2001 From: Zsolt Felfoldi Date: Wed, 9 Jan 2019 14:19:21 +0100 Subject: [PATCH] les: simplified handling of trusted nodes in the server pool --- les/serverpool.go | 165 +++++++++++++++++++++++----------------------- 1 file changed, 83 insertions(+), 82 deletions(-) diff --git a/les/serverpool.go b/les/serverpool.go index 05d63a917a..35da86f3b6 100644 --- a/les/serverpool.go +++ b/les/serverpool.go @@ -126,18 +126,18 @@ type serverPool struct { discNodes chan *enode.Node discLookups chan bool - trustedNodes []string + trustedNodes map[enode.ID]*enode.Node entries map[enode.ID]*poolEntry timeout, enableRetry chan *poolEntry adjustStats chan poolStatAdjust - knownQueue, newQueue, trustedQueue poolEntryQueue - knownSelect, newSelect *weightedRandomSelect - knownSelected, newSelected int - fastDiscover bool - connCh chan *connReq - disconnCh chan *disconnReq - registerCh chan *registerReq + knownQueue, newQueue poolEntryQueue + knownSelect, newSelect *weightedRandomSelect + knownSelected, newSelected int + fastDiscover bool + connCh chan *connReq + disconnCh chan *disconnReq + registerCh chan *registerReq } // newServerPool creates a new serverPool instance @@ -156,10 +156,9 @@ func newServerPool(db ethdb.Database, quit chan struct{}, wg *sync.WaitGroup, tr knownSelect: newWeightedRandomSelect(), newSelect: newWeightedRandomSelect(), fastDiscover: true, - trustedNodes: trustedNodes, + trustedNodes: parseTrustedNodes(trustedNodes), } - pool.trustedQueue = newPoolEntryQueue(maxKnownEntries, pool.removeEntry) pool.knownQueue = newPoolEntryQueue(maxKnownEntries, pool.removeEntry) pool.newQueue = newPoolEntryQueue(maxNewEntries, pool.removeEntry) 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, // disconnect should also always be called. 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)} select { 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 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 := ®isterReq{entry: entry, done: make(chan struct{})} select { case pool.registerCh <- req: @@ -237,7 +236,7 @@ func (pool *serverPool) disconnect(entry *poolEntry) { stopped = true 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{})} // Block until disconnection request is served. @@ -305,6 +304,9 @@ func (pool *serverPool) eventLoop() { } } entry.state = psNotConnected + if pool.trustedNodes[entry.node.ID()] == nil { + pool.server.RemovePeer(entry.node) + } if entry.knownSelected { pool.knownSelected-- @@ -341,8 +343,10 @@ func (pool *serverPool) eventLoop() { } case node := <-pool.discNodes: - entry := pool.findOrNewNode(node) - pool.updateCheckDial(entry) + if pool.trustedNodes[node.ID()] == nil { + entry := pool.findOrNewNode(node) + pool.updateCheckDial(entry) + } case conv := <-pool.discLookups: if conv { @@ -359,29 +363,34 @@ func (pool *serverPool) eventLoop() { } case req := <-pool.connCh: - // Handle peer connection requests. - entry := pool.entries[req.p.ID()] - if entry == nil { - entry = pool.findOrNewNode(req.node) - } - if entry.state == psConnected || entry.state == psRegistered { + if pool.trustedNodes[req.p.ID()] != nil { + // ignore trusted nodes req.result <- nil - continue + } else { + // Handle peer connection requests. + entry := pool.entries[req.p.ID()] + if entry == nil { + entry = pool.findOrNewNode(req.node) + } + if entry.state == psConnected || entry.state == psRegistered { + req.result <- nil + continue + } + pool.connWg.Add(1) + entry.peer = req.p + entry.state = psConnected + addr := &poolEntryAddress{ + ip: req.node.IP(), + port: uint16(req.node.TCP()), + lastSeen: mclock.Now(), + } + entry.lastConnected = addr + entry.addr = make(map[string]*poolEntryAddress) + entry.addr[addr.strKey()] = addr + entry.addrSelect = *newWeightedRandomSelect() + entry.addrSelect.update(addr) + req.result <- entry } - pool.connWg.Add(1) - entry.peer = req.p - entry.state = psConnected - addr := &poolEntryAddress{ - ip: req.node.IP(), - port: uint16(req.node.TCP()), - lastSeen: mclock.Now(), - } - entry.lastConnected = addr - entry.addr = make(map[string]*poolEntryAddress) - entry.addr[addr.strKey()] = addr - entry.addrSelect = *newWeightedRandomSelect() - entry.addrSelect.update(addr) - req.result <- entry case req := <-pool.registerCh: // Handle peer registration requests. @@ -426,7 +435,7 @@ func (pool *serverPool) findOrNewNode(node *enode.Node) *poolEntry { now := mclock.Now() entry := pool.entries[node.ID()] if entry == nil { - log.Debug("Discovered new entry", "id", node.ID()) + log.Info("Discovered new entry", "id", node.ID()) entry = &poolEntry{ node: node, addr: make(map[string]*poolEntryAddress), @@ -449,7 +458,7 @@ func (pool *serverPool) findOrNewNode(node *enode.Node) *poolEntry { } addr.lastSeen = now entry.addrSelect.update(addr) - if !entry.known || !entry.trusted { + if !entry.known { pool.newQueue.setLatest(entry) } return entry @@ -464,45 +473,48 @@ func (pool *serverPool) loadNodes() { var list []*poolEntry err = rlp.DecodeBytes(enc, &list) if err != nil { - log.Debug("Failed to decode node list", "err", err) + log.Info("Failed to decode node list", "err", err) return } 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), "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), "timeout", fmt.Sprintf("%v/%v", e.timeoutStats.avg, e.timeoutStats.weight)) pool.entries[e.node.ID()] = e - pool.knownQueue.setLatest(e) - pool.knownSelect.update((*knownEntry)(e)) - } -} - -func (pool *serverPool) connectToTrustedNodes() { - //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) + if pool.trustedNodes[e.node.ID()] == nil { + pool.knownQueue.setLatest(e) + pool.knownSelect.update((*knownEntry)(e)) } } } -// parseTrustedServers returns valid and parsed by discovery enodes. -func (pool *serverPool) parseTrustedServers() []*enode.Node { - nodes := make([]*enode.Node, 0, len(pool.trustedNodes)) - +// 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() { + //connect to trusted nodes 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) if err != nil { log.Warn("Trusted node URL invalid", "enode", node, "err", err) continue } - nodes = append(nodes, node) + nodes[node.ID()] = node } 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 // based on good statistics and recent discovery. func (pool *serverPool) checkDial() { - for _, e := range pool.trustedQueue.queue { - pool.dial(e, false) - } - fillWithKnownSelects := !pool.fastDiscover for pool.knownSelected < targetKnownSelect { entry := pool.knownSelect.choose() @@ -602,24 +610,15 @@ func (pool *serverPool) dial(entry *poolEntry, knownSelected bool) { return } entry.state = psDialed - - if !entry.trusted { - entry.knownSelected = knownSelected - if knownSelected { - pool.knownSelected++ - } else { - pool.newSelected++ - } - addr := entry.addrSelect.choose().(*poolEntryAddress) - entry.dialed = addr + entry.knownSelected = knownSelected + if knownSelected { + pool.knownSelected++ + } else { + pool.newSelected++ } - - 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) - + 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 go func() { pool.server.AddPeer(entry.node) select { @@ -639,8 +638,11 @@ func (pool *serverPool) checkDialTimeout(entry *poolEntry) { if entry.state != psDialed { 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 + if pool.trustedNodes[entry.node.ID()] == nil { + pool.server.RemovePeer(entry.node) + } if entry.knownSelected { pool.knownSelected-- } else { @@ -669,7 +671,6 @@ type poolEntry struct { lastDiscovered mclock.AbsTime known, knownSelected bool - trusted bool connectStats, delayStats poolStats responseStats, timeoutStats poolStats state int