From 14c25ee238a9cfdd4345d1687cf58f22adf883a5 Mon Sep 17 00:00:00 2001 From: sneak Date: Tue, 29 Sep 2026 11:05:28 +0000 Subject: [PATCH] Serve the /api/v1/stats prefix distribution from memory (closes #30) 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, made only for new and removed routes, within the same write. The single-route delete now runs its delete and that lookup in one transaction, so a failed lookup cannot leave the counts off. Model: opus-5-5 --- TODO.md | 7 +- internal/database/counts.go | 75 ++++++++-- internal/database/counts_test.go | 191 +++++++++++++++++++++++++- internal/database/database.go | 226 ++++++++++++++++++++++++------- internal/database/utils.go | 13 ++ internal/server/handlers_test.go | 77 +++++++++++ 6 files changed, 526 insertions(+), 63 deletions(-) diff --git a/TODO.md b/TODO.md index d3915ca..d0b7920 100644 --- a/TODO.md +++ b/TODO.md @@ -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) diff --git a/internal/database/counts.go b/internal/database/counts.go index a305ba2..71ce82f 100644 --- a/internal/database/counts.go +++ b/internal/database/counts.go @@ -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 } diff --git a/internal/database/counts_test.go b/internal/database/counts_test.go index 65cfec8..8922619 100644 --- a/internal/database/counts_test.go +++ b/internal/database/counts_test.go @@ -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) + } +} diff --git a/internal/database/database.go b/internal/database/database.go index 6dfddbf..5393f74 100644 --- a/internal/database/database.go +++ b/internal/database/database.go @@ -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 } diff --git a/internal/database/utils.go b/internal/database/utils.go index cf1ff34..8094c78 100644 --- a/internal/database/utils.go +++ b/internal/database/utils.go @@ -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 +} diff --git a/internal/server/handlers_test.go b/internal/server/handlers_test.go index 46756bd..5bff6d4 100644 --- a/internal/server/handlers_test.go +++ b/internal/server/handlers_test.go @@ -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 { -- 2.54.0