From 1716cdbf74af07dd16c2c1dbbebdaf91ee2c3a7f Mon Sep 17 00:00:00 2001 From: Calin Martinconi Date: Mon, 22 Jun 2026 23:35:23 +0300 Subject: [PATCH 1/6] feat(addressbook): prune stale entries by last-seen time --- pkg/addressbook/addressbook.go | 63 ++++++- pkg/addressbook/addressbook_test.go | 112 ++++++++++++ pkg/addressbook/export_test.go | 14 ++ pkg/node/node.go | 7 + pkg/statestore/storeadapter/export_test.go | 13 +- pkg/statestore/storeadapter/migration.go | 83 ++++++++- pkg/statestore/storeadapter/migration_test.go | 166 ++++++++++++++++++ pkg/topology/kademlia/kademlia.go | 7 + pkg/topology/kademlia/lastseen_test.go | 94 ++++++++++ 9 files changed, 545 insertions(+), 14 deletions(-) create mode 100644 pkg/topology/kademlia/lastseen_test.go diff --git a/pkg/addressbook/addressbook.go b/pkg/addressbook/addressbook.go index 32b14e0a771..7f31a4f4e2b 100644 --- a/pkg/addressbook/addressbook.go +++ b/pkg/addressbook/addressbook.go @@ -9,6 +9,7 @@ import ( "errors" "fmt" "strings" + "time" "github.com/ethersphere/bee/v2/pkg/bzz" "github.com/ethersphere/bee/v2/pkg/storage" @@ -22,22 +23,30 @@ var _ Interface = (*store)(nil) var ErrNotFound = errors.New("addressbook: not found") // verifiedAddress pairs a bzz.Address with a flag indicating whether the peer -// has been verified. +// has been verified, and the last time the overlay was seen. type verifiedAddress struct { Address *bzz.Address `json:"address"` Verified bool `json:"verified"` + // LastSeen is the Unix timestamp (seconds) of the last time the overlay + // was seen over hive or in kademlia. Used to prune stale entries. + LastSeen int64 `json:"last_seen"` } // Interface is the AddressBook interface. type Interface interface { GetPutter Remover + // UpdateLastSeen marks the overlay as seen at the current time. + UpdateLastSeen(overlay swarm.Address) error // Overlays returns a list of all overlay addresses saved in addressbook. Overlays() ([]swarm.Address, error) // IterateOverlays exposes overlays in a form of an iterator. IterateOverlays(func(swarm.Address) (bool, error)) error // Addresses returns a list of all bzz.Address-es saved in addressbook. Addresses() ([]bzz.Address, error) + // Prune removes all entries whose overlay has not been seen since the + // given time. + Prune(before time.Time) error } type GetPutter interface { @@ -63,12 +72,14 @@ type Remover interface { type store struct { store storage.StateStorer + now func() time.Time } // New creates new addressbook for state storer. func New(storer storage.StateStorer) Interface { return &store{ store: storer, + now: time.Now, } } @@ -90,13 +101,63 @@ func (s *store) Put(overlay swarm.Address, addr bzz.Address, verified bool) (err return s.store.Put(key, &verifiedAddress{ Address: &addr, Verified: verified, + LastSeen: s.now().Unix(), }) } +// UpdateLastSeen marks the overlay as seen at the current time. It is a no-op +// if the overlay is not present in the addressbook. +func (s *store) UpdateLastSeen(overlay swarm.Address) error { + key := keyPrefix + overlay.String() + v := &verifiedAddress{} + if err := s.store.Get(key, v); err != nil { + if errors.Is(err, storage.ErrNotFound) { + return nil + } + return err + } + v.LastSeen = s.now().Unix() + return s.store.Put(key, v) +} + func (s *store) Remove(overlay swarm.Address) error { return s.store.Delete(keyPrefix + overlay.String()) } +// Prune removes all entries whose overlay has not been seen since before. +// Entries without a recorded last-seen time (LastSeen == 0) are kept, leaving +// pruning to a later run once they have been observed. +func (s *store) Prune(before time.Time) error { + cutoff := before.Unix() + + var stale []swarm.Address + err := s.store.Iterate(keyPrefix, func(key, value []byte) (stop bool, err error) { + entry := &verifiedAddress{} + if err := json.Unmarshal(value, entry); err != nil { + return true, err + } + if entry.LastSeen != 0 && entry.LastSeen < cutoff { + addr, err := swarm.ParseHexAddress(strings.TrimPrefix(string(key), keyPrefix)) + if err != nil { + return true, err + } + stale = append(stale, addr) + } + return false, nil + }) + if err != nil { + return err + } + + for _, addr := range stale { + if err := s.Remove(addr); err != nil { + return err + } + } + + return nil +} + func (s *store) IterateOverlays(cb func(swarm.Address) (bool, error)) error { return s.store.Iterate(keyPrefix, func(key, _ []byte) (stop bool, err error) { k := string(key) diff --git a/pkg/addressbook/addressbook_test.go b/pkg/addressbook/addressbook_test.go index f515527d21e..ee355e9cde0 100644 --- a/pkg/addressbook/addressbook_test.go +++ b/pkg/addressbook/addressbook_test.go @@ -7,6 +7,7 @@ package addressbook_test import ( "errors" "testing" + "time" "github.com/ethereum/go-ethereum/common" "github.com/ethersphere/bee/v2/pkg/addressbook" @@ -17,6 +18,24 @@ import ( ma "github.com/multiformats/go-multiaddr" ) +func newTestAddr(t *testing.T, overlay swarm.Address) bzz.Address { + t.Helper() + + multiaddr, err := ma.NewMultiaddr("/ip4/1.1.1.1") + if err != nil { + t.Fatal(err) + } + pk, err := crypto.GenerateSecp256k1Key() + if err != nil { + t.Fatal(err) + } + bzzAddr, err := bzz.NewAddress(crypto.NewDefaultSigner(pk), []ma.Multiaddr{multiaddr}, overlay, 1, common.HexToHash("0x1").Bytes(), 1, common.Address{}) + if err != nil { + t.Fatal(err) + } + return *bzzAddr +} + type bookFunc func() (book addressbook.Interface) func TestInMem(t *testing.T) { @@ -96,3 +115,96 @@ func run(t *testing.T, f bookFunc) { t.Fatalf("expected addresses len %v, got %v", 1, len(addresses)) } } + +func TestUpdateLastSeen(t *testing.T) { + t.Parallel() + + now := time.Unix(1000, 0) + store := addressbook.NewWithClock(mock.NewStateStore(), func() time.Time { return now }) + + overlay := swarm.NewAddress([]byte{0, 1, 2, 3}) + + // UpdateLastSeen on a missing overlay is a no-op and must not create an entry. + if err := store.UpdateLastSeen(overlay); err != nil { + t.Fatal(err) + } + if _, _, err := store.Get(overlay); !errors.Is(err, addressbook.ErrNotFound) { + t.Fatalf("expected ErrNotFound, got %v", err) + } + + if err := store.Put(overlay, newTestAddr(t, overlay), true); err != nil { + t.Fatal(err) + } + + // advance the clock and bump last-seen; the entry must survive a prune at + // the original time. + now = time.Unix(5000, 0) + if err := store.UpdateLastSeen(overlay); err != nil { + t.Fatal(err) + } + + if err := store.Prune(time.Unix(4000, 0)); err != nil { + t.Fatal(err) + } + if _, _, err := store.Get(overlay); err != nil { + t.Fatalf("entry pruned despite recent last-seen: %v", err) + } +} + +func TestPrune(t *testing.T) { + t.Parallel() + + now := time.Unix(0, 0) + store := addressbook.NewWithClock(mock.NewStateStore(), func() time.Time { return now }) + + stale := swarm.NewAddress([]byte{0, 1, 2, 3}) + fresh := swarm.NewAddress([]byte{0, 1, 2, 4}) + + now = time.Unix(1000, 0) + if err := store.Put(stale, newTestAddr(t, stale), true); err != nil { + t.Fatal(err) + } + + now = time.Unix(9000, 0) + if err := store.Put(fresh, newTestAddr(t, fresh), true); err != nil { + t.Fatal(err) + } + + if err := store.Prune(time.Unix(5000, 0)); err != nil { + t.Fatal(err) + } + + if _, _, err := store.Get(stale); !errors.Is(err, addressbook.ErrNotFound) { + t.Fatalf("stale entry should have been pruned, got err=%v", err) + } + if _, _, err := store.Get(fresh); err != nil { + t.Fatalf("fresh entry should survive: %v", err) + } +} + +func TestPruneKeepsEntriesWithoutLastSeen(t *testing.T) { + t.Parallel() + + mockStore := mock.NewStateStore() + overlay := swarm.NewAddress([]byte{0, 1, 2, 3}) + + // Seed an entry without a last_seen field, mirroring records that predate + // pruning before the stamping migration runs. + if err := mockStore.Put("addressbook_entry_"+overlay.String(), &addressbook.VerifiedAddress{ + Address: addrPtr(newTestAddr(t, overlay)), + Verified: true, + }); err != nil { + t.Fatal(err) + } + + store := addressbook.New(mockStore) + if err := store.Prune(time.Unix(5000, 0)); err != nil { + t.Fatal(err) + } + + if _, _, err := store.Get(overlay); err != nil { + t.Fatalf("entry without last_seen must not be pruned: %v", err) + } +} + +func addrPtr(a bzz.Address) *bzz.Address { return &a } diff --git a/pkg/addressbook/export_test.go b/pkg/addressbook/export_test.go index 9db017e1f27..325621e1266 100644 --- a/pkg/addressbook/export_test.go +++ b/pkg/addressbook/export_test.go @@ -4,4 +4,18 @@ package addressbook +import ( + "time" + + "github.com/ethersphere/bee/v2/pkg/storage" +) + type VerifiedAddress = verifiedAddress + +// NewWithClock creates an addressbook with an overridable clock, for testing. +func NewWithClock(storer storage.StateStorer, now func() time.Time) Interface { + return &store{ + store: storer, + now: now, + } +} diff --git a/pkg/node/node.go b/pkg/node/node.go index c6a1dc54981..2a8069763ff 100644 --- a/pkg/node/node.go +++ b/pkg/node/node.go @@ -216,6 +216,7 @@ const ( reserveMinEvictCount = 1_000 cacheMinEvictCount = 10_000 maxAllowedDoubling = 1 + addressbookPruneAfter = 30 * 24 * time.Hour // remove addressbook entries not seen for this long ) func NewBee( @@ -382,6 +383,12 @@ func NewBee( addressbook := addressbook.New(stateStore) + // Prune addressbook entries whose overlays have not been seen recently, so + // the address book does not accumulate stale peers indefinitely. + if err := addressbook.Prune(time.Now().Add(-addressbookPruneAfter)); err != nil { + logger.Warning("addressbook prune failed", "error", err) + } + logger.Info("using overlay address", "address", swarmAddress) // this will set overlay if it was not set before diff --git a/pkg/statestore/storeadapter/export_test.go b/pkg/statestore/storeadapter/export_test.go index d857e0a66d5..bffe33c70ff 100644 --- a/pkg/statestore/storeadapter/export_test.go +++ b/pkg/statestore/storeadapter/export_test.go @@ -4,8 +4,13 @@ package storeadapter -var RewriteAddressbookEnvelope = rewriteAddressbookEnvelope +var ( + RewriteAddressbookEnvelope = rewriteAddressbookEnvelope + StampAddressbookLastSeen = stampAddressbookLastSeen +) -type LegacyEntry = legacyEntry -type MigratedEntry = migratedEntry -type MigratedAddress = migratedAddress +type ( + LegacyEntry = legacyEntry + MigratedEntry = migratedEntry + MigratedAddress = migratedAddress +) diff --git a/pkg/statestore/storeadapter/migration.go b/pkg/statestore/storeadapter/migration.go index 6542a5f3c78..9f3a83ff804 100644 --- a/pkg/statestore/storeadapter/migration.go +++ b/pkg/statestore/storeadapter/migration.go @@ -7,6 +7,7 @@ package storeadapter import ( "encoding/json" "fmt" + "time" "github.com/ethersphere/bee/v2/pkg/storage" "github.com/ethersphere/bee/v2/pkg/storage/migration" @@ -23,15 +24,16 @@ func allSteps(st storage.Store) migration.Steps { // and never execute newly added migrations. noop := func() error { return nil } return map[uint64]migration.StepFn{ - 1: noop, - 2: noop, - 3: noop, - 4: noop, - 5: noop, - 6: noop, - 7: noop, - 8: noop, - 9: rewriteAddressbookEnvelope(st), + 1: noop, + 2: noop, + 3: noop, + 4: noop, + 5: noop, + 6: noop, + 7: noop, + 8: noop, + 9: rewriteAddressbookEnvelope(st), + 10: stampAddressbookLastSeen(st), } } @@ -55,6 +57,7 @@ type migratedAddress struct { type migratedEntry struct { Address migratedAddress `json:"address"` Verified bool `json:"verified"` + LastSeen int64 `json:"last_seen"` } // rewriteAddressbookEnvelope wraps each "addressbook_entry_*" legacy @@ -118,3 +121,65 @@ func rewriteAddressbookEnvelope(s storage.Store) migration.StepFn { return nil } } + +// stampAddressbookLastSeen sets "last_seen" to the current time on every +// "addressbook_entry_*" record that lacks it, so that addresses carried over +// from before pruning was introduced are not immediately pruned. Entries that +// already carry a non-zero last_seen are left untouched. The whole record is +// preserved by merging into the decoded JSON object rather than re-encoding a +// typed struct. +func stampAddressbookLastSeen(s storage.Store) migration.StepFn { + return func() error { + store := &StateStorerAdapter{s} + + type item struct { + key string + val []byte + } + + var batch []item + if err := store.Iterate("addressbook_entry_", func(key, val []byte) (stop bool, err error) { + batch = append(batch, item{ + key: string(key), + val: append([]byte(nil), val...), + }) + return false, nil + }); err != nil { + return fmt.Errorf("iterate addressbook entries: %w", err) + } + + now := time.Now().Unix() + + for _, e := range batch { + var fields map[string]json.RawMessage + if err := json.Unmarshal(e.val, &fields); err != nil { + _ = store.Delete(e.key) + continue + } + + if raw, ok := fields["last_seen"]; ok { + var ls int64 + if json.Unmarshal(raw, &ls) == nil && ls != 0 { + continue + } + } + + stamp, err := json.Marshal(now) + if err != nil { + return fmt.Errorf("marshal last_seen: %w", err) + } + fields["last_seen"] = stamp + + out, err := json.Marshal(fields) + if err != nil { + return fmt.Errorf("marshal addressbook entry %q: %w", e.key, err) + } + + if err := store.Put(e.key, json.RawMessage(out)); err != nil { + return fmt.Errorf("stamp addressbook entry %q: %w", e.key, err) + } + } + + return nil + } +} diff --git a/pkg/statestore/storeadapter/migration_test.go b/pkg/statestore/storeadapter/migration_test.go index 3a74eb2a90b..51182dc58b5 100644 --- a/pkg/statestore/storeadapter/migration_test.go +++ b/pkg/statestore/storeadapter/migration_test.go @@ -7,6 +7,7 @@ package storeadapter_test import ( "errors" "testing" + "time" "github.com/ethereum/go-ethereum/common" "github.com/ethersphere/bee/v2/pkg/addressbook" @@ -191,6 +192,171 @@ func TestRewriteAddressbookEnvelope_AddressbookConsumes(t *testing.T) { } } +func TestStampAddressbookLastSeen(t *testing.T) { + t.Parallel() + + raw := newTestStore(t) + store, err := storeadapter.NewStateStorerAdapter(raw) + if err != nil { + t.Fatalf("NewStateStorerAdapter: %v", err) + } + + const prefix = "addressbook_entry_" + + // entry carried over from before pruning: no last_seen. + stampKey := prefix + "aabb" + if err := store.Put(stampKey, &storeadapter.MigratedEntry{ + Address: storeadapter.MigratedAddress{ + Overlay: "aabb", + Underlays: []string{"/ip4/1.1.1.1"}, + Signature: "sig==", + Nonce: "deadbeef", + Timestamp: 12345, + }, + Verified: true, + }); err != nil { + t.Fatalf("seed stamp: %v", err) + } + + // entry that already has a last_seen must not be touched. + keepKey := prefix + "ccdd" + const existingLastSeen = int64(42) + if err := store.Put(keepKey, &storeadapter.MigratedEntry{ + Address: storeadapter.MigratedAddress{Overlay: "ccdd"}, + LastSeen: existingLastSeen, + }); err != nil { + t.Fatalf("seed keep: %v", err) + } + + if err := storeadapter.StampAddressbookLastSeen(raw)(); err != nil { + t.Fatalf("migration: %v", err) + } + + var stamped storeadapter.MigratedEntry + if err := store.Get(stampKey, &stamped); err != nil { + t.Fatalf("get stamped: %v", err) + } + if stamped.LastSeen == 0 { + t.Fatal("last_seen was not stamped") + } + // other fields must survive the merge. + if !stamped.Verified || stamped.Address.Overlay != "aabb" || stamped.Address.Timestamp != 12345 { + t.Fatalf("entry mutated unexpectedly: %+v", stamped) + } + if len(stamped.Address.Underlays) != 1 || stamped.Address.Underlays[0] != "/ip4/1.1.1.1" { + t.Fatalf("underlays lost: %v", stamped.Address.Underlays) + } + + var kept storeadapter.MigratedEntry + if err := store.Get(keepKey, &kept); err != nil { + t.Fatalf("get kept: %v", err) + } + if kept.LastSeen != existingLastSeen { + t.Fatalf("existing last_seen overwritten: got %d want %d", kept.LastSeen, existingLastSeen) + } +} + +// TestAddressbookPruneRealStore drives addressbook.Prune over the production +// storage path (leveldb behind StateStorerAdapter), whose iterator key +// semantics differ from the in-memory mock used in the addressbook package's +// own tests. It confirms that overlays are correctly reconstructed from the +// iterated keys so the right entries are pruned and the survivors remain +// readable through the addressbook. +func TestAddressbookPruneRealStore(t *testing.T) { + t.Parallel() + + const ( + prefix = "addressbook_entry_" + validSignature = "c2lnbmF0dXJl" // base64("signature") + validUnderlay = "/ip4/127.0.0.1/tcp/1634" + nonceHex = "deadbeefdeadbeefdeadbeefdeadbeefdeadbeefdeadbeefdeadbeefdeadbeef" + ) + + store, err := storeadapter.NewStateStorerAdapter(newTestStore(t)) + if err != nil { + t.Fatalf("NewStateStorerAdapter: %v", err) + } + + seed := func(overlayHex string, lastSeen int64) { + t.Helper() + if err := store.Put(prefix+overlayHex, &storeadapter.MigratedEntry{ + Address: storeadapter.MigratedAddress{ + Overlay: overlayHex, + Underlays: []string{validUnderlay}, + Signature: validSignature, + Nonce: nonceHex, + }, + Verified: true, + LastSeen: lastSeen, + }); err != nil { + t.Fatalf("seed %s: %v", overlayHex, err) + } + } + + stale := swarm.MustParseHexAddress("aabb") + fresh := swarm.MustParseHexAddress("ccdd") + seed("aabb", 1000) + seed("ccdd", 9000) + + book := addressbook.New(store) + if err := book.Prune(time.Unix(5000, 0)); err != nil { + t.Fatalf("prune: %v", err) + } + + if _, _, err := book.Get(stale); !errors.Is(err, addressbook.ErrNotFound) { + t.Fatalf("stale entry should have been pruned, got err=%v", err) + } + + got, _, err := book.Get(fresh) + if err != nil { + t.Fatalf("fresh entry should survive prune: %v", err) + } + if !got.Overlay.Equal(fresh) { + t.Fatalf("survivor overlay mismatch: got %s want %s", got.Overlay, fresh) + } +} + +func TestStampAddressbookLastSeen_Idempotent(t *testing.T) { + t.Parallel() + + raw := newTestStore(t) + store, err := storeadapter.NewStateStorerAdapter(raw) + if err != nil { + t.Fatalf("NewStateStorerAdapter: %v", err) + } + + key := "addressbook_entry_aabb" + if err := store.Put(key, &storeadapter.MigratedEntry{ + Address: storeadapter.MigratedAddress{Overlay: "aabb"}, + Verified: true, + }); err != nil { + t.Fatalf("seed: %v", err) + } + + if err := storeadapter.StampAddressbookLastSeen(raw)(); err != nil { + t.Fatalf("first run: %v", err) + } + + var first storeadapter.MigratedEntry + if err := store.Get(key, &first); err != nil { + t.Fatalf("get after first run: %v", err) + } + + for i := 0; i < 2; i++ { + if err := storeadapter.StampAddressbookLastSeen(raw)(); err != nil { + t.Fatalf("rerun %d: %v", i, err) + } + } + + var got storeadapter.MigratedEntry + if err := store.Get(key, &got); err != nil { + t.Fatalf("get after repeated runs: %v", err) + } + if got.LastSeen != first.LastSeen { + t.Fatalf("last_seen changed across reruns: got %d want %d", got.LastSeen, first.LastSeen) + } +} + func TestRewriteAddressbookEnvelope_Idempotent(t *testing.T) { t.Parallel() diff --git a/pkg/topology/kademlia/kademlia.go b/pkg/topology/kademlia/kademlia.go index 57c4c41a950..6f0d3fc6a0b 100644 --- a/pkg/topology/kademlia/kademlia.go +++ b/pkg/topology/kademlia/kademlia.go @@ -455,6 +455,10 @@ func (k *Kad) connectionAttemptsHandler(ctx context.Context, wg *sync.WaitGroup, k.connectedPeers.Add(peer.addr) + if err := k.addressBook.UpdateLastSeen(peer.addr); err != nil { + k.logger.Debug("could not update last seen for peer", "peer_address", peer.addr, "error", err) + } + k.metrics.TotalOutboundConnections.Inc() k.collector.Record(peer.addr, im.PeerLogIn(time.Now(), im.PeerConnectionDirectionOutbound)) @@ -1214,6 +1218,9 @@ func (k *Kad) onConnected(ctx context.Context, addr swarm.Address) error { k.knownPeers.Add(addr) k.connectedPeers.Add(addr) k.waitNext.Remove(addr) + if err := k.addressBook.UpdateLastSeen(addr); err != nil { + k.logger.Debug("could not update last seen for peer", "peer_address", addr, "error", err) + } k.recalcDepth() k.notifyManageLoop() k.notifyPeerSig() diff --git a/pkg/topology/kademlia/lastseen_test.go b/pkg/topology/kademlia/lastseen_test.go new file mode 100644 index 00000000000..3ceb6b2daf2 --- /dev/null +++ b/pkg/topology/kademlia/lastseen_test.go @@ -0,0 +1,94 @@ +// Copyright 2026 The Swarm Authors. All rights reserved. +// Use of this source code is governed by a BSD-style +// license that can be found in the LICENSE file. + +package kademlia_test + +import ( + "context" + "sync" + "testing" + "time" + + "github.com/ethersphere/bee/v2/pkg/addressbook" + beeCrypto "github.com/ethersphere/bee/v2/pkg/crypto" + "github.com/ethersphere/bee/v2/pkg/discovery/mock" + "github.com/ethersphere/bee/v2/pkg/log" + "github.com/ethersphere/bee/v2/pkg/stabilization" + mockstate "github.com/ethersphere/bee/v2/pkg/statestore/mock" + "github.com/ethersphere/bee/v2/pkg/swarm" + "github.com/ethersphere/bee/v2/pkg/topology/kademlia" + "github.com/ethersphere/bee/v2/pkg/util/testutil" +) + +type spyBook struct { + addressbook.Interface + mu sync.Mutex + seen map[string]int +} + +func (s *spyBook) UpdateLastSeen(o swarm.Address) error { + s.mu.Lock() + s.seen[o.String()]++ + s.mu.Unlock() + return s.Interface.UpdateLastSeen(o) +} + +func (s *spyBook) count(o swarm.Address) int { + s.mu.Lock() + defer s.mu.Unlock() + return s.seen[o.String()] +} + +func TestKademliaBumpsLastSeenOnConnect(t *testing.T) { + t.Parallel() + + detector, err := stabilization.NewDetector(stabilization.Config{ + PeriodDuration: 1 * time.Second, + NumPeriodsForStabilization: 2, + StabilizationFactor: 1, + WarmupTime: 0, + }) + if err != nil { + t.Fatal(err) + } + + var conns, failed int32 + spy := &spyBook{Interface: addressbook.New(mockstate.NewStateStore()), seen: map[string]int{}} + base := swarm.RandAddress(t) + disc := mock.NewDiscovery() + + pk, _ := beeCrypto.GenerateSecp256k1Key() + signer := beeCrypto.NewDefaultSigner(pk) + p2ps := p2pMock(t, spy, signer, &conns, &failed) + + bit := -1 + kad, err := kademlia.New(base, spy, disc, p2ps, detector, log.Noop, kademlia.Options{ + BitSuffixLength: &bit, + ExcludeFunc: defaultExcludeFunc, + }) + if err != nil { + t.Fatal(err) + } + p2ps.SetPickyNotifier(kad) + if err := kad.Start(context.Background()); err != nil { + t.Fatal(err) + } + testutil.CleanupCloser(t, kad) + kad.SetStorageRadius(0) + + // Inbound path: Connected -> onConnected -> UpdateLastSeen. + inbound := swarm.RandAddress(t) + connectOne(t, signer, kad, spy, inbound, nil) + if got := spy.count(inbound); got == 0 { + t.Fatalf("inbound connect did not bump last-seen for %s", inbound) + } + + // Outbound path: manage loop dials -> connect closure -> UpdateLastSeen. + outbound := swarm.RandAddressAt(t, base, 0) + addOne(t, signer, kad, spy, outbound) + waitConn(t, &conns) + if got := spy.count(outbound); got == 0 { + t.Fatalf("outbound connect did not bump last-seen for %s", outbound) + } +} From 5ec07018e200bdc720fe4fab627d8970e9faf020 Mon Sep 17 00:00:00 2001 From: acud <12988138+acud@users.noreply.github.com> Date: Sat, 27 Jun 2026 18:25:44 -0600 Subject: [PATCH 2/6] chore: address PR comments --- pkg/statestore/storeadapter/migration.go | 33 +++++++----------------- pkg/topology/kademlia/kademlia.go | 4 +-- 2 files changed, 12 insertions(+), 25 deletions(-) diff --git a/pkg/statestore/storeadapter/migration.go b/pkg/statestore/storeadapter/migration.go index 9f3a83ff804..a0d0341110c 100644 --- a/pkg/statestore/storeadapter/migration.go +++ b/pkg/statestore/storeadapter/migration.go @@ -57,7 +57,7 @@ type migratedAddress struct { type migratedEntry struct { Address migratedAddress `json:"address"` Verified bool `json:"verified"` - LastSeen int64 `json:"last_seen"` + LastSeen int64 `json:"last_seen,omitempty"` } // rewriteAddressbookEnvelope wraps each "addressbook_entry_*" legacy @@ -125,9 +125,9 @@ func rewriteAddressbookEnvelope(s storage.Store) migration.StepFn { // stampAddressbookLastSeen sets "last_seen" to the current time on every // "addressbook_entry_*" record that lacks it, so that addresses carried over // from before pruning was introduced are not immediately pruned. Entries that -// already carry a non-zero last_seen are left untouched. The whole record is -// preserved by merging into the decoded JSON object rather than re-encoding a -// typed struct. +// already carry a non-zero last_seen are left untouched. The record is decoded +// into migratedEntry, the current serialization shape, whose last_seen field is +// omitempty so older records that predate it round-trip unchanged. func stampAddressbookLastSeen(s storage.Store) migration.StepFn { return func() error { store := &StateStorerAdapter{s} @@ -151,31 +151,18 @@ func stampAddressbookLastSeen(s storage.Store) migration.StepFn { now := time.Now().Unix() for _, e := range batch { - var fields map[string]json.RawMessage - if err := json.Unmarshal(e.val, &fields); err != nil { + var entry migratedEntry + if err := json.Unmarshal(e.val, &entry); err != nil { _ = store.Delete(e.key) continue } - if raw, ok := fields["last_seen"]; ok { - var ls int64 - if json.Unmarshal(raw, &ls) == nil && ls != 0 { - continue - } - } - - stamp, err := json.Marshal(now) - if err != nil { - return fmt.Errorf("marshal last_seen: %w", err) - } - fields["last_seen"] = stamp - - out, err := json.Marshal(fields) - if err != nil { - return fmt.Errorf("marshal addressbook entry %q: %w", e.key, err) + if entry.LastSeen != 0 { + continue } + entry.LastSeen = now - if err := store.Put(e.key, json.RawMessage(out)); err != nil { + if err := store.Put(e.key, &entry); err != nil { return fmt.Errorf("stamp addressbook entry %q: %w", e.key, err) } } diff --git a/pkg/topology/kademlia/kademlia.go b/pkg/topology/kademlia/kademlia.go index 6f0d3fc6a0b..c5e9736397a 100644 --- a/pkg/topology/kademlia/kademlia.go +++ b/pkg/topology/kademlia/kademlia.go @@ -456,7 +456,7 @@ func (k *Kad) connectionAttemptsHandler(ctx context.Context, wg *sync.WaitGroup, k.connectedPeers.Add(peer.addr) if err := k.addressBook.UpdateLastSeen(peer.addr); err != nil { - k.logger.Debug("could not update last seen for peer", "peer_address", peer.addr, "error", err) + k.logger.Warning("could not update last seen for peer", "peer_address", peer.addr, "error", err) } k.metrics.TotalOutboundConnections.Inc() @@ -1219,7 +1219,7 @@ func (k *Kad) onConnected(ctx context.Context, addr swarm.Address) error { k.connectedPeers.Add(addr) k.waitNext.Remove(addr) if err := k.addressBook.UpdateLastSeen(addr); err != nil { - k.logger.Debug("could not update last seen for peer", "peer_address", addr, "error", err) + k.logger.Warning("could not update last seen for peer", "peer_address", addr, "error", err) } k.recalcDepth() k.notifyManageLoop() From 157d203f653e305a3a2fb585a5cc8ff8533f5e0d Mon Sep 17 00:00:00 2001 From: Calin Martinconi Date: Mon, 13 Jul 2026 12:51:08 +0300 Subject: [PATCH 3/6] fix(addressbook): refresh last-seen on hive gossip and serialize updates --- pkg/addressbook/addressbook.go | 36 ++++++- pkg/addressbook/addressbook_test.go | 122 ++++++++++++++++++++++- pkg/addressbook/export_test.go | 3 + pkg/hive/hive.go | 15 ++- pkg/hive/lastseen_test.go | 147 ++++++++++++++++++++++++++++ 5 files changed, 315 insertions(+), 8 deletions(-) create mode 100644 pkg/hive/lastseen_test.go diff --git a/pkg/addressbook/addressbook.go b/pkg/addressbook/addressbook.go index 7f31a4f4e2b..b773f040b05 100644 --- a/pkg/addressbook/addressbook.go +++ b/pkg/addressbook/addressbook.go @@ -9,6 +9,7 @@ import ( "errors" "fmt" "strings" + "sync" "time" "github.com/ethersphere/bee/v2/pkg/bzz" @@ -18,6 +19,12 @@ import ( const keyPrefix = "addressbook_entry_" +// lastSeenUpdateInterval is the resolution at which UpdateLastSeen persists. +// Hive reports the same peer on every gossip round, while pruning operates on +// a scale of weeks, so a write is skipped when the stored value is already +// this fresh. +const lastSeenUpdateInterval = 24 * time.Hour + var _ Interface = (*store)(nil) var ErrNotFound = errors.New("addressbook: not found") @@ -54,6 +61,14 @@ type GetPutter interface { Putter } +// GetPutUpdater is the addressbook surface needed by hive: it stores peers it +// learns about and refreshes the last-seen time of the ones it already knows. +type GetPutUpdater interface { + GetPutter + // UpdateLastSeen marks the overlay as seen at the current time. + UpdateLastSeen(overlay swarm.Address) error +} + type Getter interface { // Get returns the saved bzz.Address for the requested overlay together // with its verification flag. @@ -73,6 +88,10 @@ type Remover interface { type store struct { store storage.StateStorer now func() time.Time + + // mu serializes the read-modify-write in UpdateLastSeen against Put, so a + // concurrent Put is not rolled back by a stale copy of the entry. + mu sync.Mutex } // New creates new addressbook for state storer. @@ -97,6 +116,9 @@ func (s *store) Get(overlay swarm.Address) (*bzz.Address, bool, error) { } func (s *store) Put(overlay swarm.Address, addr bzz.Address, verified bool) (err error) { + s.mu.Lock() + defer s.mu.Unlock() + key := keyPrefix + overlay.String() return s.store.Put(key, &verifiedAddress{ Address: &addr, @@ -106,8 +128,12 @@ func (s *store) Put(overlay swarm.Address, addr bzz.Address, verified bool) (err } // UpdateLastSeen marks the overlay as seen at the current time. It is a no-op -// if the overlay is not present in the addressbook. +// if the overlay is not present in the addressbook, or if the recorded time is +// younger than lastSeenUpdateInterval. func (s *store) UpdateLastSeen(overlay swarm.Address) error { + s.mu.Lock() + defer s.mu.Unlock() + key := keyPrefix + overlay.String() v := &verifiedAddress{} if err := s.store.Get(key, v); err != nil { @@ -116,7 +142,13 @@ func (s *store) UpdateLastSeen(overlay swarm.Address) error { } return err } - v.LastSeen = s.now().Unix() + + now := s.now().Unix() + if now-v.LastSeen < int64(lastSeenUpdateInterval.Seconds()) { + return nil + } + + v.LastSeen = now return s.store.Put(key, v) } diff --git a/pkg/addressbook/addressbook_test.go b/pkg/addressbook/addressbook_test.go index ee355e9cde0..ed6f5670c5c 100644 --- a/pkg/addressbook/addressbook_test.go +++ b/pkg/addressbook/addressbook_test.go @@ -14,6 +14,7 @@ import ( "github.com/ethersphere/bee/v2/pkg/bzz" "github.com/ethersphere/bee/v2/pkg/crypto" "github.com/ethersphere/bee/v2/pkg/statestore/mock" + "github.com/ethersphere/bee/v2/pkg/storage" "github.com/ethersphere/bee/v2/pkg/swarm" ma "github.com/multiformats/go-multiaddr" ) @@ -136,14 +137,15 @@ func TestUpdateLastSeen(t *testing.T) { t.Fatal(err) } - // advance the clock and bump last-seen; the entry must survive a prune at - // the original time. - now = time.Unix(5000, 0) + // advance the clock past the update interval and bump last-seen; the entry + // must survive a prune at the original time. + seenAt := now.Add(2 * 24 * time.Hour) + now = seenAt if err := store.UpdateLastSeen(overlay); err != nil { t.Fatal(err) } - if err := store.Prune(time.Unix(4000, 0)); err != nil { + if err := store.Prune(seenAt.Add(-time.Hour)); err != nil { t.Fatal(err) } if _, _, err := store.Get(overlay); err != nil { @@ -151,6 +153,118 @@ func TestUpdateLastSeen(t *testing.T) { } } +// TestUpdateLastSeenThrottled asserts that a bump within lastSeenUpdateInterval +// does not write. Hive calls UpdateLastSeen on every gossip sighting, so this +// is what keeps the write rate bounded to roughly one per peer per day. +func TestUpdateLastSeenThrottled(t *testing.T) { + t.Parallel() + + base := time.Unix(1_000_000, 0) + now := base + mockStore := mock.NewStateStore() + store := addressbook.NewWithClock(mockStore, func() time.Time { return now }) + + overlay := swarm.NewAddress([]byte{0, 1, 2, 3}) + if err := store.Put(overlay, newTestAddr(t, overlay), true); err != nil { + t.Fatal(err) + } + + // well inside the interval: must not touch the record. + now = base.Add(time.Hour) + if err := store.UpdateLastSeen(overlay); err != nil { + t.Fatal(err) + } + if got := lastSeenOf(t, mockStore, overlay); got != base.Unix() { + t.Fatalf("throttled update wrote: last_seen = %d, want %d", got, base.Unix()) + } + + // past the interval: must write. + now = base.Add(addressbook.LastSeenUpdateInterval + time.Second) + if err := store.UpdateLastSeen(overlay); err != nil { + t.Fatal(err) + } + if got := lastSeenOf(t, mockStore, overlay); got != now.Unix() { + t.Fatalf("update past interval did not write: last_seen = %d, want %d", got, now.Unix()) + } +} + +// TestUpdateLastSeenKeepsConcurrentPut pins the read-modify-write in +// UpdateLastSeen against a Put that lands between its read and its write. +// Without serialization the Put's Verified flag is rolled back, which would +// also desync the addressbook from hive's chequebook registry. +func TestUpdateLastSeenKeepsConcurrentPut(t *testing.T) { + t.Parallel() + + base := time.Unix(1_000_000, 0) + now := base + hooked := &hookStore{StateStorer: mock.NewStateStore()} + book := addressbook.NewWithClock(hooked, func() time.Time { return now }) + + overlay := swarm.NewAddress([]byte{0, 1, 2, 3}) + addr := newTestAddr(t, overlay) + + // a known, not-yet-verified peer. + if err := book.Put(overlay, addr, false); err != nil { + t.Fatal(err) + } + // move past the throttle so UpdateLastSeen really writes. + now = base.Add(addressbook.LastSeenUpdateInterval + time.Second) + + // While UpdateLastSeen holds the entry it has just read, hive verifies the + // same peer and stores it with Verified=true. + started, finished := make(chan struct{}), make(chan struct{}) + hooked.onGet = func() { + go func() { + defer close(finished) + close(started) + if err := book.Put(overlay, addr, true); err != nil { + t.Error(err) + } + }() + <-started + // Give the writer time to land. Serialized, it blocks on the + // addressbook lock until UpdateLastSeen returns; unsynchronized, its + // write completes here and is then overwritten below. + time.Sleep(100 * time.Millisecond) + } + + if err := book.UpdateLastSeen(overlay); err != nil { + t.Fatal(err) + } + <-finished + + if _, verified, err := book.Get(overlay); err != nil || !verified { + t.Fatalf("concurrent Put(verified=true) was rolled back: verified=%v err=%v", verified, err) + } +} + +// hookStore fires onGet once, immediately after a Get returns, to interleave a +// concurrent writer inside UpdateLastSeen's read-modify-write. +type hookStore struct { + storage.StateStorer + onGet func() +} + +func (h *hookStore) Get(key string, i any) error { + err := h.StateStorer.Get(key, i) + if h.onGet != nil { + f := h.onGet + h.onGet = nil + f() + } + return err +} + +func lastSeenOf(t *testing.T, store storage.StateStorer, overlay swarm.Address) int64 { + t.Helper() + + v := &addressbook.VerifiedAddress{} + if err := store.Get("addressbook_entry_"+overlay.String(), v); err != nil { + t.Fatalf("get entry: %v", err) + } + return v.LastSeen +} + func TestPrune(t *testing.T) { t.Parallel() diff --git a/pkg/addressbook/export_test.go b/pkg/addressbook/export_test.go index 325621e1266..af0060b92b8 100644 --- a/pkg/addressbook/export_test.go +++ b/pkg/addressbook/export_test.go @@ -12,6 +12,9 @@ import ( type VerifiedAddress = verifiedAddress +// LastSeenUpdateInterval exposes the UpdateLastSeen write throttle. +const LastSeenUpdateInterval = lastSeenUpdateInterval + // NewWithClock creates an addressbook with an overridable clock, for testing. func NewWithClock(storer storage.StateStorer, now func() time.Time) Interface { return &store{ diff --git a/pkg/hive/hive.go b/pkg/hive/hive.go index a991dbdd58b..a3f41477142 100644 --- a/pkg/hive/hive.go +++ b/pkg/hive/hive.go @@ -70,7 +70,7 @@ type Options struct { type Service struct { streamer p2p.Streamer - addressBook addressbook.GetPutter + addressBook addressbook.GetPutUpdater addPeersHandler func(...swarm.Address) networkID uint64 logger log.Logger @@ -92,7 +92,7 @@ type Service struct { chequebookStorer ChequebookStorer } -func New(streamer p2p.Streamer, addressbook addressbook.GetPutter, networkID uint64, overlay swarm.Address, logger log.Logger, o Options) *Service { +func New(streamer p2p.Streamer, addressbook addressbook.GetPutUpdater, networkID uint64, overlay swarm.Address, logger log.Logger, o Options) *Service { svc := &Service{ streamer: streamer, logger: logger.WithName(loggerName).Register(), @@ -370,6 +370,17 @@ func (s *Service) checkAndAddPeers(ctx context.Context, peers pb.Peers) { } if err := bzz.CheckTimestamp(bzzAddress.Timestamp, existing, bzz.TimestampSourceGossip, s.now()); err != nil { + // A peer re-presenting a record we already hold has nothing new to + // store, but it is still a sighting: peers mint their bzz.Address + // once and gossip it unchanged for their whole uptime. Refresh + // last-seen so peers we keep hearing about, but never dial, do not + // look stale to the pruner. Both errors imply a known peer, since + // an unknown one short-circuits inside CheckTimestamp. + if errors.Is(err, bzz.ErrTimestampStale) || errors.Is(err, bzz.ErrTimestampTooSoon) { + if err := s.addressBook.UpdateLastSeen(overlayAddr); err != nil { + s.logger.Debug("hive gossip: update last seen", "overlay", overlayAddr.String(), "error", err) + } + } s.bumpTimestampMetric(err) s.logger.Debug("hive gossip: timestamp validation failed", "overlay", overlayAddr.String(), "error", err) continue diff --git a/pkg/hive/lastseen_test.go b/pkg/hive/lastseen_test.go new file mode 100644 index 00000000000..23740e0a373 --- /dev/null +++ b/pkg/hive/lastseen_test.go @@ -0,0 +1,147 @@ +// Copyright 2026 The Swarm Authors. All rights reserved. +// Use of this source code is governed by a BSD-style +// license that can be found in the LICENSE file. + +package hive_test + +import ( + "sync" + "testing" + "time" + + ab "github.com/ethersphere/bee/v2/pkg/addressbook" + "github.com/ethersphere/bee/v2/pkg/bzz" + "github.com/ethersphere/bee/v2/pkg/hive" + "github.com/ethersphere/bee/v2/pkg/hive/pb" + "github.com/ethersphere/bee/v2/pkg/log" + "github.com/ethersphere/bee/v2/pkg/p2p/streamtest" + "github.com/ethersphere/bee/v2/pkg/statestore/mock" + "github.com/ethersphere/bee/v2/pkg/swarm" +) + +// lastSeenSpy counts the addressbook writes hive performs, so we can tell a +// re-sighting (UpdateLastSeen) from a record update (Put). +type lastSeenSpy struct { + ab.Interface + mu sync.Mutex + puts int + updates int +} + +func (s *lastSeenSpy) Put(o swarm.Address, a bzz.Address, v bool) error { + s.mu.Lock() + s.puts++ + s.mu.Unlock() + return s.Interface.Put(o, a, v) +} + +func (s *lastSeenSpy) UpdateLastSeen(o swarm.Address) error { + s.mu.Lock() + s.updates++ + s.mu.Unlock() + return s.Interface.UpdateLastSeen(o) +} + +func (s *lastSeenSpy) counts() (puts, updates int) { + s.mu.Lock() + defer s.mu.Unlock() + return s.puts, s.updates +} + +func newHiveWithSpy(t *testing.T, networkID uint64, now *time.Time) (*hive.Service, *lastSeenSpy) { + t.Helper() + + spy := &lastSeenSpy{Interface: ab.New(mock.NewStateStore())} + svc := hive.New(streamtest.New(), spy, networkID, swarm.RandAddress(t), log.Noop, hive.Options{ + AllowPrivateCIDRs: true, + }) + svc.SetTimeFunc(func() time.Time { return *now }) + t.Cleanup(func() { _ = svc.Close() }) + return svc, spy +} + +// TestGossipRefreshesLastSeenOnRepeatSighting covers the case that keeps +// gossip-only peers alive. A peer mints its bzz.Address once and re-presents +// that same signed record for its whole uptime, so CheckTimestamp rejects it +// as stale and there is nothing to Put. It is still a sighting, and last-seen +// must move, or the pruner eventually evicts a peer we hear about constantly. +func TestGossipRefreshesLastSeenOnRepeatSighting(t *testing.T) { + t.Parallel() + + const networkID = uint64(1) + + base := time.Unix(1_700_000_000, 0) + now := base + svc, spy := newHiveWithSpy(t, networkID, &now) + + id := newPeerIdentity(t, networkID, "/ip4/10.0.0.1/tcp/1634") + rec := id.protoAt(t, networkID, base.Unix()) + + // first sighting: unknown peer, gets stored. + svc.CheckAndAddPeers(pb.Peers{Peers: []*pb.BzzAddress{rec}}) + if puts, updates := spy.counts(); puts != 1 || updates != 0 { + t.Fatalf("first sighting: puts=%d updates=%d, want 1/0", puts, updates) + } + + // ten days later, the same peer is still being gossiped to us with the very + // same record. + now = base.Add(10 * 24 * time.Hour) + svc.CheckAndAddPeers(pb.Peers{Peers: []*pb.BzzAddress{rec}}) + + puts, updates := spy.counts() + if puts != 1 { + t.Fatalf("re-sighting stored the record again: puts=%d, want 1", puts) + } + if updates != 1 { + t.Fatalf("re-sighting did not refresh last-seen: updates=%d, want 1", updates) + } +} + +// TestGossipUpdatesRecordWhenGenuinelyNewer asserts that a record newer than +// existing.Timestamp+MinimumUpdateInterval still takes the Put path, where +// last-seen is stamped as part of the write. +func TestGossipUpdatesRecordWhenGenuinelyNewer(t *testing.T) { + t.Parallel() + + const networkID = uint64(1) + + base := time.Unix(1_700_000_000, 0) + now := base + svc, spy := newHiveWithSpy(t, networkID, &now) + + id := newPeerIdentity(t, networkID, "/ip4/10.0.0.1/tcp/1634") + svc.CheckAndAddPeers(pb.Peers{Peers: []*pb.BzzAddress{id.protoAt(t, networkID, base.Unix())}}) + + // re-minted beyond the minimum update interval. + newer := base.Add(bzz.MinimumUpdateInterval + time.Second) + now = newer + svc.CheckAndAddPeers(pb.Peers{Peers: []*pb.BzzAddress{id.protoAt(t, networkID, newer.Unix())}}) + + if puts, updates := spy.counts(); puts != 2 || updates != 0 { + t.Fatalf("newer record: puts=%d updates=%d, want 2/0", puts, updates) + } +} + +// TestGossipDoesNotRefreshLastSeenOnInvalidRecord makes sure the refresh is +// scoped to records that merely are not newer. A record rejected for any other +// reason is not evidence that the peer is alive. +func TestGossipDoesNotRefreshLastSeenOnInvalidRecord(t *testing.T) { + t.Parallel() + + const networkID = uint64(1) + + base := time.Unix(1_700_000_000, 0) + now := base + svc, spy := newHiveWithSpy(t, networkID, &now) + + id := newPeerIdentity(t, networkID, "/ip4/10.0.0.1/tcp/1634") + svc.CheckAndAddPeers(pb.Peers{Peers: []*pb.BzzAddress{id.protoAt(t, networkID, base.Unix())}}) + + // a record timestamped far in the future is rejected as invalid, not stale. + future := base.Add(24 * time.Hour) + svc.CheckAndAddPeers(pb.Peers{Peers: []*pb.BzzAddress{id.protoAt(t, networkID, future.Unix())}}) + + if puts, updates := spy.counts(); puts != 1 || updates != 0 { + t.Fatalf("invalid record touched the addressbook: puts=%d updates=%d, want 1/0", puts, updates) + } +} From 699b4d7a920347a847cf85b4a447df694146c20a Mon Sep 17 00:00:00 2001 From: Calin Martinconi Date: Wed, 15 Jul 2026 18:33:54 +0300 Subject: [PATCH 4/6] fix(addressbook): mark peers seen on every sighting and refresh connected peers --- pkg/addressbook/addressbook.go | 93 ++++++----- pkg/addressbook/addressbook_test.go | 152 +++++++++--------- pkg/addressbook/export_test.go | 8 +- pkg/hive/hive.go | 27 ++-- pkg/hive/lastseen_test.go | 77 +++++---- pkg/node/node.go | 11 +- pkg/statestore/storeadapter/migration_test.go | 9 +- pkg/topology/kademlia/export_test.go | 6 + pkg/topology/kademlia/kademlia.go | 44 ++++- pkg/topology/kademlia/lastseen_test.go | 32 ++-- 10 files changed, 246 insertions(+), 213 deletions(-) diff --git a/pkg/addressbook/addressbook.go b/pkg/addressbook/addressbook.go index b773f040b05..59e08ae89a8 100644 --- a/pkg/addressbook/addressbook.go +++ b/pkg/addressbook/addressbook.go @@ -19,12 +19,6 @@ import ( const keyPrefix = "addressbook_entry_" -// lastSeenUpdateInterval is the resolution at which UpdateLastSeen persists. -// Hive reports the same peer on every gossip round, while pruning operates on -// a scale of weeks, so a write is skipped when the stored value is already -// this fresh. -const lastSeenUpdateInterval = 24 * time.Hour - var _ Interface = (*store)(nil) var ErrNotFound = errors.New("addressbook: not found") @@ -43,17 +37,13 @@ type verifiedAddress struct { type Interface interface { GetPutter Remover - // UpdateLastSeen marks the overlay as seen at the current time. - UpdateLastSeen(overlay swarm.Address) error + Seener // Overlays returns a list of all overlay addresses saved in addressbook. Overlays() ([]swarm.Address, error) // IterateOverlays exposes overlays in a form of an iterator. IterateOverlays(func(swarm.Address) (bool, error)) error // Addresses returns a list of all bzz.Address-es saved in addressbook. Addresses() ([]bzz.Address, error) - // Prune removes all entries whose overlay has not been seen since the - // given time. - Prune(before time.Time) error } type GetPutter interface { @@ -61,12 +51,11 @@ type GetPutter interface { Putter } -// GetPutUpdater is the addressbook surface needed by hive: it stores peers it -// learns about and refreshes the last-seen time of the ones it already knows. -type GetPutUpdater interface { +// GetPutSeener is the addressbook surface needed by hive: it stores the peers +// it learns about and marks the ones it already knows as seen. +type GetPutSeener interface { GetPutter - // UpdateLastSeen marks the overlay as seen at the current time. - UpdateLastSeen(overlay swarm.Address) error + Seener } type Getter interface { @@ -85,20 +74,29 @@ type Remover interface { Remove(overlay swarm.Address) error } +type Seener interface { + // Seen marks the overlays as seen at the current time. + Seen(overlays ...swarm.Address) error +} + type store struct { store storage.StateStorer now func() time.Time - // mu serializes the read-modify-write in UpdateLastSeen against Put, so a + // mu serializes the read-modify-write in Seen against Put, so that a // concurrent Put is not rolled back by a stale copy of the entry. mu sync.Mutex } // New creates new addressbook for state storer. func New(storer storage.StateStorer) Interface { + return newStore(storer, time.Now) +} + +func newStore(storer storage.StateStorer, now func() time.Time) *store { return &store{ store: storer, - now: time.Now, + now: now, } } @@ -127,53 +125,54 @@ func (s *store) Put(overlay swarm.Address, addr bzz.Address, verified bool) (err }) } -// UpdateLastSeen marks the overlay as seen at the current time. It is a no-op -// if the overlay is not present in the addressbook, or if the recorded time is -// younger than lastSeenUpdateInterval. -func (s *store) UpdateLastSeen(overlay swarm.Address) error { +// Seen marks the overlays as seen at the current time. An overlay that is not +// present in the addressbook is skipped. +func (s *store) Seen(overlays ...swarm.Address) error { s.mu.Lock() defer s.mu.Unlock() - key := keyPrefix + overlay.String() - v := &verifiedAddress{} - if err := s.store.Get(key, v); err != nil { - if errors.Is(err, storage.ErrNotFound) { - return nil + now := s.now().Unix() + + for _, overlay := range overlays { + key := keyPrefix + overlay.String() + + v := &verifiedAddress{} + if err := s.store.Get(key, v); err != nil { + if errors.Is(err, storage.ErrNotFound) { + continue + } + return err } - return err - } - now := s.now().Unix() - if now-v.LastSeen < int64(lastSeenUpdateInterval.Seconds()) { - return nil + v.LastSeen = now + if err := s.store.Put(key, v); err != nil { + return err + } } - v.LastSeen = now - return s.store.Put(key, v) + return nil } func (s *store) Remove(overlay swarm.Address) error { return s.store.Delete(keyPrefix + overlay.String()) } -// Prune removes all entries whose overlay has not been seen since before. -// Entries without a recorded last-seen time (LastSeen == 0) are kept, leaving -// pruning to a later run once they have been observed. -func (s *store) Prune(before time.Time) error { +// Prune removes all entries from storer whose overlay has not been seen since +// before. Entries without a recorded last-seen time (LastSeen == 0) are kept, +// leaving them to a later run once they have been observed. It is a one-shot +// maintenance step meant to run at startup, before the address book is shared +// with its writers, so it takes no lock. +func Prune(storer storage.StateStorer, before time.Time) error { cutoff := before.Unix() - var stale []swarm.Address - err := s.store.Iterate(keyPrefix, func(key, value []byte) (stop bool, err error) { + var stale []string + err := storer.Iterate(keyPrefix, func(key, value []byte) (stop bool, err error) { entry := &verifiedAddress{} if err := json.Unmarshal(value, entry); err != nil { return true, err } if entry.LastSeen != 0 && entry.LastSeen < cutoff { - addr, err := swarm.ParseHexAddress(strings.TrimPrefix(string(key), keyPrefix)) - if err != nil { - return true, err - } - stale = append(stale, addr) + stale = append(stale, string(key)) } return false, nil }) @@ -181,8 +180,8 @@ func (s *store) Prune(before time.Time) error { return err } - for _, addr := range stale { - if err := s.Remove(addr); err != nil { + for _, key := range stale { + if err := storer.Delete(key); err != nil { return err } } diff --git a/pkg/addressbook/addressbook_test.go b/pkg/addressbook/addressbook_test.go index ed6f5670c5c..9be9ef4db8b 100644 --- a/pkg/addressbook/addressbook_test.go +++ b/pkg/addressbook/addressbook_test.go @@ -117,101 +117,101 @@ func run(t *testing.T, f bookFunc) { } } -func TestUpdateLastSeen(t *testing.T) { +// TestSeen covers a sighting of a peer we already hold: the last-seen time +// moves, and the entry then survives a prune that would otherwise catch it. An +// overlay we do not know is skipped rather than created. +func TestSeen(t *testing.T) { t.Parallel() - now := time.Unix(1000, 0) - store := addressbook.NewWithClock(mock.NewStateStore(), func() time.Time { return now }) + base := time.Unix(1_000_000, 0) + now := base + state := mock.NewStateStore() + book := addressbook.NewWithClock(state, func() time.Time { return now }) overlay := swarm.NewAddress([]byte{0, 1, 2, 3}) - // UpdateLastSeen on a missing overlay is a no-op and must not create an entry. - if err := store.UpdateLastSeen(overlay); err != nil { + // an unknown overlay is skipped, not created. + if err := book.Seen(overlay); err != nil { t.Fatal(err) } - if _, _, err := store.Get(overlay); !errors.Is(err, addressbook.ErrNotFound) { + if _, _, err := book.Get(overlay); !errors.Is(err, addressbook.ErrNotFound) { t.Fatalf("expected ErrNotFound, got %v", err) } - if err := store.Put(overlay, newTestAddr(t, overlay), true); err != nil { + if err := book.Put(overlay, newTestAddr(t, overlay), true); err != nil { t.Fatal(err) } - // advance the clock past the update interval and bump last-seen; the entry - // must survive a prune at the original time. - seenAt := now.Add(2 * 24 * time.Hour) - now = seenAt - if err := store.UpdateLastSeen(overlay); err != nil { + // a much later sighting moves last-seen forward, so a prune whose cutoff + // predates it keeps the entry. + now = base.Add(90 * 24 * time.Hour) + if err := book.Seen(overlay); err != nil { t.Fatal(err) } + if got := lastSeenOf(t, state, overlay); got != now.Unix() { + t.Fatalf("seen did not move last seen: got %d, want %d", got, now.Unix()) + } - if err := store.Prune(seenAt.Add(-time.Hour)); err != nil { + if err := addressbook.Prune(state, now.Add(-time.Hour)); err != nil { t.Fatal(err) } - if _, _, err := store.Get(overlay); err != nil { - t.Fatalf("entry pruned despite recent last-seen: %v", err) + if _, verified, err := book.Get(overlay); err != nil || !verified { + t.Fatalf("entry pruned despite a recent sighting: verified=%v err=%v", verified, err) } } -// TestUpdateLastSeenThrottled asserts that a bump within lastSeenUpdateInterval -// does not write. Hive calls UpdateLastSeen on every gossip sighting, so this -// is what keeps the write rate bounded to roughly one per peer per day. -func TestUpdateLastSeenThrottled(t *testing.T) { +// TestSeenVariadic marks several overlays in one call, which is how kademlia +// refreshes everything it is connected to. +func TestSeenVariadic(t *testing.T) { t.Parallel() - base := time.Unix(1_000_000, 0) - now := base - mockStore := mock.NewStateStore() - store := addressbook.NewWithClock(mockStore, func() time.Time { return now }) + now := time.Unix(1_000_000, 0) + state := mock.NewStateStore() + book := addressbook.NewWithClock(state, func() time.Time { return now }) - overlay := swarm.NewAddress([]byte{0, 1, 2, 3}) - if err := store.Put(overlay, newTestAddr(t, overlay), true); err != nil { - t.Fatal(err) + overlays := []swarm.Address{ + swarm.NewAddress([]byte{0, 1, 2, 3}), + swarm.NewAddress([]byte{0, 1, 2, 4}), } - - // well inside the interval: must not touch the record. - now = base.Add(time.Hour) - if err := store.UpdateLastSeen(overlay); err != nil { - t.Fatal(err) - } - if got := lastSeenOf(t, mockStore, overlay); got != base.Unix() { - t.Fatalf("throttled update wrote: last_seen = %d, want %d", got, base.Unix()) + for _, overlay := range overlays { + if err := book.Put(overlay, newTestAddr(t, overlay), true); err != nil { + t.Fatal(err) + } } - // past the interval: must write. - now = base.Add(addressbook.LastSeenUpdateInterval + time.Second) - if err := store.UpdateLastSeen(overlay); err != nil { + now = now.Add(time.Hour) + if err := book.Seen(overlays...); err != nil { t.Fatal(err) } - if got := lastSeenOf(t, mockStore, overlay); got != now.Unix() { - t.Fatalf("update past interval did not write: last_seen = %d, want %d", got, now.Unix()) + + for _, overlay := range overlays { + if got := lastSeenOf(t, state, overlay); got != now.Unix() { + t.Fatalf("overlay %s not marked seen: got %d, want %d", overlay, got, now.Unix()) + } } } -// TestUpdateLastSeenKeepsConcurrentPut pins the read-modify-write in -// UpdateLastSeen against a Put that lands between its read and its write. -// Without serialization the Put's Verified flag is rolled back, which would +// TestSeenKeepsConcurrentPut pins Seen's read-modify-write against a Put that +// lands between its read and its write. Seen rewrites the whole record, so +// without serialization the Put's verified flag is rolled back, which would // also desync the addressbook from hive's chequebook registry. -func TestUpdateLastSeenKeepsConcurrentPut(t *testing.T) { +func TestSeenKeepsConcurrentPut(t *testing.T) { t.Parallel() - base := time.Unix(1_000_000, 0) - now := base + now := time.Unix(1_000_000, 0) hooked := &hookStore{StateStorer: mock.NewStateStore()} book := addressbook.NewWithClock(hooked, func() time.Time { return now }) overlay := swarm.NewAddress([]byte{0, 1, 2, 3}) addr := newTestAddr(t, overlay) - // a known, not-yet-verified peer. + // a known, not yet verified peer. if err := book.Put(overlay, addr, false); err != nil { t.Fatal(err) } - // move past the throttle so UpdateLastSeen really writes. - now = base.Add(addressbook.LastSeenUpdateInterval + time.Second) - // While UpdateLastSeen holds the entry it has just read, hive verifies the - // same peer and stores it with Verified=true. + // While Seen holds the entry it has just read, hive verifies the same peer + // and stores it with Verified=true. started, finished := make(chan struct{}), make(chan struct{}) hooked.onGet = func() { go func() { @@ -223,12 +223,12 @@ func TestUpdateLastSeenKeepsConcurrentPut(t *testing.T) { }() <-started // Give the writer time to land. Serialized, it blocks on the - // addressbook lock until UpdateLastSeen returns; unsynchronized, its - // write completes here and is then overwritten below. + // addressbook lock until Seen returns; unsynchronized, its write + // completes here and is then overwritten below. time.Sleep(100 * time.Millisecond) } - if err := book.UpdateLastSeen(overlay); err != nil { + if err := book.Seen(overlay); err != nil { t.Fatal(err) } <-finished @@ -239,7 +239,7 @@ func TestUpdateLastSeenKeepsConcurrentPut(t *testing.T) { } // hookStore fires onGet once, immediately after a Get returns, to interleave a -// concurrent writer inside UpdateLastSeen's read-modify-write. +// concurrent writer inside Seen's read-modify-write. type hookStore struct { storage.StateStorer onGet func() @@ -255,69 +255,71 @@ func (h *hookStore) Get(key string, i any) error { return err } -func lastSeenOf(t *testing.T, store storage.StateStorer, overlay swarm.Address) int64 { +func lastSeenOf(t *testing.T, state storage.StateStorer, overlay swarm.Address) int64 { t.Helper() v := &addressbook.VerifiedAddress{} - if err := store.Get("addressbook_entry_"+overlay.String(), v); err != nil { + if err := state.Get("addressbook_entry_"+overlay.String(), v); err != nil { t.Fatalf("get entry: %v", err) } return v.LastSeen } +// TestPrune drops overlays last seen before the cutoff and keeps the rest. func TestPrune(t *testing.T) { t.Parallel() - now := time.Unix(0, 0) - store := addressbook.NewWithClock(mock.NewStateStore(), func() time.Time { return now }) + now := time.Unix(1_000_000_000, 0) + state := mock.NewStateStore() + book := addressbook.NewWithClock(state, func() time.Time { return now }) stale := swarm.NewAddress([]byte{0, 1, 2, 3}) - fresh := swarm.NewAddress([]byte{0, 1, 2, 4}) - - now = time.Unix(1000, 0) - if err := store.Put(stale, newTestAddr(t, stale), true); err != nil { + if err := book.Put(stale, newTestAddr(t, stale), true); err != nil { t.Fatal(err) } - now = time.Unix(9000, 0) - if err := store.Put(fresh, newTestAddr(t, fresh), true); err != nil { + now = now.Add(48 * time.Hour) + fresh := swarm.NewAddress([]byte{0, 1, 2, 4}) + if err := book.Put(fresh, newTestAddr(t, fresh), true); err != nil { t.Fatal(err) } - if err := store.Prune(time.Unix(5000, 0)); err != nil { + // cutoff sits between the two puts: stale is before it, fresh is after. + if err := addressbook.Prune(state, now.Add(-time.Hour)); err != nil { t.Fatal(err) } - if _, _, err := store.Get(stale); !errors.Is(err, addressbook.ErrNotFound) { + if _, _, err := book.Get(stale); !errors.Is(err, addressbook.ErrNotFound) { t.Fatalf("stale entry should have been pruned, got err=%v", err) } - if _, _, err := store.Get(fresh); err != nil { + if _, _, err := book.Get(fresh); err != nil { t.Fatalf("fresh entry should survive: %v", err) } } +// TestPruneKeepsEntriesWithoutLastSeen covers records that predate pruning and +// have not been stamped by the migration. They are kept, and stamped on their +// next sighting. func TestPruneKeepsEntriesWithoutLastSeen(t *testing.T) { t.Parallel() - mockStore := mock.NewStateStore() + state := mock.NewStateStore() overlay := swarm.NewAddress([]byte{0, 1, 2, 3}) - // Seed an entry without a last_seen field, mirroring records that predate - // pruning before the stamping migration runs. - if err := mockStore.Put("addressbook_entry_"+overlay.String(), &addressbook.VerifiedAddress{ + if err := state.Put("addressbook_entry_"+overlay.String(), &addressbook.VerifiedAddress{ Address: addrPtr(newTestAddr(t, overlay)), Verified: true, }); err != nil { t.Fatal(err) } - store := addressbook.New(mockStore) - if err := store.Prune(time.Unix(5000, 0)); err != nil { + if err := addressbook.Prune(state, time.Unix(5_000_000_000, 0)); err != nil { t.Fatal(err) } - if _, _, err := store.Get(overlay); err != nil { - t.Fatalf("entry without last_seen must not be pruned: %v", err) + book := addressbook.New(state) + if _, _, err := book.Get(overlay); err != nil { + t.Fatalf("entry without a last-seen time must not be pruned: %v", err) } } diff --git a/pkg/addressbook/export_test.go b/pkg/addressbook/export_test.go index af0060b92b8..f26d03bd661 100644 --- a/pkg/addressbook/export_test.go +++ b/pkg/addressbook/export_test.go @@ -12,13 +12,7 @@ import ( type VerifiedAddress = verifiedAddress -// LastSeenUpdateInterval exposes the UpdateLastSeen write throttle. -const LastSeenUpdateInterval = lastSeenUpdateInterval - // NewWithClock creates an addressbook with an overridable clock, for testing. func NewWithClock(storer storage.StateStorer, now func() time.Time) Interface { - return &store{ - store: storer, - now: now, - } + return newStore(storer, now) } diff --git a/pkg/hive/hive.go b/pkg/hive/hive.go index a3f41477142..192467d5339 100644 --- a/pkg/hive/hive.go +++ b/pkg/hive/hive.go @@ -70,7 +70,7 @@ type Options struct { type Service struct { streamer p2p.Streamer - addressBook addressbook.GetPutUpdater + addressBook addressbook.GetPutSeener addPeersHandler func(...swarm.Address) networkID uint64 logger log.Logger @@ -92,7 +92,7 @@ type Service struct { chequebookStorer ChequebookStorer } -func New(streamer p2p.Streamer, addressbook addressbook.GetPutUpdater, networkID uint64, overlay swarm.Address, logger log.Logger, o Options) *Service { +func New(streamer p2p.Streamer, addressbook addressbook.GetPutSeener, networkID uint64, overlay swarm.Address, logger log.Logger, o Options) *Service { svc := &Service{ streamer: streamer, logger: logger.WithName(loggerName).Register(), @@ -369,18 +369,19 @@ func (s *Service) checkAndAddPeers(ctx context.Context, peers pb.Peers) { continue } - if err := bzz.CheckTimestamp(bzzAddress.Timestamp, existing, bzz.TimestampSourceGossip, s.now()); err != nil { - // A peer re-presenting a record we already hold has nothing new to - // store, but it is still a sighting: peers mint their bzz.Address - // once and gossip it unchanged for their whole uptime. Refresh - // last-seen so peers we keep hearing about, but never dial, do not - // look stale to the pruner. Both errors imply a known peer, since - // an unknown one short-circuits inside CheckTimestamp. - if errors.Is(err, bzz.ErrTimestampStale) || errors.Is(err, bzz.ErrTimestampTooSoon) { - if err := s.addressBook.UpdateLastSeen(overlayAddr); err != nil { - s.logger.Debug("hive gossip: update last seen", "overlay", overlayAddr.String(), "error", err) - } + // Hearing about a peer we already know is a sighting in its own right, + // whether or not the record it carries is newer than the one we hold. + // Peers mint their bzz.Address once and gossip it unchanged for their + // whole uptime, so for a known peer there is almost never anything new + // to store, and the pruner would evict peers we are told about + // constantly. + if existing != nil { + if err := s.addressBook.Seen(overlayAddr); err != nil { + s.logger.Debug("hive gossip: mark peer seen", "overlay", overlayAddr.String(), "error", err) } + } + + if err := bzz.CheckTimestamp(bzzAddress.Timestamp, existing, bzz.TimestampSourceGossip, s.now()); err != nil { s.bumpTimestampMetric(err) s.logger.Debug("hive gossip: timestamp validation failed", "overlay", overlayAddr.String(), "error", err) continue diff --git a/pkg/hive/lastseen_test.go b/pkg/hive/lastseen_test.go index 23740e0a373..2f3c3cdaba4 100644 --- a/pkg/hive/lastseen_test.go +++ b/pkg/hive/lastseen_test.go @@ -19,13 +19,13 @@ import ( "github.com/ethersphere/bee/v2/pkg/swarm" ) -// lastSeenSpy counts the addressbook writes hive performs, so we can tell a -// re-sighting (UpdateLastSeen) from a record update (Put). +// lastSeenSpy counts the addressbook writes hive performs, so that a sighting +// can be told apart from a record update. type lastSeenSpy struct { ab.Interface - mu sync.Mutex - puts int - updates int + mu sync.Mutex + puts int + seens int } func (s *lastSeenSpy) Put(o swarm.Address, a bzz.Address, v bool) error { @@ -35,17 +35,17 @@ func (s *lastSeenSpy) Put(o swarm.Address, a bzz.Address, v bool) error { return s.Interface.Put(o, a, v) } -func (s *lastSeenSpy) UpdateLastSeen(o swarm.Address) error { +func (s *lastSeenSpy) Seen(o ...swarm.Address) error { s.mu.Lock() - s.updates++ + s.seens++ s.mu.Unlock() - return s.Interface.UpdateLastSeen(o) + return s.Interface.Seen(o...) } -func (s *lastSeenSpy) counts() (puts, updates int) { +func (s *lastSeenSpy) counts() (puts, seens int) { s.mu.Lock() defer s.mu.Unlock() - return s.puts, s.updates + return s.puts, s.seens } func newHiveWithSpy(t *testing.T, networkID uint64, now *time.Time) (*hive.Service, *lastSeenSpy) { @@ -60,12 +60,12 @@ func newHiveWithSpy(t *testing.T, networkID uint64, now *time.Time) (*hive.Servi return svc, spy } -// TestGossipRefreshesLastSeenOnRepeatSighting covers the case that keeps -// gossip-only peers alive. A peer mints its bzz.Address once and re-presents -// that same signed record for its whole uptime, so CheckTimestamp rejects it -// as stale and there is nothing to Put. It is still a sighting, and last-seen -// must move, or the pruner eventually evicts a peer we hear about constantly. -func TestGossipRefreshesLastSeenOnRepeatSighting(t *testing.T) { +// TestSeenOnRepeatGossip covers the sighting that keeps gossip-only peers +// alive. A peer mints its bzz.Address once and re-presents that same signed +// record for its whole uptime, so there is nothing new to Put. It is still a +// sighting, and last-seen must move, or the pruner evicts a peer we are told +// about constantly. +func TestSeenOnRepeatGossip(t *testing.T) { t.Parallel() const networkID = uint64(1) @@ -77,30 +77,29 @@ func TestGossipRefreshesLastSeenOnRepeatSighting(t *testing.T) { id := newPeerIdentity(t, networkID, "/ip4/10.0.0.1/tcp/1634") rec := id.protoAt(t, networkID, base.Unix()) - // first sighting: unknown peer, gets stored. + // first sighting: an unknown peer, stored by Put, which stamps last-seen. svc.CheckAndAddPeers(pb.Peers{Peers: []*pb.BzzAddress{rec}}) - if puts, updates := spy.counts(); puts != 1 || updates != 0 { - t.Fatalf("first sighting: puts=%d updates=%d, want 1/0", puts, updates) + if puts, seens := spy.counts(); puts != 1 || seens != 0 { + t.Fatalf("first sighting: puts=%d seens=%d, want 1/0", puts, seens) } - // ten days later, the same peer is still being gossiped to us with the very - // same record. + // ten days on, the same peer is still gossiped to us with the very same + // record. now = base.Add(10 * 24 * time.Hour) svc.CheckAndAddPeers(pb.Peers{Peers: []*pb.BzzAddress{rec}}) - puts, updates := spy.counts() + puts, seens := spy.counts() if puts != 1 { t.Fatalf("re-sighting stored the record again: puts=%d, want 1", puts) } - if updates != 1 { - t.Fatalf("re-sighting did not refresh last-seen: updates=%d, want 1", updates) + if seens != 1 { + t.Fatalf("re-sighting did not refresh last-seen: seens=%d, want 1", seens) } } -// TestGossipUpdatesRecordWhenGenuinelyNewer asserts that a record newer than -// existing.Timestamp+MinimumUpdateInterval still takes the Put path, where -// last-seen is stamped as part of the write. -func TestGossipUpdatesRecordWhenGenuinelyNewer(t *testing.T) { +// TestSeenOnNewerRecord asserts that a genuinely newer record still takes the +// Put path, where last-seen is stamped as part of the write. +func TestSeenOnNewerRecord(t *testing.T) { t.Parallel() const networkID = uint64(1) @@ -117,15 +116,15 @@ func TestGossipUpdatesRecordWhenGenuinelyNewer(t *testing.T) { now = newer svc.CheckAndAddPeers(pb.Peers{Peers: []*pb.BzzAddress{id.protoAt(t, networkID, newer.Unix())}}) - if puts, updates := spy.counts(); puts != 2 || updates != 0 { - t.Fatalf("newer record: puts=%d updates=%d, want 2/0", puts, updates) + if puts, _ := spy.counts(); puts != 2 { + t.Fatalf("newer record was not stored: puts=%d, want 2", puts) } } -// TestGossipDoesNotRefreshLastSeenOnInvalidRecord makes sure the refresh is -// scoped to records that merely are not newer. A record rejected for any other -// reason is not evidence that the peer is alive. -func TestGossipDoesNotRefreshLastSeenOnInvalidRecord(t *testing.T) { +// TestSeenSkippedOnInvalidRecord makes sure a record we refuse to parse is not +// taken for a sighting. Only a record carrying the peer's own signature is +// evidence that we heard about that peer at all. +func TestSeenSkippedOnInvalidRecord(t *testing.T) { t.Parallel() const networkID = uint64(1) @@ -137,11 +136,11 @@ func TestGossipDoesNotRefreshLastSeenOnInvalidRecord(t *testing.T) { id := newPeerIdentity(t, networkID, "/ip4/10.0.0.1/tcp/1634") svc.CheckAndAddPeers(pb.Peers{Peers: []*pb.BzzAddress{id.protoAt(t, networkID, base.Unix())}}) - // a record timestamped far in the future is rejected as invalid, not stale. - future := base.Add(24 * time.Hour) - svc.CheckAndAddPeers(pb.Peers{Peers: []*pb.BzzAddress{id.protoAt(t, networkID, future.Unix())}}) + // the record is signed for a different network, so it does not verify + // against ours. + svc.CheckAndAddPeers(pb.Peers{Peers: []*pb.BzzAddress{id.protoAt(t, networkID+1, base.Unix())}}) - if puts, updates := spy.counts(); puts != 1 || updates != 0 { - t.Fatalf("invalid record touched the addressbook: puts=%d updates=%d, want 1/0", puts, updates) + if puts, seens := spy.counts(); puts != 1 || seens != 0 { + t.Fatalf("invalid record touched the addressbook: puts=%d seens=%d, want 1/0", puts, seens) } } diff --git a/pkg/node/node.go b/pkg/node/node.go index 2a8069763ff..eb34733258e 100644 --- a/pkg/node/node.go +++ b/pkg/node/node.go @@ -381,14 +381,15 @@ func NewBee( return nil, fmt.Errorf("batchstore: exists: %w", err) } - addressbook := addressbook.New(stateStore) - - // Prune addressbook entries whose overlays have not been seen recently, so - // the address book does not accumulate stale peers indefinitely. - if err := addressbook.Prune(time.Now().Add(-addressbookPruneAfter)); err != nil { + // Drop addressbook entries whose overlays have not been seen recently, so + // the address book does not accumulate stale peers indefinitely. Non-fatal: + // a failed prune must never block boot. + if err := addressbook.Prune(stateStore, time.Now().Add(-addressbookPruneAfter)); err != nil { logger.Warning("addressbook prune failed", "error", err) } + addressbook := addressbook.New(stateStore) + logger.Info("using overlay address", "address", swarmAddress) // this will set overlay if it was not set before diff --git a/pkg/statestore/storeadapter/migration_test.go b/pkg/statestore/storeadapter/migration_test.go index 51182dc58b5..cff48a88afd 100644 --- a/pkg/statestore/storeadapter/migration_test.go +++ b/pkg/statestore/storeadapter/migration_test.go @@ -259,9 +259,8 @@ func TestStampAddressbookLastSeen(t *testing.T) { // TestAddressbookPruneRealStore drives addressbook.Prune over the production // storage path (leveldb behind StateStorerAdapter), whose iterator key // semantics differ from the in-memory mock used in the addressbook package's -// own tests. It confirms that overlays are correctly reconstructed from the -// iterated keys so the right entries are pruned and the survivors remain -// readable through the addressbook. +// own tests. It confirms that the right entries are pruned and the survivors +// remain readable through the addressbook. func TestAddressbookPruneRealStore(t *testing.T) { t.Parallel() @@ -298,11 +297,11 @@ func TestAddressbookPruneRealStore(t *testing.T) { seed("aabb", 1000) seed("ccdd", 9000) - book := addressbook.New(store) - if err := book.Prune(time.Unix(5000, 0)); err != nil { + if err := addressbook.Prune(store, time.Unix(5000, 0)); err != nil { t.Fatalf("prune: %v", err) } + book := addressbook.New(store) if _, _, err := book.Get(stale); !errors.Is(err, addressbook.ErrNotFound) { t.Fatalf("stale entry should have been pruned, got err=%v", err) } diff --git a/pkg/topology/kademlia/export_test.go b/pkg/topology/kademlia/export_test.go index 4d4587b7b0a..bca662c9a7b 100644 --- a/pkg/topology/kademlia/export_test.go +++ b/pkg/topology/kademlia/export_test.go @@ -18,6 +18,12 @@ var ( GenerateCommonBinPrefixes = generateCommonBinPrefixes ) +// MarkConnectedPeersSeen runs the sweep the manage loop performs on every +// lastSeenRefreshInterval tick. +func (k *Kad) MarkConnectedPeersSeen() error { + return k.markConnectedPeersSeen() +} + const ( DefaultBitSuffixLength = defaultBitSuffixLength DefaultSaturationPeers = defaultSaturationPeers diff --git a/pkg/topology/kademlia/kademlia.go b/pkg/topology/kademlia/kademlia.go index c5e9736397a..ca2ca4db610 100644 --- a/pkg/topology/kademlia/kademlia.go +++ b/pkg/topology/kademlia/kademlia.go @@ -45,6 +45,12 @@ const ( // Each underlay address gets up to 15s for connection (in libp2p.Connect). // This budget allows multiple addresses to be tried sequentially per peer. peerConnectionAttemptTimeout = 45 * time.Second // timeout for establishing a new connection with peer. + + // lastSeenRefreshInterval is how often the peers we are connected to are + // marked as seen in the addressbook. A peer we hold a connection to is + // seen continuously, so marking it on the connect event alone would let + // the addressbook pruner evict our longest-lived, most valuable peers. + lastSeenRefreshInterval = 15 * time.Minute ) // Default option values @@ -455,10 +461,6 @@ func (k *Kad) connectionAttemptsHandler(ctx context.Context, wg *sync.WaitGroup, k.connectedPeers.Add(peer.addr) - if err := k.addressBook.UpdateLastSeen(peer.addr); err != nil { - k.logger.Warning("could not update last seen for peer", "peer_address", peer.addr, "error", err) - } - k.metrics.TotalOutboundConnections.Inc() k.collector.Record(peer.addr, im.PeerLogIn(time.Now(), im.PeerConnectionDirectionOutbound)) @@ -517,6 +519,22 @@ func (k *Kad) notifyManageLoop() { } } +// markConnectedPeersSeen marks every currently connected peer as seen in the +// addressbook. +func (k *Kad) markConnectedPeersSeen() error { + var peers []swarm.Address + _ = k.connectedPeers.EachBin(func(addr swarm.Address, _ uint8) (bool, bool, error) { + peers = append(peers, addr) + return false, false, nil + }) + + if len(peers) == 0 { + return nil + } + + return k.addressBook.Seen(peers...) +} + // manage is a forever loop that manages the connection to new peers // once they get added or once others leave. func (k *Kad) manage() { @@ -573,6 +591,21 @@ func (k *Kad) manage() { } }) + k.wg.Go(func() { + for { + select { + case <-k.halt: + return + case <-k.quit: + return + case <-time.After(lastSeenRefreshInterval): + if err := k.markConnectedPeersSeen(); err != nil { + k.logger.Warning("could not mark connected peers as seen", "error", err) + } + } + } + }) + // tell each neighbor about other neighbors periodically k.wg.Go(func() { for { @@ -1218,9 +1251,6 @@ func (k *Kad) onConnected(ctx context.Context, addr swarm.Address) error { k.knownPeers.Add(addr) k.connectedPeers.Add(addr) k.waitNext.Remove(addr) - if err := k.addressBook.UpdateLastSeen(addr); err != nil { - k.logger.Warning("could not update last seen for peer", "peer_address", addr, "error", err) - } k.recalcDepth() k.notifyManageLoop() k.notifyPeerSig() diff --git a/pkg/topology/kademlia/lastseen_test.go b/pkg/topology/kademlia/lastseen_test.go index 3ceb6b2daf2..e87e3f6d1f6 100644 --- a/pkg/topology/kademlia/lastseen_test.go +++ b/pkg/topology/kademlia/lastseen_test.go @@ -27,11 +27,13 @@ type spyBook struct { seen map[string]int } -func (s *spyBook) UpdateLastSeen(o swarm.Address) error { +func (s *spyBook) Seen(overlays ...swarm.Address) error { s.mu.Lock() - s.seen[o.String()]++ + for _, o := range overlays { + s.seen[o.String()]++ + } s.mu.Unlock() - return s.Interface.UpdateLastSeen(o) + return s.Interface.Seen(overlays...) } func (s *spyBook) count(o swarm.Address) int { @@ -40,7 +42,11 @@ func (s *spyBook) count(o swarm.Address) int { return s.seen[o.String()] } -func TestKademliaBumpsLastSeenOnConnect(t *testing.T) { +// TestMarkConnectedPeersSeen covers the sweep the manage loop runs on every +// lastSeenRefreshInterval tick. A peer we hold a connection to is seen +// continuously, so marking it on the connect event alone would let the +// addressbook pruner evict our longest-lived peers. +func TestMarkConnectedPeersSeen(t *testing.T) { t.Parallel() detector, err := stabilization.NewDetector(stabilization.Config{ @@ -77,18 +83,14 @@ func TestKademliaBumpsLastSeenOnConnect(t *testing.T) { testutil.CleanupCloser(t, kad) kad.SetStorageRadius(0) - // Inbound path: Connected -> onConnected -> UpdateLastSeen. - inbound := swarm.RandAddress(t) - connectOne(t, signer, kad, spy, inbound, nil) - if got := spy.count(inbound); got == 0 { - t.Fatalf("inbound connect did not bump last-seen for %s", inbound) + connected := swarm.RandAddress(t) + connectOne(t, signer, kad, spy, connected, nil) + + if err := kad.MarkConnectedPeersSeen(); err != nil { + t.Fatal(err) } - // Outbound path: manage loop dials -> connect closure -> UpdateLastSeen. - outbound := swarm.RandAddressAt(t, base, 0) - addOne(t, signer, kad, spy, outbound) - waitConn(t, &conns) - if got := spy.count(outbound); got == 0 { - t.Fatalf("outbound connect did not bump last-seen for %s", outbound) + if got := spy.count(connected); got != 1 { + t.Fatalf("connected peer marked seen %d times, want 1", got) } } From 59b94af94939b762c2af8c1c2a910084fecb8542 Mon Sep 17 00:00:00 2001 From: Calin Martinconi Date: Fri, 17 Jul 2026 12:47:01 +0300 Subject: [PATCH 5/6] refactor(addressbook): prune when the store is opened --- pkg/addressbook/addressbook.go | 39 +++++++++++++------ pkg/addressbook/addressbook_test.go | 31 +++++++-------- pkg/addressbook/export_test.go | 3 ++ pkg/node/node.go | 8 ---- pkg/statestore/storeadapter/migration_test.go | 11 +++--- 5 files changed, 51 insertions(+), 41 deletions(-) diff --git a/pkg/addressbook/addressbook.go b/pkg/addressbook/addressbook.go index 59e08ae89a8..5ee25a9e55b 100644 --- a/pkg/addressbook/addressbook.go +++ b/pkg/addressbook/addressbook.go @@ -17,7 +17,13 @@ import ( "github.com/ethersphere/bee/v2/pkg/swarm" ) -const keyPrefix = "addressbook_entry_" +const ( + keyPrefix = "addressbook_entry_" + + // pruneAfter is how long an overlay may go unseen before its entry is + // dropped when the addressbook is opened. + pruneAfter = 30 * 24 * time.Hour +) var _ Interface = (*store)(nil) @@ -94,10 +100,18 @@ func New(storer storage.StateStorer) Interface { } func newStore(storer storage.StateStorer, now func() time.Time) *store { - return &store{ + s := &store{ store: storer, now: now, } + + // Drop entries whose overlays have not been seen recently, so the address + // book does not accumulate stale peers indefinitely. Best-effort: this is + // garbage collection, and failing it only leaves the stale entries in place + // for another run, which must not stop a node from starting. + _ = s.prune(s.now().Add(-pruneAfter)) + + return s } func (s *store) Get(overlay swarm.Address) (*bzz.Address, bool, error) { @@ -157,19 +171,22 @@ func (s *store) Remove(overlay swarm.Address) error { return s.store.Delete(keyPrefix + overlay.String()) } -// Prune removes all entries from storer whose overlay has not been seen since -// before. Entries without a recorded last-seen time (LastSeen == 0) are kept, -// leaving them to a later run once they have been observed. It is a one-shot -// maintenance step meant to run at startup, before the address book is shared -// with its writers, so it takes no lock. -func Prune(storer storage.StateStorer, before time.Time) error { +// prune removes all entries whose overlay has not been seen since before. +// Entries without a recorded last-seen time (LastSeen == 0) are kept, leaving +// them to a later run once they have been observed, as are entries that cannot +// be unmarshaled: a single unreadable record must not cost us the sweep. +// +// It runs from newStore, before the store is shared with its writers, so it +// takes no lock. +func (s *store) prune(before time.Time) error { cutoff := before.Unix() var stale []string - err := storer.Iterate(keyPrefix, func(key, value []byte) (stop bool, err error) { + err := s.store.Iterate(keyPrefix, func(key, value []byte) (stop bool, err error) { entry := &verifiedAddress{} if err := json.Unmarshal(value, entry); err != nil { - return true, err + //nolint:nilerr // an unreadable record is skipped, not fatal: it must not cost us the rest of the sweep + return false, nil } if entry.LastSeen != 0 && entry.LastSeen < cutoff { stale = append(stale, string(key)) @@ -181,7 +198,7 @@ func Prune(storer storage.StateStorer, before time.Time) error { } for _, key := range stale { - if err := storer.Delete(key); err != nil { + if err := s.store.Delete(key); err != nil { return err } } diff --git a/pkg/addressbook/addressbook_test.go b/pkg/addressbook/addressbook_test.go index 9be9ef4db8b..69ef1c8caaa 100644 --- a/pkg/addressbook/addressbook_test.go +++ b/pkg/addressbook/addressbook_test.go @@ -152,10 +152,10 @@ func TestSeen(t *testing.T) { t.Fatalf("seen did not move last seen: got %d, want %d", got, now.Unix()) } - if err := addressbook.Prune(state, now.Add(-time.Hour)); err != nil { - t.Fatal(err) - } - if _, verified, err := book.Get(overlay); err != nil || !verified { + // reopening the book prunes whatever has gone unseen for PruneAfter; the + // sighting above is what keeps this entry. + reopened := addressbook.NewWithClock(state, func() time.Time { return now }) + if _, verified, err := reopened.Get(overlay); err != nil || !verified { t.Fatalf("entry pruned despite a recent sighting: verified=%v err=%v", verified, err) } } @@ -269,7 +269,8 @@ func lastSeenOf(t *testing.T, state storage.StateStorer, overlay swarm.Address) func TestPrune(t *testing.T) { t.Parallel() - now := time.Unix(1_000_000_000, 0) + base := time.Unix(1_000_000_000, 0) + now := base state := mock.NewStateStore() book := addressbook.NewWithClock(state, func() time.Time { return now }) @@ -284,15 +285,15 @@ func TestPrune(t *testing.T) { t.Fatal(err) } - // cutoff sits between the two puts: stale is before it, fresh is after. - if err := addressbook.Prune(state, now.Add(-time.Hour)); err != nil { - t.Fatal(err) - } + // Reopen with a clock that puts the cutoff between the two puts: the first + // overlay has gone unseen for longer than PruneAfter, the second has not. + reopenAt := base.Add(addressbook.PruneAfter + time.Hour) + reopened := addressbook.NewWithClock(state, func() time.Time { return reopenAt }) - if _, _, err := book.Get(stale); !errors.Is(err, addressbook.ErrNotFound) { + if _, _, err := reopened.Get(stale); !errors.Is(err, addressbook.ErrNotFound) { t.Fatalf("stale entry should have been pruned, got err=%v", err) } - if _, _, err := book.Get(fresh); err != nil { + if _, _, err := reopened.Get(fresh); err != nil { t.Fatalf("fresh entry should survive: %v", err) } } @@ -313,11 +314,9 @@ func TestPruneKeepsEntriesWithoutLastSeen(t *testing.T) { t.Fatal(err) } - if err := addressbook.Prune(state, time.Unix(5_000_000_000, 0)); err != nil { - t.Fatal(err) - } - - book := addressbook.New(state) + // open the book far past any plausible cutoff: the entry still survives, + // because it carries no last-seen time to judge it by. + book := addressbook.NewWithClock(state, func() time.Time { return time.Unix(5_000_000_000, 0) }) if _, _, err := book.Get(overlay); err != nil { t.Fatalf("entry without a last-seen time must not be pruned: %v", err) } diff --git a/pkg/addressbook/export_test.go b/pkg/addressbook/export_test.go index f26d03bd661..f57faa3a93c 100644 --- a/pkg/addressbook/export_test.go +++ b/pkg/addressbook/export_test.go @@ -12,6 +12,9 @@ import ( type VerifiedAddress = verifiedAddress +// PruneAfter is how long an overlay may go unseen before newStore drops it. +const PruneAfter = pruneAfter + // NewWithClock creates an addressbook with an overridable clock, for testing. func NewWithClock(storer storage.StateStorer, now func() time.Time) Interface { return newStore(storer, now) diff --git a/pkg/node/node.go b/pkg/node/node.go index eb34733258e..c6a1dc54981 100644 --- a/pkg/node/node.go +++ b/pkg/node/node.go @@ -216,7 +216,6 @@ const ( reserveMinEvictCount = 1_000 cacheMinEvictCount = 10_000 maxAllowedDoubling = 1 - addressbookPruneAfter = 30 * 24 * time.Hour // remove addressbook entries not seen for this long ) func NewBee( @@ -381,13 +380,6 @@ func NewBee( return nil, fmt.Errorf("batchstore: exists: %w", err) } - // Drop addressbook entries whose overlays have not been seen recently, so - // the address book does not accumulate stale peers indefinitely. Non-fatal: - // a failed prune must never block boot. - if err := addressbook.Prune(stateStore, time.Now().Add(-addressbookPruneAfter)); err != nil { - logger.Warning("addressbook prune failed", "error", err) - } - addressbook := addressbook.New(stateStore) logger.Info("using overlay address", "address", swarmAddress) diff --git a/pkg/statestore/storeadapter/migration_test.go b/pkg/statestore/storeadapter/migration_test.go index cff48a88afd..42e0adc5687 100644 --- a/pkg/statestore/storeadapter/migration_test.go +++ b/pkg/statestore/storeadapter/migration_test.go @@ -292,14 +292,13 @@ func TestAddressbookPruneRealStore(t *testing.T) { } } + // the addressbook prunes when it is opened, against its own clock, which + // this package cannot override — so seed relative to the wall clock. + now := time.Now() stale := swarm.MustParseHexAddress("aabb") fresh := swarm.MustParseHexAddress("ccdd") - seed("aabb", 1000) - seed("ccdd", 9000) - - if err := addressbook.Prune(store, time.Unix(5000, 0)); err != nil { - t.Fatalf("prune: %v", err) - } + seed("aabb", now.Add(-90*24*time.Hour).Unix()) + seed("ccdd", now.Add(-24*time.Hour).Unix()) book := addressbook.New(store) if _, _, err := book.Get(stale); !errors.Is(err, addressbook.ErrNotFound) { From b3ba2f7b35351e956bec66e63521fb36c85f2c97 Mon Sep 17 00:00:00 2001 From: Calin Martinconi Date: Thu, 23 Jul 2026 15:41:33 +0300 Subject: [PATCH 6/6] perf(addressbook): throttle last-seen writes to once per 24h --- pkg/addressbook/addressbook.go | 16 ++++++++++++++-- pkg/addressbook/addressbook_test.go | 18 ++++++++++++++++-- 2 files changed, 30 insertions(+), 4 deletions(-) diff --git a/pkg/addressbook/addressbook.go b/pkg/addressbook/addressbook.go index 52a001a27ee..c2b50d635a9 100644 --- a/pkg/addressbook/addressbook.go +++ b/pkg/addressbook/addressbook.go @@ -23,6 +23,12 @@ const ( // pruneAfter is how long an overlay may go unseen before its entry is // dropped when the addressbook is opened. pruneAfter = 30 * 24 * time.Hour + + // seenInterval throttles Seen's disk writes: a sighting of an overlay + // already seen within this window is not written back. Sightings are + // frequent (every hive gossip and kademlia manage tick) and almost always + // redundant, given that pruning acts on a much coarser scale. + seenInterval = 24 * time.Hour ) var _ Interface = (*store)(nil) @@ -81,7 +87,8 @@ type Remover interface { } type Seener interface { - // Seen marks the overlays as seen at the current time. + // Seen marks the overlays as seen at the current time. Writes are + // throttled: an overlay already marked seen recently is left untouched. Seen(overlays ...swarm.Address) error } @@ -144,7 +151,8 @@ func (s *store) Put(overlay swarm.Address, addr bzz.Address, verified bool) (err } // Seen marks the overlays as seen at the current time. An overlay that is not -// present in the addressbook is skipped. +// present in the addressbook is skipped, as is one already seen within +// seenInterval, to keep redundant sightings off the disk. func (s *store) Seen(overlays ...swarm.Address) error { s.mu.Lock() defer s.mu.Unlock() @@ -162,6 +170,10 @@ func (s *store) Seen(overlays ...swarm.Address) error { return err } + if now-v.LastSeen < int64(seenInterval/time.Second) { + continue + } + v.LastSeen = now if err := s.store.Put(key, v); err != nil { return err diff --git a/pkg/addressbook/addressbook_test.go b/pkg/addressbook/addressbook_test.go index 3951c24354a..270d7a506c8 100644 --- a/pkg/addressbook/addressbook_test.go +++ b/pkg/addressbook/addressbook_test.go @@ -120,7 +120,8 @@ func run(t *testing.T, f bookFunc) { // TestSeen covers a sighting of a peer we already hold: the last-seen time // moves, and the entry then survives a prune that would otherwise catch it. An -// overlay we do not know is skipped rather than created. +// overlay we do not know is skipped rather than created, and a sighting soon +// after the last one is throttled rather than written. func TestSeen(t *testing.T) { t.Parallel() @@ -143,6 +144,16 @@ func TestSeen(t *testing.T) { t.Fatal(err) } + // a sighting within the throttle window is not written back; last-seen + // stays at the time of the Put. + now = base.Add(time.Hour) + if err := book.Seen(overlay); err != nil { + t.Fatal(err) + } + if got := lastSeenOf(t, state, overlay); got != base.Unix() { + t.Fatalf("throttled sighting moved last seen: got %d, want %d", got, base.Unix()) + } + // a much later sighting moves last-seen forward, so a prune whose cutoff // predates it keeps the entry. now = base.Add(90 * 24 * time.Hour) @@ -180,7 +191,7 @@ func TestSeenVariadic(t *testing.T) { } } - now = now.Add(time.Hour) + now = now.Add(25 * time.Hour) if err := book.Seen(overlays...); err != nil { t.Fatal(err) } @@ -211,6 +222,9 @@ func TestSeenKeepsConcurrentPut(t *testing.T) { t.Fatal(err) } + // move past the throttle window, so that Seen takes its write path. + now = now.Add(25 * time.Hour) + // While Seen holds the entry it has just read, hive verifies the same peer // and stores it with Verified=true. started, finished := make(chan struct{}), make(chan struct{})