diff --git a/internal/routewatch/peeringhandler.go b/internal/routewatch/peeringhandler.go index 050cb23..820ccde 100644 --- a/internal/routewatch/peeringhandler.go +++ b/internal/routewatch/peeringhandler.go @@ -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 - asPaths map[string]time.Time // key is JSON-encoded AS path + 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() { diff --git a/internal/routewatch/peeringhandler_test.go b/internal/routewatch/peeringhandler_test.go new file mode 100644 index 0000000..364fa46 --- /dev/null +++ b/internal/routewatch/peeringhandler_test.go @@ -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) + } +}