Bound peering AS-path map and swap it instead of copying (closes #10)
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
This commit was merged in pull request #16.
This commit is contained in:
@@ -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() {
|
||||
|
||||
@@ -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