2 Commits
Author SHA1 Message Date
clawbot a9b68608c6 Serve the /api/v1/stats prefix distribution from memory (closes #30)
check / check (push) Successful in 3m19s
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 every live-route write. 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, within the same write. The single-route delete now runs its delete and that lookup in one transaction.

Model: opus-5-5
2026-10-03 18:43:58 +02:00
clawbot 44a5f4cdca Create new databases with auto_vacuum incremental (closes #43)
check / check (push) Successful in 2m50s
SQLite only accepts auto_vacuum before the database file is first written. The connection switched to WAL first, which writes the file, so the PRAGMA auto_vacuum in Initialize came too late and was ignored. The setting now goes in the connection string, which the driver applies on open before the journal mode, and the late PRAGMA is removed.

Vacuum now reads every row PRAGMA incremental_vacuum returns: SQLite frees one page per row, and the single step ExecContext took freed only one page per call.

Tests check that every pooled connection sees auto_vacuum incremental on a new database and that one Vacuum call frees every page left by deleting routes.

Model: opus-5-5
2026-10-03 17:09:29 +02:00
6 changed files with 526 additions and 63 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-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 - 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 incremental, through the connection string so it is set before the file
is first written, and each periodic incremental vacuum now returns up to 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 - 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
} }
+187 -4
View File
@@ -2,6 +2,7 @@ package database
import ( import (
"context" "context"
"slices"
"sync" "sync"
"testing" "testing"
"time" "time"
@@ -11,14 +12,20 @@ import (
"github.com/google/uuid" "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 { func mkV4Route(t *testing.T, prefix string, asn int, ts time.Time) *LiveRoute {
t.Helper() t.Helper()
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",
@@ -136,8 +143,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()}
@@ -177,6 +184,171 @@ 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 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 // 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 { 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 {
@@ -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) 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
View File
@@ -33,6 +33,8 @@ const (
dirPermissions = 0750 // rwxr-x--- dirPermissions = 0750 // rwxr-x---
ipVersionV4 = 4 ipVersionV4 = 4
ipVersionV6 = 6 ipVersionV6 = 6
ipv4Bits = 32
ipv6Bits = 128
) )
// SQLite memory tuning. cache_size and busy_timeout go in the DSN so every // 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 (?, ?, ?, ?, ?, ?, ?, ?)` 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) {
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
} }
// 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
@@ -332,7 +356,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 {
@@ -340,24 +378,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 {
@@ -365,6 +409,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
} }
@@ -415,8 +460,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
@@ -459,6 +518,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 {
@@ -466,6 +550,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
} }
@@ -1032,16 +1117,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.
@@ -1065,17 +1150,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
} }
@@ -1147,9 +1221,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
@@ -1166,11 +1240,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)
@@ -1183,6 +1263,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
} }
@@ -1201,21 +1288,33 @@ func (d *Database) DeleteLiveRoute(prefix string, originASN int, peerIP string)
isV4 := ipnet.IP.To4() != nil 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) // 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 // so the delete carries no dynamically built SQL. A delete with no origin
// ASN can remove several rows. // ASN can remove several rows.
var res sql.Result var res sql.Result
switch { switch {
case isV4 && originASN == 0: 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: case isV4:
res, err = d.db.Exec( res, err = tx.Exec(
`DELETE FROM live_routes_v4 WHERE prefix = ? AND origin_asn = ? AND peer_ip = ?`, `DELETE FROM live_routes_v4 WHERE prefix = ? AND origin_asn = ? AND peer_ip = ?`,
prefix, originASN, peerIP) prefix, originASN, peerIP)
case originASN == 0: 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: default:
res, err = d.db.Exec( res, err = tx.Exec(
`DELETE FROM live_routes_v6 WHERE prefix = ? AND origin_asn = ? AND peer_ip = ?`, `DELETE FROM live_routes_v6 WHERE prefix = ? AND origin_asn = ? AND peer_ip = ?`,
prefix, originASN, peerIP) prefix, originASN, peerIP)
} }
@@ -1227,11 +1326,37 @@ func (d *Database) DeleteLiveRoute(prefix string, originASN int, peerIP string)
if err != nil { if err != nil {
return err 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 { if isV4 {
d.counts.addRoutes(-int(affected), 0) d.counts.addRoutes(-int(affected), 0)
} else { } else {
d.counts.addRoutes(0, -int(affected)) 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 return nil
} }
@@ -1241,7 +1366,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
@@ -1268,6 +1395,10 @@ func (d *Database) GetPrefixDistributionContext(ctx context.Context) (
} }
ipv4 = append(ipv4, dist) 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 // IPv6 distribution - count unique prefixes from v6 table
query = ` query = `
@@ -1293,6 +1424,9 @@ func (d *Database) GetPrefixDistributionContext(ctx context.Context) (
} }
ipv6 = append(ipv6, dist) 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 return ipv4, ipv6, nil
} }
+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
}
+77
View File
@@ -2,9 +2,11 @@ package server
import ( import (
"context" "context"
"encoding/json"
"net/http" "net/http"
"net/http/httptest" "net/http/httptest"
"runtime" "runtime"
"slices"
"testing" "testing"
"time" "time"
@@ -13,6 +15,7 @@ import (
"git.eeqj.de/sneak/routewatch/internal/logger" "git.eeqj.de/sneak/routewatch/internal/logger"
"git.eeqj.de/sneak/routewatch/internal/metrics" "git.eeqj.de/sneak/routewatch/internal/metrics"
"git.eeqj.de/sneak/routewatch/internal/streamer" "git.eeqj.de/sneak/routewatch/internal/streamer"
"github.com/google/uuid"
) )
// blockingStatsDB embeds database.Store (left nil) and overrides only // 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 // settledGoroutineCount lets transient goroutines finish, then reports the
// current count. // current count.
func settledGoroutineCount() int { func settledGoroutineCount() int {