diff --git a/TODO.md b/TODO.md index b967d9a..18b3936 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-09-29: `/api/v1/stats` serves the prefix distribution from memory, + seeded at startup and adjusted on every live-route write, so a request no + longer reads every live route (closes #30) - 2026-09-29: the entrypoint creates the data directory if it is missing and stops the start if a step fails; README "Running under upaas" no longer asks for the host directory to be created first (closes #42) @@ -103,6 +105,3 @@ The other open issue is https://git.eeqj.de/sneak/routewatch/issues/30. - Production memory under 5 GiB: whether to test under a real 5 GiB 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..49f8e77 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,145 @@ 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 write, through both the +// batch and the single-route methods, and that after every step it equals what +// the distribution query reads from the route tables. +func TestPrefixDistributionTracksWrites(t *testing.T) { + cfg := &config.Config{StateDir: t.TempDir()} + + db, err := New(cfg, logger.New()) + if err != nil { + t.Fatalf("failed to create database: %v", err) + } + defer func() { _ = db.Close() }() + + ts := time.Date(2026, 1, 2, 3, 4, 5, 0, time.UTC) + shared := mkV4Route(t, "198.51.100.0/24", 64500, ts) + sharedSecondPeer := mkV4Route(t, "198.51.100.0/24", 64500, ts) + sharedSecondPeer.PeerIP = "192.0.2.2" + other := mkV4Route(t, "203.0.113.0/24", 64501, ts) + wide := mkV4Route(t, "172.16.0.0/16", 64502, ts) + v6 := mkV6Route("2001:db8::/32", 64503, ts) + + all := []PrefixDistribution{{MaskLength: 16, Count: 1}, {MaskLength: 24, Count: 2}} + v6Only := []PrefixDistribution{{MaskLength: 32, Count: 1}} + + steps := []struct { + name string + write func() error + wantV4 []PrefixDistribution + wantV6 []PrefixDistribution + }{ + { + name: "empty database", + write: func() error { return nil }, + }, + { + name: "new prefixes", + write: func() error { return db.UpsertLiveRouteBatch([]*LiveRoute{shared, other, wide, v6}) }, + wantV4: all, wantV6: v6Only, + }, + { + name: "re-announcement", + write: func() error { return db.UpsertLiveRouteBatch([]*LiveRoute{shared, other, wide, v6}) }, + wantV4: all, wantV6: v6Only, + }, + { + name: "second peer announces a prefix that already has a route", + write: func() error { return db.UpsertLiveRoute(sharedSecondPeer) }, + wantV4: all, wantV6: v6Only, + }, + { + name: "withdrawal of a route that is not the last for its prefix", + write: func() error { + return db.DeleteLiveRouteBatch([]LiveRouteDeletion{ + {Prefix: shared.Prefix, OriginASN: shared.OriginASN, PeerIP: shared.PeerIP, IPVersion: ipVersionV4}, + }) + }, + wantV4: all, wantV6: v6Only, + }, + { + name: "withdrawal of the last route for a prefix", + write: func() error { + return db.DeleteLiveRoute(sharedSecondPeer.Prefix, sharedSecondPeer.OriginASN, sharedSecondPeer.PeerIP) + }, + wantV4: []PrefixDistribution{{MaskLength: 16, Count: 1}, {MaskLength: 24, Count: 1}}, + wantV6: v6Only, + }, + { + name: "withdrawal of every remaining route, one without an origin ASN", + write: func() error { + return db.DeleteLiveRouteBatch([]LiveRouteDeletion{ + {Prefix: other.Prefix, PeerIP: other.PeerIP, IPVersion: ipVersionV4}, + {Prefix: wide.Prefix, OriginASN: wide.OriginASN, PeerIP: wide.PeerIP, IPVersion: ipVersionV4}, + {Prefix: v6.Prefix, OriginASN: v6.OriginASN, PeerIP: v6.PeerIP, IPVersion: ipVersionV6}, + }) + }, + }, + } + + ctx := context.Background() + for _, step := range steps { + if err := step.write(); err != nil { + t.Fatalf("%s: %v", step.name, err) + } + + stats, err := db.GetStatsContext(ctx) + if err != nil { + t.Fatalf("%s: GetStatsContext: %v", step.name, err) + } + queryV4, queryV6, err := db.GetPrefixDistributionContext(ctx) + if err != nil { + t.Fatalf("%s: GetPrefixDistributionContext: %v", step.name, err) + } + + assertDistribution(t, step.name+": IPv4 distribution", stats.IPv4PrefixDistribution, step.wantV4) + assertDistribution(t, step.name+": IPv6 distribution", stats.IPv6PrefixDistribution, step.wantV6) + assertDistribution(t, step.name+": IPv4 distribution query", queryV4, step.wantV4) + assertDistribution(t, step.name+": IPv6 distribution query", queryV6, step.wantV6) + } +} + +// TestStatsDistributionNeedsNoQuery checks that the stats read serves the prefix +// distribution from memory. The read's deadline has already passed, so any +// query it ran would fail at once. The distribution query it used to run read +// every live route and passed the /api/v1/stats deadline on a large database +// (issue 30). +func TestStatsDistributionNeedsNoQuery(t *testing.T) { + cfg := &config.Config{StateDir: t.TempDir()} + + db, err := New(cfg, logger.New()) + if err != nil { + t.Fatalf("failed to create database: %v", err) + } + defer func() { _ = db.Close() }() + + ts := time.Date(2026, 1, 2, 3, 4, 5, 0, time.UTC) + if err := db.UpsertLiveRouteBatch([]*LiveRoute{ + mkV4Route(t, "198.51.100.0/24", 64500, ts), + mkV6Route("2001:db8::/32", 64501, ts), + }); err != nil { + t.Fatalf("UpsertLiveRouteBatch: %v", err) + } + + ctx, cancel := context.WithCancel(context.Background()) + cancel() + + stats, err := db.GetStatsContext(ctx) + if err != nil { + t.Fatalf("GetStatsContext: %v", err) + } + assertDistribution(t, "IPv4 distribution", stats.IPv4PrefixDistribution, + []PrefixDistribution{{MaskLength: 24, Count: 1}}) + assertDistribution(t, "IPv6 distribution", stats.IPv6PrefixDistribution, + []PrefixDistribution{{MaskLength: 32, Count: 1}}) } // TestStatsRouteTimestamps checks the oldest/newest route timestamps are read @@ -296,6 +441,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 +481,11 @@ func assertCounts(t *testing.T, when string, got Stats, want wantCounts) { t.Errorf("%s: LiveRoutes = %d, want %d", when, got.LiveRoutes, want.liveRoutes) } } + +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..99ba85e 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 } @@ -1235,6 +1321,27 @@ func (d *Database) DeleteLiveRoute(prefix string, originASN int, peerIP string) } else { d.counts.addRoutes(0, -int(affected)) } + if affected == 0 { + return nil + } + + // The prefix leaves the distribution when no live route has it any more. + lookupSQL := prefixHasLiveRouteV6SQL + if isV4 { + lookupSQL = prefixHasLiveRouteV4SQL + } + var prefixHasRoute bool + if err := d.db.QueryRow(lookupSQL, prefix).Scan(&prefixHasRoute); err != nil { + return err + } + if !prefixHasRoute { + maskLength, _ := ipnet.Mask.Size() + if isV4 { + d.counts.addToDistribution([]int{maskLength}, nil, -1) + } else { + d.counts.addToDistribution(nil, []int{maskLength}, -1) + } + } return nil } @@ -1244,7 +1351,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 +}