Serve the /api/v1/stats prefix distribution from memory (closes #30)
check / check (push) Successful in 3m52s

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
This commit is contained in:
2026-10-03 14:08:49 +00:00
parent 32704c601e
commit ec99b98d2b
6 changed files with 524 additions and 63 deletions
+180 -46
View File
@@ -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
@@ -234,54 +236,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
@@ -328,7 +352,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 {
@@ -336,24 +374,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 {
@@ -361,6 +405,7 @@ func (d *Database) UpsertLiveRouteBatch(routes []*LiveRoute) error {
}
d.counts.addRoutes(newV4, newV6)
d.counts.addToDistribution(newPrefixMaskLengthsV4, newPrefixMaskLengthsV6, 1)
return nil
}
@@ -411,8 +456,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
@@ -455,6 +514,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 {
@@ -462,6 +546,7 @@ func (d *Database) DeleteLiveRouteBatch(deletions []LiveRouteDeletion) error {
}
d.counts.addRoutes(-int(deletedV4), -int(deletedV6))
d.counts.addToDistribution(gonePrefixMaskLengthsV4, gonePrefixMaskLengthsV6, -1)
return nil
}
@@ -1028,16 +1113,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.
@@ -1061,17 +1146,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
}
@@ -1143,9 +1217,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
@@ -1162,11 +1236,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)
@@ -1179,6 +1259,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
}
@@ -1197,21 +1284,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)
}
@@ -1223,11 +1322,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
}
@@ -1237,7 +1362,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
@@ -1264,6 +1391,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 = `
@@ -1289,6 +1420,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
}