Compare commits
1
Commits
next
...
14c25ee238
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
14c25ee238 |
@@ -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
|
||||
(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).
|
||||
The other open issue is https://git.eeqj.de/sneak/routewatch/issues/30.
|
||||
|
||||
# Completed Steps
|
||||
|
||||
- 2026-10-03: `/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-10-03: a new database is created with `auto_vacuum` set to
|
||||
incremental, through the connection string so it is set before the file
|
||||
is first written, and each periodic incremental vacuum now returns up to
|
||||
@@ -118,6 +120,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
|
||||
container limit on fsn1app1 is open for sneak
|
||||
(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)
|
||||
|
||||
@@ -6,11 +6,12 @@ import (
|
||||
"sync"
|
||||
)
|
||||
|
||||
// liveCounts holds the running row counts that the stats endpoints report. They
|
||||
// are seeded once at startup from the tables and then adjusted on every write,
|
||||
// so a stats read serves them from memory instead of running a COUNT(*) over
|
||||
// each table. Those scans, once the database passed a few GiB, took the whole
|
||||
// request timeout and made /api/v1/stats return 500 (issue 27).
|
||||
// liveCounts holds the running row counts and the prefix distribution that the
|
||||
// stats endpoints report. They are seeded once at startup from the tables and
|
||||
// then adjusted on every write, so a stats read serves them from memory instead
|
||||
// of running a query over the tables. The COUNT(*) scans (issue 27) and then
|
||||
// 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
|
||||
// snapshot at one instant and writers, which already run under the database
|
||||
@@ -24,11 +25,16 @@ type liveCounts struct {
|
||||
peers int
|
||||
routesV4 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
|
||||
// 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()
|
||||
defer c.mu.Unlock()
|
||||
|
||||
@@ -39,6 +45,12 @@ func (c *liveCounts) seed(asns, prefixesV4, prefixesV6, peerings, peers, routesV
|
||||
c.peers = peers
|
||||
c.routesV4 = routesV4
|
||||
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.
|
||||
@@ -79,6 +91,30 @@ func (c *liveCounts) addRoutes(v4, v6 int) {
|
||||
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
|
||||
// single read lock so the reader sees one consistent snapshot.
|
||||
func (c *liveCounts) fill(s *Stats) {
|
||||
@@ -94,6 +130,21 @@ func (c *liveCounts) fill(s *Stats) {
|
||||
s.IPv4Routes = c.routesV4
|
||||
s.IPv6Routes = 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
|
||||
@@ -108,8 +159,9 @@ func (d *Database) countRows(ctx context.Context, table string) (int, error) {
|
||||
return n, nil
|
||||
}
|
||||
|
||||
// seedCounts reads the current row counts from the tables into the in-memory
|
||||
// counters. It runs once at startup, before the streamer begins writing.
|
||||
// seedCounts reads the current row counts and prefix distribution from the
|
||||
// tables into the in-memory counters. It runs once at startup, before the
|
||||
// streamer begins writing.
|
||||
func (d *Database) seedCounts(ctx context.Context) error {
|
||||
asns, err := d.countRows(ctx, "asns")
|
||||
if err != nil {
|
||||
@@ -139,8 +191,13 @@ func (d *Database) seedCounts(ctx context.Context) error {
|
||||
if err != nil {
|
||||
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
|
||||
}
|
||||
|
||||
@@ -2,6 +2,7 @@ package database
|
||||
|
||||
import (
|
||||
"context"
|
||||
"slices"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
@@ -11,14 +12,20 @@ import (
|
||||
"github.com/google/uuid"
|
||||
)
|
||||
|
||||
// mkV4Route builds an IPv4 live route.
|
||||
// mkV4Route builds an IPv4 live route with its mask length taken from the
|
||||
// prefix.
|
||||
func mkV4Route(t *testing.T, prefix string, asn int, ts time.Time) *LiveRoute {
|
||||
t.Helper()
|
||||
|
||||
maskLength, err := prefixMaskLength(prefix)
|
||||
if err != nil {
|
||||
t.Fatalf("prefixMaskLength(%s): %v", prefix, err)
|
||||
}
|
||||
|
||||
return &LiveRoute{
|
||||
ID: uuid.New(),
|
||||
Prefix: prefix,
|
||||
MaskLength: 24,
|
||||
MaskLength: maskLength,
|
||||
IPVersion: ipVersionV4,
|
||||
OriginASN: asn,
|
||||
PeerIP: "192.0.2.1",
|
||||
@@ -136,8 +143,8 @@ func TestLiveCountsTrackWritesInRealtime(t *testing.T) {
|
||||
}
|
||||
|
||||
// TestLiveCountsSeededFromDatabaseAtStartup writes rows, reopens the same
|
||||
// database file, and checks the counts come back from the seed scan rather than
|
||||
// starting at zero.
|
||||
// database file, and checks the counts and the prefix distribution come back
|
||||
// from the seed scan rather than starting at zero.
|
||||
func TestLiveCountsSeededFromDatabaseAtStartup(t *testing.T) {
|
||||
cfg := &config.Config{StateDir: t.TempDir()}
|
||||
|
||||
@@ -177,6 +184,171 @@ func TestLiveCountsSeededFromDatabaseAtStartup(t *testing.T) {
|
||||
t.Errorf("seeded routes = (v4 %d, v6 %d, total %d), want (1, 1, 2)",
|
||||
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 live-route write, and that
|
||||
// after every step it equals what the distribution query reads from the route
|
||||
// tables. The steps run once through the batch methods the prefix handler uses
|
||||
// and once through the single-route methods.
|
||||
func TestPrefixDistributionTracksWrites(t *testing.T) {
|
||||
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)
|
||||
v6SecondPeer := mkV6Route("2001:db8::/32", 64503, ts)
|
||||
v6SecondPeer.PeerIP = "2001:db8::2"
|
||||
|
||||
all := []PrefixDistribution{{MaskLength: 16, Count: 1}, {MaskLength: 24, Count: 2}}
|
||||
v6Only := []PrefixDistribution{{MaskLength: 32, Count: 1}}
|
||||
|
||||
steps := []struct {
|
||||
name string
|
||||
announce []*LiveRoute
|
||||
withdraw []*LiveRoute
|
||||
wantV4 []PrefixDistribution
|
||||
wantV6 []PrefixDistribution
|
||||
}{
|
||||
{
|
||||
name: "new routes",
|
||||
announce: []*LiveRoute{shared, other, wide, v6},
|
||||
wantV4: all, wantV6: v6Only,
|
||||
},
|
||||
{
|
||||
name: "re-announcement",
|
||||
announce: []*LiveRoute{shared, other, wide, v6},
|
||||
wantV4: all, wantV6: v6Only,
|
||||
},
|
||||
{
|
||||
name: "second peer announces a prefix that has a live route",
|
||||
announce: []*LiveRoute{sharedSecondPeer, v6SecondPeer},
|
||||
wantV4: all, wantV6: v6Only,
|
||||
},
|
||||
{
|
||||
name: "withdrawal of a route that is not the last for its prefix",
|
||||
withdraw: []*LiveRoute{shared, v6},
|
||||
wantV4: all, wantV6: v6Only,
|
||||
},
|
||||
{
|
||||
name: "withdrawal of the last route for a prefix",
|
||||
withdraw: []*LiveRoute{sharedSecondPeer, v6SecondPeer},
|
||||
wantV4: []PrefixDistribution{{MaskLength: 16, Count: 1}, {MaskLength: 24, Count: 1}},
|
||||
},
|
||||
{
|
||||
name: "withdrawal of every remaining route",
|
||||
withdraw: []*LiveRoute{other, wide},
|
||||
},
|
||||
{
|
||||
name: "two peers announce a new prefix together",
|
||||
announce: []*LiveRoute{shared, sharedSecondPeer, v6, v6SecondPeer},
|
||||
wantV4: []PrefixDistribution{{MaskLength: 24, Count: 1}},
|
||||
wantV6: v6Only,
|
||||
},
|
||||
{
|
||||
name: "both routes for a prefix withdrawn together",
|
||||
withdraw: []*LiveRoute{shared, sharedSecondPeer, v6, v6SecondPeer},
|
||||
},
|
||||
// The feed often withdraws a route that is not live. That must not take
|
||||
// the prefix out of the distribution. A count wrongly taken below zero is
|
||||
// left out of the answer, so the next step shows it: its announcement
|
||||
// at the same mask length would then not be counted.
|
||||
{
|
||||
name: "withdrawal of routes that are not live",
|
||||
withdraw: []*LiveRoute{other, v6},
|
||||
},
|
||||
{
|
||||
name: "announcement after a withdrawal of routes that are not live",
|
||||
announce: []*LiveRoute{other, v6},
|
||||
wantV4: []PrefixDistribution{{MaskLength: 24, Count: 1}},
|
||||
wantV6: v6Only,
|
||||
},
|
||||
}
|
||||
|
||||
for _, batch := range []bool{true, false} {
|
||||
name := "single-route writes"
|
||||
if batch {
|
||||
name = "batch writes"
|
||||
}
|
||||
|
||||
t.Run(name, func(t *testing.T) {
|
||||
db, err := New(&config.Config{StateDir: t.TempDir()}, logger.New())
|
||||
if err != nil {
|
||||
t.Fatalf("failed to create database: %v", err)
|
||||
}
|
||||
defer func() { _ = db.Close() }()
|
||||
|
||||
ctx := context.Background()
|
||||
for _, step := range steps {
|
||||
if err := announceRoutes(db, batch, step.announce); err != nil {
|
||||
t.Fatalf("%s: announce: %v", step.name, err)
|
||||
}
|
||||
if err := withdrawRoutes(db, batch, step.withdraw); err != nil {
|
||||
t.Fatalf("%s: withdraw: %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)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// announceRoutes writes routes in one UpsertLiveRouteBatch, or with one
|
||||
// UpsertLiveRoute each.
|
||||
func announceRoutes(db *Database, batch bool, routes []*LiveRoute) error {
|
||||
if batch {
|
||||
return db.UpsertLiveRouteBatch(routes)
|
||||
}
|
||||
for _, route := range routes {
|
||||
if err := db.UpsertLiveRoute(route); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// withdrawRoutes removes routes in one DeleteLiveRouteBatch, or with one
|
||||
// DeleteLiveRoute each. It names each route by prefix and peer only, with no
|
||||
// origin ASN, as a withdrawal from the feed does when its message carries no AS
|
||||
// path.
|
||||
func withdrawRoutes(db *Database, batch bool, routes []*LiveRoute) error {
|
||||
if !batch {
|
||||
for _, route := range routes {
|
||||
if err := db.DeleteLiveRoute(route.Prefix, 0, route.PeerIP); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
deletions := make([]LiveRouteDeletion, 0, len(routes))
|
||||
for _, route := range routes {
|
||||
deletions = append(deletions, LiveRouteDeletion{
|
||||
Prefix: route.Prefix, PeerIP: route.PeerIP, IPVersion: route.IPVersion,
|
||||
})
|
||||
}
|
||||
|
||||
return db.DeleteLiveRouteBatch(deletions)
|
||||
}
|
||||
|
||||
// TestStatsRouteTimestamps checks the oldest/newest route timestamps are read
|
||||
@@ -289,6 +461,9 @@ func TestLiveCountsConcurrentReadWrite(t *testing.T) {
|
||||
if want := writers * 25; 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 {
|
||||
@@ -326,3 +501,11 @@ func assertCounts(t *testing.T, when string, got Stats, want wantCounts) {
|
||||
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)
|
||||
}
|
||||
}
|
||||
|
||||
+180
-46
@@ -33,6 +33,8 @@ const (
|
||||
dirPermissions = 0750 // rwxr-x---
|
||||
ipVersionV4 = 4
|
||||
ipVersionV6 = 6
|
||||
ipv4Bits = 32
|
||||
ipv6Bits = 128
|
||||
)
|
||||
|
||||
// SQLite memory tuning. cache_size and busy_timeout go in the DSN so every
|
||||
@@ -238,54 +240,76 @@ const (
|
||||
as_path, next_hop, last_updated) VALUES (?, ?, ?, ?, ?, ?, ?, ?)`
|
||||
)
|
||||
|
||||
// upsertRouteRowV4 updates an IPv4 live route, inserting it when no row matched,
|
||||
// and reports whether a new row was inserted.
|
||||
func upsertRouteRowV4(upd, ins *sql.Stmt, route *LiveRoute, pathJSON string) (inserted bool, err error) {
|
||||
// Before a new route is inserted and after a route is deleted, the write looks
|
||||
// up whether any live route has that prefix, to keep the in-memory prefix
|
||||
// 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) {
|
||||
res, err := upd.Exec(route.MaskLength, pathJSON, route.NextHop, route.LastUpdated,
|
||||
route.Prefix, route.OriginASN, route.PeerIP)
|
||||
if err != nil {
|
||||
return false, err
|
||||
return false, false, err
|
||||
}
|
||||
affected, err := res.RowsAffected()
|
||||
if err != nil {
|
||||
return false, err
|
||||
return false, false, err
|
||||
}
|
||||
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,
|
||||
route.PeerIP, pathJSON, route.NextHop, route.LastUpdated)
|
||||
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,
|
||||
// and reports whether a new row was inserted.
|
||||
func upsertRouteRowV6(upd, ins *sql.Stmt, route *LiveRoute, pathJSON string) (inserted bool, err error) {
|
||||
// upsertRouteRowV6 updates an IPv6 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 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,
|
||||
route.Prefix, route.OriginASN, route.PeerIP)
|
||||
if err != nil {
|
||||
return false, err
|
||||
return false, false, err
|
||||
}
|
||||
affected, err := res.RowsAffected()
|
||||
if err != nil {
|
||||
return false, err
|
||||
return false, false, err
|
||||
}
|
||||
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,
|
||||
route.PeerIP, pathJSON, route.NextHop, route.LastUpdated)
|
||||
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
|
||||
@@ -332,7 +356,21 @@ func (d *Database) UpsertLiveRouteBatch(routes []*LiveRoute) error {
|
||||
}
|
||||
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
|
||||
// Mask lengths of the prefixes that get their first live route in this batch.
|
||||
var newPrefixMaskLengthsV4, newPrefixMaskLengthsV6 []int
|
||||
for _, route := range routes {
|
||||
pathJSON, err := json.Marshal(route.ASPath)
|
||||
if err != nil {
|
||||
@@ -340,24 +378,30 @@ func (d *Database) UpsertLiveRouteBatch(routes []*LiveRoute) error {
|
||||
}
|
||||
|
||||
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 {
|
||||
return fmt.Errorf("failed to upsert route %s: %w", route.Prefix, err)
|
||||
}
|
||||
if inserted {
|
||||
newV4++
|
||||
}
|
||||
if newPrefix {
|
||||
newPrefixMaskLengthsV4 = append(newPrefixMaskLengthsV4, route.MaskLength)
|
||||
}
|
||||
|
||||
continue
|
||||
}
|
||||
|
||||
inserted, err := upsertRouteRowV6(updV6, insV6, route, string(pathJSON))
|
||||
inserted, newPrefix, err := upsertRouteRowV6(updV6, insV6, hasV6, route, string(pathJSON))
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to upsert route %s: %w", route.Prefix, err)
|
||||
}
|
||||
if inserted {
|
||||
newV6++
|
||||
}
|
||||
if newPrefix {
|
||||
newPrefixMaskLengthsV6 = append(newPrefixMaskLengthsV6, route.MaskLength)
|
||||
}
|
||||
}
|
||||
|
||||
if err = tx.Commit(); err != nil {
|
||||
@@ -365,6 +409,7 @@ func (d *Database) UpsertLiveRouteBatch(routes []*LiveRoute) error {
|
||||
}
|
||||
|
||||
d.counts.addRoutes(newV4, newV6)
|
||||
d.counts.addToDistribution(newPrefixMaskLengthsV4, newPrefixMaskLengthsV6, 1)
|
||||
|
||||
return nil
|
||||
}
|
||||
@@ -415,8 +460,22 @@ func (d *Database) DeleteLiveRouteBatch(deletions []LiveRouteDeletion) error {
|
||||
}
|
||||
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
|
||||
var deletedV4, deletedV6 int64
|
||||
// Mask lengths of the prefixes this batch leaves with no live route.
|
||||
var gonePrefixMaskLengthsV4, gonePrefixMaskLengthsV6 []int
|
||||
for _, del := range deletions {
|
||||
var stmt *sql.Stmt
|
||||
|
||||
@@ -459,6 +518,31 @@ func (d *Database) DeleteLiveRouteBatch(deletions []LiveRouteDeletion) error {
|
||||
} else {
|
||||
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 {
|
||||
@@ -466,6 +550,7 @@ func (d *Database) DeleteLiveRouteBatch(deletions []LiveRouteDeletion) error {
|
||||
}
|
||||
|
||||
d.counts.addRoutes(-int(deletedV4), -int(deletedV6))
|
||||
d.counts.addToDistribution(gonePrefixMaskLengthsV4, gonePrefixMaskLengthsV6, -1)
|
||||
|
||||
return nil
|
||||
}
|
||||
@@ -1032,16 +1117,16 @@ func (d *Database) GetStats() (Stats, error) {
|
||||
|
||||
// GetStatsContext returns database statistics with context support.
|
||||
//
|
||||
// The row counts (ASNs, prefixes, peerings, peers, live routes) come from the
|
||||
// in-memory counters, seeded at startup and kept current on every write, so a
|
||||
// read runs no COUNT(*) over the tables. The oldest/newest route timestamps are
|
||||
// read from the ends of the last_updated index, and the file size from a
|
||||
// stat(); neither is a table scan. The only remaining query is the prefix
|
||||
// distribution.
|
||||
// The row counts (ASNs, prefixes, peerings, peers, live routes) and the prefix
|
||||
// distribution come from the in-memory counters, seeded at startup and kept
|
||||
// current on every write. The oldest/newest route timestamps are read from the
|
||||
// ends of the last_updated index, and the file size from a stat(). No part of
|
||||
// the read scans a table or a whole index.
|
||||
func (d *Database) GetStatsContext(ctx context.Context) (Stats, error) {
|
||||
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)
|
||||
|
||||
// Database file size is a cheap stat() on the file.
|
||||
@@ -1065,17 +1150,6 @@ func (d *Database) GetStatsContext(ctx context.Context) (Stats, error) {
|
||||
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
|
||||
}
|
||||
|
||||
@@ -1147,9 +1221,9 @@ func (d *Database) UpsertLiveRoute(route *LiveRoute) error {
|
||||
return fmt.Errorf("failed to encode AS path: %w", err)
|
||||
}
|
||||
|
||||
updateSQL, insertSQL := updateLiveRouteV4SQL, insertLiveRouteV4SQL
|
||||
updateSQL, insertSQL, lookupSQL := updateLiveRouteV4SQL, insertLiveRouteV4SQL, prefixHasLiveRouteV4SQL
|
||||
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
|
||||
@@ -1166,11 +1240,17 @@ func (d *Database) UpsertLiveRoute(route *LiveRoute) error {
|
||||
}
|
||||
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 {
|
||||
inserted, err = upsertRouteRowV4(upd, ins, route, string(pathJSON))
|
||||
inserted, newPrefix, err = upsertRouteRowV4(upd, ins, has, route, string(pathJSON))
|
||||
} else {
|
||||
inserted, err = upsertRouteRowV6(upd, ins, route, string(pathJSON))
|
||||
inserted, newPrefix, err = upsertRouteRowV6(upd, ins, has, route, string(pathJSON))
|
||||
}
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to upsert route %s: %w", route.Prefix, err)
|
||||
@@ -1183,6 +1263,13 @@ func (d *Database) UpsertLiveRoute(route *LiveRoute) error {
|
||||
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
|
||||
}
|
||||
@@ -1201,21 +1288,33 @@ func (d *Database) DeleteLiveRoute(prefix string, originASN int, peerIP string)
|
||||
|
||||
isV4 := ipnet.IP.To4() != nil
|
||||
|
||||
// The delete and the prefix lookup after it share one transaction, so the
|
||||
// in-memory counts change only once both have succeeded.
|
||||
tx, err := d.beginTx()
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to begin transaction: %w", err)
|
||||
}
|
||||
defer func() {
|
||||
if err := tx.Rollback(); err != nil && err != sql.ErrTxDone {
|
||||
d.logger.Error("Failed to rollback transaction", "error", err)
|
||||
}
|
||||
}()
|
||||
|
||||
// Literal per-table queries (rather than one formatted with the table name)
|
||||
// so the delete carries no dynamically built SQL. A delete with no origin
|
||||
// ASN can remove several rows.
|
||||
var res sql.Result
|
||||
switch {
|
||||
case isV4 && originASN == 0:
|
||||
res, err = d.db.Exec(`DELETE FROM live_routes_v4 WHERE prefix = ? AND peer_ip = ?`, prefix, peerIP)
|
||||
res, err = tx.Exec(`DELETE FROM live_routes_v4 WHERE prefix = ? AND peer_ip = ?`, prefix, peerIP)
|
||||
case isV4:
|
||||
res, err = d.db.Exec(
|
||||
res, err = tx.Exec(
|
||||
`DELETE FROM live_routes_v4 WHERE prefix = ? AND origin_asn = ? AND peer_ip = ?`,
|
||||
prefix, originASN, peerIP)
|
||||
case originASN == 0:
|
||||
res, err = d.db.Exec(`DELETE FROM live_routes_v6 WHERE prefix = ? AND peer_ip = ?`, prefix, peerIP)
|
||||
res, err = tx.Exec(`DELETE FROM live_routes_v6 WHERE prefix = ? AND peer_ip = ?`, prefix, peerIP)
|
||||
default:
|
||||
res, err = d.db.Exec(
|
||||
res, err = tx.Exec(
|
||||
`DELETE FROM live_routes_v6 WHERE prefix = ? AND origin_asn = ? AND peer_ip = ?`,
|
||||
prefix, originASN, peerIP)
|
||||
}
|
||||
@@ -1227,11 +1326,37 @@ func (d *Database) DeleteLiveRoute(prefix string, originASN int, peerIP string)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
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 := tx.QueryRow(lookupSQL, prefix).Scan(&prefixHasRoute); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if err := tx.Commit(); err != nil {
|
||||
return fmt.Errorf("failed to commit transaction: %w", err)
|
||||
}
|
||||
|
||||
if isV4 {
|
||||
d.counts.addRoutes(-int(affected), 0)
|
||||
} else {
|
||||
d.counts.addRoutes(0, -int(affected))
|
||||
}
|
||||
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
|
||||
}
|
||||
@@ -1241,7 +1366,9 @@ func (d *Database) GetPrefixDistribution() (ipv4 []PrefixDistribution, ipv6 []Pr
|
||||
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) (
|
||||
ipv4 []PrefixDistribution, ipv6 []PrefixDistribution, err error) {
|
||||
// IPv4 distribution - count unique prefixes from v4 table
|
||||
@@ -1268,6 +1395,10 @@ func (d *Database) GetPrefixDistributionContext(ctx context.Context) (
|
||||
}
|
||||
ipv4 = append(ipv4, dist)
|
||||
}
|
||||
// A read that stops partway ends the loop without an error of its own.
|
||||
if err := rows4.Err(); err != nil {
|
||||
return nil, nil, fmt.Errorf("failed to read IPv4 distribution: %w", err)
|
||||
}
|
||||
|
||||
// IPv6 distribution - count unique prefixes from v6 table
|
||||
query = `
|
||||
@@ -1293,6 +1424,9 @@ func (d *Database) GetPrefixDistributionContext(ctx context.Context) (
|
||||
}
|
||||
ipv6 = append(ipv6, dist)
|
||||
}
|
||||
if err := rows6.Err(); err != nil {
|
||||
return nil, nil, fmt.Errorf("failed to read IPv6 distribution: %w", err)
|
||||
}
|
||||
|
||||
return ipv4, ipv6, nil
|
||||
}
|
||||
|
||||
@@ -1,6 +1,8 @@
|
||||
package database
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"net"
|
||||
"strings"
|
||||
|
||||
"github.com/google/uuid"
|
||||
@@ -18,3 +20,14 @@ func detectIPVersion(prefix string) int {
|
||||
|
||||
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
|
||||
}
|
||||
|
||||
@@ -2,9 +2,11 @@ package server
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"runtime"
|
||||
"slices"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
@@ -13,6 +15,7 @@ import (
|
||||
"git.eeqj.de/sneak/routewatch/internal/logger"
|
||||
"git.eeqj.de/sneak/routewatch/internal/metrics"
|
||||
"git.eeqj.de/sneak/routewatch/internal/streamer"
|
||||
"github.com/google/uuid"
|
||||
)
|
||||
|
||||
// blockingStatsDB embeds database.Store (left nil) and overrides only
|
||||
@@ -71,6 +74,80 @@ func TestStatsHandlersDoNotLeakOnTimeout(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// TestStatsHandlersAnswerFromMemory checks that both stats handlers answer 200
|
||||
// with the live route counts and the prefix distribution while the database is
|
||||
// closed, so that any query would fail: the request path reads them from
|
||||
// memory. The prefix distribution query it used to run read every live route
|
||||
// and, on a large database, took the whole 4-second deadline, so
|
||||
// /api/v1/stats answered 500 (https://git.eeqj.de/sneak/routewatch/issues/30).
|
||||
// The oldest and newest route times still come from one-row lookups at the ends
|
||||
// of an index; with the database closed they are left out of the answer.
|
||||
func TestStatsHandlersAnswerFromMemory(t *testing.T) {
|
||||
db, err := database.New(&config.Config{StateDir: t.TempDir()}, logger.New())
|
||||
if err != nil {
|
||||
t.Fatalf("database.New: %v", err)
|
||||
}
|
||||
|
||||
ts := time.Date(2026, 1, 2, 3, 4, 5, 0, time.UTC)
|
||||
if err := db.UpsertLiveRouteBatch([]*database.LiveRoute{
|
||||
{
|
||||
ID: uuid.New(), Prefix: "198.51.100.0/24", MaskLength: 24, IPVersion: 4,
|
||||
OriginASN: 64500, PeerIP: "192.0.2.1", ASPath: []int{64500}, NextHop: "192.0.2.1",
|
||||
LastUpdated: ts,
|
||||
},
|
||||
{
|
||||
ID: uuid.New(), Prefix: "2001:db8::/32", MaskLength: 32, IPVersion: 6,
|
||||
OriginASN: 64501, PeerIP: "2001:db8::1", ASPath: []int{64501}, NextHop: "2001:db8::1",
|
||||
LastUpdated: ts,
|
||||
},
|
||||
}); err != nil {
|
||||
t.Fatalf("UpsertLiveRouteBatch: %v", err)
|
||||
}
|
||||
if err := db.Close(); err != nil {
|
||||
t.Fatalf("Close: %v", err)
|
||||
}
|
||||
|
||||
s := New(db, streamer.New(logger.New(), metrics.New()), logger.New(), &config.Config{})
|
||||
handlers := map[string]http.HandlerFunc{
|
||||
"status.json": s.handleStatusJSON(),
|
||||
"stats": s.handleStats(),
|
||||
}
|
||||
|
||||
for name, handler := range handlers {
|
||||
rec := httptest.NewRecorder()
|
||||
handler(rec, httptest.NewRequest(http.MethodGet, "/", nil))
|
||||
if rec.Code != http.StatusOK {
|
||||
t.Errorf("%s: status %d, want %d; body %s", name, rec.Code, http.StatusOK, rec.Body)
|
||||
|
||||
continue
|
||||
}
|
||||
|
||||
var body struct {
|
||||
Data struct {
|
||||
IPv4Routes int `json:"ipv4_routes"`
|
||||
IPv6Routes int `json:"ipv6_routes"`
|
||||
IPv4PrefixDistribution []database.PrefixDistribution `json:"ipv4_prefix_distribution"`
|
||||
IPv6PrefixDistribution []database.PrefixDistribution `json:"ipv6_prefix_distribution"`
|
||||
} `json:"data"`
|
||||
}
|
||||
if err := json.Unmarshal(rec.Body.Bytes(), &body); err != nil {
|
||||
t.Fatalf("%s: decoding the answer: %v", name, err)
|
||||
}
|
||||
|
||||
if body.Data.IPv4Routes != 1 || body.Data.IPv6Routes != 1 {
|
||||
t.Errorf("%s: routes = (v4 %d, v6 %d), want (1, 1)", name, body.Data.IPv4Routes, body.Data.IPv6Routes)
|
||||
}
|
||||
wantV4 := []database.PrefixDistribution{{MaskLength: 24, Count: 1}}
|
||||
if !slices.Equal(body.Data.IPv4PrefixDistribution, wantV4) {
|
||||
t.Errorf("%s: IPv4 distribution = %v, want %v", name, body.Data.IPv4PrefixDistribution, wantV4)
|
||||
}
|
||||
wantV6 := []database.PrefixDistribution{{MaskLength: 32, Count: 1}}
|
||||
if !slices.Equal(body.Data.IPv6PrefixDistribution, wantV6) {
|
||||
t.Errorf("%s: IPv6 distribution = %v, want %v", name, body.Data.IPv6PrefixDistribution, wantV6)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// settledGoroutineCount lets transient goroutines finish, then reports the
|
||||
// current count.
|
||||
func settledGoroutineCount() int {
|
||||
|
||||
Reference in New Issue
Block a user