Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions .github/workflows/go_build.yml
Original file line number Diff line number Diff line change
Expand Up @@ -47,3 +47,6 @@ jobs:

- name: Sealevel interpreter, syscall and vote ownership regressions
run: go test -race -count=1 ./pkg/sealevel

- name: Gossip and networking regressions
run: go test -race -p 2 -count=1 ./pkg/gossip ./pkg/repair ./pkg/blockstream
46 changes: 34 additions & 12 deletions pkg/gossip/client.go
Original file line number Diff line number Diff line change
Expand Up @@ -54,6 +54,7 @@ type Client struct {

contactMu sync.RWMutex
contact *ContactInfo
relay contactRelay

identityConflictMu sync.Mutex
identityConflictGossip contactEndpoint
Expand Down Expand Up @@ -475,13 +476,23 @@ func (c *Client) pushContact(conn *net.UDPConn) error {
if err != nil {
return err
}
for _, peer := range c.currentPeers() {
if err := sendUDP(conn, packet, peer); err != nil {
c.txErrors.Add(1)
packets := [][]byte{packet}
if values := c.relay.values(wallclockMillis()); len(values) != 0 {
packet, err := encodePushMessage(c.pubkey, values)
if err != nil {
return err
}
c.recordTx()
c.txPushMessages.Add(1)
packets = append(packets, packet)
}
for _, peer := range c.currentPeers() {
for _, packet := range packets {
if err := sendUDP(conn, packet, peer); err != nil {
c.txErrors.Add(1)
break
}
c.recordTx()
c.txPushMessages.Add(1)
}
}
return nil
}
Expand Down Expand Up @@ -626,6 +637,10 @@ func (c *Client) handleContactRecord(record contactRecord, shredVersion uint16)
if record.ShredVer != shredVersion || !record.GossipAddr.ok {
return
}
addr := record.GossipAddr.UDPAddr()
if addr == nil || addr.Port == 0 || addr.IP.IsUnspecified() || addr.IP.IsMulticast() {
return
}
// CRDS is keyed by validator identity. If another process publishes that
// same identity with different sockets, its newer ContactInfo replaces ours
// cluster-wide: turbine and Votor traffic then move to the other process in
Expand All @@ -637,6 +652,11 @@ func (c *Client) handleContactRecord(record contactRecord, shredVersion uint16)
c.observeOwnContactRecord(record)
return
}
// The decoder verified the original signed bytes. Retain those bytes for
// relay, and do not let an older relayed contact replace newer endpoints.
if !c.relay.accept(record, wallclockMillis()) {
return
}
// The entrypoint is commonly a validator itself. Do not add its gossip
// socket to the rotating pull-peer table, but retain the service endpoints
// from its signed ContactInfo so Turbine, repair, and Votor can route to it.
Expand Down Expand Up @@ -742,24 +762,24 @@ func (c *Client) recordRepairPeer(contact *ContactInfo) {
}

func (c *Client) recordRepairPeerRecord(record contactRecord) {
c.repairPeerMu.Lock()
defer c.repairPeerMu.Unlock()
if !record.ServeRepairAddr.ok || record.ServeRepairAddr.port == 0 {
delete(c.repairPeers, record.Pubkey)
return
}
now := time.Now()
key := record.Pubkey
c.repairPeerMu.Lock()
if existing, ok := c.repairPeers[key]; ok && sameEndpointUDPAddr(record.ServeRepairAddr, existing.Addr) {
existing.LastSeen = now
c.repairPeers[key] = existing
c.repairPeerMu.Unlock()
return
}
c.repairPeers[key] = RepairPeer{
Pubkey: record.Pubkey,
Addr: record.ServeRepairAddr.UDPAddr(),
LastSeen: now,
}
c.repairPeerMu.Unlock()
}

// TVUPeers returns non-expired TVU endpoints learned from gossip.
Expand Down Expand Up @@ -841,30 +861,32 @@ func (c *Client) LookupAlpenglow(pubkey solana.PublicKey) (*net.UDPAddr, bool) {

func (c *Client) recordTVUPeerRecord(record contactRecord) {
addr := record.TVUAddr.UDPAddr()
c.tvuPeerMu.Lock()
defer c.tvuPeerMu.Unlock()
if addr == nil {
delete(c.tvuPeers, record.Pubkey)
return
}
c.tvuPeerMu.Lock()
c.tvuPeers[record.Pubkey] = TVUPeer{
Pubkey: record.Pubkey,
TVUAddr: addr,
LastSeen: time.Now(),
}
c.tvuPeerMu.Unlock()
}

func (c *Client) recordAlpenglowPeerRecord(record contactRecord) {
addr := record.Sockets[socketTagAlpenglow].UDPAddr()
c.alpenglowPeerMu.Lock()
defer c.alpenglowPeerMu.Unlock()
if addr == nil {
delete(c.alpenglowPeers, record.Pubkey)
return
}
c.alpenglowPeerMu.Lock()
c.alpenglowPeers[record.Pubkey] = AlpenglowPeer{
Pubkey: record.Pubkey,
AlpenglowAddr: addr,
LastSeen: time.Now(),
}
c.alpenglowPeerMu.Unlock()
}

func (c *Client) recordTx() {
Expand Down
3 changes: 2 additions & 1 deletion pkg/gossip/contact_info.go
Original file line number Diff line number Diff line change
Expand Up @@ -87,6 +87,7 @@ type contactEndpoint struct {
type contactRecord struct {
Pubkey Pubkey
Wallclock uint64
Outset uint64
ShredVer uint16
GossipAddr contactEndpoint
ServeRepairAddr contactEndpoint
Expand Down Expand Up @@ -370,7 +371,7 @@ func decodeContactRecord(d *decoder) (contactRecord, error) {
if record.Wallclock, err = d.varint(10); err != nil {
return contactRecord{}, err
}
if _, err := d.u64(); err != nil {
if record.Outset, err = d.u64(); err != nil {
return contactRecord{}, err
}
if record.ShredVer, err = d.u16(); err != nil {
Expand Down
116 changes: 116 additions & 0 deletions pkg/gossip/contact_relay.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,116 @@
package gossip

import (
"bytes"
"crypto/sha256"
"sort"
"sync"
"time"
)

// Agave accepts pushed values within a 15-second wallclock window.
const contactPushWindow = uint64(15 * time.Second / time.Millisecond)

const crdsMessageHeaderSize = 4 + 32 + 8

type relayContact struct {
record contactRecord
hash [32]byte
forwarded uint64
}

// contactRelay retains at most the existing gossip peer limit. Each push tick
// relays at most one additional datagram per peer, rotating through contacts.
type contactRelay struct {
mu sync.Mutex
contacts map[Pubkey]relayContact
sequence uint64
}

func recentContact(wallclock, now, maxAge uint64) bool {
if wallclock > now {
return wallclock-now < contactPushWindow
}
return now-wallclock <= maxAge
}

// accept is called only after signature and shred-version verification.
func (r *contactRelay) accept(record contactRecord, now uint64) bool {
if !recentContact(record.Wallclock, now, uint64(peerExpirationWindow/time.Millisecond)) {
return false
}
if len(record.data)+len(record.signature)+crdsMessageHeaderSize > packetDataSize {
return false
}
h := sha256.New()
h.Write(record.signature[:])
h.Write(record.data)
var hash [32]byte
copy(hash[:], h.Sum(nil))
r.mu.Lock()
defer r.mu.Unlock()
old, exists := r.contacts[record.Pubkey]
if exists {
// ContactInfo orders restarts first, then wallclock, then the hash of
// the complete signed value, matching Agave's CRDS replacement rule.
if record.Outset < old.record.Outset ||
(record.Outset == old.record.Outset && record.Wallclock < old.record.Wallclock) ||
(record.Outset == old.record.Outset && record.Wallclock == old.record.Wallclock && bytes.Compare(hash[:], old.hash[:]) <= 0) {
return false
}
} else if len(r.contacts) >= maxKnownGossipPeers {
var oldest Pubkey
oldestWallclock := ^uint64(0)
for key, contact := range r.contacts {
if contact.record.Wallclock < oldestWallclock {
oldest, oldestWallclock = key, contact.record.Wallclock
}
}
delete(r.contacts, oldest)
}
if r.contacts == nil {
r.contacts = make(map[Pubkey]relayContact)
}
// The receive loop reuses its datagram buffer on the next packet.
record.data = bytes.Clone(record.data)
r.contacts[record.Pubkey] = relayContact{record: record, hash: hash, forwarded: old.forwarded}
return true
}

func (r *contactRelay) values(now uint64) []CrdsValue {
r.mu.Lock()
defer r.mu.Unlock()
var keys []Pubkey
for key, contact := range r.contacts {
if !recentContact(contact.record.Wallclock, now, uint64(peerExpirationWindow/time.Millisecond)) {
delete(r.contacts, key)
continue
}
if recentContact(contact.record.Wallclock, now, contactPushWindow) {
keys = append(keys, key)
}
}
sort.Slice(keys, func(i, j int) bool {
a, b := r.contacts[keys[i]], r.contacts[keys[j]]
if a.forwarded != b.forwarded {
return a.forwarded < b.forwarded
}
return bytes.Compare(keys[i][:], keys[j][:]) < 0
})
var values []CrdsValue
size := crdsMessageHeaderSize
for _, key := range keys {
contact := r.contacts[key]
record := contact.record
n := len(record.signature) + len(record.data)
if size+n > packetDataSize {
break
}
values = append(values, CrdsValue{Signature: record.signature, Data: record.data})
size += n
r.sequence++
contact.forwarded = r.sequence
r.contacts[key] = contact
}
return values
}
Loading
Loading