Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
19 commits
Select commit Hold shift + click to select a range
0faded3
p2p/discover: evict unresponsive nodes in topic RPC response handling
srene Jun 23, 2026
a94a288
p2p/discover: skip topic RPC failure handling on errClosed shutdown
srene Jul 13, 2026
9e3b89f
p2p/discover/topicindex: free IP-limit slot when removing an asked se…
srene Jul 13, 2026
86a6c09
p2p/discover: decouple DHT-removal feed drain from topic eviction work
srene Jul 13, 2026
2934174
p2p/discover: evict failed search nodes from registration tables too
srene Jul 13, 2026
965448b
p2p/discover: trim comments in topic eviction code
srene Jul 13, 2026
abc6128
p2p/discover: simplify evictRemovedNodes to a best-effort worker hand…
srene Jul 13, 2026
c275e1a
p2p/discover/topicindex: trim RemoveNode doc comment
srene Jul 14, 2026
8b5f927
p2p/discover/topicindex: trim TestRegistrationRemoveNode doc comment
srene Jul 14, 2026
7a6ab3e
p2p/discover/topicindex: trim search error-response doc comments
srene Jul 14, 2026
ef85298
p2p/discover: trim evictRemovedNodes comments
srene Jul 14, 2026
b3f1658
p2p/discover: drop issue references from topic eviction comments
srene Jul 14, 2026
afbf3ee
p2p/discover: log partial-response TOPICQUERY errors
srene Jul 20, 2026
8988d71
p2p/discover: merge topic eviction tests into v5_topic_test.go
srene Jul 20, 2026
191eb25
p2p/discover: drain DHT removals on a single buffered goroutine
srene Jul 20, 2026
62bb322
p2p/discover: clarify topic eviction comments
srene Jul 20, 2026
0fc3afd
go docs
srene Jul 20, 2026
286f15a
p2p/discover: remove test-only registration node-count scaffolding
srene Jul 20, 2026
e1392c2
go docs
srene Jul 20, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 8 additions & 0 deletions p2p/discover/table.go
Original file line number Diff line number Diff line change
Expand Up @@ -87,6 +87,7 @@ type Table struct {
closed chan struct{}

nodeFeed event.FeedOf[*enode.Node]
removedFeed event.FeedOf[enode.ID]
nodeAddedHook func(*bucket, *tableNode)
nodeRemovedHook func(*bucket, *tableNode)
}
Expand Down Expand Up @@ -625,6 +626,7 @@ func (tab *Table) nodeAdded(b *bucket, n *tableNode) {

func (tab *Table) nodeRemoved(b *bucket, n *tableNode) {
tab.revalidation.nodeRemoved(n)
tab.removedFeed.Send(n.ID())
if tab.nodeRemovedHook != nil {
tab.nodeRemovedHook(b, n)
}
Expand Down Expand Up @@ -749,6 +751,12 @@ func (tab *Table) subscribeNodes(ch chan *enode.Node) event.Subscription {
return tab.nodeFeed.Subscribe(ch)
}

// subscribeRemovedNodes subscribes to node removal events. The ID of each node
// removed from the table (e.g. on failed revalidation) is sent on the channel.
func (tab *Table) subscribeRemovedNodes(ch chan enode.ID) event.Subscription {
return tab.removedFeed.Subscribe(ch)
}

// allNodes returns all live nodes in the table.
func (tab *Table) allNodes() []*enode.Node {
tab.mutex.Lock()
Expand Down
15 changes: 15 additions & 0 deletions p2p/discover/topicindex/registration.go
Original file line number Diff line number Diff line change
Expand Up @@ -214,6 +214,21 @@ func (r *Registration) AddNodes(src *enode.Node, nodes []*enode.Node) {
}
}

// RemoveNode drops registration attempts of a node evicted as dead
func (r *Registration) RemoveNode(id enode.ID) {
b := &r.buckets[r.bucketIndex(id)]
att, ok := b.att[id]
if !ok {
return
}
//An in-flight attempt (index == -2) is left for its pending response to remove.
if att.index == -2 {
return
}
r.removeAttempt(att, "evicted")
r.refillAttempts(b)
}

func (r *Registration) setAttemptState(att *RegAttempt, state RegAttemptState) {
att.bucket.count[att.State]--
att.bucket.count[state]++
Expand Down
48 changes: 48 additions & 0 deletions p2p/discover/topicindex/registration_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -317,6 +317,54 @@ func TestRegistrationExpiry(t *testing.T) {
}
}

// TestRegistrationRemoveNode checks that RemoveNode drops a parked attempt but
// leaves an in-flight one for its pending response.
func TestRegistrationRemoveNode(t *testing.T) {
cfg := testConfig(t)
r := NewRegistration(topic1, cfg)

// Waiting attempt: removed.
waiting := nodesAtDistance(enode.ID(r.Topic()), 30, 1)
r.AddNodes(nil, waiting)
if att := r.Update(); att == nil || att.State != Waiting {
t.Fatal("no waiting attempt scheduled")
}
r.RemoveNode(waiting[0].ID())
if r.buckets[r.bucketIndex(waiting[0].ID())].att[waiting[0].ID()] != nil {
t.Fatal("waiting attempt not removed")
}

// Registered attempt (ad not yet expired): removed before expiry.
registered := nodesAtDistance(enode.ID(r.Topic()), 40, 1)
r.AddNodes(nil, registered)
att := r.Update()
if att == nil {
t.Fatal("no attempt scheduled")
}
r.StartRequest(att)
r.HandleRegistered(att, cfg.AdLifetime)
r.RemoveNode(registered[0].ID())
if r.buckets[r.bucketIndex(registered[0].ID())].att[registered[0].ID()] != nil {
t.Fatal("registered attempt not removed before ad expiry")
}

// Unknown node: no-op.
r.RemoveNode(enode.ID{42})

// In-flight attempt: left in place, the pending response handles it.
inflight := nodesAtDistance(enode.ID(r.Topic()), 50, 1)
r.AddNodes(nil, inflight)
att = r.Update()
if att == nil {
t.Fatal("no attempt scheduled")
}
r.StartRequest(att)
r.RemoveNode(inflight[0].ID())
if r.buckets[r.bucketIndex(inflight[0].ID())].att[inflight[0].ID()] == nil {
t.Fatal("in-flight attempt removed; must be left for response handling")
}
}

// TestRegistrationHandleTicketResponseDropAboveBudget verifies that a
// registrar quoting a wait time above RegAttemptTimeout causes the attempt
// to be dropped, freeing the bucket slot for another registrar. Without
Expand Down
31 changes: 28 additions & 3 deletions p2p/discover/topicindex/search.go
Original file line number Diff line number Diff line change
Expand Up @@ -53,7 +53,7 @@ type Search struct {
type searchBucket struct {
dist int
new map[enode.ID]*enode.Node
asked map[enode.ID]struct{}
asked map[enode.ID]*enode.Node
numRequests int

ips netutil.DistinctNetSet
Expand All @@ -74,7 +74,7 @@ func NewSearch(topic TopicID, cfg Config) *Search {
s.buckets[i] = searchBucket{
dist: dist,
new: make(map[enode.ID]*enode.Node, cfg.SearchBucketSize),
asked: make(map[enode.ID]struct{}, cfg.SearchBucketSize),
asked: make(map[enode.ID]*enode.Node, cfg.SearchBucketSize),
ips: netutil.DistinctNetSet{
Subnet: searchBucketSubnet,
Limit: searchBucketIPLimit,
Expand Down Expand Up @@ -155,6 +155,31 @@ func (s *Search) AddNodes(src *enode.Node, nodes []*enode.Node) {
}
}

// HandleErrorResponse drops a failed node from the table, freeing its slot and
// IP-limit entry.
func (s *Search) HandleErrorResponse(from *enode.Node, err error) {
s.log.Debug("Topic query failed", "id", from.ID(), "err", err)
s.removeNode(from.ID())
}

// removeNode drops a node from the search table. The node is removed from both
// the unasked ('new') and asked sets of its bucket, and its IP-limit entry is
// released regardless of which set it was in.
func (s *Search) removeNode(id enode.ID) {
b := s.bucket(id)
n, ok := b.new[id]
if !ok {
n, ok = b.asked[id]
}
if ok {
if ip := n.IP(); ip != nil && !netutil.IsLAN(ip) {
b.ips.Remove(ip)
}
}
delete(b.new, id)
delete(b.asked, id)
}

// QueryTarget returns a random node to which a topic query should be sent.
// Random nodes are collected from buckets progressively: only buckets with unasked nodes
// that have received at least one response, plus the next unqueried bucket
Expand Down Expand Up @@ -241,6 +266,6 @@ func (b *searchBucket) count() int {
}

func (b *searchBucket) setAsked(n *enode.Node) {
b.asked[n.ID()] = struct{}{}
b.asked[n.ID()] = n
delete(b.new, n.ID())
}
83 changes: 83 additions & 0 deletions p2p/discover/topicindex/search_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@
package topicindex

import (
"errors"
"testing"

"github.com/ethereum/go-ethereum/p2p/enode"
Expand Down Expand Up @@ -169,6 +170,88 @@ func TestSearchAddNodesOnePerBucketRule(t *testing.T) {
}
}

// TestSearchHandleErrorResponse checks that a failed topic query drops the
// queried node from the table and frees its bucket slot and IP-limit entry.
func TestSearchHandleErrorResponse(t *testing.T) {
config := testConfig(t)
s := NewSearch(topic1, config)

// Two nodes in bucket 0 (logdist 256), one in bucket 5 (logdist 251).
far := nodesAtDistanceFrom(enode.ID(topic1), 256, 2, 1)
mid := nodesAtDistanceFrom(enode.ID(topic1), 251, 1, 10)
s.AddNodes(nil, far)
s.AddNodes(nil, mid)

// The query to far[0] fails: the node must leave the table entirely.
s.HandleErrorResponse(far[0], errors.New("timeout"))
if s.buckets[0].contains(far[0].ID()) {
t.Fatal("failed node still present in its bucket")
}
if got := s.buckets[0].count(); got != 1 {
t.Fatalf("bucket[0] count is %d after eviction, want 1", got)
}

// The IP-limit slot is freed: a replacement in the same /24 is accepted.
replacement := nodeAtDistance(enode.ID(topic1), 256, far[0].IP())
s.AddNodes(nil, []*enode.Node{replacement})
if !s.buckets[0].contains(replacement.ID()) {
t.Fatal("replacement with the failed node's IP was not admitted")
}

// The failure did not count as a response: bucket 0 still has candidates
// and no response yet, so it keeps gating the walk.
if n := s.QueryTarget(); n == nil || !s.buckets[0].contains(n.ID()) {
t.Fatalf("QueryTarget should keep picking from unwarmed bucket[0], got %v", n)
}

// Failing all of bucket 0 empties it; the empty bucket no longer blocks
// the walk, so QueryTarget advances to bucket 5.
s.HandleErrorResponse(far[1], errors.New("timeout"))
s.HandleErrorResponse(replacement, errors.New("timeout"))
target := s.QueryTarget()
if target == nil {
t.Fatal("QueryTarget returned nil after bucket[0] failed out; want bucket[5] node")
}
if !s.buckets[5].contains(target.ID()) {
t.Fatalf("QueryTarget returned %v, want the bucket[5] node", target.ID())
}

// When every node has failed, the search is done and rolls over.
s.HandleErrorResponse(mid[0], errors.New("timeout"))
if !s.IsDone() {
t.Fatal("IsDone should report true once every node has failed out")
}
}

// TestSearchRemoveAskedNodeFreesIP verifies that removing a node that has moved
// to the 'asked' set (it responded before being removed) still releases its
// IP-limit entry, so a same-/24 replacement is admitted.
func TestSearchRemoveAskedNodeFreesIP(t *testing.T) {
config := testConfig(t)
s := NewSearch(topic1, config)

nodes := nodesAtDistanceFrom(enode.ID(topic1), 256, 1, 1)
n := nodes[0]
s.AddNodes(nil, nodes)

// Move n into the 'asked' set by recording a query response from it.
s.AddQueryResults(n, nil)
if _, inAsked := s.buckets[0].asked[n.ID()]; !inAsked {
t.Fatal("node was not moved to the asked set")
}

// Removing it now must free the IP-limit slot even though it is in 'asked'.
s.HandleErrorResponse(n, errors.New("timeout"))
if s.buckets[0].contains(n.ID()) {
t.Fatal("removed asked-node still present")
}
replacement := nodeAtDistance(enode.ID(topic1), 256, n.IP())
s.AddNodes(nil, []*enode.Node{replacement})
if !s.buckets[0].contains(replacement.ID()) {
t.Fatal("replacement with the asked node's IP was not admitted (IP slot leaked)")
}
}

// This checks (de)queueing of topic search results: results come out of the
// buffer in the order they were received.
func TestSearchResultsTracking(t *testing.T) {
Expand Down
15 changes: 15 additions & 0 deletions p2p/discover/topicindex/topictable.go
Original file line number Diff line number Diff line change
Expand Up @@ -194,6 +194,21 @@ func (tab *TopicTable) add(n *enode.Node, topic TopicID) *topicTableEntry {
return reg
}

// RemoveNode removes all registrations advertised by the given node, across all
// topics. It is used to evict ads pointing at a node that has become
// unresponsive.
func (tab *TopicTable) RemoveNode(id enode.ID) {
for e := tab.all.Front(); e != nil; {
next := e.Next()
reg := e.Value.(*topicTableEntry)
if reg.node.ID() == id {
tab.remove(reg)
tab.wt.removeReg(reg)
}
e = next
}
}

func (tab *TopicTable) remove(reg *topicTableEntry) {
tab.all.Remove(reg.allElem)
topicList := tab.reg[reg.topic]
Expand Down
Loading