diff --git a/pkg/addressbook/addressbook.go b/pkg/addressbook/addressbook.go index dd5bcd2d0b0..c2b50d635a9 100644 --- a/pkg/addressbook/addressbook.go +++ b/pkg/addressbook/addressbook.go @@ -9,29 +9,47 @@ import ( "errors" "fmt" "strings" + "sync" + "time" "github.com/ethersphere/bee/v2/pkg/bzz" "github.com/ethersphere/bee/v2/pkg/storage" "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 + + // 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) 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 + 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. @@ -45,6 +63,13 @@ type GetPutter interface { Putter } +// 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 + Seener +} + type Getter interface { // Get returns the saved bzz.Address for the requested overlay together // with its verification flag. @@ -61,15 +86,39 @@ type Remover interface { Remove(overlay swarm.Address) error } +type Seener interface { + // 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 +} + type store struct { store storage.StateStorer + now func() time.Time + + // 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 &store{ + return newStore(storer, time.Now) +} + +func newStore(storer storage.StateStorer, now func() time.Time) *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) { @@ -90,17 +139,89 @@ 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, Verified: verified, + LastSeen: s.now().Unix(), }) } +// Seen marks the overlays as seen at the current time. An overlay that is not +// 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() + + 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 + } + + if now-v.LastSeen < int64(seenInterval/time.Second) { + continue + } + + v.LastSeen = now + if err := s.store.Put(key, v); err != nil { + return err + } + } + + 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 +// 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 := s.store.Iterate(keyPrefix, func(key, value []byte) (stop bool, err error) { + entry := &verifiedAddress{} + if err := json.Unmarshal(value, entry); err != nil { + //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)) + } + return false, nil + }) + if err != nil { + return err + } + + for _, key := range stale { + if err := s.store.Delete(key); 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 36dac87fe51..270d7a506c8 100644 --- a/pkg/addressbook/addressbook_test.go +++ b/pkg/addressbook/addressbook_test.go @@ -8,6 +8,7 @@ import ( "encoding/json" "errors" "testing" + "time" "github.com/ethereum/go-ethereum/common" "github.com/ethersphere/bee/v2/pkg/addressbook" @@ -19,6 +20,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) { @@ -99,6 +118,227 @@ 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, and a sighting soon +// after the last one is throttled rather than written. +func TestSeen(t *testing.T) { + t.Parallel() + + 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}) + + // an unknown overlay is skipped, not created. + if err := book.Seen(overlay); err != nil { + t.Fatal(err) + } + if _, _, err := book.Get(overlay); !errors.Is(err, addressbook.ErrNotFound) { + t.Fatalf("expected ErrNotFound, got %v", err) + } + + if err := book.Put(overlay, newTestAddr(t, overlay), true); err != nil { + 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) + 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()) + } + + // 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) + } +} + +// TestSeenVariadic marks several overlays in one call, which is how kademlia +// refreshes everything it is connected to. +func TestSeenVariadic(t *testing.T) { + t.Parallel() + + now := time.Unix(1_000_000, 0) + state := mock.NewStateStore() + book := addressbook.NewWithClock(state, func() time.Time { return now }) + + overlays := []swarm.Address{ + swarm.NewAddress([]byte{0, 1, 2, 3}), + swarm.NewAddress([]byte{0, 1, 2, 4}), + } + for _, overlay := range overlays { + if err := book.Put(overlay, newTestAddr(t, overlay), true); err != nil { + t.Fatal(err) + } + } + + now = now.Add(25 * time.Hour) + if err := book.Seen(overlays...); err != nil { + t.Fatal(err) + } + + 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()) + } + } +} + +// 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 TestSeenKeepsConcurrentPut(t *testing.T) { + t.Parallel() + + 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. + if err := book.Put(overlay, addr, false); err != nil { + 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{}) + 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 Seen returns; unsynchronized, its write + // completes here and is then overwritten below. + time.Sleep(100 * time.Millisecond) + } + + if err := book.Seen(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 Seen'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, state storage.StateStorer, overlay swarm.Address) int64 { + t.Helper() + + v := &addressbook.VerifiedAddress{} + 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() + + base := time.Unix(1_000_000_000, 0) + now := base + state := mock.NewStateStore() + book := addressbook.NewWithClock(state, func() time.Time { return now }) + + stale := swarm.NewAddress([]byte{0, 1, 2, 3}) + if err := book.Put(stale, newTestAddr(t, stale), true); err != nil { + t.Fatal(err) + } + + 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) + } + + // 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 := reopened.Get(stale); !errors.Is(err, addressbook.ErrNotFound) { + t.Fatalf("stale entry should have been pruned, got err=%v", err) + } + if _, _, err := reopened.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() + + state := mock.NewStateStore() + overlay := swarm.NewAddress([]byte{0, 1, 2, 3}) + + if err := state.Put("addressbook_entry_"+overlay.String(), &addressbook.VerifiedAddress{ + Address: addrPtr(newTestAddr(t, overlay)), + Verified: true, + }); err != nil { + t.Fatal(err) + } + + // 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) + } +} + +func addrPtr(a bzz.Address) *bzz.Address { return &a } + type mockCorruptedStore struct{} func (m *mockCorruptedStore) Get(key string, i any) error { diff --git a/pkg/addressbook/export_test.go b/pkg/addressbook/export_test.go index 9db017e1f27..f57faa3a93c 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 + +// 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/hive/hive.go b/pkg/hive/hive.go index a991dbdd58b..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.GetPutter + 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.GetPutter, 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,6 +369,18 @@ func (s *Service) checkAndAddPeers(ctx context.Context, peers pb.Peers) { continue } + // 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) diff --git a/pkg/hive/lastseen_test.go b/pkg/hive/lastseen_test.go new file mode 100644 index 00000000000..2f3c3cdaba4 --- /dev/null +++ b/pkg/hive/lastseen_test.go @@ -0,0 +1,146 @@ +// 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 that a sighting +// can be told apart from a record update. +type lastSeenSpy struct { + ab.Interface + mu sync.Mutex + puts int + seens 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) Seen(o ...swarm.Address) error { + s.mu.Lock() + s.seens++ + s.mu.Unlock() + return s.Interface.Seen(o...) +} + +func (s *lastSeenSpy) counts() (puts, seens int) { + s.mu.Lock() + defer s.mu.Unlock() + return s.puts, s.seens +} + +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 +} + +// 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) + + 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: an unknown peer, stored by Put, which stamps last-seen. + svc.CheckAndAddPeers(pb.Peers{Peers: []*pb.BzzAddress{rec}}) + if puts, seens := spy.counts(); puts != 1 || seens != 0 { + t.Fatalf("first sighting: puts=%d seens=%d, want 1/0", puts, seens) + } + + // 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, seens := spy.counts() + if puts != 1 { + t.Fatalf("re-sighting stored the record again: puts=%d, want 1", puts) + } + if seens != 1 { + t.Fatalf("re-sighting did not refresh last-seen: seens=%d, want 1", seens) + } +} + +// 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) + + 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, _ := spy.counts(); puts != 2 { + t.Fatalf("newer record was not stored: puts=%d, want 2", puts) + } +} + +// 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) + + 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())}}) + + // 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, 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/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..a0d0341110c 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,omitempty"` } // rewriteAddressbookEnvelope wraps each "addressbook_entry_*" legacy @@ -118,3 +121,52 @@ 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 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} + + 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 entry migratedEntry + if err := json.Unmarshal(e.val, &entry); err != nil { + _ = store.Delete(e.key) + continue + } + + if entry.LastSeen != 0 { + continue + } + entry.LastSeen = now + + if err := store.Put(e.key, &entry); 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..42e0adc5687 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,169 @@ 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 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) + } + } + + // 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", 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) { + 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/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 5f792c38e85..d66b669b786 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 @@ -515,6 +521,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() { @@ -571,6 +593,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 { diff --git a/pkg/topology/kademlia/lastseen_test.go b/pkg/topology/kademlia/lastseen_test.go new file mode 100644 index 00000000000..e87e3f6d1f6 --- /dev/null +++ b/pkg/topology/kademlia/lastseen_test.go @@ -0,0 +1,96 @@ +// 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) Seen(overlays ...swarm.Address) error { + s.mu.Lock() + for _, o := range overlays { + s.seen[o.String()]++ + } + s.mu.Unlock() + return s.Interface.Seen(overlays...) +} + +func (s *spyBook) count(o swarm.Address) int { + s.mu.Lock() + defer s.mu.Unlock() + return s.seen[o.String()] +} + +// 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{ + 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) + + connected := swarm.RandAddress(t) + connectOne(t, signer, kad, spy, connected, nil) + + if err := kad.MarkConnectedPeersSeen(); err != nil { + t.Fatal(err) + } + + if got := spy.count(connected); got != 1 { + t.Fatalf("connected peer marked seen %d times, want 1", got) + } +}