Bound peering AS-path map and swap it instead of copying #16
@@ -21,19 +21,15 @@ const (
|
||||
// than 2 ASNs cannot contain any peering information.
|
||||
minPathLengthForPeering = 2
|
||||
|
||||
// pathExpirationTime determines how long AS paths are kept in memory
|
||||
// before being eligible for pruning. Paths older than this are removed
|
||||
// to prevent unbounded memory growth.
|
||||
pathExpirationTime = 30 * time.Minute
|
||||
|
||||
// peeringProcessInterval controls how frequently the handler processes
|
||||
// accumulated AS paths and extracts peering relationships to store
|
||||
// in the database.
|
||||
peeringProcessInterval = 30 * time.Second
|
||||
|
||||
// pathPruneInterval determines how often the handler checks for and
|
||||
// removes expired AS paths from memory.
|
||||
pathPruneInterval = 5 * time.Minute
|
||||
// maxTrackedPaths bounds how many distinct AS paths are held in memory
|
||||
// between processing runs. Once the map is full, further new paths are
|
||||
// dropped and counted until the next run empties it.
|
||||
maxTrackedPaths = 500000
|
||||
)
|
||||
|
||||
// PeeringHandler processes BGP UPDATE messages to extract and track
|
||||
@@ -46,17 +42,17 @@ type PeeringHandler struct {
|
||||
logger *logger.Logger
|
||||
|
||||
// In-memory AS path tracking
|
||||
mu sync.RWMutex
|
||||
mu sync.Mutex
|
||||
asPaths map[string]time.Time // key is JSON-encoded AS path
|
||||
droppedPaths int // paths dropped because the map was full
|
||||
|
||||
stopCh chan struct{}
|
||||
}
|
||||
|
||||
// NewPeeringHandler creates and initializes a new PeeringHandler with the
|
||||
// provided database store and logger. It starts two background goroutines:
|
||||
// one for periodic processing of accumulated AS paths into peering records,
|
||||
// and one for pruning expired paths from memory. The handler begins
|
||||
// processing immediately upon creation.
|
||||
// provided database store and logger. It starts one background goroutine that
|
||||
// periodically processes accumulated AS paths into peering records. The
|
||||
// handler begins processing immediately upon creation.
|
||||
func NewPeeringHandler(db database.Store, logger *logger.Logger) *PeeringHandler {
|
||||
h := &PeeringHandler{
|
||||
db: db,
|
||||
@@ -65,9 +61,8 @@ func NewPeeringHandler(db database.Store, logger *logger.Logger) *PeeringHandler
|
||||
stopCh: make(chan struct{}),
|
||||
}
|
||||
|
||||
// Start the periodic processing goroutines
|
||||
// Start the periodic processing goroutine
|
||||
go h.processLoop()
|
||||
go h.pruneLoop()
|
||||
|
||||
return h
|
||||
}
|
||||
@@ -106,9 +101,18 @@ func (h *PeeringHandler) HandleMessage(msg *ristypes.RISMessage) {
|
||||
|
||||
return
|
||||
}
|
||||
key := string(pathJSON)
|
||||
|
||||
h.mu.Lock()
|
||||
h.asPaths[string(pathJSON)] = timestamp
|
||||
if _, exists := h.asPaths[key]; exists {
|
||||
// Already tracked: refresh its timestamp.
|
||||
h.asPaths[key] = timestamp
|
||||
} else if len(h.asPaths) >= maxTrackedPaths {
|
||||
// Map is full; drop this new path and count it.
|
||||
h.droppedPaths++
|
||||
} else {
|
||||
h.asPaths[key] = timestamp
|
||||
}
|
||||
h.mu.Unlock()
|
||||
}
|
||||
|
||||
@@ -130,41 +134,6 @@ func (h *PeeringHandler) processLoop() {
|
||||
}
|
||||
}
|
||||
|
||||
// pruneLoop runs periodically to remove old AS paths
|
||||
func (h *PeeringHandler) pruneLoop() {
|
||||
ticker := time.NewTicker(pathPruneInterval)
|
||||
defer ticker.Stop()
|
||||
|
||||
for {
|
||||
select {
|
||||
case <-ticker.C:
|
||||
h.prunePaths()
|
||||
case <-h.stopCh:
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// prunePaths removes AS paths older than pathExpirationTime
|
||||
func (h *PeeringHandler) prunePaths() {
|
||||
cutoff := time.Now().Add(-pathExpirationTime)
|
||||
var removed int
|
||||
|
||||
h.mu.Lock()
|
||||
for pathKey, timestamp := range h.asPaths {
|
||||
if timestamp.Before(cutoff) {
|
||||
delete(h.asPaths, pathKey)
|
||||
removed++
|
||||
}
|
||||
}
|
||||
pathCount := len(h.asPaths)
|
||||
h.mu.Unlock()
|
||||
|
||||
if removed > 0 {
|
||||
h.logger.Debug("Pruned old AS paths", "removed", removed, "remaining", pathCount)
|
||||
}
|
||||
}
|
||||
|
||||
// ProcessPeeringsNow triggers immediate processing of all accumulated AS
|
||||
// paths into peering records. This bypasses the normal periodic processing
|
||||
// schedule and is primarily intended for testing purposes.
|
||||
@@ -174,15 +143,16 @@ func (h *PeeringHandler) ProcessPeeringsNow() {
|
||||
|
||||
// processPeerings extracts peerings from AS paths and writes to database
|
||||
func (h *PeeringHandler) processPeerings() {
|
||||
// Take a snapshot of current AS paths
|
||||
h.mu.RLock()
|
||||
pathsCopy := make(map[string]time.Time, len(h.asPaths))
|
||||
for k, v := range h.asPaths {
|
||||
pathsCopy[k] = v
|
||||
}
|
||||
h.mu.RUnlock()
|
||||
// Take the accumulated paths and replace the map with a fresh empty one
|
||||
// under the lock. Each path is processed exactly once and the memory is
|
||||
// released, so the map never grows past a single interval's traffic.
|
||||
h.mu.Lock()
|
||||
paths := h.asPaths
|
||||
h.asPaths = make(map[string]time.Time)
|
||||
dropped := h.droppedPaths
|
||||
h.mu.Unlock()
|
||||
|
||||
if len(pathsCopy) == 0 {
|
||||
if len(paths) == 0 {
|
||||
return
|
||||
}
|
||||
|
||||
@@ -192,7 +162,7 @@ func (h *PeeringHandler) processPeerings() {
|
||||
}
|
||||
peerings := make(map[peeringKey]time.Time)
|
||||
|
||||
for pathJSON, timestamp := range pathsCopy {
|
||||
for pathJSON, timestamp := range paths {
|
||||
var path []int
|
||||
if err := json.Unmarshal([]byte(pathJSON), &path); err != nil {
|
||||
h.logger.Error("Failed to decode AS path", "error", err)
|
||||
@@ -241,15 +211,16 @@ func (h *PeeringHandler) processPeerings() {
|
||||
}
|
||||
|
||||
h.logger.Info("Processed AS peerings",
|
||||
"paths", len(pathsCopy),
|
||||
"paths", len(paths),
|
||||
"unique_peerings", len(peerings),
|
||||
"success", successCount,
|
||||
"dropped_paths", dropped,
|
||||
"duration", time.Since(start),
|
||||
)
|
||||
}
|
||||
|
||||
// Stop gracefully shuts down the handler by signaling the background
|
||||
// goroutines to stop and performing a final synchronous processing of
|
||||
// goroutine to stop and performing a final synchronous processing of
|
||||
// any remaining AS paths. This ensures no peering data is lost during
|
||||
// shutdown.
|
||||
func (h *PeeringHandler) Stop() {
|
||||
|
||||
@@ -0,0 +1,163 @@
|
||||
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)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user