diff --git a/TODO.md b/TODO.md index b329b6e..6547430 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-02: a plain `docker build .` stamps the commit's tag or short commit (`git describe --tags --always`) into the page footer instead of `unknown`: `.dockerignore` sends `.git` without `.git/config`, a `VERSION` @@ -108,6 +110,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 6db2859..e4e3e0f 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,7 +12,8 @@ import ( "github.com/google/uuid" ) -// mkV4Route builds an IPv4 live route with its range columns populated. +// mkV4Route builds an IPv4 live route with its mask length and range columns +// taken from the prefix. func mkV4Route(t *testing.T, prefix string, asn int, ts time.Time) *LiveRoute { t.Helper() @@ -19,11 +21,15 @@ func mkV4Route(t *testing.T, prefix string, asn int, ts time.Time) *LiveRoute { if err != nil { t.Fatalf("CalculateIPv4Range(%s): %v", prefix, err) } + 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", @@ -143,8 +149,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()} @@ -184,6 +190,155 @@ 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) + + 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}, + wantV4: all, wantV6: v6Only, + }, + { + name: "withdrawal of a route that is not the last for its prefix", + withdraw: []*LiveRoute{shared}, + wantV4: all, wantV6: v6Only, + }, + { + name: "withdrawal of the last route for a prefix", + withdraw: []*LiveRoute{sharedSecondPeer}, + wantV4: []PrefixDistribution{{MaskLength: 16, Count: 1}, {MaskLength: 24, Count: 1}}, + wantV6: v6Only, + }, + { + name: "withdrawal of every remaining route", + withdraw: []*LiveRoute{other, wide, v6}, + }, + { + name: "two peers announce a new prefix together", + announce: []*LiveRoute{shared, sharedSecondPeer}, + wantV4: []PrefixDistribution{{MaskLength: 24, Count: 1}}, + }, + { + name: "both routes for a prefix withdrawn together", + withdraw: []*LiveRoute{shared, sharedSecondPeer}, + }, + } + + 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 @@ -296,6 +451,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 { @@ -333,3 +491,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 4007b48..88074b7 100644 --- a/internal/database/database.go +++ b/internal/database/database.go @@ -35,6 +35,7 @@ const ( ipv6Length = 16 ipv4Offset = 12 ipv4Bits = 32 + ipv6Bits = 128 maxIPv4 = 0xFFFFFFFF ) @@ -237,58 +238,80 @@ 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) { if route.V4IPStart == nil || route.V4IPEnd == nil { - return false, fmt.Errorf("IPv4 route %s missing range values", route.Prefix) + return false, false, fmt.Errorf("IPv4 route %s missing range values", route.Prefix) } res, err := upd.Exec(route.MaskLength, pathJSON, route.NextHop, route.LastUpdated, *route.V4IPStart, *route.V4IPEnd, 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, *route.V4IPStart, *route.V4IPEnd) 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 @@ -335,7 +358,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 { @@ -343,24 +380,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 { @@ -368,6 +411,7 @@ func (d *Database) UpsertLiveRouteBatch(routes []*LiveRoute) error { } d.counts.addRoutes(newV4, newV6) + d.counts.addToDistribution(newPrefixMaskLengthsV4, newPrefixMaskLengthsV6, 1) return nil } @@ -418,8 +462,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 @@ -462,6 +520,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 { @@ -469,6 +552,7 @@ func (d *Database) DeleteLiveRouteBatch(deletions []LiveRouteDeletion) error { } d.counts.addRoutes(-int(deletedV4), -int(deletedV6)) + d.counts.addToDistribution(gonePrefixMaskLengthsV4, gonePrefixMaskLengthsV6, -1) return nil } @@ -1035,16 +1119,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. @@ -1068,17 +1152,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 } @@ -1150,9 +1223,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 @@ -1169,11 +1242,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) @@ -1186,6 +1265,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 } @@ -1204,21 +1290,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) } @@ -1230,11 +1328,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 } @@ -1244,7 +1368,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 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..0658d4a 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,84 @@ 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) + start, end, err := database.CalculateIPv4Range("198.51.100.0/24") + if err != nil { + t.Fatalf("CalculateIPv4Range: %v", err) + } + 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, V4IPStart: &start, V4IPEnd: &end, + }, + { + 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 {