Skip to content

Commit 157d203

Browse files
committed
fix(addressbook): refresh last-seen on hive gossip and serialize updates
1 parent 5ec0701 commit 157d203

5 files changed

Lines changed: 315 additions & 8 deletions

File tree

pkg/addressbook/addressbook.go

Lines changed: 34 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,7 @@ import (
99
"errors"
1010
"fmt"
1111
"strings"
12+
"sync"
1213
"time"
1314

1415
"github.com/ethersphere/bee/v2/pkg/bzz"
@@ -18,6 +19,12 @@ import (
1819

1920
const keyPrefix = "addressbook_entry_"
2021

22+
// lastSeenUpdateInterval is the resolution at which UpdateLastSeen persists.
23+
// Hive reports the same peer on every gossip round, while pruning operates on
24+
// a scale of weeks, so a write is skipped when the stored value is already
25+
// this fresh.
26+
const lastSeenUpdateInterval = 24 * time.Hour
27+
2128
var _ Interface = (*store)(nil)
2229

2330
var ErrNotFound = errors.New("addressbook: not found")
@@ -54,6 +61,14 @@ type GetPutter interface {
5461
Putter
5562
}
5663

64+
// GetPutUpdater is the addressbook surface needed by hive: it stores peers it
65+
// learns about and refreshes the last-seen time of the ones it already knows.
66+
type GetPutUpdater interface {
67+
GetPutter
68+
// UpdateLastSeen marks the overlay as seen at the current time.
69+
UpdateLastSeen(overlay swarm.Address) error
70+
}
71+
5772
type Getter interface {
5873
// Get returns the saved bzz.Address for the requested overlay together
5974
// with its verification flag.
@@ -73,6 +88,10 @@ type Remover interface {
7388
type store struct {
7489
store storage.StateStorer
7590
now func() time.Time
91+
92+
// mu serializes the read-modify-write in UpdateLastSeen against Put, so a
93+
// concurrent Put is not rolled back by a stale copy of the entry.
94+
mu sync.Mutex
7695
}
7796

7897
// New creates new addressbook for state storer.
@@ -97,6 +116,9 @@ func (s *store) Get(overlay swarm.Address) (*bzz.Address, bool, error) {
97116
}
98117

99118
func (s *store) Put(overlay swarm.Address, addr bzz.Address, verified bool) (err error) {
119+
s.mu.Lock()
120+
defer s.mu.Unlock()
121+
100122
key := keyPrefix + overlay.String()
101123
return s.store.Put(key, &verifiedAddress{
102124
Address: &addr,
@@ -106,8 +128,12 @@ func (s *store) Put(overlay swarm.Address, addr bzz.Address, verified bool) (err
106128
}
107129

108130
// UpdateLastSeen marks the overlay as seen at the current time. It is a no-op
109-
// if the overlay is not present in the addressbook.
131+
// if the overlay is not present in the addressbook, or if the recorded time is
132+
// younger than lastSeenUpdateInterval.
110133
func (s *store) UpdateLastSeen(overlay swarm.Address) error {
134+
s.mu.Lock()
135+
defer s.mu.Unlock()
136+
111137
key := keyPrefix + overlay.String()
112138
v := &verifiedAddress{}
113139
if err := s.store.Get(key, v); err != nil {
@@ -116,7 +142,13 @@ func (s *store) UpdateLastSeen(overlay swarm.Address) error {
116142
}
117143
return err
118144
}
119-
v.LastSeen = s.now().Unix()
145+
146+
now := s.now().Unix()
147+
if now-v.LastSeen < int64(lastSeenUpdateInterval.Seconds()) {
148+
return nil
149+
}
150+
151+
v.LastSeen = now
120152
return s.store.Put(key, v)
121153
}
122154

pkg/addressbook/addressbook_test.go

Lines changed: 118 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -14,6 +14,7 @@ import (
1414
"github.com/ethersphere/bee/v2/pkg/bzz"
1515
"github.com/ethersphere/bee/v2/pkg/crypto"
1616
"github.com/ethersphere/bee/v2/pkg/statestore/mock"
17+
"github.com/ethersphere/bee/v2/pkg/storage"
1718
"github.com/ethersphere/bee/v2/pkg/swarm"
1819
ma "github.com/multiformats/go-multiaddr"
1920
)
@@ -136,21 +137,134 @@ func TestUpdateLastSeen(t *testing.T) {
136137
t.Fatal(err)
137138
}
138139

139-
// advance the clock and bump last-seen; the entry must survive a prune at
140-
// the original time.
141-
now = time.Unix(5000, 0)
140+
// advance the clock past the update interval and bump last-seen; the entry
141+
// must survive a prune at the original time.
142+
seenAt := now.Add(2 * 24 * time.Hour)
143+
now = seenAt
142144
if err := store.UpdateLastSeen(overlay); err != nil {
143145
t.Fatal(err)
144146
}
145147

146-
if err := store.Prune(time.Unix(4000, 0)); err != nil {
148+
if err := store.Prune(seenAt.Add(-time.Hour)); err != nil {
147149
t.Fatal(err)
148150
}
149151
if _, _, err := store.Get(overlay); err != nil {
150152
t.Fatalf("entry pruned despite recent last-seen: %v", err)
151153
}
152154
}
153155

156+
// TestUpdateLastSeenThrottled asserts that a bump within lastSeenUpdateInterval
157+
// does not write. Hive calls UpdateLastSeen on every gossip sighting, so this
158+
// is what keeps the write rate bounded to roughly one per peer per day.
159+
func TestUpdateLastSeenThrottled(t *testing.T) {
160+
t.Parallel()
161+
162+
base := time.Unix(1_000_000, 0)
163+
now := base
164+
mockStore := mock.NewStateStore()
165+
store := addressbook.NewWithClock(mockStore, func() time.Time { return now })
166+
167+
overlay := swarm.NewAddress([]byte{0, 1, 2, 3})
168+
if err := store.Put(overlay, newTestAddr(t, overlay), true); err != nil {
169+
t.Fatal(err)
170+
}
171+
172+
// well inside the interval: must not touch the record.
173+
now = base.Add(time.Hour)
174+
if err := store.UpdateLastSeen(overlay); err != nil {
175+
t.Fatal(err)
176+
}
177+
if got := lastSeenOf(t, mockStore, overlay); got != base.Unix() {
178+
t.Fatalf("throttled update wrote: last_seen = %d, want %d", got, base.Unix())
179+
}
180+
181+
// past the interval: must write.
182+
now = base.Add(addressbook.LastSeenUpdateInterval + time.Second)
183+
if err := store.UpdateLastSeen(overlay); err != nil {
184+
t.Fatal(err)
185+
}
186+
if got := lastSeenOf(t, mockStore, overlay); got != now.Unix() {
187+
t.Fatalf("update past interval did not write: last_seen = %d, want %d", got, now.Unix())
188+
}
189+
}
190+
191+
// TestUpdateLastSeenKeepsConcurrentPut pins the read-modify-write in
192+
// UpdateLastSeen against a Put that lands between its read and its write.
193+
// Without serialization the Put's Verified flag is rolled back, which would
194+
// also desync the addressbook from hive's chequebook registry.
195+
func TestUpdateLastSeenKeepsConcurrentPut(t *testing.T) {
196+
t.Parallel()
197+
198+
base := time.Unix(1_000_000, 0)
199+
now := base
200+
hooked := &hookStore{StateStorer: mock.NewStateStore()}
201+
book := addressbook.NewWithClock(hooked, func() time.Time { return now })
202+
203+
overlay := swarm.NewAddress([]byte{0, 1, 2, 3})
204+
addr := newTestAddr(t, overlay)
205+
206+
// a known, not-yet-verified peer.
207+
if err := book.Put(overlay, addr, false); err != nil {
208+
t.Fatal(err)
209+
}
210+
// move past the throttle so UpdateLastSeen really writes.
211+
now = base.Add(addressbook.LastSeenUpdateInterval + time.Second)
212+
213+
// While UpdateLastSeen holds the entry it has just read, hive verifies the
214+
// same peer and stores it with Verified=true.
215+
started, finished := make(chan struct{}), make(chan struct{})
216+
hooked.onGet = func() {
217+
go func() {
218+
defer close(finished)
219+
close(started)
220+
if err := book.Put(overlay, addr, true); err != nil {
221+
t.Error(err)
222+
}
223+
}()
224+
<-started
225+
// Give the writer time to land. Serialized, it blocks on the
226+
// addressbook lock until UpdateLastSeen returns; unsynchronized, its
227+
// write completes here and is then overwritten below.
228+
time.Sleep(100 * time.Millisecond)
229+
}
230+
231+
if err := book.UpdateLastSeen(overlay); err != nil {
232+
t.Fatal(err)
233+
}
234+
<-finished
235+
236+
if _, verified, err := book.Get(overlay); err != nil || !verified {
237+
t.Fatalf("concurrent Put(verified=true) was rolled back: verified=%v err=%v", verified, err)
238+
}
239+
}
240+
241+
// hookStore fires onGet once, immediately after a Get returns, to interleave a
242+
// concurrent writer inside UpdateLastSeen's read-modify-write.
243+
type hookStore struct {
244+
storage.StateStorer
245+
onGet func()
246+
}
247+
248+
func (h *hookStore) Get(key string, i any) error {
249+
err := h.StateStorer.Get(key, i)
250+
if h.onGet != nil {
251+
f := h.onGet
252+
h.onGet = nil
253+
f()
254+
}
255+
return err
256+
}
257+
258+
func lastSeenOf(t *testing.T, store storage.StateStorer, overlay swarm.Address) int64 {
259+
t.Helper()
260+
261+
v := &addressbook.VerifiedAddress{}
262+
if err := store.Get("addressbook_entry_"+overlay.String(), v); err != nil {
263+
t.Fatalf("get entry: %v", err)
264+
}
265+
return v.LastSeen
266+
}
267+
154268
func TestPrune(t *testing.T) {
155269
t.Parallel()
156270

pkg/addressbook/export_test.go

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -12,6 +12,9 @@ import (
1212

1313
type VerifiedAddress = verifiedAddress
1414

15+
// LastSeenUpdateInterval exposes the UpdateLastSeen write throttle.
16+
const LastSeenUpdateInterval = lastSeenUpdateInterval
17+
1518
// NewWithClock creates an addressbook with an overridable clock, for testing.
1619
func NewWithClock(storer storage.StateStorer, now func() time.Time) Interface {
1720
return &store{

pkg/hive/hive.go

Lines changed: 13 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -70,7 +70,7 @@ type Options struct {
7070

7171
type Service struct {
7272
streamer p2p.Streamer
73-
addressBook addressbook.GetPutter
73+
addressBook addressbook.GetPutUpdater
7474
addPeersHandler func(...swarm.Address)
7575
networkID uint64
7676
logger log.Logger
@@ -92,7 +92,7 @@ type Service struct {
9292
chequebookStorer ChequebookStorer
9393
}
9494

95-
func New(streamer p2p.Streamer, addressbook addressbook.GetPutter, networkID uint64, overlay swarm.Address, logger log.Logger, o Options) *Service {
95+
func New(streamer p2p.Streamer, addressbook addressbook.GetPutUpdater, networkID uint64, overlay swarm.Address, logger log.Logger, o Options) *Service {
9696
svc := &Service{
9797
streamer: streamer,
9898
logger: logger.WithName(loggerName).Register(),
@@ -370,6 +370,17 @@ func (s *Service) checkAndAddPeers(ctx context.Context, peers pb.Peers) {
370370
}
371371

372372
if err := bzz.CheckTimestamp(bzzAddress.Timestamp, existing, bzz.TimestampSourceGossip, s.now()); err != nil {
373+
// A peer re-presenting a record we already hold has nothing new to
374+
// store, but it is still a sighting: peers mint their bzz.Address
375+
// once and gossip it unchanged for their whole uptime. Refresh
376+
// last-seen so peers we keep hearing about, but never dial, do not
377+
// look stale to the pruner. Both errors imply a known peer, since
378+
// an unknown one short-circuits inside CheckTimestamp.
379+
if errors.Is(err, bzz.ErrTimestampStale) || errors.Is(err, bzz.ErrTimestampTooSoon) {
380+
if err := s.addressBook.UpdateLastSeen(overlayAddr); err != nil {
381+
s.logger.Debug("hive gossip: update last seen", "overlay", overlayAddr.String(), "error", err)
382+
}
383+
}
373384
s.bumpTimestampMetric(err)
374385
s.logger.Debug("hive gossip: timestamp validation failed", "overlay", overlayAddr.String(), "error", err)
375386
continue

0 commit comments

Comments
 (0)