processPeerings now takes the accumulated AS-path map under the lock and replaces it with a fresh empty one, so each path is processed once and the map never grows past a single 30-second interval's traffic. HandleMessage stops adding new paths once maxTrackedPaths (500000) is reached and counts the drops, which are logged with each run. The 30-minute prune and its ticker are removed as dead code under the swap, and the map mutex drops from RWMutex to Mutex since the read path is gone. RecordPeering already upserts last_seen, so stored peerings are unchanged. Model: opus-4-8
164 lines
4.3 KiB
Go
164 lines
4.3 KiB
Go
package routewatch
|
|
|
|
import (
|
|
"encoding/json"
|
|
"strconv"
|
|
"sync"
|
|
"testing"
|
|
"time"
|
|
|
|
"git.eeqj.de/sneak/routewatch/internal/database"
|
|
"git.eeqj.de/sneak/routewatch/internal/logger"
|
|
"git.eeqj.de/sneak/routewatch/internal/ristypes"
|
|
)
|
|
|
|
const (
|
|
testASNA = 64500
|
|
testASNB = 64501
|
|
testASNC = 64502
|
|
)
|
|
|
|
// recordingStore wraps mockStore to count every RecordPeering call, so a
|
|
// test can tell how many times a peering was written across separate runs.
|
|
type recordingStore struct {
|
|
*mockStore
|
|
|
|
mu sync.Mutex
|
|
calls int
|
|
}
|
|
|
|
func (r *recordingStore) RecordPeering(asA, asB int, ts time.Time) error {
|
|
r.mu.Lock()
|
|
r.calls++
|
|
r.mu.Unlock()
|
|
|
|
return r.mockStore.RecordPeering(asA, asB, ts)
|
|
}
|
|
|
|
func (r *recordingStore) callCount() int {
|
|
r.mu.Lock()
|
|
defer r.mu.Unlock()
|
|
|
|
return r.calls
|
|
}
|
|
|
|
// newTestHandler builds a PeeringHandler without starting the periodic
|
|
// processing goroutine, so tests drive processing explicitly.
|
|
func newTestHandler(db database.Store) *PeeringHandler {
|
|
return &PeeringHandler{
|
|
db: db,
|
|
logger: logger.New(),
|
|
asPaths: make(map[string]time.Time),
|
|
stopCh: make(chan struct{}),
|
|
}
|
|
}
|
|
|
|
func snapshot(h *PeeringHandler) (tracked, dropped int) {
|
|
h.mu.Lock()
|
|
defer h.mu.Unlock()
|
|
|
|
return len(h.asPaths), h.droppedPaths
|
|
}
|
|
|
|
func pathKey(t *testing.T, asns ...int) string {
|
|
t.Helper()
|
|
|
|
b, err := json.Marshal(ristypes.ASPath(asns))
|
|
if err != nil {
|
|
t.Fatalf("failed to marshal path: %v", err)
|
|
}
|
|
|
|
return string(b)
|
|
}
|
|
|
|
func handle(h *PeeringHandler, ts time.Time, asns ...int) {
|
|
h.HandleMessage(&ristypes.RISMessage{
|
|
Path: ristypes.ASPath(asns),
|
|
ParsedTimestamp: ts,
|
|
})
|
|
}
|
|
|
|
// TestPeeringHandlerProcessesEachRunAndEmpties verifies that a processing run
|
|
// empties the path map (the swap) and that a path seen again after a run is
|
|
// recorded in the next run too.
|
|
func TestPeeringHandlerProcessesEachRunAndEmpties(t *testing.T) {
|
|
store := &recordingStore{mockStore: newMockStore()}
|
|
h := newTestHandler(store)
|
|
|
|
now := time.Now().UTC()
|
|
|
|
// Run 1: one path, one peering recorded, map emptied afterwards.
|
|
handle(h, now, testASNA, testASNB)
|
|
h.ProcessPeeringsNow()
|
|
|
|
if tracked, _ := snapshot(h); tracked != 0 {
|
|
t.Fatalf("map not empty after first run: %d paths remain", tracked)
|
|
}
|
|
if got := store.callCount(); got != 1 {
|
|
t.Fatalf("want 1 RecordPeering call after first run, got %d", got)
|
|
}
|
|
|
|
// Run 2: the same path again is recorded again (RecordPeering upserts).
|
|
handle(h, now.Add(time.Second), testASNA, testASNB)
|
|
h.ProcessPeeringsNow()
|
|
|
|
if tracked, _ := snapshot(h); tracked != 0 {
|
|
t.Fatalf("map not empty after second run: %d paths remain", tracked)
|
|
}
|
|
if got := store.callCount(); got != 2 {
|
|
t.Fatalf("want 2 RecordPeering calls after second run, got %d", got)
|
|
}
|
|
}
|
|
|
|
// TestPeeringHandlerCapDropsAndCounts verifies that a full map drops new paths
|
|
// and counts them, while a path already tracked is refreshed rather than
|
|
// dropped.
|
|
func TestPeeringHandlerCapDropsAndCounts(t *testing.T) {
|
|
store := &recordingStore{mockStore: newMockStore()}
|
|
h := newTestHandler(store)
|
|
|
|
now := time.Now().UTC()
|
|
|
|
// Fill the map to exactly maxTrackedPaths, including one real path key so
|
|
// the "already tracked" branch can be exercised. The filler keys are never
|
|
// processed in this test, so their contents do not matter.
|
|
existing := pathKey(t, testASNA, testASNB)
|
|
|
|
h.mu.Lock()
|
|
h.asPaths[existing] = now
|
|
for i := 0; len(h.asPaths) < maxTrackedPaths; i++ {
|
|
h.asPaths[strconv.Itoa(i)] = now
|
|
}
|
|
h.mu.Unlock()
|
|
|
|
// A new path is dropped and counted because the map is full.
|
|
handle(h, now.Add(time.Second), testASNA, testASNC)
|
|
|
|
tracked, dropped := snapshot(h)
|
|
if tracked != maxTrackedPaths {
|
|
t.Fatalf("want map size %d after drop, got %d", maxTrackedPaths, tracked)
|
|
}
|
|
if dropped != 1 {
|
|
t.Fatalf("want dropped count 1, got %d", dropped)
|
|
}
|
|
|
|
// A path already tracked is refreshed, not dropped.
|
|
refreshed := now.Add(2 * time.Second)
|
|
handle(h, refreshed, testASNA, testASNB)
|
|
|
|
tracked, dropped = snapshot(h)
|
|
if tracked != maxTrackedPaths {
|
|
t.Fatalf("want map size %d after refresh, got %d", maxTrackedPaths, tracked)
|
|
}
|
|
if dropped != 1 {
|
|
t.Fatalf("want dropped count still 1 after refresh, got %d", dropped)
|
|
}
|
|
|
|
h.mu.Lock()
|
|
gotTS := h.asPaths[existing]
|
|
h.mu.Unlock()
|
|
if !gotTS.Equal(refreshed) {
|
|
t.Fatalf("existing path timestamp not refreshed: want %v, got %v", refreshed, gotTS)
|
|
}
|
|
}
|