Skip to content
Merged
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
62 changes: 61 additions & 1 deletion internal/hub/hub.go
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@ import (
"log/slog"
"slices"
"sync/atomic"
"time"
)

// EventType identifies the kind of server-push event. These match the
Expand Down Expand Up @@ -54,6 +55,8 @@ type Event struct {
IATA string
PayloadType uint8
ChannelHash string // hex string, non-empty only for channelMessage events
// Repeat marks a later hearing of a stored observation; only IncludeRepeats clients get it.
Repeat bool
}

// Scope mirrors the client-side subscribe message. All fields are optional:
Expand Down Expand Up @@ -84,12 +87,15 @@ type Client struct {
ResolvePath bool
// IncludeObserverKey works like ResolvePath, for observerPublicKey.
IncludeObserverKey bool
// IncludeRepeats works like ResolvePath, for Repeat events.
IncludeRepeats bool
}

// ClientOptions are the connection-wide settings a "configure" message sets.
type ClientOptions struct {
ResolvePath bool
IncludeObserverKey bool
IncludeRepeats bool
}

// payloadFor falls back to a narrower variant if the event lacks one.
Expand Down Expand Up @@ -153,6 +159,12 @@ type Hub struct {

// observerKeyClients counts registered clients with IncludeObserverKey set.
observerKeyClients atomic.Int64
// repeatClients counts registered clients with IncludeRepeats set.
repeatClients atomic.Int64
// repeatDrops counts repeats refused because the broadcast channel was busy.
repeatDrops atomic.Int64
repeatDropLogAt atomic.Int64 // unix nanos of the last drop log line
sent *sentPaths
}

// subscribeMsg carries a client registration, a scope subscription, or a
Expand Down Expand Up @@ -185,6 +197,7 @@ func New() *Hub {
unsubscribe: make(chan unsubscribeMsg, 64),
remove: make(chan *Client, 64),
broadcast: make(chan Event, 512),
sent: newSentPaths(sentPathsTTL, sentPathsMax),
}
}

Expand Down Expand Up @@ -226,6 +239,17 @@ func (h *Hub) ObserverKeyWanted() bool {
return h.observerKeyClients.Load() > 0
}

// RepeatsWanted lets ingest skip repeat work when nobody wants it.
func (h *Hub) RepeatsWanted() bool {
return h.repeatClients.Load() > 0
}

// MarkSent records a hearing's path and reports whether it was not already sent recently,
// so broker copies and same-path duplicates go out once.
func (h *Hub) MarkSent(packetHash, observerID, path []byte) bool {
return h.sent.mark(packetHash, observerID, path, time.Now())
}

// Remove deregisters a client and closes its Send channel.
// Safe to call from any goroutine (e.g. the WS handler's defer).
func (h *Hub) Remove(c *Client) {
Expand All @@ -241,6 +265,31 @@ func (h *Hub) Broadcast(e Event) {
}
}

// BroadcastRepeat enqueues a repeat event, dropping it once the broadcast channel is half
// full so repeats never crowd out first hearings.
func (h *Hub) BroadcastRepeat(e Event) {
e.Repeat = true
if len(h.broadcast) >= cap(h.broadcast)/2 {
h.dropRepeat()
return
}
select {
case h.broadcast <- e:
default:
h.dropRepeat()
}
}

// dropRepeat counts a dropped repeat and logs the running total at most once a minute.
func (h *Hub) dropRepeat() {
total := h.repeatDrops.Add(1)
now := time.Now().UnixNano()
last := h.repeatDropLogAt.Load()
if now-last >= int64(time.Minute) && h.repeatDropLogAt.CompareAndSwap(last, now) {
slog.Warn("hub: broadcast channel busy, dropping repeat events", "component", "hub", "dropped_total", total)
}
}

// Run is the hub's single-goroutine event loop. Call it in a dedicated
// goroutine: go hub.Run().
//
Expand Down Expand Up @@ -270,8 +319,16 @@ func (h *Hub) Run() {
h.observerKeyClients.Add(-1)
}
}
if msg.client.IncludeRepeats != msg.options.IncludeRepeats {
if msg.options.IncludeRepeats {
h.repeatClients.Add(1)
} else {
h.repeatClients.Add(-1)
}
}
msg.client.ResolvePath = msg.options.ResolvePath
msg.client.IncludeObserverKey = msg.options.IncludeObserverKey
msg.client.IncludeRepeats = msg.options.IncludeRepeats
}
default:
// AddScope path — client must already be registered.
Expand All @@ -291,13 +348,16 @@ func (h *Hub) Run() {
if c.IncludeObserverKey {
h.observerKeyClients.Add(-1)
}
if c.IncludeRepeats {
h.repeatClients.Add(-1)
}
close(c.Send)
close(c.laggedCH)
}

