Serve the /api/v1/stats prefix distribution from memory (closes #30)
check / check (push) Waiting to run

The prefix distribution was the last query on the stats path. It read one
index entry per live route on every request, so on a large database it took
the whole 4-second deadline and /api/v1/stats answered 500.

The distribution now lives next to the in-memory counts: seeded once at
startup from the same query, then kept exact by the live-route writes. A new
route whose prefix had no live route adds one at its mask length; a delete
that leaves a prefix with no live route takes one away. Each check is one
lookup on the prefix index, made only for new and removed routes, inside the
write.

Model: opus-5-5
This commit is contained in:
2026-09-29 11:05:28 +00:00
parent 6422d9fa0c
commit 539d15c0a5
5 changed files with 394 additions and 60 deletions
+3 -4
View File
@@ -24,10 +24,12 @@ https://git.eeqj.de/sneak/routewatch/pulls/6. After that, setting
routewatch up under upaas on fsn1app1 and deploying it are his routewatch up under upaas on fsn1app1 and deploying it are his
(https://git.eeqj.de/sneak/routewatch/issues/31), and so is the run under a (https://git.eeqj.de/sneak/routewatch/issues/31), and so is the run under a
real 5 GiB limit (https://git.eeqj.de/sneak/routewatch/issues/3). real 5 GiB limit (https://git.eeqj.de/sneak/routewatch/issues/3).
The other open issue is https://git.eeqj.de/sneak/routewatch/issues/30.
# Completed Steps # Completed Steps
- 2026-09-29: `/api/v1/stats` serves the prefix distribution from memory,
seeded at startup and adjusted on every live-route write, so a request no
longer reads every live route (closes #30)
- 2026-09-29: the entrypoint creates the data directory if it is missing and - 2026-09-29: the entrypoint creates the data directory if it is missing and
stops the start if a step fails; README "Running under upaas" no longer stops the start if a step fails; README "Running under upaas" no longer
asks for the host directory to be created first (closes #42) asks for the host directory to be created first (closes #42)
@@ -103,6 +105,3 @@ The other open issue is https://git.eeqj.de/sneak/routewatch/issues/30.
- Production memory under 5 GiB: whether to test under a real 5 GiB - Production memory under 5 GiB: whether to test under a real 5 GiB
container limit on fsn1app1 is open for sneak container limit on fsn1app1 is open for sneak
(https://git.eeqj.de/sneak/routewatch/issues/3) (https://git.eeqj.de/sneak/routewatch/issues/3)
- `/api/v1/stats` answered HTTP 500 after 35 hours on the live feed, seen
on `3898daa`, which predates the 2026-09-22 in-memory statistics
(https://git.eeqj.de/sneak/routewatch/issues/30)
+66 -9
View File
@@ -6,11 +6,12 @@ import (
"sync" "sync"
) )
// liveCounts holds the running row counts that the stats endpoints report. They // liveCounts holds the running row counts and the prefix distribution that the
// are seeded once at startup from the tables and then adjusted on every write, // stats endpoints report. They are seeded once at startup from the tables and
// so a stats read serves them from memory instead of running a COUNT(*) over // then adjusted on every write, so a stats read serves them from memory instead
// each table. Those scans, once the database passed a few GiB, took the whole // of running a query over the tables. The COUNT(*) scans (issue 27) and then
// request timeout and made /api/v1/stats return 500 (issue 27). // the prefix distribution query (issue 30) each grew with the database until
// they took the whole request timeout and made /api/v1/stats return 500.
// //
// A single mutex guards all fields so the stats reader takes a consistent // A single mutex guards all fields so the stats reader takes a consistent
// snapshot at one instant and writers, which already run under the database // snapshot at one instant and writers, which already run under the database
@@ -24,11 +25,16 @@ type liveCounts struct {
peers int peers int
routesV4 int routesV4 int
routesV6 int routesV6 int
// The prefix distribution: for each mask length, the number of distinct
// prefixes that have at least one live route.
distributionV4 [ipv4Bits + 1]int
distributionV6 [ipv6Bits + 1]int
} }
// seed sets every count to the value read from the tables at startup. It runs // seed sets every count to the value read from the tables at startup. It runs
// before any writer, so it needs no coordination with the adjust methods. // before any writer, so it needs no coordination with the adjust methods.
func (c *liveCounts) seed(asns, prefixesV4, prefixesV6, peerings, peers, routesV4, routesV6 int) { func (c *liveCounts) seed(asns, prefixesV4, prefixesV6, peerings, peers, routesV4, routesV6 int,
distributionV4, distributionV6 []PrefixDistribution) {
c.mu.Lock() c.mu.Lock()
defer c.mu.Unlock() defer c.mu.Unlock()
@@ -39,6 +45,12 @@ func (c *liveCounts) seed(asns, prefixesV4, prefixesV6, peerings, peers, routesV
c.peers = peers c.peers = peers
c.routesV4 = routesV4 c.routesV4 = routesV4
c.routesV6 = routesV6 c.routesV6 = routesV6
for _, entry := range distributionV4 {
addAtMaskLength(c.distributionV4[:], entry.MaskLength, entry.Count)
}
for _, entry := range distributionV6 {
addAtMaskLength(c.distributionV6[:], entry.MaskLength, entry.Count)
}
} }
// addASNs adds n to the ASN count. // addASNs adds n to the ASN count.
@@ -79,6 +91,30 @@ func (c *liveCounts) addRoutes(v4, v6 int) {
c.mu.Unlock() c.mu.Unlock()
} }
// addToDistribution adds n to the IPv4 and IPv6 prefix distributions once for
// each listed mask length. A write lists the mask lengths of the prefixes it
// gave their first live route with n = 1, and of the prefixes it left with no
// live route with n = -1.
func (c *liveCounts) addToDistribution(maskLengthsV4, maskLengthsV6 []int, n int) {
c.mu.Lock()
defer c.mu.Unlock()
for _, maskLength := range maskLengthsV4 {
addAtMaskLength(c.distributionV4[:], maskLength, n)
}
for _, maskLength := range maskLengthsV6 {
addAtMaskLength(c.distributionV6[:], maskLength, n)
}
}
// addAtMaskLength adds n to counts[maskLength]. A mask length the array has no
// entry for is ignored, so a malformed route cannot crash the daemon.
func addAtMaskLength(counts []int, maskLength, n int) {
if maskLength >= 0 && maskLength < len(counts) {
counts[maskLength] += n
}
}
// fill copies the counts into a Stats, including the derived totals, under a // fill copies the counts into a Stats, including the derived totals, under a
// single read lock so the reader sees one consistent snapshot. // single read lock so the reader sees one consistent snapshot.
func (c *liveCounts) fill(s *Stats) { func (c *liveCounts) fill(s *Stats) {
@@ -94,6 +130,21 @@ func (c *liveCounts) fill(s *Stats) {
s.IPv4Routes = c.routesV4 s.IPv4Routes = c.routesV4
s.IPv6Routes = c.routesV6 s.IPv6Routes = c.routesV6
s.LiveRoutes = c.routesV4 + c.routesV6 s.LiveRoutes = c.routesV4 + c.routesV6
s.IPv4PrefixDistribution = distributionList(c.distributionV4[:])
s.IPv6PrefixDistribution = distributionList(c.distributionV6[:])
}
// distributionList lists the mask lengths that have at least one prefix, in
// ascending order, the way the distribution query returns them.
func distributionList(counts []int) []PrefixDistribution {
var list []PrefixDistribution
for maskLength, count := range counts {
if count > 0 {
list = append(list, PrefixDistribution{MaskLength: maskLength, Count: count})
}
}
return list
} }
// countRows returns the number of rows in the named table. It is used only at // countRows returns the number of rows in the named table. It is used only at
@@ -108,8 +159,9 @@ func (d *Database) countRows(ctx context.Context, table string) (int, error) {
return n, nil return n, nil
} }
// seedCounts reads the current row counts from the tables into the in-memory // seedCounts reads the current row counts and prefix distribution from the
// counters. It runs once at startup, before the streamer begins writing. // tables into the in-memory counters. It runs once at startup, before the
// streamer begins writing.
func (d *Database) seedCounts(ctx context.Context) error { func (d *Database) seedCounts(ctx context.Context) error {
asns, err := d.countRows(ctx, "asns") asns, err := d.countRows(ctx, "asns")
if err != nil { if err != nil {
@@ -139,8 +191,13 @@ func (d *Database) seedCounts(ctx context.Context) error {
if err != nil { if err != nil {
return err return err
} }
distributionV4, distributionV6, err := d.GetPrefixDistributionContext(ctx)
if err != nil {
return err
}
d.counts.seed(asns, prefixesV4, prefixesV6, peerings, peers, routesV4, routesV6) d.counts.seed(asns, prefixesV4, prefixesV6, peerings, peers, routesV4, routesV6,
distributionV4, distributionV6)
return nil return nil
} }
+160 -4
View File
@@ -2,6 +2,7 @@ package database
import ( import (
"context" "context"
"slices"
"sync" "sync"
"testing" "testing"
"time" "time"
@@ -11,7 +12,8 @@ import (
"github.com/google/uuid" "github.com/google/uuid"
) )
// mkV4Route builds an IPv4 live route with its range columns populated. // mkV4Route builds an IPv4 live route with its mask length and range columns
// taken from the prefix.
func mkV4Route(t *testing.T, prefix string, asn int, ts time.Time) *LiveRoute { func mkV4Route(t *testing.T, prefix string, asn int, ts time.Time) *LiveRoute {
t.Helper() t.Helper()
@@ -19,11 +21,15 @@ func mkV4Route(t *testing.T, prefix string, asn int, ts time.Time) *LiveRoute {
if err != nil { if err != nil {
t.Fatalf("CalculateIPv4Range(%s): %v", prefix, err) t.Fatalf("CalculateIPv4Range(%s): %v", prefix, err)
} }
maskLength, err := prefixMaskLength(prefix)
if err != nil {
t.Fatalf("prefixMaskLength(%s): %v", prefix, err)
}
return &LiveRoute{ return &LiveRoute{
ID: uuid.New(), ID: uuid.New(),
Prefix: prefix, Prefix: prefix,
MaskLength: 24, MaskLength: maskLength,
IPVersion: ipVersionV4, IPVersion: ipVersionV4,
OriginASN: asn, OriginASN: asn,
PeerIP: "192.0.2.1", PeerIP: "192.0.2.1",
@@ -143,8 +149,8 @@ func TestLiveCountsTrackWritesInRealtime(t *testing.T) {
} }
// TestLiveCountsSeededFromDatabaseAtStartup writes rows, reopens the same // TestLiveCountsSeededFromDatabaseAtStartup writes rows, reopens the same
// database file, and checks the counts come back from the seed scan rather than // database file, and checks the counts and the prefix distribution come back
// starting at zero. // from the seed scan rather than starting at zero.
func TestLiveCountsSeededFromDatabaseAtStartup(t *testing.T) { func TestLiveCountsSeededFromDatabaseAtStartup(t *testing.T) {
cfg := &config.Config{StateDir: t.TempDir()} cfg := &config.Config{StateDir: t.TempDir()}
@@ -184,6 +190,145 @@ func TestLiveCountsSeededFromDatabaseAtStartup(t *testing.T) {
t.Errorf("seeded routes = (v4 %d, v6 %d, total %d), want (1, 1, 2)", t.Errorf("seeded routes = (v4 %d, v6 %d, total %d), want (1, 1, 2)",
stats.IPv4Routes, stats.IPv6Routes, stats.LiveRoutes) stats.IPv4Routes, stats.IPv6Routes, stats.LiveRoutes)
} }
assertDistribution(t, "seeded IPv4 distribution", stats.IPv4PrefixDistribution,
[]PrefixDistribution{{MaskLength: 24, Count: 1}})
assertDistribution(t, "seeded IPv6 distribution", stats.IPv6PrefixDistribution,
[]PrefixDistribution{{MaskLength: 32, Count: 1}})
}
// TestPrefixDistributionTracksWrites checks that the prefix distribution the
// stats read reports stays exact across each kind of write, through both the
// batch and the single-route methods, and that after every step it equals what
// the distribution query reads from the route tables.
func TestPrefixDistributionTracksWrites(t *testing.T) {
cfg := &config.Config{StateDir: t.TempDir()}
db, err := New(cfg, logger.New())
if err != nil {
t.Fatalf("failed to create database: %v", err)
}
defer func() { _ = db.Close() }()
ts := time.Date(2026, 1, 2, 3, 4, 5, 0, time.UTC)
shared := mkV4Route(t, "198.51.100.0/24", 64500, ts)
sharedSecondPeer := mkV4Route(t, "198.51.100.0/24", 64500, ts)
sharedSecondPeer.PeerIP = "192.0.2.2"
other := mkV4Route(t, "203.0.113.0/24", 64501, ts)
wide := mkV4Route(t, "172.16.0.0/16", 64502, ts)
v6 := mkV6Route("2001:db8::/32", 64503, ts)
all := []PrefixDistribution{{MaskLength: 16, Count: 1}, {MaskLength: 24, Count: 2}}
v6Only := []PrefixDistribution{{MaskLength: 32, Count: 1}}
steps := []struct {
name string
write func() error
wantV4 []PrefixDistribution
wantV6 []PrefixDistribution
}{
{
name: "empty database",
write: func() error { return nil },
},
{
name: "new prefixes",
write: func() error { return db.UpsertLiveRouteBatch([]*LiveRoute{shared, other, wide, v6}) },
wantV4: all, wantV6: v6Only,
},
{
name: "re-announcement",
write: func() error { return db.UpsertLiveRouteBatch([]*LiveRoute{shared, other, wide, v6}) },
wantV4: all, wantV6: v6Only,
},
{
name: "second peer announces a prefix that already has a route",
write: func() error { return db.UpsertLiveRoute(sharedSecondPeer) },
wantV4: all, wantV6: v6Only,
},
{
name: "withdrawal of a route that is not the last for its prefix",
write: func() error {
return db.DeleteLiveRouteBatch([]LiveRouteDeletion{
{Prefix: shared.Prefix, OriginASN: shared.OriginASN, PeerIP: shared.PeerIP, IPVersion: ipVersionV4},
})
},
wantV4: all, wantV6: v6Only,
},
{
name: "withdrawal of the last route for a prefix",
write: func() error {
return db.DeleteLiveRoute(sharedSecondPeer.Prefix, sharedSecondPeer.OriginASN, sharedSecondPeer.PeerIP)
},
wantV4: []PrefixDistribution{{MaskLength: 16, Count: 1}, {MaskLength: 24, Count: 1}},
wantV6: v6Only,
},
{
name: "withdrawal of every remaining route, one without an origin ASN",
write: func() error {
return db.DeleteLiveRouteBatch([]LiveRouteDeletion{
{Prefix: other.Prefix, PeerIP: other.PeerIP, IPVersion: ipVersionV4},
{Prefix: wide.Prefix, OriginASN: wide.OriginASN, PeerIP: wide.PeerIP, IPVersion: ipVersionV4},
{Prefix: v6.Prefix, OriginASN: v6.OriginASN, PeerIP: v6.PeerIP, IPVersion: ipVersionV6},
})
},
},
}
ctx := context.Background()
for _, step := range steps {
if err := step.write(); err != nil {
t.Fatalf("%s: %v", step.name, err)
}
stats, err := db.GetStatsContext(ctx)
if err != nil {
t.Fatalf("%s: GetStatsContext: %v", step.name, err)
}
queryV4, queryV6, err := db.GetPrefixDistributionContext(ctx)
if err != nil {
t.Fatalf("%s: GetPrefixDistributionContext: %v", step.name, err)
}
assertDistribution(t, step.name+": IPv4 distribution", stats.IPv4PrefixDistribution, step.wantV4)
assertDistribution(t, step.name+": IPv6 distribution", stats.IPv6PrefixDistribution, step.wantV6)
assertDistribution(t, step.name+": IPv4 distribution query", queryV4, step.wantV4)
assertDistribution(t, step.name+": IPv6 distribution query", queryV6, step.wantV6)
}
}
// TestStatsDistributionNeedsNoQuery checks that the stats read serves the prefix
// distribution from memory. The read's deadline has already passed, so any
// query it ran would fail at once. The distribution query it used to run read
// every live route and passed the /api/v1/stats deadline on a large database
// (issue 30).
func TestStatsDistributionNeedsNoQuery(t *testing.T) {
cfg := &config.Config{StateDir: t.TempDir()}
db, err := New(cfg, logger.New())
if err != nil {
t.Fatalf("failed to create database: %v", err)
}
defer func() { _ = db.Close() }()
ts := time.Date(2026, 1, 2, 3, 4, 5, 0, time.UTC)
if err := db.UpsertLiveRouteBatch([]*LiveRoute{
mkV4Route(t, "198.51.100.0/24", 64500, ts),
mkV6Route("2001:db8::/32", 64501, ts),
}); err != nil {
t.Fatalf("UpsertLiveRouteBatch: %v", err)
}
ctx, cancel := context.WithCancel(context.Background())
cancel()
stats, err := db.GetStatsContext(ctx)
if err != nil {
t.Fatalf("GetStatsContext: %v", err)
}
assertDistribution(t, "IPv4 distribution", stats.IPv4PrefixDistribution,
[]PrefixDistribution{{MaskLength: 24, Count: 1}})
assertDistribution(t, "IPv6 distribution", stats.IPv6PrefixDistribution,
[]PrefixDistribution{{MaskLength: 32, Count: 1}})
} }
// TestStatsRouteTimestamps checks the oldest/newest route timestamps are read // TestStatsRouteTimestamps checks the oldest/newest route timestamps are read
@@ -296,6 +441,9 @@ func TestLiveCountsConcurrentReadWrite(t *testing.T) {
if want := writers * 25; stats.IPv6Routes != want { if want := writers * 25; stats.IPv6Routes != want {
t.Errorf("IPv6Routes = %d, want %d", stats.IPv6Routes, want) t.Errorf("IPv6Routes = %d, want %d", stats.IPv6Routes, want)
} }
// Every writer announced the same prefix, so it counts once.
assertDistribution(t, "IPv6 distribution", stats.IPv6PrefixDistribution,
[]PrefixDistribution{{MaskLength: 32, Count: 1}})
} }
type wantCounts struct { type wantCounts struct {
@@ -333,3 +481,11 @@ func assertCounts(t *testing.T, when string, got Stats, want wantCounts) {
t.Errorf("%s: LiveRoutes = %d, want %d", when, got.LiveRoutes, want.liveRoutes) t.Errorf("%s: LiveRoutes = %d, want %d", when, got.LiveRoutes, want.liveRoutes)
} }
} }
func assertDistribution(t *testing.T, what string, got, want []PrefixDistribution) {
t.Helper()
if !slices.Equal(got, want) {
t.Errorf("%s = %v, want %v", what, got, want)
}
}
+152 -43
View File
@@ -35,6 +35,7 @@ const (
ipv6Length = 16 ipv6Length = 16
ipv4Offset = 12 ipv4Offset = 12
ipv4Bits = 32 ipv4Bits = 32
ipv6Bits = 128
maxIPv4 = 0xFFFFFFFF maxIPv4 = 0xFFFFFFFF
) )
@@ -237,58 +238,80 @@ const (
as_path, next_hop, last_updated) VALUES (?, ?, ?, ?, ?, ?, ?, ?)` as_path, next_hop, last_updated) VALUES (?, ?, ?, ?, ?, ?, ?, ?)`
) )
// upsertRouteRowV4 updates an IPv4 live route, inserting it when no row matched, // Before a new route is inserted and after a route is deleted, the write looks
// and reports whether a new row was inserted. // up whether any live route has that prefix, to keep the in-memory prefix
func upsertRouteRowV4(upd, ins *sql.Stmt, route *LiveRoute, pathJSON string) (inserted bool, err error) { // distribution exact. The lookup reads one entry of the prefix index.
const (
prefixHasLiveRouteV4SQL = `SELECT EXISTS (SELECT 1 FROM live_routes_v4 WHERE prefix = ?)`
prefixHasLiveRouteV6SQL = `SELECT EXISTS (SELECT 1 FROM live_routes_v6 WHERE prefix = ?)`
)
// upsertRouteRowV4 updates an IPv4 live route, inserting it when no row matched.
// It reports whether a new row was inserted and whether that row is the first
// live route for its prefix.
func upsertRouteRowV4(upd, ins, has *sql.Stmt, route *LiveRoute, pathJSON string) (
inserted, newPrefix bool, err error) {
if route.V4IPStart == nil || route.V4IPEnd == nil { if route.V4IPStart == nil || route.V4IPEnd == nil {
return false, fmt.Errorf("IPv4 route %s missing range values", route.Prefix) return false, false, fmt.Errorf("IPv4 route %s missing range values", route.Prefix)
} }
res, err := upd.Exec(route.MaskLength, pathJSON, route.NextHop, route.LastUpdated, res, err := upd.Exec(route.MaskLength, pathJSON, route.NextHop, route.LastUpdated,
*route.V4IPStart, *route.V4IPEnd, route.Prefix, route.OriginASN, route.PeerIP) *route.V4IPStart, *route.V4IPEnd, route.Prefix, route.OriginASN, route.PeerIP)
if err != nil { if err != nil {
return false, err return false, false, err
} }
affected, err := res.RowsAffected() affected, err := res.RowsAffected()
if err != nil { if err != nil {
return false, err return false, false, err
} }
if affected > 0 { if affected > 0 {
return false, nil return false, false, nil
}
var prefixHadRoute bool
if err := has.QueryRow(route.Prefix).Scan(&prefixHadRoute); err != nil {
return false, false, err
} }
_, err = ins.Exec(route.ID.String(), route.Prefix, route.MaskLength, route.OriginASN, _, err = ins.Exec(route.ID.String(), route.Prefix, route.MaskLength, route.OriginASN,
route.PeerIP, pathJSON, route.NextHop, route.LastUpdated, *route.V4IPStart, *route.V4IPEnd) route.PeerIP, pathJSON, route.NextHop, route.LastUpdated, *route.V4IPStart, *route.V4IPEnd)
if err != nil { if err != nil {
return false, err return false, false, err
} }
return true, nil return true, !prefixHadRoute, nil
} }
// upsertRouteRowV6 updates an IPv6 live route, inserting it when no row matched, // upsertRouteRowV6 updates an IPv6 live route, inserting it when no row matched.
// and reports whether a new row was inserted. // It reports whether a new row was inserted and whether that row is the first
func upsertRouteRowV6(upd, ins *sql.Stmt, route *LiveRoute, pathJSON string) (inserted bool, err error) { // live route for its prefix.
func upsertRouteRowV6(upd, ins, has *sql.Stmt, route *LiveRoute, pathJSON string) (
inserted, newPrefix bool, err error) {
res, err := upd.Exec(route.MaskLength, pathJSON, route.NextHop, route.LastUpdated, res, err := upd.Exec(route.MaskLength, pathJSON, route.NextHop, route.LastUpdated,
route.Prefix, route.OriginASN, route.PeerIP) route.Prefix, route.OriginASN, route.PeerIP)
if err != nil { if err != nil {
return false, err return false, false, err
} }
affected, err := res.RowsAffected() affected, err := res.RowsAffected()
if err != nil { if err != nil {
return false, err return false, false, err
} }
if affected > 0 { if affected > 0 {
return false, nil return false, false, nil
}
var prefixHadRoute bool
if err := has.QueryRow(route.Prefix).Scan(&prefixHadRoute); err != nil {
return false, false, err
} }
_, err = ins.Exec(route.ID.String(), route.Prefix, route.MaskLength, route.OriginASN, _, err = ins.Exec(route.ID.String(), route.Prefix, route.MaskLength, route.OriginASN,
route.PeerIP, pathJSON, route.NextHop, route.LastUpdated) route.PeerIP, pathJSON, route.NextHop, route.LastUpdated)
if err != nil { if err != nil {
return false, err return false, false, err
} }
return true, nil return true, !prefixHadRoute, nil
} }
// UpsertLiveRouteBatch inserts or updates multiple live routes in a single transaction // UpsertLiveRouteBatch inserts or updates multiple live routes in a single transaction
@@ -335,7 +358,21 @@ func (d *Database) UpsertLiveRouteBatch(routes []*LiveRoute) error {
} }
defer func() { _ = insV6.Close() }() defer func() { _ = insV6.Close() }()
hasV4, err := tx.Prepare(prefixHasLiveRouteV4SQL)
if err != nil {
return fmt.Errorf("failed to prepare IPv4 prefix lookup statement: %w", err)
}
defer func() { _ = hasV4.Close() }()
hasV6, err := tx.Prepare(prefixHasLiveRouteV6SQL)
if err != nil {
return fmt.Errorf("failed to prepare IPv6 prefix lookup statement: %w", err)
}
defer func() { _ = hasV6.Close() }()
var newV4, newV6 int var newV4, newV6 int
// Mask lengths of the prefixes that get their first live route in this batch.
var newPrefixMaskLengthsV4, newPrefixMaskLengthsV6 []int
for _, route := range routes { for _, route := range routes {
pathJSON, err := json.Marshal(route.ASPath) pathJSON, err := json.Marshal(route.ASPath)
if err != nil { if err != nil {
@@ -343,24 +380,30 @@ func (d *Database) UpsertLiveRouteBatch(routes []*LiveRoute) error {
} }
if route.IPVersion == ipVersionV4 { if route.IPVersion == ipVersionV4 {
inserted, err := upsertRouteRowV4(updV4, insV4, route, string(pathJSON)) inserted, newPrefix, err := upsertRouteRowV4(updV4, insV4, hasV4, route, string(pathJSON))
if err != nil { if err != nil {
return fmt.Errorf("failed to upsert route %s: %w", route.Prefix, err) return fmt.Errorf("failed to upsert route %s: %w", route.Prefix, err)
} }
if inserted { if inserted {
newV4++ newV4++
} }
if newPrefix {
newPrefixMaskLengthsV4 = append(newPrefixMaskLengthsV4, route.MaskLength)
}
continue continue
} }
inserted, err := upsertRouteRowV6(updV6, insV6, route, string(pathJSON)) inserted, newPrefix, err := upsertRouteRowV6(updV6, insV6, hasV6, route, string(pathJSON))
if err != nil { if err != nil {
return fmt.Errorf("failed to upsert route %s: %w", route.Prefix, err) return fmt.Errorf("failed to upsert route %s: %w", route.Prefix, err)
} }
if inserted { if inserted {
newV6++ newV6++
} }
if newPrefix {
newPrefixMaskLengthsV6 = append(newPrefixMaskLengthsV6, route.MaskLength)
}
} }
if err = tx.Commit(); err != nil { if err = tx.Commit(); err != nil {
@@ -368,6 +411,7 @@ func (d *Database) UpsertLiveRouteBatch(routes []*LiveRoute) error {
} }
d.counts.addRoutes(newV4, newV6) d.counts.addRoutes(newV4, newV6)
d.counts.addToDistribution(newPrefixMaskLengthsV4, newPrefixMaskLengthsV6, 1)
return nil return nil
} }
@@ -418,8 +462,22 @@ func (d *Database) DeleteLiveRouteBatch(deletions []LiveRouteDeletion) error {
} }
defer func() { _ = stmtV6WithoutOrigin.Close() }() defer func() { _ = stmtV6WithoutOrigin.Close() }()
hasV4, err := tx.Prepare(prefixHasLiveRouteV4SQL)
if err != nil {
return fmt.Errorf("failed to prepare IPv4 prefix lookup statement: %w", err)
}
defer func() { _ = hasV4.Close() }()
hasV6, err := tx.Prepare(prefixHasLiveRouteV6SQL)
if err != nil {
return fmt.Errorf("failed to prepare IPv6 prefix lookup statement: %w", err)
}
defer func() { _ = hasV6.Close() }()
// Process deletions // Process deletions
var deletedV4, deletedV6 int64 var deletedV4, deletedV6 int64
// Mask lengths of the prefixes this batch leaves with no live route.
var gonePrefixMaskLengthsV4, gonePrefixMaskLengthsV6 []int
for _, del := range deletions { for _, del := range deletions {
var stmt *sql.Stmt var stmt *sql.Stmt
@@ -462,6 +520,31 @@ func (d *Database) DeleteLiveRouteBatch(deletions []LiveRouteDeletion) error {
} else { } else {
deletedV6 += affected deletedV6 += affected
} }
if affected == 0 {
continue
}
// The prefix leaves the distribution when no live route has it any more.
has := hasV4
if del.IPVersion != ipVersionV4 {
has = hasV6
}
var prefixHasRoute bool
if err := has.QueryRow(del.Prefix).Scan(&prefixHasRoute); err != nil {
return fmt.Errorf("failed to look up prefix %s: %w", del.Prefix, err)
}
if prefixHasRoute {
continue
}
maskLength, err := prefixMaskLength(del.Prefix)
if err != nil {
return err
}
if del.IPVersion == ipVersionV4 {
gonePrefixMaskLengthsV4 = append(gonePrefixMaskLengthsV4, maskLength)
} else {
gonePrefixMaskLengthsV6 = append(gonePrefixMaskLengthsV6, maskLength)
}
} }
if err = tx.Commit(); err != nil { if err = tx.Commit(); err != nil {
@@ -469,6 +552,7 @@ func (d *Database) DeleteLiveRouteBatch(deletions []LiveRouteDeletion) error {
} }
d.counts.addRoutes(-int(deletedV4), -int(deletedV6)) d.counts.addRoutes(-int(deletedV4), -int(deletedV6))
d.counts.addToDistribution(gonePrefixMaskLengthsV4, gonePrefixMaskLengthsV6, -1)
return nil return nil
} }
@@ -1035,16 +1119,16 @@ func (d *Database) GetStats() (Stats, error) {
// GetStatsContext returns database statistics with context support. // GetStatsContext returns database statistics with context support.
// //
// The row counts (ASNs, prefixes, peerings, peers, live routes) come from the // The row counts (ASNs, prefixes, peerings, peers, live routes) and the prefix
// in-memory counters, seeded at startup and kept current on every write, so a // distribution come from the in-memory counters, seeded at startup and kept
// read runs no COUNT(*) over the tables. The oldest/newest route timestamps are // current on every write. The oldest/newest route timestamps are read from the
// read from the ends of the last_updated index, and the file size from a // ends of the last_updated index, and the file size from a stat(). No part of
// stat(); neither is a table scan. The only remaining query is the prefix // the read scans a table or a whole index.
// distribution.
func (d *Database) GetStatsContext(ctx context.Context) (Stats, error) { func (d *Database) GetStatsContext(ctx context.Context) (Stats, error) {
var stats Stats var stats Stats
// Row counts from memory, as a single consistent snapshot. // Row counts and prefix distribution from memory, as a single consistent
// snapshot.
d.counts.fill(&stats) d.counts.fill(&stats)
// Database file size is a cheap stat() on the file. // Database file size is a cheap stat() on the file.
@@ -1068,17 +1152,6 @@ func (d *Database) GetStatsContext(ctx context.Context) (Stats, error) {
stats.NewestRoute = newest stats.NewestRoute = newest
} }
// Prefix distribution counts distinct prefixes per mask length. It stays a
// query over the covering (mask_length, prefix) index rather than an
// in-memory counter: maintaining distinct-prefix-per-mask in memory would
// need a per-prefix table of roughly a million entries, memory this service
// is tuned to avoid.
stats.IPv4PrefixDistribution, stats.IPv6PrefixDistribution, err = d.GetPrefixDistributionContext(ctx)
if err != nil {
// Log but don't fail.
d.logger.Warn("Failed to get prefix distribution", "error", err)
}
return stats, nil return stats, nil
} }
@@ -1150,9 +1223,9 @@ func (d *Database) UpsertLiveRoute(route *LiveRoute) error {
return fmt.Errorf("failed to encode AS path: %w", err) return fmt.Errorf("failed to encode AS path: %w", err)
} }
updateSQL, insertSQL := updateLiveRouteV4SQL, insertLiveRouteV4SQL updateSQL, insertSQL, lookupSQL := updateLiveRouteV4SQL, insertLiveRouteV4SQL, prefixHasLiveRouteV4SQL
if route.IPVersion == ipVersionV6 { if route.IPVersion == ipVersionV6 {
updateSQL, insertSQL = updateLiveRouteV6SQL, insertLiveRouteV6SQL updateSQL, insertSQL, lookupSQL = updateLiveRouteV6SQL, insertLiveRouteV6SQL, prefixHasLiveRouteV6SQL
} }
// The write lock is held, so no other writer can insert this key between the // The write lock is held, so no other writer can insert this key between the
@@ -1169,11 +1242,17 @@ func (d *Database) UpsertLiveRoute(route *LiveRoute) error {
} }
defer func() { _ = ins.Close() }() defer func() { _ = ins.Close() }()
var inserted bool has, err := d.db.Prepare(lookupSQL)
if err != nil {
return fmt.Errorf("failed to prepare prefix lookup statement: %w", err)
}
defer func() { _ = has.Close() }()
var inserted, newPrefix bool
if route.IPVersion == ipVersionV4 { if route.IPVersion == ipVersionV4 {
inserted, err = upsertRouteRowV4(upd, ins, route, string(pathJSON)) inserted, newPrefix, err = upsertRouteRowV4(upd, ins, has, route, string(pathJSON))
} else { } else {
inserted, err = upsertRouteRowV6(upd, ins, route, string(pathJSON)) inserted, newPrefix, err = upsertRouteRowV6(upd, ins, has, route, string(pathJSON))
} }
if err != nil { if err != nil {
return fmt.Errorf("failed to upsert route %s: %w", route.Prefix, err) return fmt.Errorf("failed to upsert route %s: %w", route.Prefix, err)
@@ -1186,6 +1265,13 @@ func (d *Database) UpsertLiveRoute(route *LiveRoute) error {
d.counts.addRoutes(0, 1) d.counts.addRoutes(0, 1)
} }
} }
if newPrefix {
if route.IPVersion == ipVersionV4 {
d.counts.addToDistribution([]int{route.MaskLength}, nil, 1)
} else {
d.counts.addToDistribution(nil, []int{route.MaskLength}, 1)
}
}
return nil return nil
} }
@@ -1235,6 +1321,27 @@ func (d *Database) DeleteLiveRoute(prefix string, originASN int, peerIP string)
} else { } else {
d.counts.addRoutes(0, -int(affected)) d.counts.addRoutes(0, -int(affected))
} }
if affected == 0 {
return nil
}
// The prefix leaves the distribution when no live route has it any more.
lookupSQL := prefixHasLiveRouteV6SQL
if isV4 {
lookupSQL = prefixHasLiveRouteV4SQL
}
var prefixHasRoute bool
if err := d.db.QueryRow(lookupSQL, prefix).Scan(&prefixHasRoute); err != nil {
return err
}
if !prefixHasRoute {
maskLength, _ := ipnet.Mask.Size()
if isV4 {
d.counts.addToDistribution([]int{maskLength}, nil, -1)
} else {
d.counts.addToDistribution(nil, []int{maskLength}, -1)
}
}
return nil return nil
} }
@@ -1244,7 +1351,9 @@ func (d *Database) GetPrefixDistribution() (ipv4 []PrefixDistribution, ipv6 []Pr
return d.GetPrefixDistributionContext(context.Background()) return d.GetPrefixDistributionContext(context.Background())
} }
// GetPrefixDistributionContext returns the distribution of unique prefixes by mask length with context support // GetPrefixDistributionContext returns the distribution of unique prefixes by mask length with context support.
// It reads every live route, so the stats read does not call it; it seeds the
// in-memory distribution once at startup.
func (d *Database) GetPrefixDistributionContext(ctx context.Context) ( func (d *Database) GetPrefixDistributionContext(ctx context.Context) (
ipv4 []PrefixDistribution, ipv6 []PrefixDistribution, err error) { ipv4 []PrefixDistribution, ipv6 []PrefixDistribution, err error) {
// IPv4 distribution - count unique prefixes from v4 table // IPv4 distribution - count unique prefixes from v4 table
+13
View File
@@ -1,6 +1,8 @@
package database package database
import ( import (
"fmt"
"net"
"strings" "strings"
"github.com/google/uuid" "github.com/google/uuid"
@@ -18,3 +20,14 @@ func detectIPVersion(prefix string) int {
return ipVersionV4 return ipVersionV4
} }
// prefixMaskLength returns the mask length of a prefix such as 192.0.2.0/24.
func prefixMaskLength(prefix string) (int, error) {
_, network, err := net.ParseCIDR(prefix)
if err != nil {
return 0, fmt.Errorf("invalid prefix %s: %w", prefix, err)
}
maskLength, _ := network.Mask.Size()
return maskLength, nil
}