Compare commits
12
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
cd59cb8a8d | ||
|
|
3898daad4e | ||
|
|
08f060045b | ||
|
|
658aadbb81 | ||
|
|
efaf79c4e3 | ||
|
|
63d62b7bc1 | ||
|
|
cb2374f033 | ||
|
|
6d8ae7f592 | ||
|
|
211fdad9c0 | ||
|
|
594e7a504b | ||
|
|
54014c88c8 | ||
|
|
b0cd884019 |
@@ -0,0 +1,12 @@
|
|||||||
|
root = true
|
||||||
|
|
||||||
|
[*]
|
||||||
|
indent_style = space
|
||||||
|
indent_size = 4
|
||||||
|
end_of_line = lf
|
||||||
|
charset = utf-8
|
||||||
|
trim_trailing_whitespace = true
|
||||||
|
insert_final_newline = true
|
||||||
|
|
||||||
|
[Makefile]
|
||||||
|
indent_style = tab
|
||||||
@@ -0,0 +1,9 @@
|
|||||||
|
name: check
|
||||||
|
on: [push]
|
||||||
|
jobs:
|
||||||
|
check:
|
||||||
|
runs-on: ubuntu-latest
|
||||||
|
steps:
|
||||||
|
# actions/checkout v4.2.2, 2026-02-28
|
||||||
|
- uses: actions/checkout@11bd71901bbe5b1630ceea73d27597364c9af683
|
||||||
|
- run: script/cibuild
|
||||||
+40
-2
@@ -1,5 +1,22 @@
|
|||||||
|
# Lint stage — fast feedback on formatting and lint issues.
|
||||||
|
# The golangci-lint image bundles Go, gcc and make, so it can run go vet on
|
||||||
|
# the CGO sqlite package and golangci-lint without extra installs.
|
||||||
|
# golangci/golangci-lint:v2.7.2 (Go 1.25.5), 2026-09-21
|
||||||
|
FROM golangci/golangci-lint@sha256:5d6d5c70a61f1356adfd9dd6316ce286799fefc9d743421356ff1b00842368ba AS lint
|
||||||
|
|
||||||
|
WORKDIR /src
|
||||||
|
|
||||||
|
COPY go.mod go.sum ./
|
||||||
|
RUN go mod download
|
||||||
|
|
||||||
|
COPY . .
|
||||||
|
|
||||||
|
RUN make fmt-check
|
||||||
|
RUN make lint
|
||||||
|
|
||||||
# Build stage
|
# Build stage
|
||||||
FROM golang:1.24-bookworm AS builder
|
# golang:1.24-bookworm, 2026-09-21
|
||||||
|
FROM golang@sha256:1a6d4452c65dea36aac2e2d606b01b4a029ec90cc1ae53890540ce6173ea77ac AS builder
|
||||||
|
|
||||||
# Install build dependencies (zstd for archive, gcc for CGO/sqlite3)
|
# Install build dependencies (zstd for archive, gcc for CGO/sqlite3)
|
||||||
RUN apt-get update && apt-get install -y --no-install-recommends \
|
RUN apt-get update && apt-get install -y --no-install-recommends \
|
||||||
@@ -10,12 +27,19 @@ RUN apt-get update && apt-get install -y --no-install-recommends \
|
|||||||
|
|
||||||
WORKDIR /src
|
WORKDIR /src
|
||||||
|
|
||||||
|
# Force BuildKit to run the lint stage before compiling or testing.
|
||||||
|
COPY --from=lint /src/go.sum /dev/null
|
||||||
|
|
||||||
# Copy everything
|
# Copy everything
|
||||||
COPY . .
|
COPY . .
|
||||||
|
|
||||||
# Vendor dependencies (must be after copying source)
|
# Vendor dependencies (must be after copying source)
|
||||||
RUN go mod download && go mod vendor
|
RUN go mod download && go mod vendor
|
||||||
|
|
||||||
|
# Run the test suite in the build stage: -race needs cgo and the C compiler
|
||||||
|
# installed above. The suite is offline (the live-feed test is opt-in).
|
||||||
|
RUN make test
|
||||||
|
|
||||||
# Build the binary with CGO enabled (required for sqlite3)
|
# Build the binary with CGO enabled (required for sqlite3)
|
||||||
RUN CGO_ENABLED=1 GOOS=linux go build -o /routewatch ./cmd/routewatch
|
RUN CGO_ENABLED=1 GOOS=linux go build -o /routewatch ./cmd/routewatch
|
||||||
|
|
||||||
@@ -26,7 +50,8 @@ RUN tar --zstd -cf /routewatch-source.tar.zst \
|
|||||||
.
|
.
|
||||||
|
|
||||||
# Runtime stage
|
# Runtime stage
|
||||||
FROM debian:bookworm-slim
|
# debian:bookworm-slim, 2026-09-21
|
||||||
|
FROM debian@sha256:3783cc01769c7b2b1b83a5c5ad96c815348e28ed7da68e2e3687004faa906251
|
||||||
|
|
||||||
# Install runtime dependencies
|
# Install runtime dependencies
|
||||||
# - ca-certificates: for HTTPS connections
|
# - ca-certificates: for HTTPS connections
|
||||||
@@ -53,6 +78,19 @@ RUN chown -R routewatch:routewatch /app
|
|||||||
|
|
||||||
ENV XDG_DATA_HOME=/var/lib
|
ENV XDG_DATA_HOME=/var/lib
|
||||||
|
|
||||||
|
# Cap the Go heap at 1.5 GiB so the runtime collects harder before the
|
||||||
|
# container's memory limit is reached. runuser preserves this the way it does
|
||||||
|
# XDG_DATA_HOME above.
|
||||||
|
ENV GOMEMLIMIT=1536MiB
|
||||||
|
|
||||||
|
# Cap glibc's malloc arenas. The SQLite C library allocates and frees millions
|
||||||
|
# of small page-cache chunks from many threads; glibc otherwise creates up to
|
||||||
|
# eight arenas per core (hundreds on a large host) and keeps each arena's freed
|
||||||
|
# chunks resident, so process RSS climbs far above SQLite's live heap and never
|
||||||
|
# comes back down. Two arenas keep that retained memory bounded; database writes
|
||||||
|
# are already serialized, so the lost allocator concurrency costs nothing here.
|
||||||
|
ENV MALLOC_ARENA_MAX=2
|
||||||
|
|
||||||
# Expose HTTP port
|
# Expose HTTP port
|
||||||
EXPOSE 8080
|
EXPOSE 8080
|
||||||
|
|
||||||
|
|||||||
@@ -165,14 +165,61 @@ bgp_peers(id, peer_ip, peer_asn, last_message_type, last_seen)
|
|||||||
Configuration is handled via environment variables and OS-specific paths:
|
Configuration is handled via environment variables and OS-specific paths:
|
||||||
|
|
||||||
| Variable | Default | Description |
|
| Variable | Default | Description |
|
||||||
|----------|---------|-------------|
|
|----------|----------|-------------|
|
||||||
| `PORT` | `8080` | HTTP server port |
|
| `PORT` | `8080` | HTTP server port |
|
||||||
| `DEBUG` | (empty) | Set to `routewatch` for debug logging |
|
| `DEBUG` | (empty) | Set to `routewatch` for debug logging |
|
||||||
|
| `GOMEMLIMIT` | `1536MiB` (in the Docker image) | Go soft memory limit; see Memory |
|
||||||
|
| `MALLOC_ARENA_MAX` | `2` (in the Docker image) | glibc malloc arena cap; see Memory |
|
||||||
|
|
||||||
State directory (database location):
|
State directory (database location):
|
||||||
- macOS: `~/Library/Application Support/routewatch/`
|
- macOS: `~/Library/Application Support/routewatch/`
|
||||||
- Linux: `/var/lib/routewatch/` or `~/.local/share/routewatch/`
|
- Linux: `/var/lib/routewatch/` or `~/.local/share/routewatch/`
|
||||||
|
|
||||||
|
## Memory
|
||||||
|
|
||||||
|
The daemon holds a live routing table, so its memory grows with the size of the
|
||||||
|
data it tracks. The image sets ceilings that keep it inside a 5 GiB container.
|
||||||
|
|
||||||
|
Budget:
|
||||||
|
- Go heap: a 1.5 GiB soft limit (`GOMEMLIMIT=1536MiB`, set in the image).
|
||||||
|
- SQLite: at most 640 MiB of page cache across the connection pool (64 MiB per
|
||||||
|
connection, 10 connections) and a 1.5 GiB hard heap limit for the C library.
|
||||||
|
- glibc allocator: the SQLite C library runs on glibc `malloc`, which frees
|
||||||
|
page-cache chunks back to per-arena free lists rather than to the kernel, so
|
||||||
|
process RSS tracks the high-water mark of those arenas, not SQLite's live
|
||||||
|
heap. glibc creates up to eight arenas per core, so on a many-core host the
|
||||||
|
retained memory — and thus RSS — grows with the core count. The image sets
|
||||||
|
`MALLOC_ARENA_MAX=2` to bound it; the two-arena cap costs nothing here because
|
||||||
|
database writes are already serialized.
|
||||||
|
- About 0.2 GiB for everything else in the runtime.
|
||||||
|
|
||||||
|
Run the container with a memory limit of 5 GiB and swap disabled:
|
||||||
|
|
||||||
|
```bash
|
||||||
|
docker run --memory=5g --memory-swap=5g ...
|
||||||
|
```
|
||||||
|
|
||||||
|
or the equivalent in your deployment tool. This leaves headroom above the
|
||||||
|
ceilings for spikes and the kernel page cache.
|
||||||
|
|
||||||
|
Override the Go soft limit by setting `GOMEMLIMIT` in the environment (for
|
||||||
|
example `-e GOMEMLIMIT=1GiB`); this replaces the image default. `MALLOC_ARENA_MAX`
|
||||||
|
can be overridden the same way, but raising it lets RSS climb again on a
|
||||||
|
many-core host.
|
||||||
|
|
||||||
|
What happens at each limit:
|
||||||
|
- Go soft limit: as the heap approaches `GOMEMLIMIT`, the runtime runs garbage
|
||||||
|
collection more aggressively rather than growing further.
|
||||||
|
- SQLite: at a 1 GiB soft heap limit it recycles its page cache instead of
|
||||||
|
allocating more; at the 1.5 GiB hard heap limit a statement fails with an
|
||||||
|
out-of-memory error, and the handler logs it and drops that batch. The process
|
||||||
|
keeps running.
|
||||||
|
- Handler queues: each of the four handler queues holds at most 20,000 messages.
|
||||||
|
When a queue fills, the streamer drops messages instead of blocking.
|
||||||
|
|
||||||
|
With `DEBUG=routewatch` the daemon logs a `System stats` line every 60 seconds
|
||||||
|
with the goroutine count and Go memory figures.
|
||||||
|
|
||||||
## Development
|
## Development
|
||||||
|
|
||||||
```bash
|
```bash
|
||||||
@@ -204,7 +251,9 @@ them. We provide:
|
|||||||
- `script/projectname` — print the project name (used for the Docker
|
- `script/projectname` — print the project name (used for the Docker
|
||||||
image tag)
|
image tag)
|
||||||
- `script/test` — run the test suite
|
- `script/test` — run the test suite
|
||||||
(`go test -timeout 30s -race -cover ./...`, verbose rerun on failure)
|
(`go test -short -timeout 30s -race -cover ./...`, verbose rerun on
|
||||||
|
failure). The `-short` flag skips the live-network integration test so the
|
||||||
|
default run is deterministic and offline.
|
||||||
- `script/lint` — run `go vet ./...` and `golangci-lint run`
|
- `script/lint` — run `go vet ./...` and `golangci-lint run`
|
||||||
- `script/fmt` — format all code (writes)
|
- `script/fmt` — format all code (writes)
|
||||||
- `script/fmt-check` — check formatting (read-only)
|
- `script/fmt-check` — check formatting (read-only)
|
||||||
@@ -218,6 +267,11 @@ them. We provide:
|
|||||||
- `script/install-precommit` — install the git pre-commit hook that
|
- `script/install-precommit` — install the git pre-commit hook that
|
||||||
runs `script/precommit`
|
runs `script/precommit`
|
||||||
|
|
||||||
|
The live-network integration test `TestRouteWatchLiveFeed` streams the RIPE
|
||||||
|
RIS feed for a few seconds and is skipped in short mode. To run it on demand,
|
||||||
|
invoke `go test` directly without `-short`:
|
||||||
|
`go test -run TestRouteWatchLiveFeed ./internal/routewatch/`.
|
||||||
|
|
||||||
## License
|
## License
|
||||||
|
|
||||||
See LICENSE file.
|
See LICENSE file.
|
||||||
|
|||||||
+100
-21
@@ -38,6 +38,22 @@ const (
|
|||||||
maxIPv4 = 0xFFFFFFFF
|
maxIPv4 = 0xFFFFFFFF
|
||||||
)
|
)
|
||||||
|
|
||||||
|
// SQLite memory tuning. cache_size and busy_timeout go in the DSN so every
|
||||||
|
// pooled connection gets them; the heap limits are process-wide and set once.
|
||||||
|
const (
|
||||||
|
// sqliteCacheSizeKiB is the per-connection page cache; negative means KiB.
|
||||||
|
// -65536 = 64 MiB, so at most 640 MiB across the 10-connection pool.
|
||||||
|
sqliteCacheSizeKiB = -65536
|
||||||
|
// sqliteBusyTimeoutMs is how long a connection waits on a locked database.
|
||||||
|
sqliteBusyTimeoutMs = 5000
|
||||||
|
// sqliteSoftHeapLimitBytes (1 GiB) makes SQLite recycle its cache rather
|
||||||
|
// than allocate once its C heap passes this size.
|
||||||
|
sqliteSoftHeapLimitBytes = 1073741824
|
||||||
|
// sqliteHardHeapLimitBytes (1.5 GiB) fails a statement with SQLITE_NOMEM
|
||||||
|
// instead of growing the C heap without bound.
|
||||||
|
sqliteHardHeapLimitBytes = 1610612736
|
||||||
|
)
|
||||||
|
|
||||||
// Common errors
|
// Common errors
|
||||||
var (
|
var (
|
||||||
// ErrInvalidIP is returned when an IP address is malformed
|
// ErrInvalidIP is returned when an IP address is malformed
|
||||||
@@ -71,11 +87,17 @@ func New(cfg *config.Config, logger *logger.Logger) (*Database, error) {
|
|||||||
return nil, fmt.Errorf("failed to create database directory: %w", err)
|
return nil, fmt.Errorf("failed to create database directory: %w", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
// Add connection parameters for go-sqlite3
|
// Per-connection SQLite settings go in the DSN so every pooled connection
|
||||||
// Configure SQLite connection parameters
|
// gets them, not just the one that runs the Initialize pragmas. _txlock=
|
||||||
|
// immediate makes every transaction take the write lock at BEGIN. Without it
|
||||||
|
// a transaction that reads before writing starts as a reader and, when it
|
||||||
|
// then writes while another connection holds the write lock, fails at once
|
||||||
|
// with "database is locked" without waiting for _busy_timeout.
|
||||||
dsn := fmt.Sprintf(
|
dsn := fmt.Sprintf(
|
||||||
"file:%s",
|
"file:%s?_cache_size=%d&_synchronous=OFF&_busy_timeout=%d&_journal_mode=WAL&_txlock=immediate",
|
||||||
dbPath,
|
dbPath,
|
||||||
|
sqliteCacheSizeKiB,
|
||||||
|
sqliteBusyTimeoutMs,
|
||||||
)
|
)
|
||||||
db, err := sql.Open("sqlite3", dsn)
|
db, err := sql.Open("sqlite3", dsn)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -104,15 +126,17 @@ func New(cfg *config.Config, logger *logger.Logger) (*Database, error) {
|
|||||||
|
|
||||||
// Initialize creates the database schema if it doesn't exist.
|
// Initialize creates the database schema if it doesn't exist.
|
||||||
func (d *Database) Initialize() error {
|
func (d *Database) Initialize() error {
|
||||||
// Set SQLite pragmas for performance
|
// Set SQLite pragmas for performance. Per-connection settings (cache_size,
|
||||||
|
// synchronous, busy_timeout, journal_mode) live in the DSN; temp_store is
|
||||||
|
// left at its default so DISTINCT temp B-trees spill to disk instead of C
|
||||||
|
// heap. The heap limits below are process-wide, so setting them once here is
|
||||||
|
// enough for the whole pool.
|
||||||
pragmas := []string{
|
pragmas := []string{
|
||||||
"PRAGMA journal_mode=WAL", // Write-Ahead Logging
|
"PRAGMA journal_mode=WAL", // Write-Ahead Logging
|
||||||
"PRAGMA synchronous=OFF", // Don't wait for disk writes
|
|
||||||
"PRAGMA cache_size=-3145728", // 3GB cache (upper limit for 2.4GB DB)
|
|
||||||
"PRAGMA temp_store=MEMORY", // Use memory for temp tables
|
|
||||||
"PRAGMA busy_timeout=5000", // 5 second busy timeout
|
|
||||||
"PRAGMA analysis_limit=0", // Disable automatic ANALYZE
|
"PRAGMA analysis_limit=0", // Disable automatic ANALYZE
|
||||||
"PRAGMA auto_vacuum=INCREMENTAL", // Enable incremental vacuum
|
"PRAGMA auto_vacuum=INCREMENTAL", // Enable incremental vacuum
|
||||||
|
fmt.Sprintf("PRAGMA soft_heap_limit=%d", sqliteSoftHeapLimitBytes),
|
||||||
|
fmt.Sprintf("PRAGMA hard_heap_limit=%d", sqliteHardHeapLimitBytes),
|
||||||
}
|
}
|
||||||
|
|
||||||
for _, pragma := range pragmas {
|
for _, pragma := range pragmas {
|
||||||
@@ -956,23 +980,20 @@ func (d *Database) GetStatsContext(ctx context.Context) (Stats, error) {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return stats, fmt.Errorf("failed to count IPv6 routes: %w", err)
|
return stats, fmt.Errorf("failed to count IPv6 routes: %w", err)
|
||||||
}
|
}
|
||||||
|
stats.IPv4Routes = v4Count
|
||||||
|
stats.IPv6Routes = v6Count
|
||||||
stats.LiveRoutes = v4Count + v6Count
|
stats.LiveRoutes = v4Count + v6Count
|
||||||
|
|
||||||
// Get oldest and newest route timestamps
|
// Get oldest and newest route timestamps. Each query reads a single row from
|
||||||
routeTimestampQuery := `
|
// one end of the last_updated index, so the cost is a log-time index lookup
|
||||||
SELECT MIN(last_updated), MAX(last_updated) FROM (
|
// rather than a full scan of both route tables. Selecting the last_updated
|
||||||
SELECT last_updated FROM live_routes_v4
|
// column directly (rather than MIN/MAX, whose result has no column type) lets
|
||||||
UNION ALL
|
// the driver parse the DATETIME value into time.Time; the union scan aggregate
|
||||||
SELECT last_updated FROM live_routes_v6
|
// used before returned an untyped string and logged a warning on every call.
|
||||||
)
|
stats.OldestRoute, stats.NewestRoute, err = d.routeTimestampRange(ctx)
|
||||||
`
|
|
||||||
var oldestRoute, newestRoute *time.Time
|
|
||||||
err = d.db.QueryRowContext(ctx, routeTimestampQuery).Scan(&oldestRoute, &newestRoute)
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
// Display-only fields; log but keep the rest of the stats.
|
||||||
d.logger.Warn("Failed to get route timestamps", "error", err)
|
d.logger.Warn("Failed to get route timestamps", "error", err)
|
||||||
} else {
|
|
||||||
stats.OldestRoute = oldestRoute
|
|
||||||
stats.NewestRoute = newestRoute
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// Get prefix distribution
|
// Get prefix distribution
|
||||||
@@ -985,6 +1006,64 @@ func (d *Database) GetStatsContext(ctx context.Context) (Stats, error) {
|
|||||||
return stats, nil
|
return stats, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// routeTimestampRange returns the earliest and latest last_updated across both
|
||||||
|
// live route tables, or nil values when both tables are empty. Each query reads
|
||||||
|
// one row from an end of the last_updated index rather than scanning the tables.
|
||||||
|
func (d *Database) routeTimestampRange(ctx context.Context) (oldest, newest *time.Time, err error) {
|
||||||
|
oldestV4, ok, err := d.scanRouteTimestamp(ctx,
|
||||||
|
"SELECT last_updated FROM live_routes_v4 ORDER BY last_updated ASC LIMIT 1")
|
||||||
|
if err != nil {
|
||||||
|
return nil, nil, err
|
||||||
|
}
|
||||||
|
if ok {
|
||||||
|
oldest = &oldestV4
|
||||||
|
}
|
||||||
|
|
||||||
|
oldestV6, ok, err := d.scanRouteTimestamp(ctx,
|
||||||
|
"SELECT last_updated FROM live_routes_v6 ORDER BY last_updated ASC LIMIT 1")
|
||||||
|
if err != nil {
|
||||||
|
return nil, nil, err
|
||||||
|
}
|
||||||
|
if ok && (oldest == nil || oldestV6.Before(*oldest)) {
|
||||||
|
oldest = &oldestV6
|
||||||
|
}
|
||||||
|
|
||||||
|
newestV4, ok, err := d.scanRouteTimestamp(ctx,
|
||||||
|
"SELECT last_updated FROM live_routes_v4 ORDER BY last_updated DESC LIMIT 1")
|
||||||
|
if err != nil {
|
||||||
|
return nil, nil, err
|
||||||
|
}
|
||||||
|
if ok {
|
||||||
|
newest = &newestV4
|
||||||
|
}
|
||||||
|
|
||||||
|
newestV6, ok, err := d.scanRouteTimestamp(ctx,
|
||||||
|
"SELECT last_updated FROM live_routes_v6 ORDER BY last_updated DESC LIMIT 1")
|
||||||
|
if err != nil {
|
||||||
|
return nil, nil, err
|
||||||
|
}
|
||||||
|
if ok && (newest == nil || newestV6.After(*newest)) {
|
||||||
|
newest = &newestV6
|
||||||
|
}
|
||||||
|
|
||||||
|
return oldest, newest, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// scanRouteTimestamp runs a single-row timestamp query. ok is false when the
|
||||||
|
// table is empty. The query selects the last_updated column directly so the
|
||||||
|
// driver parses the DATETIME value into a time.Time.
|
||||||
|
func (d *Database) scanRouteTimestamp(ctx context.Context, query string) (ts time.Time, ok bool, err error) {
|
||||||
|
err = d.db.QueryRowContext(ctx, query).Scan(&ts)
|
||||||
|
switch {
|
||||||
|
case errors.Is(err, sql.ErrNoRows):
|
||||||
|
return time.Time{}, false, nil
|
||||||
|
case err != nil:
|
||||||
|
return time.Time{}, false, err
|
||||||
|
default:
|
||||||
|
return ts, true, nil
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// UpsertLiveRoute inserts or updates a live route
|
// UpsertLiveRoute inserts or updates a live route
|
||||||
func (d *Database) UpsertLiveRoute(route *LiveRoute) error {
|
func (d *Database) UpsertLiveRoute(route *LiveRoute) error {
|
||||||
d.lock("UpsertLiveRoute")
|
d.lock("UpsertLiveRoute")
|
||||||
|
|||||||
@@ -1,8 +1,36 @@
|
|||||||
package database
|
package database
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"context"
|
||||||
|
"database/sql"
|
||||||
"net"
|
"net"
|
||||||
|
"sync"
|
||||||
"testing"
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"git.eeqj.de/sneak/routewatch/internal/config"
|
||||||
|
"git.eeqj.de/sneak/routewatch/internal/logger"
|
||||||
|
"github.com/google/uuid"
|
||||||
|
)
|
||||||
|
|
||||||
|
// tempStoreMemory is the PRAGMA temp_store value meaning "hold temp B-trees in
|
||||||
|
// memory"; the DSN change must leave temp_store below this so they spill to disk.
|
||||||
|
const tempStoreMemory = 2
|
||||||
|
|
||||||
|
// heldConnections is how many pooled connections the pragma test holds open at
|
||||||
|
// once so each is a distinct SQLite connection that parsed the DSN.
|
||||||
|
const heldConnections = 5
|
||||||
|
|
||||||
|
// Parameters for the checkpoint-contention regression test.
|
||||||
|
const (
|
||||||
|
// contentionIterations is how many batch writes race the checkpoint loop.
|
||||||
|
contentionIterations = 400
|
||||||
|
// contendedASNCount is the small set of ASNs the batches reuse, so most
|
||||||
|
// batches update existing rows and exercise the read-before-write path.
|
||||||
|
contendedASNCount = 16
|
||||||
|
// asnSecondBand offsets a second ASN per batch so each batch writes more
|
||||||
|
// than one row.
|
||||||
|
asnSecondBand = 100
|
||||||
)
|
)
|
||||||
|
|
||||||
func TestIPToUint32(t *testing.T) {
|
func TestIPToUint32(t *testing.T) {
|
||||||
@@ -282,6 +310,221 @@ func TestIPv4RangeIntegration(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// TestConnectionPoolPragmas holds several pooled connections open at once and
|
||||||
|
// checks each one carries the per-connection settings from the DSN, plus the
|
||||||
|
// process-wide hard heap limit.
|
||||||
|
func TestConnectionPoolPragmas(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() }()
|
||||||
|
|
||||||
|
ctx := context.Background()
|
||||||
|
|
||||||
|
// Hold distinct connections open simultaneously so the pool must open a new
|
||||||
|
// one (each parsing the DSN) rather than hand back the same connection.
|
||||||
|
conns := make([]*sql.Conn, 0, heldConnections)
|
||||||
|
defer func() {
|
||||||
|
for _, c := range conns {
|
||||||
|
_ = c.Close()
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
|
||||||
|
for i := 0; i < heldConnections; i++ {
|
||||||
|
c, err := db.db.Conn(ctx)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("failed to open connection %d: %v", i, err)
|
||||||
|
}
|
||||||
|
conns = append(conns, c)
|
||||||
|
}
|
||||||
|
|
||||||
|
for i, c := range conns {
|
||||||
|
var cacheSize int
|
||||||
|
if err := c.QueryRowContext(ctx, "PRAGMA cache_size").Scan(&cacheSize); err != nil {
|
||||||
|
t.Fatalf("conn %d: failed to read cache_size: %v", i, err)
|
||||||
|
}
|
||||||
|
if cacheSize != sqliteCacheSizeKiB {
|
||||||
|
t.Errorf("conn %d: cache_size = %d, want %d", i, cacheSize, sqliteCacheSizeKiB)
|
||||||
|
}
|
||||||
|
|
||||||
|
var busyTimeout int
|
||||||
|
if err := c.QueryRowContext(ctx, "PRAGMA busy_timeout").Scan(&busyTimeout); err != nil {
|
||||||
|
t.Fatalf("conn %d: failed to read busy_timeout: %v", i, err)
|
||||||
|
}
|
||||||
|
if busyTimeout != sqliteBusyTimeoutMs {
|
||||||
|
t.Errorf("conn %d: busy_timeout = %d, want %d", i, busyTimeout, sqliteBusyTimeoutMs)
|
||||||
|
}
|
||||||
|
|
||||||
|
var tempStore int
|
||||||
|
if err := c.QueryRowContext(ctx, "PRAGMA temp_store").Scan(&tempStore); err != nil {
|
||||||
|
t.Fatalf("conn %d: failed to read temp_store: %v", i, err)
|
||||||
|
}
|
||||||
|
if tempStore == tempStoreMemory {
|
||||||
|
t.Errorf("conn %d: temp_store = %d, want anything but %d (MEMORY)", i, tempStore, tempStoreMemory)
|
||||||
|
}
|
||||||
|
|
||||||
|
var hardHeapLimit int64
|
||||||
|
if err := c.QueryRowContext(ctx, "PRAGMA hard_heap_limit").Scan(&hardHeapLimit); err != nil {
|
||||||
|
t.Fatalf("conn %d: failed to read hard_heap_limit: %v", i, err)
|
||||||
|
}
|
||||||
|
if hardHeapLimit != sqliteHardHeapLimitBytes {
|
||||||
|
t.Errorf("conn %d: hard_heap_limit = %d, want %d", i, hardHeapLimit, sqliteHardHeapLimitBytes)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestBatchWriteDuringCheckpoint reproduces issue #25. A batch write reads
|
||||||
|
// (SELECT) before it writes (INSERT/UPDATE). Under the default deferred locking
|
||||||
|
// the transaction begins as a reader and, when the maintainer's WAL checkpoint
|
||||||
|
// holds the write lock, its upgrade to writer fails immediately with "database
|
||||||
|
// is locked" without honouring busy_timeout, dropping the batch. With
|
||||||
|
// _txlock=immediate the transaction takes the write lock at BEGIN and waits, so
|
||||||
|
// no batch is dropped. The checkpoint runs without the Database mutex, exactly
|
||||||
|
// as the background maintainer does in production.
|
||||||
|
func TestBatchWriteDuringCheckpoint(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() }()
|
||||||
|
|
||||||
|
ctx, cancel := context.WithCancel(context.Background())
|
||||||
|
|
||||||
|
var wg sync.WaitGroup
|
||||||
|
wg.Add(1)
|
||||||
|
go func() {
|
||||||
|
defer wg.Done()
|
||||||
|
for {
|
||||||
|
select {
|
||||||
|
case <-ctx.Done():
|
||||||
|
return
|
||||||
|
default:
|
||||||
|
_ = db.Checkpoint(ctx) // errors are the checkpoint's own to absorb
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
|
||||||
|
ts := time.Now().UTC()
|
||||||
|
for i := 0; i < contentionIterations; i++ {
|
||||||
|
asns := map[int]time.Time{
|
||||||
|
i % contendedASNCount: ts,
|
||||||
|
(i % contendedASNCount) + asnSecondBand: ts,
|
||||||
|
}
|
||||||
|
if err := db.GetOrCreateASNBatch(asns); err != nil {
|
||||||
|
cancel()
|
||||||
|
wg.Wait()
|
||||||
|
t.Fatalf("batch write failed under checkpoint contention: %v", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
cancel()
|
||||||
|
wg.Wait()
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestStatsRouteTimestampsAndCounts checks GetStatsContext reports the correct
|
||||||
|
// route counts and the oldest/newest last_updated across both route tables. The
|
||||||
|
// old union-scan query read the aggregate result into *time.Time, which the
|
||||||
|
// driver could not parse, so it logged a warning every call and left both
|
||||||
|
// timestamps nil; this asserts they are populated from the right rows.
|
||||||
|
func TestStatsRouteTimestampsAndCounts(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() }()
|
||||||
|
|
||||||
|
ctx := context.Background()
|
||||||
|
|
||||||
|
// Empty database: no routes, so both timestamps are nil and no error.
|
||||||
|
empty, err := db.GetStatsContext(ctx)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("GetStatsContext on empty database: %v", err)
|
||||||
|
}
|
||||||
|
if empty.OldestRoute != nil || empty.NewestRoute != nil {
|
||||||
|
t.Fatalf("empty database timestamps = (%v, %v), want (nil, nil)",
|
||||||
|
empty.OldestRoute, empty.NewestRoute)
|
||||||
|
}
|
||||||
|
if empty.LiveRoutes != 0 {
|
||||||
|
t.Fatalf("empty database LiveRoutes = %d, want 0", empty.LiveRoutes)
|
||||||
|
}
|
||||||
|
|
||||||
|
base := time.Date(2026, 1, 2, 3, 4, 5, 0, time.UTC)
|
||||||
|
oldest := base
|
||||||
|
middle := base.Add(time.Minute)
|
||||||
|
newest := base.Add(2 * time.Minute)
|
||||||
|
|
||||||
|
mkV4 := func(prefix string, asn int, ts time.Time) *LiveRoute {
|
||||||
|
start, end, rerr := CalculateIPv4Range(prefix)
|
||||||
|
if rerr != nil {
|
||||||
|
t.Fatalf("CalculateIPv4Range(%s): %v", prefix, rerr)
|
||||||
|
}
|
||||||
|
|
||||||
|
return &LiveRoute{
|
||||||
|
ID: uuid.New(),
|
||||||
|
Prefix: prefix,
|
||||||
|
MaskLength: 24,
|
||||||
|
IPVersion: ipVersionV4,
|
||||||
|
OriginASN: asn,
|
||||||
|
PeerIP: "192.0.2.1",
|
||||||
|
ASPath: []int{asn},
|
||||||
|
NextHop: "192.0.2.254",
|
||||||
|
LastUpdated: ts,
|
||||||
|
V4IPStart: &start,
|
||||||
|
V4IPEnd: &end,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// Two IPv4 routes (one oldest, one middle) and one IPv6 route (newest).
|
||||||
|
routes := []*LiveRoute{
|
||||||
|
mkV4("198.51.100.0/24", 64500, middle),
|
||||||
|
mkV4("203.0.113.0/24", 64501, oldest),
|
||||||
|
{
|
||||||
|
ID: uuid.New(),
|
||||||
|
Prefix: "2001:db8::/32",
|
||||||
|
MaskLength: 32,
|
||||||
|
IPVersion: ipVersionV6,
|
||||||
|
OriginASN: 64502,
|
||||||
|
PeerIP: "2001:db8::1",
|
||||||
|
ASPath: []int{64502},
|
||||||
|
NextHop: "2001:db8::ffff",
|
||||||
|
LastUpdated: newest,
|
||||||
|
},
|
||||||
|
}
|
||||||
|
for _, route := range routes {
|
||||||
|
if err := db.UpsertLiveRoute(route); err != nil {
|
||||||
|
t.Fatalf("UpsertLiveRoute(%s): %v", route.Prefix, err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
stats, err := db.GetStatsContext(ctx)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("GetStatsContext: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
if stats.IPv4Routes != 2 {
|
||||||
|
t.Errorf("IPv4Routes = %d, want 2", stats.IPv4Routes)
|
||||||
|
}
|
||||||
|
if stats.IPv6Routes != 1 {
|
||||||
|
t.Errorf("IPv6Routes = %d, want 1", stats.IPv6Routes)
|
||||||
|
}
|
||||||
|
if stats.LiveRoutes != 3 {
|
||||||
|
t.Errorf("LiveRoutes = %d, want 3", stats.LiveRoutes)
|
||||||
|
}
|
||||||
|
if stats.OldestRoute == nil || !stats.OldestRoute.Equal(oldest) {
|
||||||
|
t.Errorf("OldestRoute = %v, want %v", stats.OldestRoute, oldest)
|
||||||
|
}
|
||||||
|
if stats.NewestRoute == nil || !stats.NewestRoute.Equal(newest) {
|
||||||
|
t.Errorf("NewestRoute = %v, want %v", stats.NewestRoute, newest)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
func BenchmarkIPToUint32(b *testing.B) {
|
func BenchmarkIPToUint32(b *testing.B) {
|
||||||
ip := net.ParseIP("192.168.1.1")
|
ip := net.ParseIP("192.168.1.1")
|
||||||
b.ResetTimer()
|
b.ResetTimer()
|
||||||
|
|||||||
@@ -18,6 +18,8 @@ type Stats struct {
|
|||||||
Peers int
|
Peers int
|
||||||
FileSizeBytes int64
|
FileSizeBytes int64
|
||||||
LiveRoutes int
|
LiveRoutes int
|
||||||
|
IPv4Routes int
|
||||||
|
IPv6Routes int
|
||||||
OldestRoute *time.Time
|
OldestRoute *time.Time
|
||||||
NewestRoute *time.Time
|
NewestRoute *time.Time
|
||||||
IPv4PrefixDistribution []PrefixDistribution
|
IPv4PrefixDistribution []PrefixDistribution
|
||||||
|
|||||||
+22
-18
@@ -63,24 +63,28 @@ type RISLiveMessage struct {
|
|||||||
// the actual BGP update data including AS path, communities, announcements,
|
// the actual BGP update data including AS path, communities, announcements,
|
||||||
// and withdrawals.
|
// and withdrawals.
|
||||||
type RISMessage struct {
|
type RISMessage struct {
|
||||||
Type string `json:"type"`
|
Type string `json:"type"`
|
||||||
Timestamp float64 `json:"timestamp"`
|
Timestamp float64 `json:"timestamp"`
|
||||||
ParsedTimestamp time.Time `json:"-"` // Parsed from Timestamp field
|
ParsedTimestamp time.Time `json:"-"` // Parsed from Timestamp field
|
||||||
Peer string `json:"peer"`
|
Peer string `json:"peer"`
|
||||||
PeerASN string `json:"peer_asn"`
|
PeerASN string `json:"peer_asn"`
|
||||||
ID string `json:"id"`
|
ID string `json:"id"`
|
||||||
Host string `json:"host"`
|
Host string `json:"host"`
|
||||||
RRC string `json:"rrc,omitempty"`
|
RRC string `json:"rrc,omitempty"`
|
||||||
MrtTime float64 `json:"mrt_time,omitempty"`
|
MrtTime float64 `json:"mrt_time,omitempty"`
|
||||||
SocketTime float64 `json:"socket_time,omitempty"`
|
SocketTime float64 `json:"socket_time,omitempty"`
|
||||||
Path ASPath `json:"path,omitempty"`
|
Path ASPath `json:"path,omitempty"`
|
||||||
Community [][]int `json:"community,omitempty"`
|
// Community and Raw are present in the feed but read by no handler.
|
||||||
Origin string `json:"origin,omitempty"`
|
// They are the largest fields on a message that lives in up to four
|
||||||
MED *int `json:"med,omitempty"`
|
// handler queues, so json:"-" keeps them out of the decoded message
|
||||||
LocalPref *int `json:"local_pref,omitempty"`
|
// to save queue memory. Do not decode them without a consumer.
|
||||||
Announcements []RISAnnouncement `json:"announcements,omitempty"`
|
Community [][]int `json:"-"`
|
||||||
Withdrawals []string `json:"withdrawals,omitempty"`
|
Origin string `json:"origin,omitempty"`
|
||||||
Raw string `json:"raw,omitempty"`
|
MED *int `json:"med,omitempty"`
|
||||||
|
LocalPref *int `json:"local_pref,omitempty"`
|
||||||
|
Announcements []RISAnnouncement `json:"announcements,omitempty"`
|
||||||
|
Withdrawals []string `json:"withdrawals,omitempty"`
|
||||||
|
Raw string `json:"-"`
|
||||||
}
|
}
|
||||||
|
|
||||||
// RISAnnouncement represents a BGP route announcement within a RIS message.
|
// RISAnnouncement represents a BGP route announcement within a RIS message.
|
||||||
|
|||||||
@@ -0,0 +1,67 @@
|
|||||||
|
package ristypes
|
||||||
|
|
||||||
|
import (
|
||||||
|
"bufio"
|
||||||
|
"encoding/json"
|
||||||
|
"os"
|
||||||
|
"testing"
|
||||||
|
)
|
||||||
|
|
||||||
|
// messageExamplesPath is the captured RIS Live feed used as decode fixtures,
|
||||||
|
// one JSON message per line.
|
||||||
|
const messageExamplesPath = "../../docs/message-examples.json"
|
||||||
|
|
||||||
|
// TestDecodeDropsCommunityAndRaw decodes every captured message the way the
|
||||||
|
// streamer does and checks that the fields no handler reads (Community, Raw)
|
||||||
|
// stay empty while the fields handlers use (Path, Announcements) still decode.
|
||||||
|
func TestDecodeDropsCommunityAndRaw(t *testing.T) {
|
||||||
|
f, err := os.Open(messageExamplesPath)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("open fixtures: %v", err)
|
||||||
|
}
|
||||||
|
defer f.Close()
|
||||||
|
|
||||||
|
var messages, withPath, withAnnouncements int
|
||||||
|
|
||||||
|
scanner := bufio.NewScanner(f)
|
||||||
|
scanner.Buffer(make([]byte, 0, 64*1024), 1024*1024)
|
||||||
|
for scanner.Scan() {
|
||||||
|
line := scanner.Bytes()
|
||||||
|
if len(line) == 0 {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
|
||||||
|
var wrapper RISLiveMessage
|
||||||
|
if err := json.Unmarshal(line, &wrapper); err != nil {
|
||||||
|
t.Fatalf("unmarshal message %d: %v", messages+1, err)
|
||||||
|
}
|
||||||
|
messages++
|
||||||
|
|
||||||
|
msg := wrapper.Data
|
||||||
|
if msg.Community != nil {
|
||||||
|
t.Errorf("message %d: Community decoded, want empty: %v", messages, msg.Community)
|
||||||
|
}
|
||||||
|
if msg.Raw != "" {
|
||||||
|
t.Errorf("message %d: Raw decoded, want empty", messages)
|
||||||
|
}
|
||||||
|
|
||||||
|
if len(msg.Path) > 0 {
|
||||||
|
withPath++
|
||||||
|
}
|
||||||
|
if len(msg.Announcements) > 0 {
|
||||||
|
withAnnouncements++
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if err := scanner.Err(); err != nil {
|
||||||
|
t.Fatalf("scan fixtures: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
if messages == 0 {
|
||||||
|
t.Fatal("no messages decoded from fixtures")
|
||||||
|
}
|
||||||
|
// The fixtures include announcement messages; a used field must still decode,
|
||||||
|
// otherwise an empty Community/Raw would prove nothing.
|
||||||
|
if withPath == 0 || withAnnouncements == 0 {
|
||||||
|
t.Fatalf("used fields did not decode: withPath=%d withAnnouncements=%d", withPath, withAnnouncements)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -426,6 +426,9 @@ func (m *mockStore) Ping(ctx context.Context) error {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func TestRouteWatchLiveFeed(t *testing.T) {
|
func TestRouteWatchLiveFeed(t *testing.T) {
|
||||||
|
if testing.Short() {
|
||||||
|
t.Skip("skipping live RIPE RIS network feed test in short mode; run without -short to include it")
|
||||||
|
}
|
||||||
|
|
||||||
// Create mock database
|
// Create mock database
|
||||||
mockDB := newMockStore()
|
mockDB := newMockStore()
|
||||||
|
|||||||
@@ -10,9 +10,11 @@ import (
|
|||||||
)
|
)
|
||||||
|
|
||||||
const (
|
const (
|
||||||
// asHandlerQueueSize is the queue capacity for ASN operations
|
// asHandlerQueueSize is the queue capacity for ASN operations, about 4
|
||||||
// DO NOT set this higher than 100000 without explicit instructions
|
// seconds of feed at peak. The streamer drops rather than blocks when a
|
||||||
asHandlerQueueSize = 100000
|
// queue is full, so this bounds memory. Batches still flush on a timer
|
||||||
|
// (asnBatchTimeout), so a queue smaller than asnBatchSize is fine.
|
||||||
|
asHandlerQueueSize = 20000
|
||||||
|
|
||||||
// asnBatchSize is the number of ASN operations to batch together
|
// asnBatchSize is the number of ASN operations to batch together
|
||||||
asnBatchSize = 30000
|
asnBatchSize = 30000
|
||||||
|
|||||||
@@ -14,8 +14,11 @@ import (
|
|||||||
)
|
)
|
||||||
|
|
||||||
const (
|
const (
|
||||||
// peerHandlerQueueSize is the queue capacity for peer tracking operations
|
// peerHandlerQueueSize is the queue capacity for peer tracking operations,
|
||||||
peerHandlerQueueSize = 100000
|
// about 4 seconds of feed at peak. The streamer drops rather than blocks
|
||||||
|
// when a queue is full, so this bounds memory. Batches still flush on a
|
||||||
|
// timer (peerBatchTimeout).
|
||||||
|
peerHandlerQueueSize = 20000
|
||||||
|
|
||||||
// peerBatchSize is the number of peer updates to batch together
|
// peerBatchSize is the number of peer updates to batch together
|
||||||
peerBatchSize = 10000
|
peerBatchSize = 10000
|
||||||
|
|||||||
@@ -11,29 +11,25 @@ import (
|
|||||||
)
|
)
|
||||||
|
|
||||||
const (
|
const (
|
||||||
// peeringHandlerQueueSize defines the buffer capacity for the peering
|
// peeringHandlerQueueSize is the buffer capacity for the peering handler's
|
||||||
// handler's message queue. This should be large enough to handle bursts
|
// message queue, about 4 seconds of feed at peak. The streamer drops
|
||||||
// of BGP UPDATE messages without blocking.
|
// rather than blocks when a queue is full, so this bounds memory.
|
||||||
peeringHandlerQueueSize = 100000
|
peeringHandlerQueueSize = 20000
|
||||||
|
|
||||||
// minPathLengthForPeering specifies the minimum number of ASNs required
|
// minPathLengthForPeering specifies the minimum number of ASNs required
|
||||||
// in a BGP AS path to extract peering relationships. A path with fewer
|
// in a BGP AS path to extract peering relationships. A path with fewer
|
||||||
// than 2 ASNs cannot contain any peering information.
|
// than 2 ASNs cannot contain any peering information.
|
||||||
minPathLengthForPeering = 2
|
minPathLengthForPeering = 2
|
||||||
|
|
||||||
// pathExpirationTime determines how long AS paths are kept in memory
|
|
||||||
// before being eligible for pruning. Paths older than this are removed
|
|
||||||
// to prevent unbounded memory growth.
|
|
||||||
pathExpirationTime = 30 * time.Minute
|
|
||||||
|
|
||||||
// peeringProcessInterval controls how frequently the handler processes
|
// peeringProcessInterval controls how frequently the handler processes
|
||||||
// accumulated AS paths and extracts peering relationships to store
|
// accumulated AS paths and extracts peering relationships to store
|
||||||
// in the database.
|
// in the database.
|
||||||
peeringProcessInterval = 30 * time.Second
|
peeringProcessInterval = 30 * time.Second
|
||||||
|
|
||||||
// pathPruneInterval determines how often the handler checks for and
|
// maxTrackedPaths bounds how many distinct AS paths are held in memory
|
||||||
// removes expired AS paths from memory.
|
// between processing runs. Once the map is full, further new paths are
|
||||||
pathPruneInterval = 5 * time.Minute
|
// dropped and counted until the next run empties it.
|
||||||
|
maxTrackedPaths = 500000
|
||||||
)
|
)
|
||||||
|
|
||||||
// PeeringHandler processes BGP UPDATE messages to extract and track
|
// PeeringHandler processes BGP UPDATE messages to extract and track
|
||||||
@@ -46,17 +42,17 @@ type PeeringHandler struct {
|
|||||||
logger *logger.Logger
|
logger *logger.Logger
|
||||||
|
|
||||||
// In-memory AS path tracking
|
// In-memory AS path tracking
|
||||||
mu sync.RWMutex
|
mu sync.Mutex
|
||||||
asPaths map[string]time.Time // key is JSON-encoded AS path
|
asPaths map[string]time.Time // key is JSON-encoded AS path
|
||||||
|
droppedPaths int // paths dropped because the map was full
|
||||||
|
|
||||||
stopCh chan struct{}
|
stopCh chan struct{}
|
||||||
}
|
}
|
||||||
|
|
||||||
// NewPeeringHandler creates and initializes a new PeeringHandler with the
|
// NewPeeringHandler creates and initializes a new PeeringHandler with the
|
||||||
// provided database store and logger. It starts two background goroutines:
|
// provided database store and logger. It starts one background goroutine that
|
||||||
// one for periodic processing of accumulated AS paths into peering records,
|
// periodically processes accumulated AS paths into peering records. The
|
||||||
// and one for pruning expired paths from memory. The handler begins
|
// handler begins processing immediately upon creation.
|
||||||
// processing immediately upon creation.
|
|
||||||
func NewPeeringHandler(db database.Store, logger *logger.Logger) *PeeringHandler {
|
func NewPeeringHandler(db database.Store, logger *logger.Logger) *PeeringHandler {
|
||||||
h := &PeeringHandler{
|
h := &PeeringHandler{
|
||||||
db: db,
|
db: db,
|
||||||
@@ -65,9 +61,8 @@ func NewPeeringHandler(db database.Store, logger *logger.Logger) *PeeringHandler
|
|||||||
stopCh: make(chan struct{}),
|
stopCh: make(chan struct{}),
|
||||||
}
|
}
|
||||||
|
|
||||||
// Start the periodic processing goroutines
|
// Start the periodic processing goroutine
|
||||||
go h.processLoop()
|
go h.processLoop()
|
||||||
go h.pruneLoop()
|
|
||||||
|
|
||||||
return h
|
return h
|
||||||
}
|
}
|
||||||
@@ -106,9 +101,18 @@ func (h *PeeringHandler) HandleMessage(msg *ristypes.RISMessage) {
|
|||||||
|
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
key := string(pathJSON)
|
||||||
|
|
||||||
h.mu.Lock()
|
h.mu.Lock()
|
||||||
h.asPaths[string(pathJSON)] = timestamp
|
if _, exists := h.asPaths[key]; exists {
|
||||||
|
// Already tracked: refresh its timestamp.
|
||||||
|
h.asPaths[key] = timestamp
|
||||||
|
} else if len(h.asPaths) >= maxTrackedPaths {
|
||||||
|
// Map is full; drop this new path and count it.
|
||||||
|
h.droppedPaths++
|
||||||
|
} else {
|
||||||
|
h.asPaths[key] = timestamp
|
||||||
|
}
|
||||||
h.mu.Unlock()
|
h.mu.Unlock()
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -130,41 +134,6 @@ func (h *PeeringHandler) processLoop() {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// pruneLoop runs periodically to remove old AS paths
|
|
||||||
func (h *PeeringHandler) pruneLoop() {
|
|
||||||
ticker := time.NewTicker(pathPruneInterval)
|
|
||||||
defer ticker.Stop()
|
|
||||||
|
|
||||||
for {
|
|
||||||
select {
|
|
||||||
case <-ticker.C:
|
|
||||||
h.prunePaths()
|
|
||||||
case <-h.stopCh:
|
|
||||||
return
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// prunePaths removes AS paths older than pathExpirationTime
|
|
||||||
func (h *PeeringHandler) prunePaths() {
|
|
||||||
cutoff := time.Now().Add(-pathExpirationTime)
|
|
||||||
var removed int
|
|
||||||
|
|
||||||
h.mu.Lock()
|
|
||||||
for pathKey, timestamp := range h.asPaths {
|
|
||||||
if timestamp.Before(cutoff) {
|
|
||||||
delete(h.asPaths, pathKey)
|
|
||||||
removed++
|
|
||||||
}
|
|
||||||
}
|
|
||||||
pathCount := len(h.asPaths)
|
|
||||||
h.mu.Unlock()
|
|
||||||
|
|
||||||
if removed > 0 {
|
|
||||||
h.logger.Debug("Pruned old AS paths", "removed", removed, "remaining", pathCount)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// ProcessPeeringsNow triggers immediate processing of all accumulated AS
|
// ProcessPeeringsNow triggers immediate processing of all accumulated AS
|
||||||
// paths into peering records. This bypasses the normal periodic processing
|
// paths into peering records. This bypasses the normal periodic processing
|
||||||
// schedule and is primarily intended for testing purposes.
|
// schedule and is primarily intended for testing purposes.
|
||||||
@@ -174,15 +143,16 @@ func (h *PeeringHandler) ProcessPeeringsNow() {
|
|||||||
|
|
||||||
// processPeerings extracts peerings from AS paths and writes to database
|
// processPeerings extracts peerings from AS paths and writes to database
|
||||||
func (h *PeeringHandler) processPeerings() {
|
func (h *PeeringHandler) processPeerings() {
|
||||||
// Take a snapshot of current AS paths
|
// Take the accumulated paths and replace the map with a fresh empty one
|
||||||
h.mu.RLock()
|
// under the lock. Each path is processed exactly once and the memory is
|
||||||
pathsCopy := make(map[string]time.Time, len(h.asPaths))
|
// released, so the map never grows past a single interval's traffic.
|
||||||
for k, v := range h.asPaths {
|
h.mu.Lock()
|
||||||
pathsCopy[k] = v
|
paths := h.asPaths
|
||||||
}
|
h.asPaths = make(map[string]time.Time)
|
||||||
h.mu.RUnlock()
|
dropped := h.droppedPaths
|
||||||
|
h.mu.Unlock()
|
||||||
|
|
||||||
if len(pathsCopy) == 0 {
|
if len(paths) == 0 {
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -192,7 +162,7 @@ func (h *PeeringHandler) processPeerings() {
|
|||||||
}
|
}
|
||||||
peerings := make(map[peeringKey]time.Time)
|
peerings := make(map[peeringKey]time.Time)
|
||||||
|
|
||||||
for pathJSON, timestamp := range pathsCopy {
|
for pathJSON, timestamp := range paths {
|
||||||
var path []int
|
var path []int
|
||||||
if err := json.Unmarshal([]byte(pathJSON), &path); err != nil {
|
if err := json.Unmarshal([]byte(pathJSON), &path); err != nil {
|
||||||
h.logger.Error("Failed to decode AS path", "error", err)
|
h.logger.Error("Failed to decode AS path", "error", err)
|
||||||
@@ -241,15 +211,16 @@ func (h *PeeringHandler) processPeerings() {
|
|||||||
}
|
}
|
||||||
|
|
||||||
h.logger.Info("Processed AS peerings",
|
h.logger.Info("Processed AS peerings",
|
||||||
"paths", len(pathsCopy),
|
"paths", len(paths),
|
||||||
"unique_peerings", len(peerings),
|
"unique_peerings", len(peerings),
|
||||||
"success", successCount,
|
"success", successCount,
|
||||||
|
"dropped_paths", dropped,
|
||||||
"duration", time.Since(start),
|
"duration", time.Since(start),
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
// Stop gracefully shuts down the handler by signaling the background
|
// Stop gracefully shuts down the handler by signaling the background
|
||||||
// goroutines to stop and performing a final synchronous processing of
|
// goroutine to stop and performing a final synchronous processing of
|
||||||
// any remaining AS paths. This ensures no peering data is lost during
|
// any remaining AS paths. This ensures no peering data is lost during
|
||||||
// shutdown.
|
// shutdown.
|
||||||
func (h *PeeringHandler) Stop() {
|
func (h *PeeringHandler) Stop() {
|
||||||
|
|||||||
@@ -0,0 +1,163 @@
|
|||||||
|
package routewatch
|
||||||
|
|
||||||
|
import (
|
||||||
|
"encoding/json"
|
||||||
|
"strconv"
|
||||||
|
"sync"
|
||||||
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"git.eeqj.de/sneak/routewatch/internal/database"
|
||||||
|
"git.eeqj.de/sneak/routewatch/internal/logger"
|
||||||
|
"git.eeqj.de/sneak/routewatch/internal/ristypes"
|
||||||
|
)
|
||||||
|
|
||||||
|
const (
|
||||||
|
testASNA = 64500
|
||||||
|
testASNB = 64501
|
||||||
|
testASNC = 64502
|
||||||
|
)
|
||||||
|
|
||||||
|
// recordingStore wraps mockStore to count every RecordPeering call, so a
|
||||||
|
// test can tell how many times a peering was written across separate runs.
|
||||||
|
type recordingStore struct {
|
||||||
|
*mockStore
|
||||||
|
|
||||||
|
mu sync.Mutex
|
||||||
|
calls int
|
||||||
|
}
|
||||||
|
|
||||||
|
func (r *recordingStore) RecordPeering(asA, asB int, ts time.Time) error {
|
||||||
|
r.mu.Lock()
|
||||||
|
r.calls++
|
||||||
|
r.mu.Unlock()
|
||||||
|
|
||||||
|
return r.mockStore.RecordPeering(asA, asB, ts)
|
||||||
|
}
|
||||||
|
|
||||||
|
func (r *recordingStore) callCount() int {
|
||||||
|
r.mu.Lock()
|
||||||
|
defer r.mu.Unlock()
|
||||||
|
|
||||||
|
return r.calls
|
||||||
|
}
|
||||||
|
|
||||||
|
// newTestHandler builds a PeeringHandler without starting the periodic
|
||||||
|
// processing goroutine, so tests drive processing explicitly.
|
||||||
|
func newTestHandler(db database.Store) *PeeringHandler {
|
||||||
|
return &PeeringHandler{
|
||||||
|
db: db,
|
||||||
|
logger: logger.New(),
|
||||||
|
asPaths: make(map[string]time.Time),
|
||||||
|
stopCh: make(chan struct{}),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func snapshot(h *PeeringHandler) (tracked, dropped int) {
|
||||||
|
h.mu.Lock()
|
||||||
|
defer h.mu.Unlock()
|
||||||
|
|
||||||
|
return len(h.asPaths), h.droppedPaths
|
||||||
|
}
|
||||||
|
|
||||||
|
func pathKey(t *testing.T, asns ...int) string {
|
||||||
|
t.Helper()
|
||||||
|
|
||||||
|
b, err := json.Marshal(ristypes.ASPath(asns))
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("failed to marshal path: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
return string(b)
|
||||||
|
}
|
||||||
|
|
||||||
|
func handle(h *PeeringHandler, ts time.Time, asns ...int) {
|
||||||
|
h.HandleMessage(&ristypes.RISMessage{
|
||||||
|
Path: ristypes.ASPath(asns),
|
||||||
|
ParsedTimestamp: ts,
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestPeeringHandlerProcessesEachRunAndEmpties verifies that a processing run
|
||||||
|
// empties the path map (the swap) and that a path seen again after a run is
|
||||||
|
// recorded in the next run too.
|
||||||
|
func TestPeeringHandlerProcessesEachRunAndEmpties(t *testing.T) {
|
||||||
|
store := &recordingStore{mockStore: newMockStore()}
|
||||||
|
h := newTestHandler(store)
|
||||||
|
|
||||||
|
now := time.Now().UTC()
|
||||||
|
|
||||||
|
// Run 1: one path, one peering recorded, map emptied afterwards.
|
||||||
|
handle(h, now, testASNA, testASNB)
|
||||||
|
h.ProcessPeeringsNow()
|
||||||
|
|
||||||
|
if tracked, _ := snapshot(h); tracked != 0 {
|
||||||
|
t.Fatalf("map not empty after first run: %d paths remain", tracked)
|
||||||
|
}
|
||||||
|
if got := store.callCount(); got != 1 {
|
||||||
|
t.Fatalf("want 1 RecordPeering call after first run, got %d", got)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Run 2: the same path again is recorded again (RecordPeering upserts).
|
||||||
|
handle(h, now.Add(time.Second), testASNA, testASNB)
|
||||||
|
h.ProcessPeeringsNow()
|
||||||
|
|
||||||
|
if tracked, _ := snapshot(h); tracked != 0 {
|
||||||
|
t.Fatalf("map not empty after second run: %d paths remain", tracked)
|
||||||
|
}
|
||||||
|
if got := store.callCount(); got != 2 {
|
||||||
|
t.Fatalf("want 2 RecordPeering calls after second run, got %d", got)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestPeeringHandlerCapDropsAndCounts verifies that a full map drops new paths
|
||||||
|
// and counts them, while a path already tracked is refreshed rather than
|
||||||
|
// dropped.
|
||||||
|
func TestPeeringHandlerCapDropsAndCounts(t *testing.T) {
|
||||||
|
store := &recordingStore{mockStore: newMockStore()}
|
||||||
|
h := newTestHandler(store)
|
||||||
|
|
||||||
|
now := time.Now().UTC()
|
||||||
|
|
||||||
|
// Fill the map to exactly maxTrackedPaths, including one real path key so
|
||||||
|
// the "already tracked" branch can be exercised. The filler keys are never
|
||||||
|
// processed in this test, so their contents do not matter.
|
||||||
|
existing := pathKey(t, testASNA, testASNB)
|
||||||
|
|
||||||
|
h.mu.Lock()
|
||||||
|
h.asPaths[existing] = now
|
||||||
|
for i := 0; len(h.asPaths) < maxTrackedPaths; i++ {
|
||||||
|
h.asPaths[strconv.Itoa(i)] = now
|
||||||
|
}
|
||||||
|
h.mu.Unlock()
|
||||||
|
|
||||||
|
// A new path is dropped and counted because the map is full.
|
||||||
|
handle(h, now.Add(time.Second), testASNA, testASNC)
|
||||||
|
|
||||||
|
tracked, dropped := snapshot(h)
|
||||||
|
if tracked != maxTrackedPaths {
|
||||||
|
t.Fatalf("want map size %d after drop, got %d", maxTrackedPaths, tracked)
|
||||||
|
}
|
||||||
|
if dropped != 1 {
|
||||||
|
t.Fatalf("want dropped count 1, got %d", dropped)
|
||||||
|
}
|
||||||
|
|
||||||
|
// A path already tracked is refreshed, not dropped.
|
||||||
|
refreshed := now.Add(2 * time.Second)
|
||||||
|
handle(h, refreshed, testASNA, testASNB)
|
||||||
|
|
||||||
|
tracked, dropped = snapshot(h)
|
||||||
|
if tracked != maxTrackedPaths {
|
||||||
|
t.Fatalf("want map size %d after refresh, got %d", maxTrackedPaths, tracked)
|
||||||
|
}
|
||||||
|
if dropped != 1 {
|
||||||
|
t.Fatalf("want dropped count still 1 after refresh, got %d", dropped)
|
||||||
|
}
|
||||||
|
|
||||||
|
h.mu.Lock()
|
||||||
|
gotTS := h.asPaths[existing]
|
||||||
|
h.mu.Unlock()
|
||||||
|
if !gotTS.Equal(refreshed) {
|
||||||
|
t.Fatalf("existing path timestamp not refreshed: want %v, got %v", refreshed, gotTS)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -14,9 +14,11 @@ import (
|
|||||||
)
|
)
|
||||||
|
|
||||||
const (
|
const (
|
||||||
// prefixHandlerQueueSize is the queue capacity for prefix tracking operations
|
// prefixHandlerQueueSize is the queue capacity for prefix tracking
|
||||||
// DO NOT set this higher than 100000 without explicit instructions
|
// operations, about 4 seconds of feed at peak. The streamer drops rather
|
||||||
prefixHandlerQueueSize = 100000
|
// than blocks when a queue is full, so this bounds memory. Batches still
|
||||||
|
// flush on a timer (prefixBatchTimeout).
|
||||||
|
prefixHandlerQueueSize = 20000
|
||||||
|
|
||||||
// prefixBatchSize is the number of prefix updates to batch together
|
// prefixBatchSize is the number of prefix updates to batch together
|
||||||
prefixBatchSize = 25000
|
prefixBatchSize = 25000
|
||||||
|
|||||||
+12
-67
@@ -179,35 +179,14 @@ func (s *Server) handleStatusJSON() http.HandlerFunc {
|
|||||||
|
|
||||||
metrics := s.streamer.GetMetrics()
|
metrics := s.streamer.GetMetrics()
|
||||||
|
|
||||||
// Get database stats with timeout
|
// Serve database statistics from the cache, which runs the table scans at
|
||||||
statsChan := make(chan database.Stats)
|
// most once per interval so this request does not.
|
||||||
errChan := make(chan error)
|
dbStats, err := s.stats.get()
|
||||||
|
if err != nil {
|
||||||
go func() {
|
|
||||||
dbStats, err := s.db.GetStatsContext(ctx)
|
|
||||||
if err != nil {
|
|
||||||
s.logger.Debug("Database stats query failed", "error", err)
|
|
||||||
errChan <- err
|
|
||||||
|
|
||||||
return
|
|
||||||
}
|
|
||||||
statsChan <- dbStats
|
|
||||||
}()
|
|
||||||
|
|
||||||
var dbStats database.Stats
|
|
||||||
select {
|
|
||||||
case <-ctx.Done():
|
|
||||||
s.logger.Error("Database stats timeout in status.json")
|
|
||||||
writeJSONError(w, http.StatusRequestTimeout, "Database timeout")
|
|
||||||
|
|
||||||
return
|
|
||||||
case err := <-errChan:
|
|
||||||
s.logger.Error("Failed to get database stats", "error", err)
|
s.logger.Error("Failed to get database stats", "error", err)
|
||||||
writeJSONError(w, http.StatusInternalServerError, err.Error())
|
writeJSONError(w, http.StatusInternalServerError, err.Error())
|
||||||
|
|
||||||
return
|
return
|
||||||
case dbStats = <-statsChan:
|
|
||||||
// Success
|
|
||||||
}
|
}
|
||||||
|
|
||||||
uptime := time.Since(metrics.ConnectedSince).Truncate(time.Second).String()
|
uptime := time.Since(metrics.ConnectedSince).Truncate(time.Second).String()
|
||||||
@@ -217,13 +196,6 @@ func (s *Server) handleStatusJSON() http.HandlerFunc {
|
|||||||
|
|
||||||
const bitsPerMegabit = 1000000.0
|
const bitsPerMegabit = 1000000.0
|
||||||
|
|
||||||
// Get route counts from database
|
|
||||||
ipv4Routes, ipv6Routes, err := s.db.GetLiveRouteCountsContext(ctx)
|
|
||||||
if err != nil {
|
|
||||||
s.logger.Warn("Failed to get live route counts", "error", err)
|
|
||||||
// Continue with zero counts
|
|
||||||
}
|
|
||||||
|
|
||||||
// Get route update metrics
|
// Get route update metrics
|
||||||
routeMetrics := s.streamer.GetMetricsTracker().GetRouteMetrics()
|
routeMetrics := s.streamer.GetMetricsTracker().GetRouteMetrics()
|
||||||
|
|
||||||
@@ -257,8 +229,8 @@ func (s *Server) handleStatusJSON() http.HandlerFunc {
|
|||||||
Peers: dbStats.Peers,
|
Peers: dbStats.Peers,
|
||||||
DatabaseSizeBytes: dbStats.FileSizeBytes,
|
DatabaseSizeBytes: dbStats.FileSizeBytes,
|
||||||
LiveRoutes: dbStats.LiveRoutes,
|
LiveRoutes: dbStats.LiveRoutes,
|
||||||
IPv4Routes: ipv4Routes,
|
IPv4Routes: dbStats.IPv4Routes,
|
||||||
IPv6Routes: ipv6Routes,
|
IPv6Routes: dbStats.IPv6Routes,
|
||||||
OldestRoute: dbStats.OldestRoute,
|
OldestRoute: dbStats.OldestRoute,
|
||||||
NewestRoute: dbStats.NewestRoute,
|
NewestRoute: dbStats.NewestRoute,
|
||||||
IPv4UpdatesPerSec: routeMetrics.IPv4UpdatesPerSec,
|
IPv4UpdatesPerSec: routeMetrics.IPv4UpdatesPerSec,
|
||||||
@@ -398,34 +370,14 @@ func (s *Server) handleStats() http.HandlerFunc {
|
|||||||
|
|
||||||
metrics := s.streamer.GetMetrics()
|
metrics := s.streamer.GetMetrics()
|
||||||
|
|
||||||
// Get database stats with timeout
|
// Serve database statistics from the cache, which runs the table scans at
|
||||||
statsChan := make(chan database.Stats)
|
// most once per interval so this request does not.
|
||||||
errChan := make(chan error)
|
dbStats, err := s.stats.get()
|
||||||
|
if err != nil {
|
||||||
go func() {
|
|
||||||
dbStats, err := s.db.GetStatsContext(ctx)
|
|
||||||
if err != nil {
|
|
||||||
s.logger.Debug("Database stats query failed", "error", err)
|
|
||||||
errChan <- err
|
|
||||||
|
|
||||||
return
|
|
||||||
}
|
|
||||||
statsChan <- dbStats
|
|
||||||
}()
|
|
||||||
|
|
||||||
var dbStats database.Stats
|
|
||||||
select {
|
|
||||||
case <-ctx.Done():
|
|
||||||
s.logger.Error("Database stats timeout")
|
|
||||||
// Don't write response here - timeout middleware already handles it
|
|
||||||
return
|
|
||||||
case err := <-errChan:
|
|
||||||
s.logger.Error("Failed to get database stats", "error", err)
|
s.logger.Error("Failed to get database stats", "error", err)
|
||||||
writeJSONError(w, http.StatusInternalServerError, err.Error())
|
writeJSONError(w, http.StatusInternalServerError, err.Error())
|
||||||
|
|
||||||
return
|
return
|
||||||
case dbStats = <-statsChan:
|
|
||||||
// Success
|
|
||||||
}
|
}
|
||||||
|
|
||||||
uptime := time.Since(metrics.ConnectedSince).Truncate(time.Second).String()
|
uptime := time.Since(metrics.ConnectedSince).Truncate(time.Second).String()
|
||||||
@@ -435,13 +387,6 @@ func (s *Server) handleStats() http.HandlerFunc {
|
|||||||
|
|
||||||
const bitsPerMegabit = 1000000.0
|
const bitsPerMegabit = 1000000.0
|
||||||
|
|
||||||
// Get route counts from database
|
|
||||||
ipv4Routes, ipv6Routes, err := s.db.GetLiveRouteCountsContext(ctx)
|
|
||||||
if err != nil {
|
|
||||||
s.logger.Warn("Failed to get live route counts", "error", err)
|
|
||||||
// Continue with zero counts
|
|
||||||
}
|
|
||||||
|
|
||||||
// Get route update metrics
|
// Get route update metrics
|
||||||
routeMetrics := s.streamer.GetMetricsTracker().GetRouteMetrics()
|
routeMetrics := s.streamer.GetMetricsTracker().GetRouteMetrics()
|
||||||
|
|
||||||
@@ -533,8 +478,8 @@ func (s *Server) handleStats() http.HandlerFunc {
|
|||||||
Peers: dbStats.Peers,
|
Peers: dbStats.Peers,
|
||||||
DatabaseSizeBytes: dbStats.FileSizeBytes,
|
DatabaseSizeBytes: dbStats.FileSizeBytes,
|
||||||
LiveRoutes: dbStats.LiveRoutes,
|
LiveRoutes: dbStats.LiveRoutes,
|
||||||
IPv4Routes: ipv4Routes,
|
IPv4Routes: dbStats.IPv4Routes,
|
||||||
IPv6Routes: ipv6Routes,
|
IPv6Routes: dbStats.IPv6Routes,
|
||||||
OldestRoute: dbStats.OldestRoute,
|
OldestRoute: dbStats.OldestRoute,
|
||||||
NewestRoute: dbStats.NewestRoute,
|
NewestRoute: dbStats.NewestRoute,
|
||||||
IPv4UpdatesPerSec: routeMetrics.IPv4UpdatesPerSec,
|
IPv4UpdatesPerSec: routeMetrics.IPv4UpdatesPerSec,
|
||||||
|
|||||||
@@ -0,0 +1,58 @@
|
|||||||
|
package server
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"net/http"
|
||||||
|
"net/http/httptest"
|
||||||
|
"sync/atomic"
|
||||||
|
"testing"
|
||||||
|
|
||||||
|
"git.eeqj.de/sneak/routewatch/internal/database"
|
||||||
|
"git.eeqj.de/sneak/routewatch/internal/logger"
|
||||||
|
"git.eeqj.de/sneak/routewatch/internal/metrics"
|
||||||
|
"git.eeqj.de/sneak/routewatch/internal/streamer"
|
||||||
|
)
|
||||||
|
|
||||||
|
// countingStatsDB embeds database.Store (left nil) and overrides only
|
||||||
|
// GetStatsContext, counting how many times it runs. The stats handlers read
|
||||||
|
// their database statistics through the cache, which calls this; every other
|
||||||
|
// Store method is unused on the stats path and would panic if called.
|
||||||
|
type countingStatsDB struct {
|
||||||
|
database.Store
|
||||||
|
calls *atomic.Int64
|
||||||
|
}
|
||||||
|
|
||||||
|
func (d countingStatsDB) GetStatsContext(_ context.Context) (database.Stats, error) {
|
||||||
|
d.calls.Add(1)
|
||||||
|
|
||||||
|
return database.Stats{}, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestStatsHandlersServeFromCache drives both stats handlers many times and
|
||||||
|
// checks that they answer 200 while the database statistics are computed at most
|
||||||
|
// once within the refresh interval. Before the fix each request ran the counts
|
||||||
|
// and MIN/MAX scans itself, which took the full timeout and returned 500 once
|
||||||
|
// the database grew large.
|
||||||
|
func TestStatsHandlersServeFromCache(t *testing.T) {
|
||||||
|
var calls atomic.Int64
|
||||||
|
db := countingStatsDB{calls: &calls}
|
||||||
|
s := New(db, streamer.New(logger.New(), metrics.New()), logger.New())
|
||||||
|
|
||||||
|
handlers := []http.HandlerFunc{s.handleStatusJSON(), s.handleStats()}
|
||||||
|
|
||||||
|
const iterations = 20
|
||||||
|
for _, handler := range handlers {
|
||||||
|
for range iterations {
|
||||||
|
req := httptest.NewRequest(http.MethodGet, "/", nil)
|
||||||
|
rec := httptest.NewRecorder()
|
||||||
|
handler(rec, req)
|
||||||
|
if rec.Code != http.StatusOK {
|
||||||
|
t.Fatalf("handler returned %d, want %d", rec.Code, http.StatusOK)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
if got := calls.Load(); got != 1 {
|
||||||
|
t.Fatalf("GetStatsContext ran %d times, want 1 within the interval", got)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -35,6 +35,7 @@ type Server struct {
|
|||||||
logger *logger.Logger
|
logger *logger.Logger
|
||||||
srv *http.Server
|
srv *http.Server
|
||||||
asnFetcher ASNFetcher
|
asnFetcher ASNFetcher
|
||||||
|
stats *statsCache
|
||||||
}
|
}
|
||||||
|
|
||||||
// New creates a new HTTP server
|
// New creates a new HTTP server
|
||||||
@@ -44,6 +45,9 @@ func New(db database.Store, streamer *streamer.Streamer, logger *logger.Logger)
|
|||||||
streamer: streamer,
|
streamer: streamer,
|
||||||
logger: logger,
|
logger: logger,
|
||||||
}
|
}
|
||||||
|
s.stats = newStatsCache(func(ctx context.Context) (database.Stats, error) {
|
||||||
|
return s.db.GetStatsContext(ctx)
|
||||||
|
})
|
||||||
|
|
||||||
s.setupRoutes()
|
s.setupRoutes()
|
||||||
|
|
||||||
|
|||||||
@@ -0,0 +1,112 @@
|
|||||||
|
package server
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"sync"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"git.eeqj.de/sneak/routewatch/internal/database"
|
||||||
|
)
|
||||||
|
|
||||||
|
const (
|
||||||
|
// statsRefreshInterval is how often the cached database statistics are
|
||||||
|
// recomputed. The scans behind GetStatsContext grow with the tables, so a
|
||||||
|
// request serves the cached copy instead of running them.
|
||||||
|
statsRefreshInterval = 30 * time.Second
|
||||||
|
|
||||||
|
// statsComputeTimeout bounds a single statistics computation so a stuck scan
|
||||||
|
// cannot block the refresh forever.
|
||||||
|
statsComputeTimeout = 20 * time.Second
|
||||||
|
)
|
||||||
|
|
||||||
|
// statsFetch computes fresh statistics. It is the expensive database scan that
|
||||||
|
// the cache runs at most once per interval.
|
||||||
|
type statsFetch func(ctx context.Context) (database.Stats, error)
|
||||||
|
|
||||||
|
// statsCache serves the most recent database statistics and recomputes them at
|
||||||
|
// most once per interval. The first request computes synchronously so it has
|
||||||
|
// real data to return; afterwards requests serve the cached copy immediately
|
||||||
|
// and a stale copy triggers a single background refresh, so no request waits on
|
||||||
|
// the scans.
|
||||||
|
type statsCache struct {
|
||||||
|
fetch statsFetch
|
||||||
|
interval time.Duration
|
||||||
|
now func() time.Time
|
||||||
|
|
||||||
|
mu sync.Mutex
|
||||||
|
stats database.Stats
|
||||||
|
haveStats bool
|
||||||
|
fetchedAt time.Time
|
||||||
|
refreshing bool
|
||||||
|
}
|
||||||
|
|
||||||
|
// newStatsCache returns a cache that recomputes statistics with fetch no more
|
||||||
|
// than once per statsRefreshInterval.
|
||||||
|
func newStatsCache(fetch statsFetch) *statsCache {
|
||||||
|
return &statsCache{
|
||||||
|
fetch: fetch,
|
||||||
|
interval: statsRefreshInterval,
|
||||||
|
now: time.Now,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// get returns the cached statistics. On the first call it computes them
|
||||||
|
// synchronously and returns any error. Later calls return the cached copy, and
|
||||||
|
// when that copy is older than the interval they start one background refresh.
|
||||||
|
func (c *statsCache) get() (database.Stats, error) {
|
||||||
|
c.mu.Lock()
|
||||||
|
|
||||||
|
if !c.haveStats {
|
||||||
|
// Cold start: compute once under the lock so concurrent first callers
|
||||||
|
// wait for this single computation rather than each starting their own.
|
||||||
|
stats, err := c.compute()
|
||||||
|
if err != nil {
|
||||||
|
c.mu.Unlock()
|
||||||
|
|
||||||
|
return database.Stats{}, err
|
||||||
|
}
|
||||||
|
c.store(stats)
|
||||||
|
c.mu.Unlock()
|
||||||
|
|
||||||
|
return stats, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
if c.now().Sub(c.fetchedAt) >= c.interval && !c.refreshing {
|
||||||
|
c.refreshing = true
|
||||||
|
go c.refresh()
|
||||||
|
}
|
||||||
|
|
||||||
|
stats := c.stats
|
||||||
|
c.mu.Unlock()
|
||||||
|
|
||||||
|
return stats, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// refresh recomputes the statistics in the background and replaces the cached
|
||||||
|
// copy. A failed computation leaves the previous copy in place.
|
||||||
|
func (c *statsCache) refresh() {
|
||||||
|
stats, err := c.compute()
|
||||||
|
|
||||||
|
c.mu.Lock()
|
||||||
|
defer c.mu.Unlock()
|
||||||
|
c.refreshing = false
|
||||||
|
if err == nil {
|
||||||
|
c.store(stats)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// compute runs the fetch with its own bounded context, independent of any
|
||||||
|
// request, so one request's cancellation cannot abort a shared refresh.
|
||||||
|
func (c *statsCache) compute() (database.Stats, error) {
|
||||||
|
ctx, cancel := context.WithTimeout(context.Background(), statsComputeTimeout)
|
||||||
|
defer cancel()
|
||||||
|
|
||||||
|
return c.fetch(ctx)
|
||||||
|
}
|
||||||
|
|
||||||
|
// store records a fresh result. The caller must hold the mutex.
|
||||||
|
func (c *statsCache) store(stats database.Stats) {
|
||||||
|
c.stats = stats
|
||||||
|
c.haveStats = true
|
||||||
|
c.fetchedAt = c.now()
|
||||||
|
}
|
||||||
@@ -0,0 +1,173 @@
|
|||||||
|
package server
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"errors"
|
||||||
|
"runtime"
|
||||||
|
"sync/atomic"
|
||||||
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"git.eeqj.de/sneak/routewatch/internal/database"
|
||||||
|
)
|
||||||
|
|
||||||
|
// testClock is a concurrency-safe clock the cache tests advance by hand, so the
|
||||||
|
// interval boundary is exercised without waiting real time.
|
||||||
|
type testClock struct {
|
||||||
|
ns atomic.Int64
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *testClock) now() time.Time { return time.Unix(0, c.ns.Load()) }
|
||||||
|
func (c *testClock) advance(d time.Duration) { c.ns.Add(int64(d)) }
|
||||||
|
|
||||||
|
// waitForCalls waits until calls reaches want, giving a background refresh time
|
||||||
|
// to finish.
|
||||||
|
func waitForCalls(calls *atomic.Int64, want int64) bool {
|
||||||
|
const attempts = 200
|
||||||
|
for range attempts {
|
||||||
|
if calls.Load() >= want {
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
time.Sleep(5 * time.Millisecond)
|
||||||
|
}
|
||||||
|
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestStatsCacheComputesOncePerInterval is the core guarantee: many reads in a
|
||||||
|
// row run the expensive fetch at most once per interval, and crossing the
|
||||||
|
// interval boundary allows exactly one more computation.
|
||||||
|
func TestStatsCacheComputesOncePerInterval(t *testing.T) {
|
||||||
|
clk := &testClock{}
|
||||||
|
clk.ns.Store(int64(time.Hour)) // start at a non-zero instant
|
||||||
|
|
||||||
|
var calls atomic.Int64
|
||||||
|
c := newStatsCache(func(_ context.Context) (database.Stats, error) {
|
||||||
|
calls.Add(1)
|
||||||
|
|
||||||
|
return database.Stats{}, nil
|
||||||
|
})
|
||||||
|
c.now = clk.now
|
||||||
|
|
||||||
|
const reads = 50
|
||||||
|
for range reads {
|
||||||
|
if _, err := c.get(); err != nil {
|
||||||
|
t.Fatalf("get returned error: %v", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if got := calls.Load(); got != 1 {
|
||||||
|
t.Fatalf("fetch ran %d times within the interval, want 1", got)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Cross the interval: the next read serves the stale copy and starts one
|
||||||
|
// background refresh.
|
||||||
|
clk.advance(c.interval)
|
||||||
|
if _, err := c.get(); err != nil {
|
||||||
|
t.Fatalf("get after interval returned error: %v", err)
|
||||||
|
}
|
||||||
|
if !waitForCalls(&calls, 2) {
|
||||||
|
t.Fatalf("background refresh did not run, fetch ran %d times", calls.Load())
|
||||||
|
}
|
||||||
|
|
||||||
|
for range reads {
|
||||||
|
if _, err := c.get(); err != nil {
|
||||||
|
t.Fatalf("get returned error: %v", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if got := calls.Load(); got != 2 {
|
||||||
|
t.Fatalf("fetch ran %d times across one interval boundary, want 2", got)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestStatsCacheColdStartReturnsError checks the first computation's error
|
||||||
|
// reaches the caller, since there is no cached copy to serve instead.
|
||||||
|
func TestStatsCacheColdStartReturnsError(t *testing.T) {
|
||||||
|
wantErr := errors.New("boom")
|
||||||
|
c := newStatsCache(func(_ context.Context) (database.Stats, error) {
|
||||||
|
return database.Stats{}, wantErr
|
||||||
|
})
|
||||||
|
|
||||||
|
if _, err := c.get(); !errors.Is(err, wantErr) {
|
||||||
|
t.Fatalf("get returned %v, want %v", err, wantErr)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestStatsCacheServesLastGoodCopyOnRefreshError checks that once a copy exists,
|
||||||
|
// a later failing refresh does not surface an error or drop the good data.
|
||||||
|
func TestStatsCacheServesLastGoodCopyOnRefreshError(t *testing.T) {
|
||||||
|
clk := &testClock{}
|
||||||
|
clk.ns.Store(int64(time.Hour))
|
||||||
|
|
||||||
|
const wantASNs = 7
|
||||||
|
var calls atomic.Int64
|
||||||
|
var failing atomic.Bool
|
||||||
|
c := newStatsCache(func(_ context.Context) (database.Stats, error) {
|
||||||
|
calls.Add(1)
|
||||||
|
if failing.Load() {
|
||||||
|
return database.Stats{}, errors.New("boom")
|
||||||
|
}
|
||||||
|
|
||||||
|
return database.Stats{ASNs: wantASNs}, nil
|
||||||
|
})
|
||||||
|
c.now = clk.now
|
||||||
|
|
||||||
|
got, err := c.get()
|
||||||
|
if err != nil || got.ASNs != wantASNs {
|
||||||
|
t.Fatalf("cold start returned (%+v, %v), want ASNs=%d, nil", got, err, wantASNs)
|
||||||
|
}
|
||||||
|
|
||||||
|
failing.Store(true)
|
||||||
|
clk.advance(c.interval)
|
||||||
|
|
||||||
|
got, err = c.get()
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("get during failing refresh returned error: %v", err)
|
||||||
|
}
|
||||||
|
if got.ASNs != wantASNs {
|
||||||
|
t.Fatalf("get returned ASNs=%d, want the last good copy %d", got.ASNs, wantASNs)
|
||||||
|
}
|
||||||
|
if !waitForCalls(&calls, 2) {
|
||||||
|
t.Fatalf("refresh was not attempted, fetch ran %d times", calls.Load())
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestStatsCacheBackgroundRefreshDoesNotLeak forces many stale refreshes and
|
||||||
|
// checks the goroutine count returns to its starting value.
|
||||||
|
func TestStatsCacheBackgroundRefreshDoesNotLeak(t *testing.T) {
|
||||||
|
clk := &testClock{}
|
||||||
|
clk.ns.Store(int64(time.Hour))
|
||||||
|
|
||||||
|
var calls atomic.Int64
|
||||||
|
c := newStatsCache(func(_ context.Context) (database.Stats, error) {
|
||||||
|
calls.Add(1)
|
||||||
|
|
||||||
|
return database.Stats{}, nil
|
||||||
|
})
|
||||||
|
c.now = clk.now
|
||||||
|
|
||||||
|
if _, err := c.get(); err != nil {
|
||||||
|
t.Fatalf("cold start returned error: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
baseline := runtime.NumGoroutine()
|
||||||
|
|
||||||
|
const rounds = 20
|
||||||
|
for i := range rounds {
|
||||||
|
clk.advance(c.interval)
|
||||||
|
if _, err := c.get(); err != nil {
|
||||||
|
t.Fatalf("get returned error: %v", err)
|
||||||
|
}
|
||||||
|
if !waitForCalls(&calls, int64(i+2)) {
|
||||||
|
t.Fatalf("refresh %d did not run", i)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
const settleAttempts = 100
|
||||||
|
for range settleAttempts {
|
||||||
|
if runtime.NumGoroutine() <= baseline {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
time.Sleep(10 * time.Millisecond)
|
||||||
|
}
|
||||||
|
t.Fatalf("goroutines did not settle to baseline %d, got %d", baseline, runtime.NumGoroutine())
|
||||||
|
}
|
||||||
@@ -106,6 +106,7 @@ type handlerInfo struct {
|
|||||||
type Streamer struct {
|
type Streamer struct {
|
||||||
logger *logger.Logger
|
logger *logger.Logger
|
||||||
client *http.Client
|
client *http.Client
|
||||||
|
url string
|
||||||
handlers []*handlerInfo
|
handlers []*handlerInfo
|
||||||
rawHandler RawMessageHandler
|
rawHandler RawMessageHandler
|
||||||
mu sync.RWMutex
|
mu sync.RWMutex
|
||||||
@@ -124,6 +125,7 @@ type Streamer struct {
|
|||||||
func New(logger *logger.Logger, metrics *metrics.Tracker) *Streamer {
|
func New(logger *logger.Logger, metrics *metrics.Tracker) *Streamer {
|
||||||
return &Streamer{
|
return &Streamer{
|
||||||
logger: logger,
|
logger: logger,
|
||||||
|
url: risLiveURL,
|
||||||
client: &http.Client{
|
client: &http.Client{
|
||||||
Timeout: 0, // No timeout for streaming
|
Timeout: 0, // No timeout for streaming
|
||||||
Transport: &http.Transport{
|
Transport: &http.Transport{
|
||||||
@@ -463,7 +465,14 @@ func (s *Streamer) streamWithReconnect(ctx context.Context) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (s *Streamer) stream(ctx context.Context) error {
|
func (s *Streamer) stream(ctx context.Context) error {
|
||||||
req, err := http.NewRequestWithContext(ctx, "GET", risLiveURL, nil)
|
// connCtx is scoped to this single connection: cancelling it when stream
|
||||||
|
// returns stops the ticker goroutines below, so a reconnect does not leak
|
||||||
|
// them. Without this they would live until the streamer's lifetime context
|
||||||
|
// is cancelled, leaking two per reconnect.
|
||||||
|
connCtx, connCancel := context.WithCancel(ctx)
|
||||||
|
defer connCancel()
|
||||||
|
|
||||||
|
req, err := http.NewRequestWithContext(ctx, "GET", s.url, nil)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return fmt.Errorf("failed to create request: %w", err)
|
return fmt.Errorf("failed to create request: %w", err)
|
||||||
}
|
}
|
||||||
@@ -516,7 +525,7 @@ func (s *Streamer) stream(ctx context.Context) error {
|
|||||||
select {
|
select {
|
||||||
case <-metricsTicker.C:
|
case <-metricsTicker.C:
|
||||||
s.logMetrics()
|
s.logMetrics()
|
||||||
case <-ctx.Done():
|
case <-connCtx.Done():
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -536,7 +545,7 @@ func (s *Streamer) stream(ctx context.Context) error {
|
|||||||
s.metrics.RecordWireBytes(delta)
|
s.metrics.RecordWireBytes(delta)
|
||||||
lastWireBytes = currentBytes
|
lastWireBytes = currentBytes
|
||||||
}
|
}
|
||||||
case <-ctx.Done():
|
case <-connCtx.Done():
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,7 +1,12 @@
|
|||||||
package streamer
|
package streamer
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"context"
|
||||||
|
"net/http"
|
||||||
|
"net/http/httptest"
|
||||||
|
"runtime"
|
||||||
"testing"
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
"git.eeqj.de/sneak/routewatch/internal/logger"
|
"git.eeqj.de/sneak/routewatch/internal/logger"
|
||||||
"git.eeqj.de/sneak/routewatch/internal/metrics"
|
"git.eeqj.de/sneak/routewatch/internal/metrics"
|
||||||
@@ -32,3 +37,68 @@ func TestNewStreamer(t *testing.T) {
|
|||||||
t.Error("metrics tracker not set correctly")
|
t.Error("metrics tracker not set correctly")
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// TestStreamDoesNotLeakTickersAcrossReconnects drives many short-lived
|
||||||
|
// connections (each stream call is one reconnect cycle) and asserts the
|
||||||
|
// goroutine count returns to its starting value. Each connection starts two
|
||||||
|
// ticker goroutines; before the fix they lived until the streamer's lifetime
|
||||||
|
// context was cancelled, so every reconnect leaked two.
|
||||||
|
func TestStreamDoesNotLeakTickersAcrossReconnects(t *testing.T) {
|
||||||
|
// The handler returns immediately, so the response body is empty and each
|
||||||
|
// stream call ends at once, standing in for a dropped connection.
|
||||||
|
srv := httptest.NewServer(http.HandlerFunc(func(_ http.ResponseWriter, _ *http.Request) {}))
|
||||||
|
defer srv.Close()
|
||||||
|
|
||||||
|
s := New(logger.New(), metrics.New())
|
||||||
|
s.url = srv.URL
|
||||||
|
|
||||||
|
// One warm-up connection so any persistent HTTP transport goroutine exists
|
||||||
|
// before we take the baseline.
|
||||||
|
if err := s.stream(context.Background()); err != nil {
|
||||||
|
t.Fatalf("warm-up stream returned error: %v", err)
|
||||||
|
}
|
||||||
|
s.client.CloseIdleConnections()
|
||||||
|
|
||||||
|
baseline := settledGoroutineCount()
|
||||||
|
|
||||||
|
const reconnects = 20
|
||||||
|
for range reconnects {
|
||||||
|
if err := s.stream(context.Background()); err != nil {
|
||||||
|
t.Fatalf("stream returned error: %v", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
s.client.CloseIdleConnections()
|
||||||
|
|
||||||
|
if !waitForGoroutines(baseline) {
|
||||||
|
t.Fatalf("goroutines did not return to baseline %d after %d reconnects, got %d",
|
||||||
|
baseline, reconnects, runtime.NumGoroutine())
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// settledGoroutineCount lets transient goroutines finish, then reports the
|
||||||
|
// current count.
|
||||||
|
func settledGoroutineCount() int {
|
||||||
|
prev := runtime.NumGoroutine()
|
||||||
|
for range 20 {
|
||||||
|
time.Sleep(10 * time.Millisecond)
|
||||||
|
cur := runtime.NumGoroutine()
|
||||||
|
if cur == prev {
|
||||||
|
return cur
|
||||||
|
}
|
||||||
|
prev = cur
|
||||||
|
}
|
||||||
|
|
||||||
|
return prev
|
||||||
|
}
|
||||||
|
|
||||||
|
// waitForGoroutines waits until the goroutine count drops to target or below.
|
||||||
|
func waitForGoroutines(target int) bool {
|
||||||
|
for range 100 {
|
||||||
|
if runtime.NumGoroutine() <= target {
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
time.Sleep(10 * time.Millisecond)
|
||||||
|
}
|
||||||
|
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
|||||||
+2
-2
@@ -7,9 +7,9 @@ ROOT="$(cd "$(dirname "$0")/.." && pwd -P)"
|
|||||||
|
|
||||||
main() {
|
main() {
|
||||||
cd "$ROOT"
|
cd "$ROOT"
|
||||||
go test -timeout 30s -race -cover ./... || {
|
go test -short -timeout 30s -race -cover ./... || {
|
||||||
echo "--- Rerunning with -v for details ---"
|
echo "--- Rerunning with -v for details ---"
|
||||||
go test -timeout 30s -race -v ./...
|
go test -short -timeout 30s -race -v ./...
|
||||||
exit 1
|
exit 1
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user