case evt := <-h.broadcast:
for c := range clients {
if !c.matches(evt) {
if (evt.Repeat && !c.IncludeRepeats) || !c.matches(evt) {
continue
}
outEvt := evt
Expand Down
109 changes: 109 additions & 0 deletions internal/hub/hub_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -434,3 +434,112 @@ func TestHub_ObserverKeyWanted_TracksOptIns(t *testing.T) {
h.Remove(b)
waitFor(false)
}

func TestHub_Repeat_OnlyReachesOptedInClients(t *testing.T) {
h := runHub(t)
plain, opted := h.NewClient(), h.NewClient()
h.AddScope(plain, "all", Scope{})
h.AddScope(opted, "all", Scope{})
h.Configure(opted, ClientOptions{IncludeRepeats: true})
time.Sleep(10 * time.Millisecond)

read := func(c *Client) (Event, bool) {
select {
case evt := <-c.Send:
return evt, true
case <-time.After(50 * time.Millisecond):
return Event{}, false
}
}
h.BroadcastRepeat(Event{Type: EventPacketObservation})
if evt, ok := read(opted); !ok || !evt.Repeat {
t.Fatalf("opted-in client: got %+v, %t", evt, ok)
}
if evt, ok := read(plain); ok {
t.Fatalf("plain client got repeat %+v", evt)
}
h.Broadcast(Event{Type: EventPacketObservation})
for _, c := range []*Client{plain, opted} {
if evt, ok := read(c); !ok || evt.Repeat {
t.Fatalf("normal event: got %+v, %t", evt, ok)
}
}
}

func TestHub_RepeatsWanted_TracksOptIns(t *testing.T) {
h := runHub(t)
waitFor := func(want bool) {
t.Helper()
deadline := time.Now().Add(time.Second)
for h.RepeatsWanted() != want {
if time.Now().After(deadline) {
t.Fatalf("RepeatsWanted = %t, want %t", !want, want)
}
time.Sleep(time.Millisecond)
}
}
a, b := h.NewClient(), h.NewClient()
if h.RepeatsWanted() {
t.Fatal("no client opted in yet")
}
h.Configure(a, ClientOptions{IncludeRepeats: true})
h.Configure(a, ClientOptions{IncludeRepeats: true}) // repeat must not double count
h.Configure(b, ClientOptions{IncludeRepeats: true})
waitFor(true)
h.Configure(a, ClientOptions{ResolvePath: true})
h.Remove(b)
waitFor(false)
}

func TestHub_BroadcastRepeat_DropsFirstWhenBusy(t *testing.T) {
h := New() // not running, so the broadcast channel only fills
h.BroadcastRepeat(Event{Type: EventPacketObservation})
if len(h.broadcast) != 1 || h.repeatDrops.Load() != 0 {
t.Fatalf("repeat on an idle hub: queued %d, dropped %d", len(h.broadcast), h.repeatDrops.Load())
}
for len(h.broadcast) < cap(h.broadcast)/2 {
h.Broadcast(Event{Type: EventPacketObservation})
}
queued := len(h.broadcast)
h.BroadcastRepeat(Event{Type: EventPacketObservation})
h.BroadcastRepeat(Event{Type: EventPacketObservation})
if len(h.broadcast) != queued || h.repeatDrops.Load() != 2 {
t.Fatalf("repeats past half full: queued %d (want %d), dropped %d", len(h.broadcast), queued, h.repeatDrops.Load())
}
h.Broadcast(Event{Type: EventPacketObservation})
if len(h.broadcast) != queued+1 {
t.Fatal("normal event not enqueued past half full")
}
}

func TestSentPaths(t *testing.T) {
start := time.Now()
s := newSentPaths(time.Minute, 3)
hash, observer := []byte{1, 2}, []byte{3}
if !s.mark(hash, observer, []byte{0xaa}, start) {
t.Fatal("first hearing refused")
}
if s.mark(hash, observer, []byte{0xaa}, start.Add(time.Second)) {
t.Fatal("exact copy within the window allowed")
}
if !s.mark(hash, observer, []byte{0xaa, 0xbb}, start.Add(time.Second)) {
t.Fatal("new path refused")
}
if !s.mark(hash, []byte{4}, []byte{0xaa}, start.Add(time.Second)) {
t.Fatal("other observer refused")
}
if !s.mark(hash, observer, []byte{0xaa}, start.Add(time.Minute)) {
t.Fatal("copy after the window refused")
}

s = newSentPaths(time.Hour, 3)
for i := range 4 {
s.mark(hash, observer, []byte{byte(i)}, start)
}
if len(s.at) != 3 || len(s.order) != 3 {
t.Fatalf("size bound: %d keys, %d queued", len(s.at), len(s.order))
}
if !s.mark(hash, observer, []byte{0}, start) {
t.Fatal("oldest key not evicted at the size bound")
}
}
59 changes: 59 additions & 0 deletions internal/hub/sent.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,59 @@
// Copyright 2026 Beacon Contributors
// SPDX-License-Identifier: AGPL-3.0-or-later

package hub

import (
"hash/maphash"
"sync"
"time"
)

const (
sentPathsTTL = 2 * time.Minute
sentPathsMax = 200_000
)

// sentPaths remembers recently sent (packet, observer, path) hearings. It is shared by all
// broker workers; keys are hashed so memory stays small at the size bound.
type sentPaths struct {
mu sync.Mutex
seed maphash.Seed
ttl time.Duration
max int
at map[uint64]time.Time
order []sentPath // oldest first
}

type sentPath struct {
key uint64
at time.Time
}

func newSentPaths(ttl time.Duration, max int) *sentPaths {
return &sentPaths{seed: maphash.MakeSeed(), ttl: ttl, max: max, at: make(map[uint64]time.Time)}
}

// mark records the hearing and reports whether it was not already recorded within the TTL.
func (s *sentPaths) mark(packetHash, observerID, path []byte, now time.Time) bool {
var h maphash.Hash
h.SetSeed(s.seed)
for _, part := range [][]byte{packetHash, observerID, path} {
h.WriteByte(byte(len(part)))
h.Write(part)
}
key := h.Sum64()

s.mu.Lock()
defer s.mu.Unlock()
for len(s.order) > 0 && (now.Sub(s.order[0].at) >= s.ttl || len(s.order) >= s.max) {
delete(s.at, s.order[0].key)
s.order = s.order[1:]
}
if _, ok := s.at[key]; ok {
return false
}
s.at[key] = now
s.order = append(s.order, sentPath{key, now})
return true
}
7 changes: 6 additions & 1 deletion internal/ingest/ingest.go
Original file line number Diff line number Diff line change
Expand Up @@ -385,7 +385,8 @@ func (w *Worker) broadcast(eventType hub.EventType, iata string, payloadType uin
// passed in rather than computed here because the caller already has the
// path-hash resolution results in hand from other per-packet work (known
// route detection, capability detection) — this adds no extra DB calls.
func (w *Worker) broadcastPacketObservation(iata string, payloadType uint8, evt packetObservationEvent, resolvedPath []api.ResolvedHop, observerKey string) {
// Repeats go through the hub's drop-first path.
func (w *Worker) broadcastPacketObservation(iata string, payloadType uint8, evt packetObservationEvent, resolvedPath []api.ResolvedHop, observerKey string, repeat bool) {
base, err := json.Marshal(evt)
if err != nil {
w.log.Error("failed to marshal packetObservation event", "error", err)
Expand Down Expand Up @@ -414,6 +415,10 @@ func (w *Worker) broadcastPacketObservation(iata string, payloadType uint8, evt
out.PayloadWithKey = nil
}
}
if repeat {
w.hub.BroadcastRepeat(out)
return
}
w.hub.Broadcast(out)
}

Expand Down
2 changes: 1 addition & 1 deletion internal/ingest/observer_key_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,7 @@ func TestBroadcastPacketObservation_ObserverKeyOptIn(t *testing.T) {

broadcast := func() hub.Event {
t.Helper()
w.broadcastPacketObservation("YVR", 4, packetObservationEvent{}, []api.ResolvedHop{{}}, key)
w.broadcastPacketObservation("YVR", 4, packetObservationEvent{}, []api.ResolvedHop{{}}, key, false)
for {
select {
case evt := <-client.Send:
Expand Down
Loading
Loading