Author SHA1 Message Date
sneak a2c9d68bfe Create new databases with auto_vacuum incremental (closes #43)
check / check (push) Successful in 3m15s
SQLite only accepts auto_vacuum before the database file is first
written. The connection switched to WAL first, which writes the file,
so the PRAGMA auto_vacuum in Initialize came too late and was ignored.
The setting now goes in the connection string, which the driver applies
on open before the journal mode, and the late PRAGMA is removed.
Vacuum now reads every row PRAGMA incremental_vacuum returns: SQLite
frees one page per row, and the single step ExecContext takes freed
only one page per call. Tests check that every pooled connection sees
auto_vacuum incremental on a new database and that one Vacuum call
frees every page left by deleting routes.

Model: opus-5-5
2026-10-03 14:08:08 +00:00
clawbot 32704c601e Look up IP addresses by prefix at each mask length (closes #48)
check / check (push) Successful in 3m13s
Looking up an address on /ip/ and /api/v1/ip/ read every IPv6 live route and about half of the IPv4 range index, so on a day-sized database it passed the 30-second request timeout. The lookup is now one function for both families: from the longest mask length down, it looks up the address's network prefix on the existing prefix index, and the first live route wins.

The feed sends IPv6 withdrawals uncompressed, so they never matched the stored compressed prefix and never removed a route. Prefixes are now stored in the text form net/netip prints, so IPv6 withdrawals take effect.

ip_start, ip_end and GetASInfoForIP are removed; nothing read them. A database from before this change must be deleted.

Model: opus-5-5
2026-10-03 16:04:38 +02:00
clawbot 7322f936e4 Stamp the git tag or short commit in a plain docker build (closes #46)
check / check (push) Successful in 3m6s
A plain `docker build .` stamped `unknown` into the page footer:
.dockerignore left out .git, and the Dockerfile built without the -X
flags. The context now carries .git without .git/config, which can hold
a credential. The build stage takes the VERSION build argument when one
is given, otherwise `git describe --tags --always`, fails the build if
.git is present and no version comes out, and stamps it together with
the full commit the footer links to. `make build` now uses
`git describe` as well, and script/docker is the current shared copy.

Model: opus-5-5
2026-10-02 08:17:35 +02:00
clawbot 6422d9fa0c Container makes its data directory usable itself (closes #42)
check / check (push) Successful in 8s
sneak's standing rule: the container makes its data directory usable itself, with no step on the host. entrypoint.sh now creates /var/lib/berlin.sneak.app.routewatch if it is missing and stops the start when any step fails (set -euo pipefail); before, a failed cd went on to change the ownership of whatever directory the script was in, and a failed chown still started the daemon. Taking ownership of the directory and switching to the routewatch user through setpriv are unchanged. The README's upaas volume line now says only which path to mount.

The empty-directory and other-uid cases were run by hand on the built image with upaas-style bind mounts, not added as an automated test.

Model: opus-5-5
2026-09-29 12:23:42 +02:00
clawbot 057e0bd9a9 Add a .dockerignore (closes #39)
check / check (push) Successful in 6s
Adds a .dockerignore, which REPO_POLICIES.md lists among the files every repo has. Until now every docker build sent the whole working tree as its build context, and the image's source archive is made from that context. It now keeps out .git, local build and test output, archives, local databases, .env and a local Go workspace. Every tracked file stays in, so the format check, the linter, the tests, the build and the source archive read the same files as before.

Without .git in the context the binary's build info records no git revision; nothing in routewatch reads it.

Model: opus-5-5
2026-09-29 07:42:20 +02:00
clawbot 87a3bc115e README: license and author up front; bring TODO.md up to date (closes #38)
check / check (push) Successful in 3m14s
Brings README.md and TODO.md in line with REPO_POLICIES.md and with what is on next. The README's first line now names routewatch, what it is, its MIT licence and its author, @sneak; the License section says MIT and links LICENSE; a new Author section names @sneak.

TODO.md's Status, Next Step and Future Steps now describe next as it is: it waits for sneak to merge the milestone PR, after which the upaas deploy and the run under a real 5 GiB limit are his. Completed Steps gains a line for each issue that landed without one.

Docs only. make fmt formats only Go here, so the markdown was wrapped by hand.

Model: opus-5-5
2026-09-29 07:06:30 +02:00
clawbot d28a59023a Add MIT LICENSE (closes #1)
check / check (push) Successful in 4m5s
Adds LICENSE with the MIT licence sneak chose for routewatch, copied byte for byte from sneak/webhooker's LICENSE: copyright 2026, Jeffrey Paul. The rest of the issue, .editorconfig and the Gitea check workflow, was already on next, and the README already points at LICENSE.

Judgement call: our repos write the copyright holder in several different ways and sneak named no form; webhooker's gives his full name and email address, the most explicit of them.

Model: opus-5-5
2026-09-29 05:05:35 +02:00
clawbot 1ac24669d3 Stop the streamer without sending on closed queues (closes #34)
check / check (push) Successful in 2m38s
Stopping the daemon while the RIS Live feed was flowing could panic with "send on closed channel" and skip the rest of the shutdown. Stop cancels the stream and closes the handler queues under the streamer's write lock, but the read loop checked for a stop only before parsing each line. It now checks again under the read lock it already takes just before handing a message to the queues, so it never sends to a closed queue. Stop also clears its cancel function and returns early when there is none, so a second call no longer closes the queues again. A test forces both cases.

Behaviour change: Stop before Start now does nothing.

Model: opus-5-5
2026-09-28 21:42:33 +02:00
clawbot 6187ac8503 Let the daemon receive docker stop's signal itself (closes #33)
check / check (push) Successful in 2m42s
entrypoint.sh now switches to the routewatch user (UID 1000) with setpriv instead of runuser. setpriv replaces itself with the daemon, so the daemon is the container's main process and receives docker stop's signal itself. runuser stayed in between, passed the signal on and killed the daemon 2 seconds later, so every stop ended with exit 143. The daemon now gets the whole wait the caller allows, up to its own 60-second limit, and a clean stop exits 0. Taking ownership of the state directory and the MALLOC_ARENA_MAX check still run as root first.

Not fixed here: a stop while the feed is flowing can still panic (#34).

Model: opus-5-5
2026-09-28 21:08:33 +02:00
clawbot f2a9e90625 Refuse invalid settings at start, health check follows PORT (closes #31)
check / check (push) Successful in 3m17s
A set but invalid PORT (anything but plain digits from 1 to 65535) or a relative XDG_DATA_HOME now stops the start before the database opens; before, a bad PORT left the daemon running without HTTP. entrypoint.sh refuses a MALLOC_ARENA_MAX that is not a positive whole number, since glibc ignores a bad one silently. The HEALTHCHECK probes the port PORT names, 8080 when unset. PORT is now read in internal/config, so server.New takes the config.

The README gives the real Linux state directory, lists XDG_DATA_HOME, and adds "Running under upaas": port, volume, environment, the 5g memory limit and the health check.

Unverified: the 5g memory limit could not be exercised on the build host.

Model: opus-5-5
2026-09-28 20:08:42 +02:00
clawbot df9e23d503 Serve /api/v1/stats counts from realtime in-memory counters (closes #27)
check / check (push) Successful in 2m36s
Realtime in-memory counters seeded at startup and adjusted on every insert, update and delete; no periodic recompute. Independent review passed: #29 (comment)

model: claude-opus-4-8 (implementation and review); merged by claude-fable-5
2026-09-22 09:41:16 +02:00
26 changed files with 1524 additions and 960 deletions
+32
View File
@@ -0,0 +1,32 @@
# Docker does not read .gitignore, and a pattern here matches from the root of
# the build context only: a pattern meant for every directory needs `**/`.
# .git is sent without its config. Without a VERSION build argument the
# stage that compiles runs `git describe --tags --always` on .git, which
# does not need .git/config; that file can hold a credential, such as a
# password in a remote URL or the token the CI checkout step stores there.
.git/config
# Local build and debug output: `make build`, `make run`, `make asupdate`, test
# binaries, coverage profiles and source archives.
/bin
/log.txt
/out
/pkg/asinfo/asdata.json
**/*.tar.zst
**/*.test
**/*.out
**/*.tmp
# Local databases and secrets. The image carries a source archive of the whole
# build context, so these would otherwise ship inside it.
**/*.db
**/*.db-journal
**/*.db-wal
**/.env
# A local Go workspace points at directories outside the build context.
/go.work
/go.work.sum
**/.DS_Store
+21 -6
View File
@@ -40,8 +40,23 @@ RUN go mod download && go mod vendor
# installed above. The suite is offline (the live-feed test is opt-in). # installed above. The suite is offline (the live-feed test is opt-in).
RUN make test RUN make test
# Build the binary with CGO enabled (required for sqlite3) # Build the binary with CGO enabled (required for sqlite3). The version the
RUN CGO_ENABLED=1 GOOS=linux go build -o /routewatch ./cmd/routewatch # page footer shows is the VERSION build argument when one is given, otherwise
# `git describe --tags --always` of the .git in the build context (git comes
# with this image): the tag on a tagged commit, tag-N-gHASH after one, the
# short commit when no tag is reachable. A context that carries .git and still
# yields no version fails the build. The footer links to the full commit.
ARG VERSION
RUN version="${VERSION:-$(git describe --tags --always || echo unknown)}"; \
if [ -e .git ] && { [ -z "$version" ] || [ "$version" = dev ] || \
[ "$version" = unknown ]; }; then \
echo "no version could be derived although the build context carries .git" >&2; \
exit 1; \
fi; \
CGO_ENABLED=1 GOOS=linux go build -o /routewatch -ldflags "\
-X git.eeqj.de/sneak/routewatch/internal/version.GitRevision=$(git rev-parse --verify HEAD || echo unknown) \
-X git.eeqj.de/sneak/routewatch/internal/version.GitRevisionShort=$version" \
./cmd/routewatch
# Create source archive with vendored dependencies # Create source archive with vendored dependencies
RUN tar --zstd -cf /routewatch-source.tar.zst \ RUN tar --zstd -cf /routewatch-source.tar.zst \
@@ -79,8 +94,8 @@ 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 # 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 # container's memory limit is reached. setpriv in the entrypoint preserves this
# XDG_DATA_HOME above. # the way it does XDG_DATA_HOME above.
ENV GOMEMLIMIT=1536MiB ENV GOMEMLIMIT=1536MiB
# Cap glibc's malloc arenas. The SQLite C library allocates and frees millions # Cap glibc's malloc arenas. The SQLite C library allocates and frees millions
@@ -96,8 +111,8 @@ EXPOSE 8080
COPY ./entrypoint.sh /entrypoint.sh COPY ./entrypoint.sh /entrypoint.sh
# Health check using the health endpoint # Health check using the health endpoint, on the port PORT names
HEALTHCHECK --interval=30s --timeout=5s --start-period=10s --retries=3 \ HEALTHCHECK --interval=30s --timeout=5s --start-period=10s --retries=3 \
CMD curl -sf http://localhost:8080/.well-known/healthcheck.json || exit 1 CMD curl -sf "http://localhost:${PORT:-8080}/.well-known/healthcheck.json" || exit 1
ENTRYPOINT ["/bin/bash", "/entrypoint.sh" ] ENTRYPOINT ["/bin/bash", "/entrypoint.sh" ]
+21
View File
@@ -0,0 +1,21 @@
MIT License
Copyright (c) 2026 Jeffrey Paul <sneak@sneak.berlin>
Permission is hereby granted, free of charge, to any person obtaining a copy
of this software and associated documentation files (the "Software"), to deal
in the Software without restriction, including without limitation the rights
to use, copy, modify, merge, publish, distribute, sublicense, and/or sell
copies of the Software, and to permit persons to whom the Software is
furnished to do so, subject to the following conditions:
The above copyright notice and this permission notice shall be included in all
copies or substantial portions of the Software.
THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,
FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE
AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER
LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,
OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE
SOFTWARE.
+1 -1
View File
@@ -2,7 +2,7 @@ export DEBUG = routewatch
# Git revision for version embedding # Git revision for version embedding
GIT_REVISION := $(shell git rev-parse HEAD 2>/dev/null || echo "unknown") GIT_REVISION := $(shell git rev-parse HEAD 2>/dev/null || echo "unknown")
GIT_REVISION_SHORT := $(shell git rev-parse --short HEAD 2>/dev/null || echo "unknown") GIT_REVISION_SHORT := $(shell git describe --tags --always 2>/dev/null || echo "unknown")
VERSION_PKG := git.eeqj.de/sneak/routewatch/internal/version VERSION_PKG := git.eeqj.de/sneak/routewatch/internal/version
LDFLAGS := -X $(VERSION_PKG).GitRevision=$(GIT_REVISION) -X $(VERSION_PKG).GitRevisionShort=$(GIT_REVISION_SHORT) LDFLAGS := -X $(VERSION_PKG).GitRevision=$(GIT_REVISION) -X $(VERSION_PKG).GitRevisionShort=$(GIT_REVISION_SHORT)
+40 -7
View File
@@ -1,6 +1,9 @@
# RouteWatch # RouteWatch
RouteWatch is a real-time BGP routing table monitor that streams BGP UPDATE messages from the RIPE RIS Live service, maintains a live routing table in SQLite, and provides HTTP APIs for querying routing information. RouteWatch is an MIT-licensed Go daemon by @sneak that monitors the BGP routing
table in real time: it streams BGP UPDATE messages from the RIPE RIS Live
service, maintains a live routing table in SQLite, and provides HTTP APIs for
querying routing information.
## Features ## Features
@@ -139,7 +142,10 @@ routewatch/
- **Backpressure**: Probabilistic message dropping when queues exceed 50% capacity - **Backpressure**: Probabilistic message dropping when queues exceed 50% capacity
- **Graceful Shutdown**: 60-second timeout, flushes all pending batches - **Graceful Shutdown**: 60-second timeout, flushes all pending batches
- **Reconnection**: Exponential backoff (5s-320s) with reset after 30s of stable connection - **Reconnection**: Exponential backoff (5s-320s) with reset after 30s of stable connection
- **IPv4 Optimization**: IP ranges stored as uint32 for O(1) lookups - **IP Lookup**: the most specific live route is found by looking up the
address's prefix at each mask length, longest first, in the prefix index
(at most 33 lookups for IPv4, 129 for IPv6); prefixes are stored in the
text form Go's `net/netip` prints
### Database Schema ### Database Schema
@@ -151,7 +157,7 @@ prefixes_v6(id, prefix, mask_length, first_seen, last_seen)
-- Live routing tables (one per IP version) -- Live routing tables (one per IP version)
live_routes_v4(id, prefix, mask_length, origin_asn, peer_ip, as_path, live_routes_v4(id, prefix, mask_length, origin_asn, peer_ip, as_path,
next_hop, last_updated, v4_ip_start, v4_ip_end) next_hop, last_updated)
live_routes_v6(id, prefix, mask_length, origin_asn, peer_ip, as_path, live_routes_v6(id, prefix, mask_length, origin_asn, peer_ip, as_path,
next_hop, last_updated) next_hop, last_updated)
@@ -166,14 +172,21 @@ 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, a whole number from 1 to 65535 |
| `DEBUG` | (empty) | Set to `routewatch` for debug logging | | `DEBUG` | (empty) | Set to `routewatch` for debug logging |
| `XDG_DATA_HOME` | `/var/lib` (in the Docker image) | Base of the state directory; must be an absolute path |
| `GOMEMLIMIT` | `1536MiB` (in the Docker image) | Go soft memory limit; see Memory | | `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 | | `MALLOC_ARENA_MAX` | `2` (in the Docker image) | glibc malloc arena cap, a positive whole number; see Memory |
A variable that is set to an invalid value stops the start with an error and a
non-zero exit. An empty variable counts as unset.
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/berlin.sneak.app.routewatch/` when running as root,
otherwise `$XDG_DATA_HOME/berlin.sneak.app.routewatch/` (with `XDG_DATA_HOME`
unset, `~/.local/share/berlin.sneak.app.routewatch/`). In the Docker image
this is `/var/lib/berlin.sneak.app.routewatch/`.
## Memory ## Memory
@@ -220,6 +233,22 @@ What happens at each limit:
With `DEBUG=routewatch` the daemon logs a `System stats` line every 60 seconds With `DEBUG=routewatch` the daemon logs a `System stats` line every 60 seconds
with the goroutine count and Go memory figures. with the goroutine count and Go memory figures.
## Running under upaas
What the [upaas](https://git.eeqj.de/sneak/upaas) app needs:
- Container port: `8080`.
- Volume: one, at container path `/var/lib/berlin.sneak.app.routewatch`.
- Environment: nothing is required. Leave `XDG_DATA_HOME`, `GOMEMLIMIT` and
`MALLOC_ARENA_MAX` at the image's values. `DEBUG=routewatch` is optional and
adds the `System stats` memory line to the log.
- Memory Limit: `5g`, the 5 GiB limit from Memory above. upaas sets no swap
limit, so on a host with swap Docker allows the same amount of swap again.
- Health check: the image's `HEALTHCHECK` requests
`/.well-known/healthcheck.json` on the container port. upaas reads the
container's health 60 seconds after a deploy and fails the deploy unless it
is `healthy`.
## Development ## Development
```bash ```bash
@@ -274,4 +303,8 @@ invoke `go test` directly without `-short`:
## License ## License
See LICENSE file. MIT. See [`LICENSE`](LICENSE).
## Author
[@sneak](https://sneak.berlin)
+87 -12
View File
@@ -10,19 +10,93 @@
# Status # Status
pre-1.0. No git tags. Runs in production-style Docker deployment, but pre-1.0. No git tags. The Docker build runs the format check, the linter
the policy compliance branch (repo-policies-compliance, make check and the tests, and the Gitea workflow runs that build on every push. The
passing, clean tree) is unmerged to main and the CI workflow is missing. image sets memory ceilings for a 5 GiB container (README "Memory") and the
README says how to run it under upaas (README "Running under upaas"). A
35-hour run of `3898daa` on the live feed peaked at about 1 GiB, without a
container memory limit.
# Next Step # Next Step
Merge repo-policies-compliance into main (3 commits: policy files and `next` waits for sneak to merge it to `main` through
.gitignore, Makefile targets fmt-check/check/docker/hooks, gofmt pass), https://git.eeqj.de/sneak/routewatch/pulls/6. After that, setting
then add .gitea/workflows/check.yml as a small follow-up commit so CI routewatch up under upaas on fsn1app1 and deploying it are his
runs make check on main. (https://git.eeqj.de/sneak/routewatch/issues/31), and so is the run under a
real 5 GiB limit (https://git.eeqj.de/sneak/routewatch/issues/3).
The other open issue is https://git.eeqj.de/sneak/routewatch/issues/30.
# Completed Steps # Completed Steps
- 2026-10-03: a new database is created with `auto_vacuum` set to
incremental, through the connection string so it is set before the file
is first written, and each periodic incremental vacuum now returns up to
1000 free pages; the late `PRAGMA auto_vacuum` in `Initialize`, which
SQLite ignored, is gone (closes #43)
- 2026-10-03: looking up an IP address no longer reads every IPv6 route:
both families find the most specific live route with at most 33 or 129
lookups on the prefix index, and the IPv4 range columns are gone. Prefixes
from the feed are stored in one text form, so an IPv6 withdrawal, which the
feed sends uncompressed, now removes its route (closes #48)
- 2026-10-02: a plain `docker build .` stamps the commit's tag or short
commit (`git describe --tags --always`) into the page footer instead of
`unknown`: `.dockerignore` sends `.git` without `.git/config`, a `VERSION`
build argument takes precedence, and `make build` stamps the same value
(closes #46)
- 2026-09-29: the entrypoint creates the data directory if it is missing and
stops the start if a step fails; README "Running under upaas" no longer
asks for the host directory to be created first (closes #42)
- 2026-09-29: `.dockerignore` keeps `.git`, local build output, local
databases and `.env` out of the Docker build context, and so out of the
source archive in the image (closes #39)
- 2026-09-29: README first line names the MIT license and the author;
the License section now says MIT and links `LICENSE`, and an Author
section was added; this file brought up to date (closes #38)
- 2026-09-29: MIT `LICENSE` (closes #1)
- 2026-09-28: stopping the daemon while the feed is flowing no longer
panics with "send on closed channel": the read loop checks for a stop
just before handing a message to the handler queues, and a second
`Stop` no longer closes the queues again (closes #34)
- 2026-09-28: `docker stop` no longer kills the daemon 2 seconds after the
stop signal: the entrypoint switches to the `routewatch` user with
`setpriv` instead of `runuser`, so the daemon receives the signal itself
and gets the whole wait `docker stop` allows, up to its own 60-second
limit (closes #33)
- 2026-09-28: ready to run under upaas: a set but invalid `PORT`,
`XDG_DATA_HOME` or `MALLOC_ARENA_MAX` stops the start, the health
check follows `PORT`, README "Running under upaas" section (closes
#31)
- 2026-09-22: realtime in-memory database statistics: counts seeded at
startup and adjusted on every write, oldest/newest route timestamps via
index-end lookups; `/api/v1/stats` no longer scans the tables (closes
#27)
- 2026-09-21: batch writes take the write lock when their transaction
begins (`_txlock=immediate`), so they wait out a WAL checkpoint instead
of failing with "database is locked" (closes #25)
- 2026-09-21: `MALLOC_ARENA_MAX=2` in the image caps glibc malloc arenas,
so memory outside the Go runtime no longer grows with the core count
(closes #23)
- 2026-09-21: `GOMEMLIMIT=1536MiB` in the image; README Memory section
with the memory budget and the 5 GiB container limit (closes #13)
- 2026-09-21: the four handler queues hold at most 20,000 messages each,
down from 100,000 (closes #11)
- 2026-09-21: two goroutine leaks fixed: the stats handlers after a
timeout and the streamer's tickers on every reconnect (closes #12)
- 2026-09-21: parsed RIS messages no longer keep the unused `Community`
and `Raw` fields (closes #9)
- 2026-09-21: `.editorconfig`, and a Gitea workflow that runs
`script/cibuild` on every push (closes #14)
- 2026-09-21: the peering handler's AS-path map holds at most 500,000
paths and is swapped for an empty one every 30 seconds instead of copied
(closes #10)
- 2026-09-21: SQLite memory bounded across the whole connection pool: a
64 MiB page cache on each connection, 1 GiB soft and 1.5 GiB hard heap
limits (closes #8)
- 2026-09-21: the Docker build runs the format check, the linter and the
tests, linting in a separate stage on a golangci-lint image pinned by
digest (closes #5)
- 2026-09-21: `make test` skips the live-network feed test (`-short`), so
`make check` no longer depends on the network (closes #2)
- 2026-07-07 Adopted scripts-to-rule-them-all: `script/` entrypoints, - 2026-07-07 Adopted scripts-to-rule-them-all: `script/` entrypoints,
Makefile shims, README Entrypoints section Makefile shims, README Entrypoints section
- 2026-02-22: repo policy compliance: required policy files, .gitignore - 2026-02-22: repo policy compliance: required policy files, .gitignore
@@ -41,8 +115,9 @@ runs make check on main.
# Future Steps # Future Steps
- Verify main is green after the merge: make check locally and the new - Production memory under 5 GiB: whether to test under a real 5 GiB
CI workflow passing container limit on fsn1app1 is open for sneak
- Review stale remote branches fix-min-time-calculation and (https://git.eeqj.de/sneak/routewatch/issues/3)
optimize-sqlite-settings: land or delete - `/api/v1/stats` answered HTTP 500 after 35 hours on the live feed, seen
- Clean up the tmp/ directory at the repo root: gitignore or remove on `3898daa`, which predates the 2026-09-22 in-memory statistics
(https://git.eeqj.de/sneak/routewatch/issues/30)
+14 -1
View File
@@ -1,7 +1,20 @@
#!/bin/bash #!/bin/bash
set -euo pipefail
# glibc silently ignores a malformed MALLOC_ARENA_MAX, so refuse it here.
if [[ -n "${MALLOC_ARENA_MAX:-}" && ! "$MALLOC_ARENA_MAX" =~ ^[1-9][0-9]*$ ]]; then
echo "MALLOC_ARENA_MAX must be a positive whole number, got '$MALLOC_ARENA_MAX'" >&2
exit 1
fi
# Give the data directory to the routewatch user before the daemon starts,
# whether it is missing, an empty root-owned mount, or holds another uid's files.
mkdir -p /var/lib/berlin.sneak.app.routewatch
cd /var/lib/berlin.sneak.app.routewatch cd /var/lib/berlin.sneak.app.routewatch
chown -R routewatch:routewatch . chown -R routewatch:routewatch .
chmod 700 . chmod 700 .
exec runuser -u routewatch -- /app/routewatch # setpriv replaces itself with the daemon, so the daemon receives the stop
# signal directly. runuser would stay in between and kill the daemon 2 seconds
# after passing the signal on.
exec setpriv --reuid=routewatch --regid=routewatch --init-groups -- /app/routewatch
+40 -1
View File
@@ -6,6 +6,7 @@ import (
"os" "os"
"path/filepath" "path/filepath"
"runtime" "runtime"
"strconv"
"time" "time"
) )
@@ -18,6 +19,12 @@ const (
// defaultRouteExpirationMinutes is the default route expiration timeout in minutes // defaultRouteExpirationMinutes is the default route expiration timeout in minutes
defaultRouteExpirationMinutes = 5 defaultRouteExpirationMinutes = 5
// defaultPort is the HTTP port used when PORT is not set
defaultPort = 8080
// maxPort is the highest TCP port number
maxPort = 65535
) )
// Config holds configuration for the entire application // Config holds configuration for the entire application
@@ -25,6 +32,9 @@ type Config struct {
// StateDir is the directory for all application state (database, snapshots) // StateDir is the directory for all application state (database, snapshots)
StateDir string StateDir string
// Port is the TCP port the HTTP server listens on
Port int
// MaxRuntime is the maximum runtime (0 = run forever) // MaxRuntime is the maximum runtime (0 = run forever)
MaxRuntime time.Duration MaxRuntime time.Duration
@@ -43,8 +53,14 @@ func New() (*Config, error) {
return nil, fmt.Errorf("failed to determine state directory: %w", err) return nil, fmt.Errorf("failed to determine state directory: %w", err)
} }
port, err := getPort()
if err != nil {
return nil, err
}
return &Config{ return &Config{
StateDir: stateDir, StateDir: stateDir,
Port: port,
MaxRuntime: 0, // Run forever by default MaxRuntime: 0, // Run forever by default
EnableBatchedDatabaseWrites: true, // Enable batching by default EnableBatchedDatabaseWrites: true, // Enable batching by default
RouteExpirationTimeout: defaultRouteExpirationMinutes * time.Minute, // For active route monitoring RouteExpirationTimeout: defaultRouteExpirationMinutes * time.Minute, // For active route monitoring
@@ -69,13 +85,20 @@ func getStateDirectory() (string, error) {
return filepath.Join(home, "Library", "Application Support", AppIdentifier), nil return filepath.Join(home, "Library", "Application Support", AppIdentifier), nil
case "linux", "freebsd", "openbsd", "netbsd": case "linux", "freebsd", "openbsd", "netbsd":
// The XDG spec requires an absolute path; a relative one would put
// the database somewhere unexpected.
xdgData := os.Getenv("XDG_DATA_HOME")
if xdgData != "" && !filepath.IsAbs(xdgData) {
return "", fmt.Errorf("XDG_DATA_HOME must be an absolute path, got %q", xdgData)
}
// Unix-like: /var/lib/berlin.sneak.app.routewatch if root, else XDG_DATA_HOME // Unix-like: /var/lib/berlin.sneak.app.routewatch if root, else XDG_DATA_HOME
if os.Geteuid() == 0 { if os.Geteuid() == 0 {
return filepath.Join("/var/lib", AppIdentifier), nil return filepath.Join("/var/lib", AppIdentifier), nil
} }
// Check XDG_DATA_HOME first // Check XDG_DATA_HOME first
if xdgData := os.Getenv("XDG_DATA_HOME"); xdgData != "" { if xdgData != "" {
return filepath.Join(xdgData, AppIdentifier), nil return filepath.Join(xdgData, AppIdentifier), nil
} }
@@ -92,6 +115,22 @@ func getStateDirectory() (string, error) {
} }
} }
// getPort returns the HTTP port from PORT, or defaultPort when PORT is not set
func getPort() (int, error) {
value := os.Getenv("PORT")
if value == "" {
return defaultPort, nil
}
// ParseUint, unlike Atoi, refuses a sign: the health check URL cannot use "+9090"
port, err := strconv.ParseUint(value, 10, 0)
if err != nil || port < 1 || port > maxPort {
return 0, fmt.Errorf("PORT must be a whole number from 1 to %d, got %q", maxPort, value)
}
return int(port), nil
}
// EnsureDirectories creates all necessary directories if they don't exist // EnsureDirectories creates all necessary directories if they don't exist
func (c *Config) EnsureDirectories() error { func (c *Config) EnsureDirectories() error {
// Ensure state directory exists // Ensure state directory exists
+62
View File
@@ -0,0 +1,62 @@
package config
import (
"runtime"
"testing"
)
func TestNewReadsPort(t *testing.T) {
tests := map[string]int{
"": defaultPort,
"1": 1,
"9090": 9090,
"65535": 65535,
}
for value, want := range tests {
t.Run(value, func(t *testing.T) {
t.Setenv("PORT", value)
t.Setenv("XDG_DATA_HOME", "")
cfg, err := New()
if err != nil {
t.Fatalf("New() with PORT=%q: %v", value, err)
}
if cfg.Port != want {
t.Errorf("New() with PORT=%q: Port = %d, want %d", value, cfg.Port, want)
}
})
}
}
func TestNewRefusesInvalidPort(t *testing.T) {
for _, value := range []string{"0", "65536", "-1", "+9090", "http", "80.5"} {
t.Run(value, func(t *testing.T) {
t.Setenv("PORT", value)
t.Setenv("XDG_DATA_HOME", "")
if _, err := New(); err == nil {
t.Errorf("New() with PORT=%q returned no error", value)
}
})
}
}
func TestNewRefusesRelativeXDGDataHome(t *testing.T) {
if runtime.GOOS == "darwin" {
t.Skip("macOS does not read XDG_DATA_HOME")
}
t.Setenv("PORT", "")
t.Setenv("XDG_DATA_HOME", "relative/path")
if _, err := New(); err == nil {
t.Error("New() with a relative XDG_DATA_HOME returned no error")
}
t.Setenv("XDG_DATA_HOME", "/var/lib")
if _, err := New(); err != nil {
t.Errorf("New() with XDG_DATA_HOME=/var/lib: %v", err)
}
}
+146
View File
@@ -0,0 +1,146 @@
package database
import (
"context"
"fmt"
"sync"
)
// liveCounts holds the running row counts that the stats endpoints report. They
// are seeded once at startup from the tables and then adjusted on every write,
// so a stats read serves them from memory instead of running a COUNT(*) over
// each table. Those scans, once the database passed a few GiB, took the whole
// request timeout and made /api/v1/stats return 500 (issue 27).
//
// A single mutex guards all fields so the stats reader takes a consistent
// snapshot at one instant and writers, which already run under the database
// write lock, adjust the counts after their transaction commits.
type liveCounts struct {
mu sync.RWMutex
asns int
prefixesV4 int
prefixesV6 int
peerings int
peers int
routesV4 int
routesV6 int
}
// seed sets every count to the value read from the tables at startup. It runs
// before any writer, so it needs no coordination with the adjust methods.
func (c *liveCounts) seed(asns, prefixesV4, prefixesV6, peerings, peers, routesV4, routesV6 int) {
c.mu.Lock()
defer c.mu.Unlock()
c.asns = asns
c.prefixesV4 = prefixesV4
c.prefixesV6 = prefixesV6
c.peerings = peerings
c.peers = peers
c.routesV4 = routesV4
c.routesV6 = routesV6
}
// addASNs adds n to the ASN count.
func (c *liveCounts) addASNs(n int) {
c.mu.Lock()
c.asns += n
c.mu.Unlock()
}
// addPrefixes adds to the IPv4 and IPv6 prefix counts.
func (c *liveCounts) addPrefixes(v4, v6 int) {
c.mu.Lock()
c.prefixesV4 += v4
c.prefixesV6 += v6
c.mu.Unlock()
}
// addPeerings adds n to the peering count.
func (c *liveCounts) addPeerings(n int) {
c.mu.Lock()
c.peerings += n
c.mu.Unlock()
}
// addPeers adds n to the BGP peer count.
func (c *liveCounts) addPeers(n int) {
c.mu.Lock()
c.peers += n
c.mu.Unlock()
}
// addRoutes adds to the IPv4 and IPv6 live-route counts. Deletions pass
// negative values.
func (c *liveCounts) addRoutes(v4, v6 int) {
c.mu.Lock()
c.routesV4 += v4
c.routesV6 += v6
c.mu.Unlock()
}
// fill copies the counts into a Stats, including the derived totals, under a
// single read lock so the reader sees one consistent snapshot.
func (c *liveCounts) fill(s *Stats) {
c.mu.RLock()
defer c.mu.RUnlock()
s.ASNs = c.asns
s.IPv4Prefixes = c.prefixesV4
s.IPv6Prefixes = c.prefixesV6
s.Prefixes = c.prefixesV4 + c.prefixesV6
s.Peerings = c.peerings
s.Peers = c.peers
s.IPv4Routes = c.routesV4
s.IPv6Routes = c.routesV6
s.LiveRoutes = c.routesV4 + c.routesV6
}
// countRows returns the number of rows in the named table. It is used only at
// startup to seed the in-memory counters, so a full COUNT(*) is acceptable.
func (d *Database) countRows(ctx context.Context, table string) (int, error) {
var n int
// table is one of a fixed set of literals below, never external input.
if err := d.db.QueryRowContext(ctx, "SELECT COUNT(*) FROM "+table).Scan(&n); err != nil {
return 0, fmt.Errorf("failed to count %s: %w", table, err)
}
return n, nil
}
// seedCounts reads the current row counts from the tables into the in-memory
// counters. It runs once at startup, before the streamer begins writing.
func (d *Database) seedCounts(ctx context.Context) error {
asns, err := d.countRows(ctx, "asns")
if err != nil {
return err
}
prefixesV4, err := d.countRows(ctx, "prefixes_v4")
if err != nil {
return err
}
prefixesV6, err := d.countRows(ctx, "prefixes_v6")
if err != nil {
return err
}
peerings, err := d.countRows(ctx, "peerings")
if err != nil {
return err
}
peers, err := d.countRows(ctx, "bgp_peers")
if err != nil {
return err
}
routesV4, err := d.countRows(ctx, "live_routes_v4")
if err != nil {
return err
}
routesV6, err := d.countRows(ctx, "live_routes_v6")
if err != nil {
return err
}
d.counts.seed(asns, prefixesV4, prefixesV6, peerings, peers, routesV4, routesV6)
return nil
}
+328
View File
@@ -0,0 +1,328 @@
package database
import (
"context"
"sync"
"testing"
"time"
"git.eeqj.de/sneak/routewatch/internal/config"
"git.eeqj.de/sneak/routewatch/internal/logger"
"github.com/google/uuid"
)
// mkV4Route builds an IPv4 live route.
func mkV4Route(t *testing.T, prefix string, asn int, ts time.Time) *LiveRoute {
t.Helper()
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,
}
}
// mkV6Route builds an IPv6 live route.
func mkV6Route(prefix string, asn int, ts time.Time) *LiveRoute {
return &LiveRoute{
ID: uuid.New(),
Prefix: prefix,
MaskLength: 32,
IPVersion: ipVersionV6,
OriginASN: asn,
PeerIP: "2001:db8::1",
ASPath: []int{asn},
NextHop: "2001:db8::ffff",
LastUpdated: ts,
}
}
// TestLiveCountsTrackWritesInRealtime checks that the stats counts start at
// zero, reflect each write the moment it commits (no recompute, no timer), do
// not move when a route is merely re-announced, and drop when a route is
// deleted. These counts are what /api/v1/stats reports; before this change the
// endpoint recomputed them with a COUNT(*) over each table on every request.
func TestLiveCountsTrackWritesInRealtime(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, err := db.GetStatsContext(ctx)
if err != nil {
t.Fatalf("GetStatsContext on empty database: %v", err)
}
if empty.ASNs != 0 || empty.Prefixes != 0 || empty.Peerings != 0 ||
empty.Peers != 0 || empty.LiveRoutes != 0 {
t.Fatalf("empty database counts nonzero: %+v", empty)
}
ts := time.Date(2026, 1, 2, 3, 4, 5, 0, time.UTC)
if err := db.GetOrCreateASNBatch(map[int]time.Time{64500: ts, 64501: ts}); err != nil {
t.Fatalf("GetOrCreateASNBatch: %v", err)
}
if err := db.UpdatePrefixesBatch(map[string]time.Time{
"198.51.100.0/24": ts,
"2001:db8::/32": ts,
}); err != nil {
t.Fatalf("UpdatePrefixesBatch: %v", err)
}
if err := db.UpdatePeerBatch(map[string]PeerUpdate{
"192.0.2.1": {PeerIP: "192.0.2.1", PeerASN: 64500, MessageType: "UPDATE", Timestamp: ts},
}); err != nil {
t.Fatalf("UpdatePeerBatch: %v", err)
}
if err := db.RecordPeering(64500, 64501, ts); err != nil {
t.Fatalf("RecordPeering: %v", err)
}
routes := []*LiveRoute{
mkV4Route(t, "198.51.100.0/24", 64500, ts),
mkV4Route(t, "203.0.113.0/24", 64501, ts.Add(time.Minute)),
mkV6Route("2001:db8::/32", 64502, ts.Add(2*time.Minute)),
}
if err := db.UpsertLiveRouteBatch(routes); err != nil {
t.Fatalf("UpsertLiveRouteBatch: %v", err)
}
stats, err := db.GetStatsContext(ctx)
if err != nil {
t.Fatalf("GetStatsContext: %v", err)
}
assertCounts(t, "after inserts", stats, wantCounts{
asns: 2, prefixes: 2, peerings: 1, peers: 1,
ipv4Routes: 2, ipv6Routes: 1, liveRoutes: 3,
})
// Re-announcing the same routes is an update, not an insert: counts hold.
if err := db.UpsertLiveRouteBatch(routes); err != nil {
t.Fatalf("UpsertLiveRouteBatch (re-announce): %v", err)
}
stats, err = db.GetStatsContext(ctx)
if err != nil {
t.Fatalf("GetStatsContext: %v", err)
}
assertCounts(t, "after re-announce", stats, wantCounts{
asns: 2, prefixes: 2, peerings: 1, peers: 1,
ipv4Routes: 2, ipv6Routes: 1, liveRoutes: 3,
})
// A withdrawal removes one route.
if err := db.DeleteLiveRouteBatch([]LiveRouteDeletion{
{Prefix: "203.0.113.0/24", OriginASN: 64501, PeerIP: "192.0.2.1", IPVersion: ipVersionV4},
}); err != nil {
t.Fatalf("DeleteLiveRouteBatch: %v", err)
}
stats, err = db.GetStatsContext(ctx)
if err != nil {
t.Fatalf("GetStatsContext: %v", err)
}
assertCounts(t, "after delete", stats, wantCounts{
asns: 2, prefixes: 2, peerings: 1, peers: 1,
ipv4Routes: 1, ipv6Routes: 1, liveRoutes: 2,
})
}
// TestLiveCountsSeededFromDatabaseAtStartup writes rows, reopens the same
// database file, and checks the counts come back from the seed scan rather than
// starting at zero.
func TestLiveCountsSeededFromDatabaseAtStartup(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)
}
ts := time.Date(2026, 1, 2, 3, 4, 5, 0, time.UTC)
if err := db.GetOrCreateASNBatch(map[int]time.Time{64500: ts, 64501: ts, 64502: ts}); err != nil {
t.Fatalf("GetOrCreateASNBatch: %v", err)
}
if err := db.UpsertLiveRouteBatch([]*LiveRoute{
mkV4Route(t, "198.51.100.0/24", 64500, ts),
mkV6Route("2001:db8::/32", 64502, ts),
}); err != nil {
t.Fatalf("UpsertLiveRouteBatch: %v", err)
}
if err := db.Close(); err != nil {
t.Fatalf("Close: %v", err)
}
reopened, err := New(cfg, logger.New())
if err != nil {
t.Fatalf("failed to reopen database: %v", err)
}
defer func() { _ = reopened.Close() }()
stats, err := reopened.GetStatsContext(context.Background())
if err != nil {
t.Fatalf("GetStatsContext after reopen: %v", err)
}
if stats.ASNs != 3 {
t.Errorf("seeded ASNs = %d, want 3", stats.ASNs)
}
if stats.IPv4Routes != 1 || stats.IPv6Routes != 1 || stats.LiveRoutes != 2 {
t.Errorf("seeded routes = (v4 %d, v6 %d, total %d), want (1, 1, 2)",
stats.IPv4Routes, stats.IPv6Routes, stats.LiveRoutes)
}
}
// TestStatsRouteTimestamps checks the oldest/newest route timestamps are read
// from the right rows across both tables and parse into time.Time. The old
// MIN/MAX union query read its result into *time.Time, which the driver could
// not parse, so it logged a warning every call and left both timestamps nil.
func TestStatsRouteTimestamps(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, 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)
}
base := time.Date(2026, 1, 2, 3, 4, 5, 0, time.UTC)
oldest := base
newest := base.Add(2 * time.Minute)
if err := db.UpsertLiveRouteBatch([]*LiveRoute{
mkV4Route(t, "198.51.100.0/24", 64500, base.Add(time.Minute)),
mkV4Route(t, "203.0.113.0/24", 64501, oldest),
mkV6Route("2001:db8::/32", 64502, newest),
}); err != nil {
t.Fatalf("UpsertLiveRouteBatch: %v", err)
}
stats, err := db.GetStatsContext(ctx)
if err != nil {
t.Fatalf("GetStatsContext: %v", err)
}
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)
}
}
// TestLiveCountsConcurrentReadWrite runs writers and stats readers at once so
// the race detector proves the counters are safe under concurrent use.
func TestLiveCountsConcurrentReadWrite(t *testing.T) {
cfg := &config.Config{StateDir: t.TempDir()}
db, err := New(cfg, logger.New())
if err != nil {
t.Fatalf("failed to create database: %v", err)
}
defer func() { _ = db.Close() }()
ts := time.Date(2026, 1, 2, 3, 4, 5, 0, time.UTC)
const writers = 4
var wg sync.WaitGroup
wg.Add(writers)
for w := range writers {
go func(base int) {
defer wg.Done()
for i := range 25 {
asn := 65000 + base*100 + i
route := mkV6Route("2001:db8::/32", asn, ts)
route.PeerIP = "2001:db8::" + uuid.NewString()
if err := db.UpsertLiveRoute(route); err != nil {
t.Errorf("UpsertLiveRoute: %v", err)
return
}
}
}(w)
}
var readerWG sync.WaitGroup
readerWG.Add(1)
stop := make(chan struct{})
go func() {
defer readerWG.Done()
for {
select {
case <-stop:
return
default:
if _, err := db.GetStatsContext(context.Background()); err != nil {
t.Errorf("GetStatsContext: %v", err)
return
}
}
}
}()
wg.Wait()
close(stop)
readerWG.Wait()
stats, err := db.GetStatsContext(context.Background())
if err != nil {
t.Fatalf("GetStatsContext: %v", err)
}
if want := writers * 25; stats.IPv6Routes != want {
t.Errorf("IPv6Routes = %d, want %d", stats.IPv6Routes, want)
}
}
type wantCounts struct {
asns int
prefixes int
peerings int
peers int
ipv4Routes int
ipv6Routes int
liveRoutes int
}
func assertCounts(t *testing.T, when string, got Stats, want wantCounts) {
t.Helper()
if got.ASNs != want.asns {
t.Errorf("%s: ASNs = %d, want %d", when, got.ASNs, want.asns)
}
if got.Prefixes != want.prefixes {
t.Errorf("%s: Prefixes = %d, want %d", when, got.Prefixes, want.prefixes)
}
if got.Peerings != want.peerings {
t.Errorf("%s: Peerings = %d, want %d", when, got.Peerings, want.peerings)
}
if got.Peers != want.peers {
t.Errorf("%s: Peers = %d, want %d", when, got.Peers, want.peers)
}
if got.IPv4Routes != want.ipv4Routes {
t.Errorf("%s: IPv4Routes = %d, want %d", when, got.IPv4Routes, want.ipv4Routes)
}
if got.IPv6Routes != want.ipv6Routes {
t.Errorf("%s: IPv6Routes = %d, want %d", when, got.IPv6Routes, want.ipv6Routes)
}
if got.LiveRoutes != want.liveRoutes {
t.Errorf("%s: LiveRoutes = %d, want %d", when, got.LiveRoutes, want.liveRoutes)
}
}
File diff suppressed because it is too large Load Diff
+156 -281
View File
@@ -3,23 +3,36 @@ package database
import ( import (
"context" "context"
"database/sql" "database/sql"
"net" "errors"
"fmt"
"net/netip"
"sync" "sync"
"testing" "testing"
"time" "time"
"git.eeqj.de/sneak/routewatch/internal/config" "git.eeqj.de/sneak/routewatch/internal/config"
"git.eeqj.de/sneak/routewatch/internal/logger" "git.eeqj.de/sneak/routewatch/internal/logger"
"github.com/google/uuid"
) )
// tempStoreMemory is the PRAGMA temp_store value meaning "hold temp B-trees in // 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. // memory"; the DSN change must leave temp_store below this so they spill to disk.
const tempStoreMemory = 2 const tempStoreMemory = 2
// autoVacuumIncremental is the PRAGMA auto_vacuum value meaning "incremental".
const autoVacuumIncremental = 2
// vacuumTestRoutes is how many routes the vacuum test writes and then deletes,
// enough to leave many free pages in the file.
const vacuumTestRoutes = 2000
// heldConnections is how many pooled connections the pragma test holds open at // heldConnections is how many pooled connections the pragma test holds open at
// once so each is a distinct SQLite connection that parsed the DSN. // once so each is a distinct SQLite connection that parsed the DSN.
const heldConnections = 5 const heldConnections = 5
// testPeerIP is the peer every route in the IP lookup test is learned from.
const testPeerIP = "192.0.2.254"
// Parameters for the checkpoint-contention regression test. // Parameters for the checkpoint-contention regression test.
const ( const (
// contentionIterations is how many batch writes race the checkpoint loop. // contentionIterations is how many batch writes race the checkpoint loop.
@@ -32,286 +45,98 @@ const (
asnSecondBand = 100 asnSecondBand = 100
) )
func TestIPToUint32(t *testing.T) { // TestGetIPInfoFindsMostSpecificLiveRoute stores nested live prefixes for both
tests := []struct { // families and checks that a lookup returns the most specific one covering the
name string // address, ErrNoRoute when none covers it, and the next less specific prefix
ip string // once the only route of the most specific one is withdrawn.
expected uint32 func TestGetIPInfoFindsMostSpecificLiveRoute(t *testing.T) {
}{ cfg := &config.Config{StateDir: t.TempDir()}
{
name: "Simple IP", db, err := New(cfg, logger.New())
ip: "192.168.1.1", if err != nil {
expected: 3232235777, // 192<<24 + 168<<16 + 1<<8 + 1 t.Fatalf("failed to create database: %v", err)
},
{
name: "Minimum IP",
ip: "0.0.0.0",
expected: 0,
},
{
name: "Maximum IP",
ip: "255.255.255.255",
expected: 4294967295,
},
{
name: "10.0.0.0",
ip: "10.0.0.0",
expected: 167772160,
},
{
name: "172.16.0.0",
ip: "172.16.0.0",
expected: 2886729728,
},
{
name: "8.8.8.8",
ip: "8.8.8.8",
expected: 134744072,
},
{
name: "1.2.3.4",
ip: "1.2.3.4",
expected: 16909060,
},
} }
defer func() { _ = db.Close() }()
for _, tt := range tests { // Nested live prefixes, each originated by its own AS.
t.Run(tt.name, func(t *testing.T) { origins := map[string]int{
ip := net.ParseIP(tt.ip) "10.0.0.0/8": 64500,
if ip == nil { "10.1.0.0/16": 64501,
t.Fatalf("Failed to parse IP: %s", tt.ip) "10.1.2.0/24": 64502,
} "2001:db8::/32": 64500,
"2001:db8:1::/48": 64501,
result := ipToUint32(ip) "2001:db8:1:2::/64": 64502,
if result != tt.expected { }
t.Errorf("ipToUint32(%s) = %d, want %d", tt.ip, result, tt.expected) ts := time.Date(2026, 1, 2, 3, 4, 5, 0, time.UTC)
} routes := make([]*LiveRoute, 0, len(origins))
for prefix, asn := range origins {
// Test with IPv4-mapped IPv6 address routes = append(routes, &LiveRoute{
ip6 := net.ParseIP(tt.ip).To16() ID: uuid.New(),
if ip6 != nil { Prefix: prefix,
result6 := ipToUint32(ip6) MaskLength: netip.MustParsePrefix(prefix).Bits(),
if result6 != tt.expected { IPVersion: detectIPVersion(prefix),
t.Errorf("ipToUint32(%s as IPv6) = %d, want %d", tt.ip, result6, tt.expected) OriginASN: asn,
} PeerIP: testPeerIP,
} ASPath: []int{asn},
NextHop: testPeerIP,
LastUpdated: ts,
}) })
} }
} if err := db.UpsertLiveRouteBatch(routes); err != nil {
t.Fatalf("UpsertLiveRouteBatch: %v", err)
func TestCalculateIPv4Range(t *testing.T) {
tests := []struct {
name string
cidr string
wantStart uint32
wantEnd uint32
wantErr bool
}{
{
name: "Single IP /32",
cidr: "192.168.1.1/32",
wantStart: 3232235777,
wantEnd: 3232235777,
},
{
name: "Class C /24",
cidr: "192.168.1.0/24",
wantStart: 3232235776, // 192.168.1.0
wantEnd: 3232236031, // 192.168.1.255
},
{
name: "Class B /16",
cidr: "192.168.0.0/16",
wantStart: 3232235520, // 192.168.0.0
wantEnd: 3232301055, // 192.168.255.255
},
{
name: "Class A /8",
cidr: "10.0.0.0/8",
wantStart: 167772160, // 10.0.0.0
wantEnd: 184549375, // 10.255.255.255
},
{
name: "Entire IPv4 space /0",
cidr: "0.0.0.0/0",
wantStart: 0,
wantEnd: 4294967295,
},
{
name: "Small subnet /30",
cidr: "192.168.1.0/30",
wantStart: 3232235776, // 192.168.1.0
wantEnd: 3232235779, // 192.168.1.3
},
{
name: "Medium subnet /20",
cidr: "172.16.0.0/20",
wantStart: 2886729728, // 172.16.0.0
wantEnd: 2886733823, // 172.16.15.255
},
{
name: "Private range 172.16/12",
cidr: "172.16.0.0/12",
wantStart: 2886729728, // 172.16.0.0
wantEnd: 2887778303, // 172.31.255.255
},
{
name: "Google DNS /29",
cidr: "8.8.8.8/29",
wantStart: 134744072, // 8.8.8.8 (network is actually 8.8.8.8 with /29)
wantEnd: 134744079, // 8.8.8.15
},
{
name: "Non-zero host bits",
cidr: "192.168.1.5/24",
wantStart: 3232235776, // 192.168.1.0 (network address)
wantEnd: 3232236031, // 192.168.1.255
},
{
name: "Invalid CIDR",
cidr: "192.168.1.1/33",
wantErr: true,
},
{
name: "Invalid IP",
cidr: "256.256.256.256/24",
wantErr: true,
},
{
name: "IPv6 CIDR",
cidr: "2001:db8::/32",
wantErr: true,
},
{
name: "Empty CIDR",
cidr: "",
wantErr: true,
},
{
name: "Missing mask",
cidr: "192.168.1.1",
wantErr: true,
},
} }
for _, tt := range tests { // lookup checks that ip resolves to the live prefix want, or to ErrNoRoute
t.Run(tt.name, func(t *testing.T) { // when want is empty.
start, end, err := CalculateIPv4Range(tt.cidr) lookup := func(ip, want string) {
t.Helper()
if tt.wantErr { info, err := db.GetIPInfo(ip)
if err == nil { if want == "" {
t.Errorf("CalculateIPv4Range(%s) expected error, got nil", tt.cidr) if !errors.Is(err, ErrNoRoute) {
} t.Errorf("GetIPInfo(%s) = %+v, %v; want ErrNoRoute", ip, info, err)
return
} }
if err != nil { return
t.Errorf("CalculateIPv4Range(%s) unexpected error: %v", tt.cidr, err) }
return if err != nil {
} t.Errorf("GetIPInfo(%s): %v", ip, err)
if start != tt.wantStart { return
t.Errorf("CalculateIPv4Range(%s) start = %d, want %d", tt.cidr, start, tt.wantStart) }
} if info.Netblock != want || info.MaskLength != netip.MustParsePrefix(want).Bits() ||
info.ASN != origins[want] {
if end != tt.wantEnd { t.Errorf("GetIPInfo(%s) = %s (mask %d) AS%d, want %s AS%d",
t.Errorf("CalculateIPv4Range(%s) end = %d, want %d", tt.cidr, end, tt.wantEnd) ip, info.Netblock, info.MaskLength, info.ASN, want, origins[want])
} }
// Verify that start <= end
if start > end {
t.Errorf("CalculateIPv4Range(%s) start (%d) > end (%d)", tt.cidr, start, end)
}
// Verify the range size matches the CIDR mask
if !tt.wantErr && tt.cidr != "" {
_, ipNet, _ := net.ParseCIDR(tt.cidr)
if ipNet != nil {
ones, bits := ipNet.Mask.Size()
expectedSize := uint32(1) << uint(bits-ones)
actualSize := end - start + 1
if actualSize != expectedSize {
t.Errorf("CalculateIPv4Range(%s) range size = %d, want %d", tt.cidr, actualSize, expectedSize)
}
}
}
})
}
}
func TestIPv4RangeIntegration(t *testing.T) {
// Test that our functions work correctly together
tests := []struct {
name string
cidr string
testIPs []string
shouldContain []bool
}{
{
name: "192.168.1.0/24",
cidr: "192.168.1.0/24",
testIPs: []string{
"192.168.1.0",
"192.168.1.1",
"192.168.1.255",
"192.168.0.255",
"192.168.2.0",
},
shouldContain: []bool{true, true, true, false, false},
},
{
name: "10.0.0.0/8",
cidr: "10.0.0.0/8",
testIPs: []string{
"10.0.0.0",
"10.255.255.255",
"10.1.2.3",
"9.255.255.255",
"11.0.0.0",
},
shouldContain: []bool{true, true, true, false, false},
},
{
name: "172.16.0.0/12",
cidr: "172.16.0.0/12",
testIPs: []string{
"172.16.0.0",
"172.31.255.255",
"172.20.1.1",
"172.15.255.255",
"172.32.0.0",
},
shouldContain: []bool{true, true, true, false, false},
},
} }
for _, tt := range tests { lookup("10.1.2.3", "10.1.2.0/24")
t.Run(tt.name, func(t *testing.T) { lookup("::ffff:10.1.2.3", "10.1.2.0/24")
start, end, err := CalculateIPv4Range(tt.cidr) lookup("10.1.3.4", "10.1.0.0/16")
if err != nil { lookup("10.2.0.1", "10.0.0.0/8")
t.Fatalf("Failed to calculate range for %s: %v", tt.cidr, err) lookup("192.0.2.1", "")
} lookup("2001:db8:1:2::3", "2001:db8:1:2::/64")
lookup("2001:db8:1:3::4", "2001:db8:1::/48")
lookup("2001:db8:2::1", "2001:db8::/32")
lookup("2001:db9::1", "")
for i, testIP := range tt.testIPs { err = db.DeleteLiveRouteBatch([]LiveRouteDeletion{
ip := net.ParseIP(testIP) {Prefix: "10.1.2.0/24", OriginASN: 64502, PeerIP: testPeerIP, IPVersion: ipVersionV4},
if ip == nil { {Prefix: "2001:db8:1:2::/64", OriginASN: 64502, PeerIP: testPeerIP, IPVersion: ipVersionV6},
t.Fatalf("Failed to parse test IP: %s", testIP) })
} if err != nil {
t.Fatalf("DeleteLiveRouteBatch: %v", err)
ipUint := ipToUint32(ip)
contained := ipUint >= start && ipUint <= end
if contained != tt.shouldContain[i] {
t.Errorf("IP %s in range %s: got %v, want %v", testIP, tt.cidr, contained, tt.shouldContain[i])
}
}
})
} }
lookup("10.1.2.3", "10.1.0.0/16")
lookup("2001:db8:1:2::3", "2001:db8:1::/48")
} }
// TestConnectionPoolPragmas holds several pooled connections open at once and // TestConnectionPoolPragmas holds several pooled connections open at once and
// checks each one carries the per-connection settings from the DSN, plus the // checks each one carries the per-connection settings from the DSN, plus the
// process-wide hard heap limit. // process-wide hard heap limit, and sees the new file with auto_vacuum
// incremental.
func TestConnectionPoolPragmas(t *testing.T) { func TestConnectionPoolPragmas(t *testing.T) {
cfg := &config.Config{StateDir: t.TempDir()} cfg := &config.Config{StateDir: t.TempDir()}
@@ -372,6 +197,74 @@ func TestConnectionPoolPragmas(t *testing.T) {
if hardHeapLimit != sqliteHardHeapLimitBytes { if hardHeapLimit != sqliteHardHeapLimitBytes {
t.Errorf("conn %d: hard_heap_limit = %d, want %d", i, hardHeapLimit, sqliteHardHeapLimitBytes) t.Errorf("conn %d: hard_heap_limit = %d, want %d", i, hardHeapLimit, sqliteHardHeapLimitBytes)
} }
var autoVacuum int
if err := c.QueryRowContext(ctx, "PRAGMA auto_vacuum").Scan(&autoVacuum); err != nil {
t.Fatalf("conn %d: failed to read auto_vacuum: %v", i, err)
}
if autoVacuum != autoVacuumIncremental {
t.Errorf("conn %d: auto_vacuum = %d, want %d (incremental)", i, autoVacuum, autoVacuumIncremental)
}
}
}
// TestVacuumReturnsFreePages checks that after routes are deleted from a new
// database, one Vacuum call returns every page they used. The deletes leave
// fewer free pages than the 1000 Vacuum frees per call, so none may remain.
// With auto_vacuum off (issue https://git.eeqj.de/sneak/routewatch/issues/43)
// the free pages stayed in the file, and with the PRAGMA run by ExecContext
// Vacuum freed only one page per call.
func TestVacuumReturnsFreePages(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()
ts := time.Now().UTC()
const asn = 64500
routes := make([]*LiveRoute, 0, vacuumTestRoutes)
deletions := make([]LiveRouteDeletion, 0, vacuumTestRoutes)
for i := range vacuumTestRoutes {
route := mkV4Route(t, fmt.Sprintf("10.%d.%d.0/24", i/256, i%256), asn, ts)
routes = append(routes, route)
deletions = append(deletions, LiveRouteDeletion{
Prefix: route.Prefix,
OriginASN: asn,
PeerIP: route.PeerIP,
IPVersion: ipVersionV4,
})
}
if err := db.UpsertLiveRouteBatch(routes); err != nil {
t.Fatalf("failed to write routes: %v", err)
}
if err := db.DeleteLiveRouteBatch(deletions); err != nil {
t.Fatalf("failed to delete routes: %v", err)
}
var before int
if err := db.db.QueryRowContext(ctx, "PRAGMA freelist_count").Scan(&before); err != nil {
t.Fatalf("failed to read freelist_count: %v", err)
}
if before == 0 {
t.Fatalf("no free pages after deleting %d routes", vacuumTestRoutes)
}
if err := db.Vacuum(ctx); err != nil {
t.Fatalf("Vacuum failed: %v", err)
}
var after int
if err := db.db.QueryRowContext(ctx, "PRAGMA freelist_count").Scan(&after); err != nil {
t.Fatalf("failed to read freelist_count: %v", err)
}
if after != 0 {
t.Errorf("free pages after Vacuum = %d of %d, want 0", after, before)
} }
} }
@@ -424,21 +317,3 @@ func TestBatchWriteDuringCheckpoint(t *testing.T) {
cancel() cancel()
wg.Wait() wg.Wait()
} }
func BenchmarkIPToUint32(b *testing.B) {
ip := net.ParseIP("192.168.1.1")
b.ResetTimer()
for i := 0; i < b.N; i++ {
_ = ipToUint32(ip)
}
}
func BenchmarkCalculateIPv4Range(b *testing.B) {
cidr := "192.168.0.0/16"
b.ResetTimer()
for i := 0; i < b.N; i++ {
_, _, _ = CalculateIPv4Range(cidr)
}
}
+2 -2
View File
@@ -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
@@ -61,8 +63,6 @@ type Store interface {
GetLiveRouteCountsContext(ctx context.Context) (ipv4Count, ipv6Count int, err error) GetLiveRouteCountsContext(ctx context.Context) (ipv4Count, ipv6Count int, err error)
// IP lookup operations // IP lookup operations
GetASInfoForIP(ip string) (*ASInfo, error)
GetASInfoForIPContext(ctx context.Context, ip string) (*ASInfo, error)
GetIPInfo(ip string) (*IPInfo, error) GetIPInfo(ip string) (*IPInfo, error)
GetIPInfoContext(ctx context.Context, ip string) (*IPInfo, error) GetIPInfoContext(ctx context.Context, ip string) (*IPInfo, error)
-13
View File
@@ -77,9 +77,6 @@ type LiveRoute struct {
ASPath []int `json:"as_path"` ASPath []int `json:"as_path"`
NextHop string `json:"next_hop"` NextHop string `json:"next_hop"`
LastUpdated time.Time `json:"last_updated"` LastUpdated time.Time `json:"last_updated"`
// IPv4 range fields for fast lookups (nil for IPv6)
V4IPStart *uint32 `json:"v4_ip_start,omitempty"`
V4IPEnd *uint32 `json:"v4_ip_end,omitempty"`
} }
// PrefixDistribution represents the distribution of prefixes by mask length // PrefixDistribution represents the distribution of prefixes by mask length
@@ -88,16 +85,6 @@ type PrefixDistribution struct {
Count int `json:"count"` Count int `json:"count"`
} }
// ASInfo represents AS information for an IP lookup (legacy format)
type ASInfo struct {
ASN int `json:"asn"`
Handle string `json:"handle"`
Description string `json:"description"`
Prefix string `json:"prefix"`
LastUpdated time.Time `json:"last_updated"`
Age string `json:"age"`
}
// IPInfo represents comprehensive IP information for the /ip endpoint // IPInfo represents comprehensive IP information for the /ip endpoint
type IPInfo struct { type IPInfo struct {
IP string `json:"ip"` IP string `json:"ip"`
-6
View File
@@ -107,9 +107,6 @@ CREATE TABLE IF NOT EXISTS live_routes_v4 (
as_path TEXT NOT NULL, -- JSON array as_path TEXT NOT NULL, -- JSON array
next_hop TEXT NOT NULL, next_hop TEXT NOT NULL,
last_updated DATETIME NOT NULL, last_updated DATETIME NOT NULL,
-- IPv4 range columns for fast lookups
ip_start INTEGER NOT NULL, -- Start of IPv4 range as 32-bit unsigned int
ip_end INTEGER NOT NULL, -- End of IPv4 range as 32-bit unsigned int
UNIQUE(prefix, origin_asn, peer_ip) UNIQUE(prefix, origin_asn, peer_ip)
); );
@@ -123,7 +120,6 @@ CREATE TABLE IF NOT EXISTS live_routes_v6 (
as_path TEXT NOT NULL, -- JSON array as_path TEXT NOT NULL, -- JSON array
next_hop TEXT NOT NULL, next_hop TEXT NOT NULL,
last_updated DATETIME NOT NULL, last_updated DATETIME NOT NULL,
-- Note: IPv6 doesn't use integer range columns
UNIQUE(prefix, origin_asn, peer_ip) UNIQUE(prefix, origin_asn, peer_ip)
); );
@@ -132,8 +128,6 @@ CREATE INDEX IF NOT EXISTS idx_live_routes_v4_prefix ON live_routes_v4(prefix);
CREATE INDEX IF NOT EXISTS idx_live_routes_v4_mask_length ON live_routes_v4(mask_length); CREATE INDEX IF NOT EXISTS idx_live_routes_v4_mask_length ON live_routes_v4(mask_length);
CREATE INDEX IF NOT EXISTS idx_live_routes_v4_origin_asn ON live_routes_v4(origin_asn); CREATE INDEX IF NOT EXISTS idx_live_routes_v4_origin_asn ON live_routes_v4(origin_asn);
CREATE INDEX IF NOT EXISTS idx_live_routes_v4_last_updated ON live_routes_v4(last_updated); CREATE INDEX IF NOT EXISTS idx_live_routes_v4_last_updated ON live_routes_v4(last_updated);
-- Indexes for IPv4 range queries
CREATE INDEX IF NOT EXISTS idx_live_routes_v4_ip_range ON live_routes_v4(ip_start, ip_end);
-- Index to optimize prefix distribution queries -- Index to optimize prefix distribution queries
CREATE INDEX IF NOT EXISTS idx_live_routes_v4_mask_prefix ON live_routes_v4(mask_length, prefix); CREATE INDEX IF NOT EXISTS idx_live_routes_v4_mask_prefix ON live_routes_v4(mask_length, prefix);
+1 -20
View File
@@ -232,25 +232,6 @@ func (m *mockStore) GetLiveRouteCountsContext(ctx context.Context) (ipv4Count, i
return m.GetLiveRouteCounts() return m.GetLiveRouteCounts()
} }
// GetASInfoForIP mock implementation
func (m *mockStore) GetASInfoForIP(ip string) (*database.ASInfo, error) {
// Simple mock - return a test AS
now := time.Now()
return &database.ASInfo{
ASN: 15169,
Handle: "GOOGLE",
Description: "Google LLC",
Prefix: "8.8.8.0/24",
LastUpdated: now.Add(-5 * time.Minute),
Age: "5m0s",
}, nil
}
// GetASInfoForIPContext mock implementation with context support
func (m *mockStore) GetASInfoForIPContext(ctx context.Context, ip string) (*database.ASInfo, error) {
return m.GetASInfoForIP(ip)
}
// GetASDetails mock implementation // GetASDetails mock implementation
func (m *mockStore) GetASDetails(asn int) (*database.ASN, []database.LiveRoute, error) { func (m *mockStore) GetASDetails(asn int) (*database.ASN, []database.LiveRoute, error) {
m.mu.Lock() m.mu.Lock()
@@ -450,7 +431,7 @@ func TestRouteWatchLiveFeed(t *testing.T) {
} }
// Create server // Create server
srv := server.New(mockDB, s, logger) srv := server.New(mockDB, s, logger, cfg)
// Create RouteWatch with 5 second limit // Create RouteWatch with 5 second limit
deps := Dependencies{ deps := Dependencies{
+18 -44
View File
@@ -2,6 +2,7 @@ package routewatch
import ( import (
"net" "net"
"net/netip"
"strings" "strings"
"sync" "sync"
"time" "time"
@@ -108,7 +109,7 @@ func (h *PrefixHandler) HandleMessage(msg *ristypes.RISMessage) {
for _, announcement := range msg.Announcements { for _, announcement := range msg.Announcements {
for _, prefix := range announcement.Prefixes { for _, prefix := range announcement.Prefixes {
h.batch = append(h.batch, prefixUpdate{ h.batch = append(h.batch, prefixUpdate{
prefix: prefix, prefix: canonicalPrefix(prefix),
originASN: originASN, originASN: originASN,
peer: msg.Peer, peer: msg.Peer,
messageType: "announcement", messageType: "announcement",
@@ -125,7 +126,7 @@ func (h *PrefixHandler) HandleMessage(msg *ristypes.RISMessage) {
// Process withdrawals // Process withdrawals
for _, prefix := range msg.Withdrawals { for _, prefix := range msg.Withdrawals {
h.batch = append(h.batch, prefixUpdate{ h.batch = append(h.batch, prefixUpdate{
prefix: prefix, prefix: canonicalPrefix(prefix),
originASN: originASN, // Use the originASN from path if available originASN: originASN, // Use the originASN from path if available
peer: msg.Peer, peer: msg.Peer,
messageType: "withdrawal", messageType: "withdrawal",
@@ -264,6 +265,21 @@ func (h *PrefixHandler) flushBatchLocked() {
h.lastFlush = time.Now() h.lastFlush = time.Now()
} }
// canonicalPrefix returns prefix in the text form net/netip prints, the form
// the IP lookup builds when it looks a prefix up. The feed sends IPv6
// announcements compressed ("2001:db8::/32") but IPv6 withdrawals uncompressed
// ("2001:db8:0:0:0:0:0:0/32"); stored as received, a withdrawal would not match
// the route its announcement stored. A prefix that does not parse is returned
// unchanged, and the batch flush reports it.
func canonicalPrefix(prefix string) string {
p, err := netip.ParsePrefix(prefix)
if err != nil {
return prefix
}
return p.Masked().String()
}
// parseCIDR extracts the mask length and IP version from a prefix string // parseCIDR extracts the mask length and IP version from a prefix string
func parseCIDR(prefix string) (maskLength int, ipVersion int, err error) { func parseCIDR(prefix string) (maskLength int, ipVersion int, err error) {
_, ipNet, err := net.ParseCIDR(prefix) _, ipNet, err := net.ParseCIDR(prefix)
@@ -315,20 +331,6 @@ func (h *PrefixHandler) processAnnouncement(_ *database.Prefix, update prefixUpd
LastUpdated: update.timestamp, LastUpdated: update.timestamp,
} }
// For IPv4, calculate the IP range
if ipVersion == ipv4Version {
start, end, err := database.CalculateIPv4Range(update.prefix)
if err == nil {
liveRoute.V4IPStart = &start
liveRoute.V4IPEnd = &end
} else {
h.logger.Error("Failed to calculate IPv4 range",
"prefix", update.prefix,
"error", err,
)
}
}
if err := h.db.UpsertLiveRoute(liveRoute); err != nil { if err := h.db.UpsertLiveRoute(liveRoute); err != nil {
h.logger.Error("Failed to upsert live route", h.logger.Error("Failed to upsert live route",
"prefix", update.prefix, "prefix", update.prefix,
@@ -372,20 +374,6 @@ func (h *PrefixHandler) createLiveRoute(update prefixUpdate) *database.LiveRoute
LastUpdated: update.timestamp, LastUpdated: update.timestamp,
} }
// For IPv4, calculate the IP range
if ipVersion == ipv4Version {
start, end, err := database.CalculateIPv4Range(update.prefix)
if err == nil {
liveRoute.V4IPStart = &start
liveRoute.V4IPEnd = &end
} else {
h.logger.Error("Failed to calculate IPv4 range",
"prefix", update.prefix,
"error", err,
)
}
}
return liveRoute return liveRoute
} }
@@ -425,20 +413,6 @@ func (h *PrefixHandler) processAnnouncementDirect(update prefixUpdate) {
LastUpdated: update.timestamp, LastUpdated: update.timestamp,
} }
// For IPv4, calculate the IP range
if ipVersion == ipv4Version {
start, end, err := database.CalculateIPv4Range(update.prefix)
if err == nil {
liveRoute.V4IPStart = &start
liveRoute.V4IPEnd = &end
} else {
h.logger.Error("Failed to calculate IPv4 range",
"prefix", update.prefix,
"error", err,
)
}
}
if err := h.db.UpsertLiveRoute(liveRoute); err != nil { if err := h.db.UpsertLiveRoute(liveRoute); err != nil {
h.logger.Error("Failed to upsert live route", h.logger.Error("Failed to upsert live route",
"prefix", update.prefix, "prefix", update.prefix,
+74
View File
@@ -0,0 +1,74 @@
package routewatch
import (
"errors"
"testing"
"time"
"git.eeqj.de/sneak/routewatch/internal/config"
"git.eeqj.de/sneak/routewatch/internal/database"
"git.eeqj.de/sneak/routewatch/internal/logger"
"git.eeqj.de/sneak/routewatch/internal/ristypes"
)
const testPeerIP = "2001:db8:ffff::1"
// TestPrefixHandlerStoresPrefixesTheIPLookupFinds runs announcements and
// withdrawals through the prefix handler into a real database and looks the
// addresses up. The prefixes are written the way the feed sends them: IPv6
// announcements compressed, IPv6 withdrawals uncompressed. The withdrawal must
// remove the route the announcement stored, and the IP lookup must find the
// stored prefix.
func TestPrefixHandlerStoresPrefixesTheIPLookupFinds(t *testing.T) {
db, err := database.New(&config.Config{StateDir: t.TempDir()}, logger.New())
if err != nil {
t.Fatalf("failed to create database: %v", err)
}
defer func() { _ = db.Close() }()
// Built without NewPrefixHandler's flush timer, so the test flushes itself.
h := &PrefixHandler{db: db, logger: logger.New()}
flush := func() {
h.mu.Lock()
defer h.mu.Unlock()
h.flushBatchLocked()
}
ts := time.Date(2026, 1, 2, 3, 4, 5, 0, time.UTC)
h.HandleMessage(&ristypes.RISMessage{
Peer: testPeerIP,
Path: ristypes.ASPath{testASNA, testASNB},
ParsedTimestamp: ts,
Announcements: []ristypes.RISAnnouncement{{
NextHop: testPeerIP,
Prefixes: []string{"2001:db8:1::/48", "192.0.2.0/24"},
}},
})
flush()
for ip, want := range map[string]string{
"2001:db8:1::1": "2001:db8:1::/48",
"192.0.2.1": "192.0.2.0/24",
} {
info, err := db.GetIPInfo(ip)
if err != nil {
t.Fatalf("GetIPInfo(%s) after announcement: %v", ip, err)
}
if info.Netblock != want || info.ASN != testASNB {
t.Errorf("GetIPInfo(%s) = %s AS%d, want %s AS%d", ip, info.Netblock, info.ASN, want, testASNB)
}
}
h.HandleMessage(&ristypes.RISMessage{
Peer: testPeerIP,
ParsedTimestamp: ts.Add(time.Minute),
Withdrawals: []string{"2001:db8:1:0:0:0:0:0/48", "192.0.2.0/24"},
})
flush()
for _, ip := range []string{"2001:db8:1::1", "192.0.2.1"} {
if info, err := db.GetIPInfo(ip); !errors.Is(err, database.ErrNoRoute) {
t.Errorf("GetIPInfo(%s) after withdrawal = %+v, %v; want ErrNoRoute", ip, info, err)
}
}
}
+4 -18
View File
@@ -219,13 +219,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()
@@ -259,8 +252,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,
@@ -439,13 +432,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()
@@ -537,8 +523,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,
+2 -1
View File
@@ -8,6 +8,7 @@ import (
"testing" "testing"
"time" "time"
"git.eeqj.de/sneak/routewatch/internal/config"
"git.eeqj.de/sneak/routewatch/internal/database" "git.eeqj.de/sneak/routewatch/internal/database"
"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"
@@ -38,7 +39,7 @@ func (d blockingStatsDB) GetStatsContext(_ context.Context) (database.Stats, err
func TestStatsHandlersDoNotLeakOnTimeout(t *testing.T) { func TestStatsHandlersDoNotLeakOnTimeout(t *testing.T) {
release := make(chan struct{}) release := make(chan struct{})
db := blockingStatsDB{release: release} db := blockingStatsDB{release: release}
s := New(db, streamer.New(logger.New(), metrics.New()), logger.New()) s := New(db, streamer.New(logger.New(), metrics.New()), logger.New(), &config.Config{})
handlers := map[string]http.HandlerFunc{ handlers := map[string]http.HandlerFunc{
"status.json": s.handleStatusJSON(), "status.json": s.handleStatusJSON(),
+7 -9
View File
@@ -4,9 +4,10 @@ package server
import ( import (
"context" "context"
"net/http" "net/http"
"os" "strconv"
"time" "time"
"git.eeqj.de/sneak/routewatch/internal/config"
"git.eeqj.de/sneak/routewatch/internal/database" "git.eeqj.de/sneak/routewatch/internal/database"
"git.eeqj.de/sneak/routewatch/internal/logger" "git.eeqj.de/sneak/routewatch/internal/logger"
"git.eeqj.de/sneak/routewatch/internal/streamer" "git.eeqj.de/sneak/routewatch/internal/streamer"
@@ -33,16 +34,18 @@ type Server struct {
db database.Store db database.Store
streamer *streamer.Streamer streamer *streamer.Streamer
logger *logger.Logger logger *logger.Logger
port int
srv *http.Server srv *http.Server
asnFetcher ASNFetcher asnFetcher ASNFetcher
} }
// New creates a new HTTP server // New creates a new HTTP server
func New(db database.Store, streamer *streamer.Streamer, logger *logger.Logger) *Server { func New(db database.Store, streamer *streamer.Streamer, logger *logger.Logger, cfg *config.Config) *Server {
s := &Server{ s := &Server{
db: db, db: db,
streamer: streamer, streamer: streamer,
logger: logger, logger: logger,
port: cfg.Port,
} }
s.setupRoutes() s.setupRoutes()
@@ -52,11 +55,6 @@ func New(db database.Store, streamer *streamer.Streamer, logger *logger.Logger)
// Start starts the HTTP server // Start starts the HTTP server
func (s *Server) Start() error { func (s *Server) Start() error {
port := os.Getenv("PORT")
if port == "" {
port = "8080"
}
const ( const (
readHeaderTimeout = 40 * time.Second readHeaderTimeout = 40 * time.Second
readTimeout = 60 * time.Second readTimeout = 60 * time.Second
@@ -65,7 +63,7 @@ func (s *Server) Start() error {
) )
s.srv = &http.Server{ s.srv = &http.Server{
Addr: ":" + port, Addr: ":" + strconv.Itoa(s.port),
Handler: s.router, Handler: s.router,
ReadHeaderTimeout: readHeaderTimeout, ReadHeaderTimeout: readHeaderTimeout,
ReadTimeout: readTimeout, ReadTimeout: readTimeout,
@@ -73,7 +71,7 @@ func (s *Server) Start() error {
IdleTimeout: idleTimeout, IdleTimeout: idleTimeout,
} }
s.logger.Info("Starting HTTP server", "port", port, "addr", s.srv.Addr) s.logger.Info("Starting HTTP server", "port", s.port, "addr", s.srv.Addr)
// Start in goroutine but log when actually listening // Start in goroutine but log when actually listening
go func() { go func() {
+15 -3
View File
@@ -210,9 +210,14 @@ func (s *Streamer) Start() error {
// the connection status in metrics. This method is safe to call multiple times. // the connection status in metrics. This method is safe to call multiple times.
func (s *Streamer) Stop() { func (s *Streamer) Stop() {
s.mu.Lock() s.mu.Lock()
if s.cancel != nil { if s.cancel == nil {
s.cancel() // Not started, or already stopped: closing the queues again would panic.
s.mu.Unlock()
return
} }
s.cancel()
s.cancel = nil
// Close all handler queues to signal workers to stop // Close all handler queues to signal workers to stop
for _, info := range s.handlers { for _, info := range s.handlers {
close(info.queue) close(info.queue)
@@ -660,8 +665,15 @@ func (s *Streamer) stream(ctx context.Context) error {
continue continue
} }
// Dispatch to interested handlers // Dispatch to interested handlers. Stop cancels ctx and closes the
// queues under the write lock, so if ctx is cancelled here, under the
// read lock, the queues are closed and must not be sent to.
s.mu.RLock() s.mu.RLock()
if ctx.Err() != nil {
s.mu.RUnlock()
return ctx.Err()
}
for _, info := range s.handlers { for _, info := range s.handlers {
if !info.handler.WantsMessage(msg.Type) { if !info.handler.WantsMessage(msg.Type) {
continue continue
+42
View File
@@ -2,6 +2,8 @@ package streamer
import ( import (
"context" "context"
"errors"
"io"
"net/http" "net/http"
"net/http/httptest" "net/http/httptest"
"runtime" "runtime"
@@ -10,6 +12,7 @@ import (
"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"
"git.eeqj.de/sneak/routewatch/internal/ristypes"
) )
func TestNewStreamer(t *testing.T) { func TestNewStreamer(t *testing.T) {
@@ -75,6 +78,45 @@ func TestStreamDoesNotLeakTickersAcrossReconnects(t *testing.T) {
} }
} }
// updateHandler wants UPDATE messages and does nothing with them.
type updateHandler struct{}
func (updateHandler) WantsMessage(messageType string) bool { return messageType == "UPDATE" }
func (updateHandler) HandleMessage(*ristypes.RISMessage) {}
func (updateHandler) QueueCapacity() int { return 10 }
// TestStopBeforeMessageReachesQueues stops the streamer after the read loop
// has checked for cancellation but before it hands the message to the handler
// queues. That is the gap a stop from another goroutine can land in, and it
// used to end in "send on closed channel". The raw handler runs in that gap on
// the read loop itself, so calling Stop from it hits the gap every time.
func TestStopBeforeMessageReachesQueues(t *testing.T) {
const line = `{"type":"ris_message","data":{"type":"UPDATE","peer":"192.0.2.1",` +
`"peer_asn":"64496","timestamp":1700000000}}` + "\n"
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
_, _ = io.WriteString(w, line)
}))
defer srv.Close()
s := New(logger.New(), metrics.New())
s.url = srv.URL
s.RegisterHandler(updateHandler{})
s.RegisterRawHandler(func(string) { s.Stop() })
// Start would run the stream in the background, where the test cannot
// wait for it. Setting cancel as Start does lets Stop cancel the stream
// run here instead.
ctx, cancel := context.WithCancel(context.Background())
s.cancel = cancel
if err := s.stream(ctx); !errors.Is(err, context.Canceled) {
t.Fatalf("stream returned %v, want %v", err, context.Canceled)
}
// A second Stop must not close the queues again.
s.Stop()
}
// settledGoroutineCount lets transient goroutines finish, then reports the // settledGoroutineCount lets transient goroutines finish, then reports the
// current count. // current count.
func settledGoroutineCount() int { func settledGoroutineCount() int {
+2 -1
View File
@@ -7,7 +7,8 @@ package version
var ( var (
// GitRevision is the git commit hash // GitRevision is the git commit hash
GitRevision = "unknown" GitRevision = "unknown"
// GitRevisionShort is the short git commit hash (7 chars) // GitRevisionShort is the version the page footer shows: the tag or
// short commit hash from `git describe --tags --always`
GitRevisionShort = "unknown" GitRevisionShort = "unknown"
) )
+11 -2
View File
@@ -1,7 +1,8 @@
#!/bin/sh #!/bin/sh
# script/docker: build the Docker image tagged with the project name. # script/docker: build the Docker image tagged with the project name.
# Identical in all repos; the tag comes from script/projectname. # Identical in all repos; the tag comes from script/projectname.
# Generic: needs no adaptation. # --no-cache because the gate phases the final stage depends on are RUN
# steps, and a cached one is a check that did not run.
set -eu set -eu
SCRIPT_DIR="$(cd "$(dirname "$0")" && pwd -P)" SCRIPT_DIR="$(cd "$(dirname "$0")" && pwd -P)"
@@ -9,7 +10,15 @@ ROOT="$(cd "$SCRIPT_DIR/.." && pwd -P)"
main() { main() {
cd "$ROOT" cd "$ROOT"
docker build -t "$("$SCRIPT_DIR/projectname")" . # Own line: a failing command substitution inside an argument does
# not trip `set -e`, so the inline form degrades silently to an
# empty constant. The VERSION build argument takes precedence over
# the version a build stage derives from the .git in the context.
version="$(git describe --tags --always --dirty 2>/dev/null || true)"
[ -n "$version" ] || version="unknown"
docker build --no-cache \
--build-arg VERSION="$version" \
-t "$("$SCRIPT_DIR/projectname")" .
} }
main "$@" main "$@"