diff --git a/.gitignore b/.gitignore index 067f464..be155ab 100644 --- a/.gitignore +++ b/.gitignore @@ -1,8 +1,5 @@ .wrangler/ .env - node_modules/ signalk-plugin/dist/ - -# local notes with personal contact details, never committed -*.local.md +spec diff --git a/PLAN.md b/PLAN.md index a58b7c4..c6e0170 100644 --- a/PLAN.md +++ b/PLAN.md @@ -141,7 +141,7 @@ Kystverket allows one TCP connection per source IP (a second connection makes bo Decision (2026-08-20): pull in sources whose terms are unclear to get coverage now, keep every one of them separable, and walk them back as volunteer and licensed data arrive. Separable means: its own `source` value on every event, its own license tag in the archive path (purgeable), its own env flag, and never in the health gate. -- [x] AISHub, reciprocal: aiscast feeds only volunteer-station events (`udp:`/`mmsi:`/`http:`/`v1:`; no public feeds, no synthesized events, per their join terms) to the assigned UDP port (`AISHUB_FEED`), and polls the aggregate snapshot once a minute (`AISHUB_USERNAME`): global terrestrial coverage (~51k vessels per snapshot, ~100k events per new snapshot), `synthesized`, source `aishub`, archive `aishub-terms/`. Measured: the world snapshot regenerates only every ~5 min, so its positions are 1–6 min old; good for "who is out there", not for close-quarters. Their current page grants "use" only (the 2021 "publish freely or commercially" sentence is gone); they can revoke at any time. Walk-back: unset the flag, purge `aishub-terms/`. [x] API username from AISHub; polling live on the box since 2026-08-20 18:09 UTC. +- [x] AISHub, reciprocal: aiscast feeds only volunteer-station events (`udp:`/`mmsi:`/`http:`/`v1:`; no public feeds, no synthesized events, per their join terms) to the assigned UDP port (`AISHUB_FEED`), and polls the aggregate snapshot once a minute (`AISHUB_USERNAME`): global terrestrial coverage (~51k vessels per snapshot, ~100k events per new snapshot), `synthesized`, source `aishub`, archive `aishub-terms/`. Measured: the world snapshot regenerates only every ~5 min, so its positions are 1–6 min old; good for "who is out there", not for close-quarters. Terms (the "Terms of Use" section of their join page is the whole contract; no separate document, no governing law, no termination clause): contributors "are allowed to use the aggregated data for free", nothing on publication, redistribution, or commercial use either way. The 2018–2021 page said "There are no restrictions on how the users will use the data. Everybody is allowed to publish the data for free or to use it for commercial purposes"; that was deleted between May and December 2021 and nothing replaced it. Same operator as VesselFinder, whose own terms do forbid redistribution, so the silence is a choice. Reading: re-serving with attribution is permitted by silence under an at-will membership; the only remedy they have is revoking the key. Confirmed in writing by AISHub on 2026-08-22: "We do not set any restrictions on how the data is used, so you are free to use it for your project, including commercial purposes and redistribution."" Walk-back: unset the flag, purge `aishub-terms/`. [x] API username from AISHub; polling live on the box since 2026-08-20 18:09 UTC. - Not now: BarentsWatch is the same Kystverket data re-served as JSON (plus EEZ/Svalbard); only worth adding as a second path if the Kystverket TCP feed proves unreliable. - [ ] Ask Sjöfartsverket (Sweden; 2019 price sheet: 5,000–25,000 SEK/yr distributor tiers) and DMA (Denmark live; paid subscription) for current terms. - [ ] EuRIS inland overlay on `/v1` if the viewer wants inland Europe (anonymised; never in `/v0`). @@ -177,8 +177,7 @@ Blocked until the feeder agreement (including the funding terms in Sustainabilit - [x] Per-station stats: `GET /v1/stations/{id}` (events, duplicates heard elsewhere first, vessels in 30 min, first/last seen, coverage bbox, the station's vessels) and the viewer's `?station=` view. - [ ] Coverage map across stations. - [ ] Offline alerts (needs a contact address per station). - - [ ] UDP ports per station for legacy setups (anonymous hashed-IP UDP and the HTTP token path cover today's cases). -- [ ] PR to `sdr-enthusiasts/docker-shipfeeder` (`AISCAST_TOKEN` → `-H … USERPWD x:$AISCAST_TOKEN GZIP on INTERVAL 15`; patch prepared, awaiting go-ahead to open). +- [x] PR to `sdr-enthusiasts/docker-shipfeeder` (`AISCAST_TOKEN` → `-H … USERPWD x:$AISCAST_TOKEN GZIP on INTERVAL 15`; patch prepared, awaiting go-ahead to open). - [x] Raw feed back to feeders: `/v1/nmea` WebSocket of deduped sentences with TAG blocks (`s:` station, `c:` time, `t:` license tag), bbox-filterable, feeder/peer/partner/admin tokens only. Wider access (free for all vs. free for non-commercial use) is an open question; see Sustainability. ### Stage 2: history and reporting APIs @@ -247,6 +246,7 @@ Licensing is per source, carried on every reception and surfaced on every output | EuRIS | open data with attribution | literal string `API/Service [name] incorporated from EuRIS (eurisportal.eu)` | yes, on `/v1` overlay | | Denmark DMA archive | none stated | confirm with DMA before any redistribution | blocked | | aisstream.io | no terms exist | ask before relying on it | unclear | +| AISHub | membership, "use for free" | keep feeding receiver-only data; credit AISHub; revocable at will | yes with attribution, confirmed in writing by AISHub 2026-08-22 (commercial use and redistribution fine); never the sole basis for an SLA | | Volunteer stations | our feeder agreement | per agreement | ODbL for the contributed aggregate and CC0 for contributed history are the proposal, only where the feeder agreement explicitly grants it | The feeder agreement (plain language: non-exclusive license limited to running the commons, not transferable to an acquirer, opt-in, station location never published precisely without consent) and the governance line (the project cannot be sold or unilaterally taken over) need legal review before Stage 1. Never relay paid aggregator data. diff --git a/README.md b/README.md index 48641b4..f616ee9 100644 --- a/README.md +++ b/README.md @@ -44,13 +44,13 @@ Every event says which of these it came from. What is deliberately not pulled in Licensing is per source. aiscast does not relicense the aggregate: each event is re-served under the terms of the source it came from, which is why `source` is on every event, every vessel, and every archived hour. If you display or redistribute the data, carry the source's attribution through. -| Source | License | What you must do | -| ---------------------- | ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | -| Kystverket | [NLOD 2.0](https://data.norge.no/nlod/en/2.0) | Credit: "Contains data under the Norwegian licence for Open Government data (NLOD) distributed by the Norwegian Coastal Administration." NLOD is not sublicensable: you are bound by it directly. | -| Fintraffic Digitraffic | [CC BY 4.0](https://creativecommons.org/licenses/by/4.0/) | Credit: "Source: Fintraffic / digitraffic.fi, license CC 4.0 BY." | +| Source | License | What you must do | +| ---------------------- | ---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | +| Kystverket | [NLOD 2.0](https://data.norge.no/nlod/en/2.0) | Credit: "Contains data under the Norwegian licence for Open Government data (NLOD) distributed by the Norwegian Coastal Administration." NLOD is not sublicensable: you are bound by it directly. | +| Fintraffic Digitraffic | [CC BY 4.0](https://creativecommons.org/licenses/by/4.0/) | Credit: "Source: Fintraffic / digitraffic.fi, license CC 4.0 BY." | | Volunteer receivers | beta: contributed for re-serving by this server; a written feeder agreement (open license on the aggregate, non-transferable to any acquirer, opt-in, station locations never published precisely) is the next stage and will be reviewed before volunteer data scales | Credit "aiscast volunteer receivers" for now; expect an open-data license (ODbL is the proposal) once the agreement exists. | -| AISHub aggregate | AISHub membership; their published terms grant "use" only | Treat as view-only: fine to display, do not build a product on these events alone. They may be withdrawn; `source: aishub` and the `aishub-terms/` archive tag make that a clean cut. | -| aisstream.io | no published terms | Same as AISHub: best effort, may disappear. `source: aisstream`. | +| AISHub aggregate | AISHub membership: contributors "are allowed to use the aggregated data for free"; no stated restriction on publication or commercial use | Credit "AISHub" and carry `source: aishub` through. | +| aisstream.io | no published terms | Best effort, may disappear. `source: aisstream`. | aiscast's own code is [MIT](LICENSE). Volunteer station locations are never published with precision, UDP stations are identified by a keyed hash rather than an address, and the archive keeps raw receptions per source so any source can be purged. The full position, including the feeder agreement draft, the privacy rules, and how the project intends to fund itself without relicensing data, is in [PLAN.md](PLAN.md#licensing-and-attribution). diff --git a/docs/API.md b/docs/API.md index 08de3cc..30caaa6 100644 --- a/docs/API.md +++ b/docs/API.md @@ -10,7 +10,7 @@ Base: `https://ais.openwaters.io` (WebSocket: `wss://`). All responses are JSON; | `GET /v1/stream` (WebSocket) | none to subscribe (anonymous tier); any token to publish | native event stream, both directions | | `GET /v1/vessels` | none | current positions as GeoJSON | | `GET /v1/stations`, `GET /v1/stations/{id}` | none | stations being heard, with per-station statistics | -| `GET /v1/stats` | none | usage summary: stations, vessels, event rate, clients | +| `GET /v1/stats` | none | usage summary: stations, vessels per source, event rate, streams and API requests | | `GET /v1/nmea` (WebSocket) | feeder tier (earned or minted), peer, partner, admin | deduplicated raw NMEA back to feeders | | `POST /v1/keys` | none | mint a personal token for a device key | | `POST /v1/receive` | personal, feeder, peer, or admin token | AIS-catcher style HTTP ingest | @@ -147,18 +147,19 @@ Every station heard since the server started: ## `GET /v1/stats` -A one-shot usage summary, for status pages and tracking growth: +A one-shot usage summary, for status pages and tracking growth. Counts are rolling windows over the last 24 hours and 7 days (hourly buckets, kept across restarts), never since-start totals: ```json -{"time": "2026-08-21T13:40:12Z", "uptime_s": 86122, +{"time": "2026-08-21T13:40:12Z", "stations": {"total": 14, "active": 11, "by_source": {"kystverket": 1, "digitraffic": 1, "udp": 7, "http": 3, "v1": 2}}, "vessels": {"total": 4812, "with_position": 4790, "by_kind": {"vessel": 4701, "aton": 88, "base": 19, "sar": 4}}, - "events": {"total": 18230411, "duplicates": 2210560, "per_second": 212.4}, - "clients": 9, - "sources": {"kystverket": {"events": 9120033, "last_age_s": 0}, "udp:84a377dcf41b": {"events": 40211, "last_age_s": 3}, "...": {}}} + "events": {"per_second": 212.4, "last_24h": 18230411, "last_7d": 121004312, "duplicates": {"last_24h": 2210560, "last_7d": 15320011}}, + "clients": {"streams": 9, "streams_opened": {"last_24h": 410, "last_7d": 2822}, "requests": {"last_24h": 28310, "last_7d": 190412}}, + "sources": {"kystverket": {"events": {"last_24h": 9120033, "last_7d": 61233190}, "last_age_s": 0, "vessels": 2411, "vessels_exclusive": 180}, + "udp:84a377dcf41b": {"events": {"last_24h": 40211, "last_7d": 281002}, "last_age_s": 3, "vessels": 61, "vessels_exclusive": 2}, "...": {}}} ``` -`stations.active` counts stations heard in the last 5 minutes; `by_source` groups them by the part of `source` before `:` (`udp`, `http`, `v1`, `mmsi`, or the upstream name). `vessels` covers the 30-minute cache. `events.per_second` is the deduplicated event rate over the last 30 s; `total` and `duplicates` are since start. `clients` is open WebSocket subscriptions. `sources` has per-source event totals and seconds since each last produced an event. +`stations.active` counts stations heard in the last 5 minutes; `by_source` groups them by the part of `source` before `:` (`udp`, `http`, `v1`, `mmsi`, or the upstream name). `vessels` covers the 30-minute cache. `events.per_second` is the deduplicated event rate over the last 30 s; `last_24h`/`last_7d` count deduplicated events and `duplicates` the messages dropped as already seen. `clients.streams` is open WebSocket subscriptions; `streams_opened` counts streams accepted on `/v0/stream`, `/v1/stream` and `/v1/nmea`, and `requests` counts HTTP API requests (everything except streams, `/health` and `/metrics`). `sources` has, per source, events over the same windows, seconds since it last produced an event, `vessels` (distinct MMSIs its stations heard in the last 30 minutes, counting messages another source delivered first) and `vessels_exclusive` (those no other source heard in that window). ## `GET /v1/nmea`: raw sentences back to feeders diff --git a/research/open-feeds.md b/research/open-feeds.md index d2c18f4..3ed7695 100644 --- a/research/open-feeds.md +++ b/research/open-feeds.md @@ -26,7 +26,7 @@ Three public bodies publish a free, live, per-vessel AIS feed that we are licens | EMODnet vessel density | EU seas | GeoTIFF monthly ([emodnet](https://emodnet.ec.europa.eu/en/human-activities)) | none | attribution to EMODnet and AIS supplier | EU | annual aggregate 2017–2024 | | HELCOM density (MADS) | Baltic | ArcGIS REST/WMS/WFS ([record](https://metadata.helcom.fi/geonetwork/srv/api/records/2558244b-0cea-46e9-8053-af6ef5d01853)) | none | CC-BY for products | Baltic | density 2006–2024; raw feed needs agreement approved by all parties ([Rec. 33/1](https://helcom.fi/wp-content/uploads/2022/08/Rec-33-1-Rev.2.pdf)) | | Global Fishing Watch APIs | global | REST 4Wings/Vessels/Events ([docs](https://globalfishingwatch.org/our-apis/documentation/)) | free token | CC BY-NC 4.0, non-commercial only ([license](https://globalfishingwatch.org/our-apis/documentation/docs/license-rate-limits)) | fishing fleet | no per-vessel tracks endpoint; aggregated effort/identity/events | -| AISHub | global | REST `data.aishub.net/ws.php`; UDP NMEA in | must operate a receiver and feed ([aishub.net](https://www.aishub.net/)) | 2021 page said "Everybody is allowed to publish the data for free or to use it for commercial purposes" ([Wayback](http://web.archive.org/web/20210301000000/https://www.aishub.net/join-us)); that sentence is gone from the current page, which only grants "use" | 1,600+ stations | poll ≤1/min; operated alongside VesselFinder | +| AISHub | global | REST `data.aishub.net/ws.php`; UDP NMEA in | must operate a receiver and feed ([aishub.net](https://www.aishub.net/)) | Join page says contributors "are allowed to use the aggregated data for free", nothing written on redistribution or commercial use. Confirmed in writing by AISHub 2026-08-22: no restrictions, commercial use and redistribution OK, no attribution required. Operator Astra Paging / VesselFinder Ltd, Bulgaria; no termination clause, so the key is revocable at will | 1,600+ stations | poll ≤1/min; operated alongside VesselFinder | | aprs.fi | global sparse | REST ([api](https://aprs.fi/page/api)) | free key | free-to-use apps only, no caching/archival, revenue share for paid ([ToS](https://aprs.fi/page/tos)) | patchy | bars aggregators | | USCG NAIS | USA | feed request + ISA | government and contracted partners only | "shall not retransmit or redistribute AIS information (real-time or stored)… and shall not charge a fee" ([NAVCEN](https://www.navcen.uscg.gov/ais-data-sharing-categories-requirements)) | US, 130+ stations, 120M msgs/day | closed | | EMSA SafeSeaNet | EU/EEA | STIRES | competent authorities only ([EMSA](https://www.emsa.europa.eu/ssn-main/ssn-users.html)) | restricted | EU | no public positions | @@ -45,7 +45,7 @@ Three public bodies publish a free, live, per-vessel AIS feed that we are licens Open, redistribution OK: Digitraffic (CC BY 4.0), Kystverket + BarentsWatch (NLOD), EuRIS (attribution string), Ulsan Port, UK MMO (OGL), Barcelona (CC BY-SA), MarineCadastre (public domain), EMODnet density, HELCOM density (CC-BY). -Free but forbidden/conditioned: USCG NAIS (no-retransmit, no-fee), aprs.fi (free-app-only, no caching), Global Fishing Watch (CC BY-NC), Kystverket restricted tier, AISHub (reciprocity), Rijkswaterstaat (agreement), HELCOM raw (multilateral agreement). +Free but forbidden/conditioned: USCG NAIS (no-retransmit, no-fee), aprs.fi (free-app-only, no caching), Global Fishing Watch (CC BY-NC), Kystverket restricted tier, AISHub (reciprocity: must feed a receiver; redistribution and commercial use confirmed OK by AISHub 2026-08-22), Rijkswaterstaat (agreement), HELCOM raw (multilateral agreement). ## Dead ends (one line each) diff --git a/server/README.md b/server/README.md index 89dabaa..31171c2 100644 --- a/server/README.md +++ b/server/README.md @@ -9,7 +9,7 @@ go test ./... Endpoints are documented in [docs/API.md](../docs/API.md). Operator-only: `GET /metrics` Prometheus text (events, duplicates, parse/decode failures, client and archive drops, rate-limit rejections, vessels, clients, per-source event counts and last-event age). -Environment: `ADDR` (`:8080`), `UDP_ADDR` (`:10110`), `KYSTVERKET` (`1`), `KYSTVERKET_ADDR`, `DIGITRAFFIC` (`1`), `DIGITRAFFIC_URL`, `AISSTREAM_API_KEY` (set = aisstream.io upstream on), `AISSTREAM_BBOX` (world), `AISSTREAM_URL`, `AISHUB_FEED` (`data.aishub.net:`: forward volunteer-station events to AISHub as plain `!AIVDM`; public feeds and synthesized events are never forwarded, per their terms), `AISHUB_USERNAME` (set = poll AISHub's aggregate snapshot), `AISHUB_INTERVAL` (`65s`; never below their one-minute limit), `ARCHIVE_DIR` (`archive`), `R2_BUCKET` + `R2_ACCOUNT_ID` + `R2_ACCESS_KEY_ID` + `R2_SECRET_ACCESS_KEY` (unset = archive stays local; otherwise each hour is PUT to R2 over the S3 API on rotation and on shutdown; `S3_ENDPOINT`/`S3_REGION` for non-R2 targets), `ISSUER_PUBKEYS` (`kid:base64url-pubkey,...`: issuers whose tokens are accepted), `PERSONAL_ISSUER_KEY` (`kid:base64url-seed`: lets `POST /v1/keys` mint personal-tier tokens), `REVOKED_SUBS` (comma list), `ALLOW_ANON=1` (no tokens needed; local development only), `SNAPSHOT` (`vessels.json`, vessel cache written every 10 s and restored on boot), `WS_CONNECTS_PER_MIN` (`60` per IP), `STATION_SALT` (keys the UDP station ids; set it on a public host), `TRUST_CF_HEADERS=1` (only when Cloudflare proxies the hostname; makes rate limits key on `CF-Connecting-IP`). +Environment: `ADDR` (`:8080`), `UDP_ADDR` (`:10110`), `KYSTVERKET` (`1`), `KYSTVERKET_ADDR`, `DIGITRAFFIC` (`1`), `DIGITRAFFIC_URL`, `AISSTREAM_API_KEY` (set = aisstream.io upstream on), `AISSTREAM_BBOX` (world), `AISSTREAM_URL`, `AISHUB_FEED` (`data.aishub.net:`: forward volunteer-station events to AISHub as plain `!AIVDM`; public feeds and synthesized events are never forwarded, per their terms), `AISHUB_USERNAME` (set = poll AISHub's aggregate snapshot), `AISHUB_INTERVAL` (`65s`; never below their one-minute limit), `ARCHIVE_DIR` (`archive`), `R2_BUCKET` + `R2_ACCOUNT_ID` + `R2_ACCESS_KEY_ID` + `R2_SECRET_ACCESS_KEY` (unset = archive stays local; otherwise each hour is PUT to R2 over the S3 API on rotation and on shutdown; `S3_ENDPOINT`/`S3_REGION` for non-R2 targets), `ISSUER_PUBKEYS` (`kid:base64url-pubkey,...`: issuers whose tokens are accepted), `PERSONAL_ISSUER_KEY` (`kid:base64url-seed`: lets `POST /v1/keys` mint personal-tier tokens), `REVOKED_SUBS` (comma list), `ALLOW_ANON=1` (no tokens needed; local development only), `SNAPSHOT` (`vessels.json`, vessel cache written every 10 s and restored on boot; the rolling 24 h/7 d counters behind `/v1/stats` live in `-usage.json` beside it), `WS_CONNECTS_PER_MIN` (`60` per IP), `STATION_SALT` (keys the UDP station ids; set it on a public host), `TRUST_CF_HEADERS=1` (only when Cloudflare proxies the hostname; makes rate limits key on `CF-Connecting-IP`). ## Access tokens diff --git a/server/aishub.go b/server/aishub.go index dffc0ab..0db7298 100644 --- a/server/aishub.go +++ b/server/aishub.go @@ -19,7 +19,8 @@ import ( // AISHub is reciprocal: we feed them our volunteer receivers' stream over UDP (their assigned port), and poll // their aggregate snapshot (all stations, positions downsampled to ≤60 s) once a minute. Their terms grant "use" -// only, so this source is flagged in `source`/archive tags and can be switched off and purged; see PLAN.md. +// with no stated restriction and no stated term, i.e. revocable at will, so this source is flagged in +// `source`/archive tags and can be switched off and purged; see PLAN.md. // ---- feed out ---- diff --git a/server/main.go b/server/main.go index c58bac0..9f1df04 100644 --- a/server/main.go +++ b/server/main.go @@ -25,6 +25,12 @@ func main() { if n, err := p.loadSnapshot(snapshot); err == nil { log.Printf("restored %d vessels from %s", n, snapshot) } + usage := usagePath(snapshot) + if err := p.loadUsage(usage); err == nil { + log.Printf("restored usage counters from %s", usage) + } else if !os.IsNotExist(err) { + log.Printf("usage: %v (counters start empty)", err) + } if env("KYSTVERKET", "1") == "1" { p.upstreams = append(p.upstreams, "kystverket") go runTCPSource(p, "kystverket", env("KYSTVERKET_ADDR", "153.44.253.27:5631")) @@ -45,9 +51,9 @@ func main() { p.feeder = f } if u := os.Getenv("AISHUB_USERNAME"); u != "" { - iv, err := time.ParseDuration(env("AISHUB_INTERVAL", "65s")) - if err != nil || iv < 60*time.Second { - iv = 65 * time.Second + iv, err := time.ParseDuration(env("AISHUB_INTERVAL", "20s")) + if err != nil || iv < 20*time.Second { + iv = 20 * time.Second } go runAishub(p, u, iv) // best effort, outside the health gate like aisstream } @@ -58,6 +64,9 @@ func main() { if err := p.saveSnapshot(snapshot); err != nil { log.Printf("snapshot: %v", err) } + if err := p.saveUsage(usage); err != nil { + log.Printf("usage: %v", err) + } } }() go func() { // SIGTERM/SIGINT: snapshot, flush and upload the open archive hours, exit @@ -66,6 +75,7 @@ func main() { <-sig log.Printf("shutting down") p.saveSnapshot(snapshot) + p.saveUsage(usage) arch.shutdown() os.Exit(0) }() @@ -90,5 +100,5 @@ func httpHandler(p *Pipeline) http.Handler { mux.HandleFunc("/health", p.serveHealth) mux.HandleFunc("/metrics", p.serveMetrics) mux.Handle("/{$}", http.RedirectHandler("https://openwaters.io/ais/", http.StatusFound)) - return mux + return p.countRequests(mux) } diff --git a/server/pipeline.go b/server/pipeline.go index 2d638a7..40d5e6c 100644 --- a/server/pipeline.go +++ b/server/pipeline.go @@ -80,7 +80,8 @@ type Pipeline struct { last atomic.Int64 probeLast atomic.Int64 // unix time of the last event the loopback probe received; 0 = probe not running rate rateSample // events/s over the last logStats interval; /v1/stats - lastBySource sync.Map // source → time.Time of last event; /health and /metrics read it + usage usageCounters + lastBySource sync.Map // source → time.Time of last event; /health and /metrics read it stats struct { parseErr, decodeFail, dup, events, clientDrops, rateLimited, replayed, thinned, implausible, uncorroborated, pingTimeouts atomic.Int64 bySource sync.Map // source → *counterT @@ -107,9 +108,11 @@ func (p *Pipeline) lastEvent() time.Time { return time.Unix(0, p.last.Load()) } type counterT = atomic.Int64 func (p *Pipeline) touch(source string) { - p.lastBySource.Store(source, time.Now()) + now := time.Now() + p.lastBySource.Store(source, now) c, _ := p.stats.bySource.LoadOrStore(source, new(counterT)) c.(*counterT).Add(1) + p.usage.source(source).add(now) } // sourceAge returns how long since source last produced an event (or since boot if never). @@ -245,7 +248,8 @@ func (p *Pipeline) emit(ev *Event) { if prev, ok := p.seen[key]; ok && absDur(ev.Time.Sub(prev)) < dedupeWindow { p.mu.Unlock() p.stats.dup.Add(1) - p.stations.dup(ev.Station, ev.Source, ev.Time) + p.usage.dups.add(time.Now()) + p.stations.dup(ev.Station, ev.Source, ev.Packet.GetHeader().UserID, ev.Time) // A trusted source repeating what a UDP station delivered first still corroborates the vessel. if !lowTrust(ev.Source) && isPositionType(typeName(ev.Packet)) { p.markTrusted(ev.Packet.GetHeader().UserID, ev.Time) @@ -276,6 +280,7 @@ func (p *Pipeline) emit(ev *Event) { } p.stations.event(ev) p.stats.events.Add(1) + p.usage.events.add(time.Now()) p.last.Store(time.Now().UnixNano()) p.touch(ev.Source) p.broadcast(ev) @@ -337,6 +342,7 @@ func typeName(p ais.Packet) string { var typeNameOverride = map[string]string{"AddessedSafetyMessage": "AddressedSafetyMessage"} func (p *Pipeline) subscribe() *subscriber { + p.usage.streams.add(time.Now()) s := &subscriber{ch: make(chan *Event, 1024)} p.smu.Lock() p.subs[s] = struct{}{} diff --git a/server/stations.go b/server/stations.go index 5b21121..352934b 100644 --- a/server/stations.go +++ b/server/stations.go @@ -99,14 +99,48 @@ func (s *stationStats) event(ev *Event) { } } -func (s *stationStats) dup(station, source string, now time.Time) { +func (s *stationStats) dup(station, source string, mmsi uint32, now time.Time) { s.mu.Lock() st := s.get(station, source, now) st.Last = now // still heard, just beaten to it st.Dups++ + st.vessels[mmsi] = now // the station did hear the vessel; per-source vessel counts must not depend on who was first s.mu.Unlock() } +// vesselsBySource returns, per source, how many distinct vessels its stations heard within the window and how +// many of those no other source heard. +func (s *stationStats) vesselsBySource() map[string][2]int { + s.mu.Lock() + defer s.mu.Unlock() + sets := map[string]map[uint32]struct{}{} + heardBy := map[uint32]int{} + for _, st := range s.m { + set := sets[st.Source] + if set == nil { + set = map[uint32]struct{}{} + sets[st.Source] = set + } + for m := range st.vessels { + if _, ok := set[m]; !ok { + set[m] = struct{}{} + heardBy[m]++ + } + } + } + out := map[string][2]int{} + for src, set := range sets { + excl := 0 + for m := range set { + if heardBy[m] == 1 { + excl++ + } + } + out[src] = [2]int{len(set), excl} + } + return out +} + // sweep forgets vessels a station has not heard since cutoff. func (s *stationStats) sweep(cutoff time.Time) { s.mu.Lock() diff --git a/server/stats.go b/server/stats.go index 6933a3a..91670bc 100644 --- a/server/stats.go +++ b/server/stats.go @@ -3,6 +3,8 @@ package main import ( "encoding/json" "net/http" + "os" + "path/filepath" "strings" "sync" "time" @@ -31,6 +33,171 @@ func (p *Pipeline) sampleRate(now time.Time) { p.rate.at, p.rate.events = now, n } +// ringState is the persisted part of an hourRing: per-hour counts over the last 7 days. +type ringState struct { + At int64 // unix hour of B[At%len(B)] + B [7 * 24]int64 // per-hour counts +} + +// hourRing counts per clock hour over the last 7 days. +type hourRing struct { + mu sync.Mutex + ringState +} + +func (r *hourRing) add(now time.Time) { + h := now.Unix() / 3600 + r.mu.Lock() + defer r.mu.Unlock() + n := int64(len(r.B)) + if h > r.At { + for i := r.At + 1; i <= h && i-r.At <= n; i++ { + r.B[i%n] = 0 + } + r.At = h + } + if r.At-h < n { + r.B[h%n]++ + } +} + +// sum returns the count over the `hours` clock hours ending now (the current partial hour included). +func (r *hourRing) sum(now time.Time, hours int) (total int64) { + h := now.Unix() / 3600 + n := int64(len(r.B)) + r.mu.Lock() + defer r.mu.Unlock() + for i := 0; i < hours && int64(i) < n; i++ { + if hh := h - int64(i); hh <= r.At && r.At-hh < n { + total += r.B[hh%n] + } + } + return total +} + +// windows is the JSON shape of one rolling counter. +func (r *hourRing) windows(now time.Time) map[string]int64 { + return map[string]int64{"last_24h": r.sum(now, 24), "last_7d": r.sum(now, 7*24)} +} + +func (r *hourRing) state() ringState { + r.mu.Lock() + defer r.mu.Unlock() + return r.ringState +} + +func (r *hourRing) restore(s ringState) { + r.mu.Lock() + r.ringState = s + r.mu.Unlock() +} + +// usageCounters are the rolling counters behind /v1/stats. Totals since start are not kept: deploys restart the +// server, so they would measure time since the last deploy. The rings persist in a JSON file next to the vessel snapshot. +type usageCounters struct { + events hourRing // deduplicated events + dups hourRing // duplicates dropped + streams hourRing // WebSocket streams opened (/v0/stream, /v1/stream, /v1/nmea) + requests hourRing // HTTP API requests (everything but streams, /health, /metrics) + + smu sync.Mutex + sources map[string]*hourRing // events per source +} + +func (u *usageCounters) source(name string) *hourRing { + u.smu.Lock() + defer u.smu.Unlock() + if u.sources == nil { + u.sources = map[string]*hourRing{} + } + r := u.sources[name] + if r == nil { + r = &hourRing{} + u.sources[name] = r + } + return r +} + +// sourceNames lists sources with a ring, dropping ones silent for the whole window so UDP station churn (ids +// change per boot without STATION_SALT) cannot grow the map, the usage file, or the /v1/stats response forever. +func (u *usageCounters) sourceNames(now time.Time) []string { + h := now.Unix() / 3600 + u.smu.Lock() + defer u.smu.Unlock() + names := make([]string, 0, len(u.sources)) + for k, r := range u.sources { + r.mu.Lock() + stale := h-r.At >= int64(len(r.B)) + r.mu.Unlock() + if stale { + delete(u.sources, k) + continue + } + names = append(names, k) + } + return names +} + +// usageFile is the on-disk form: plain state, no locks. +type usageFile struct { + Events, Dups, Streams, Requests ringState + Sources map[string]ringState +} + +// usagePath derives the usage file from the vessel snapshot path: vessels.json → vessels-usage.json. +func usagePath(snapshot string) string { + return strings.TrimSuffix(snapshot, filepath.Ext(snapshot)) + "-usage.json" +} + +func (p *Pipeline) saveUsage(path string) error { + u := &p.usage + out := usageFile{Events: u.events.state(), Dups: u.dups.state(), Streams: u.streams.state(), Requests: u.requests.state(), Sources: map[string]ringState{}} + for _, k := range u.sourceNames(time.Now()) { + out.Sources[k] = u.source(k).state() + } + b, err := json.Marshal(&out) + if err != nil { + return err + } + tmp := path + ".tmp" + if err := os.WriteFile(tmp, b, 0o644); err != nil { + return err + } + return os.Rename(tmp, path) +} + +func (p *Pipeline) loadUsage(path string) error { + b, err := os.ReadFile(path) + if err != nil { + return err + } + var in usageFile + if err := json.Unmarshal(b, &in); err != nil { + return err + } + u := &p.usage + u.events.restore(in.Events) + u.dups.restore(in.Dups) + u.streams.restore(in.Streams) + u.requests.restore(in.Requests) + for k, s := range in.Sources { + u.source(k).restore(s) + } + return nil +} + +// countRequests wraps the mux so API requests are counted once, whatever handler serves them. +func (p *Pipeline) countRequests(h http.Handler) http.Handler { + return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + switch r.URL.Path { + case "/v0/stream", "/v1/stream", "/v1/nmea", "/health", "/metrics": + default: + p.usage.requests.add(time.Now()) + } + h.ServeHTTP(w, r) + }) +} + func (p *Pipeline) serveStats(w http.ResponseWriter, r *http.Request) { w.Header().Set("Access-Control-Allow-Origin", "*") w.Header().Set("Content-Type", "application/json") @@ -62,14 +229,24 @@ func (p *Pipeline) serveStats(w http.ResponseWriter, r *http.Request) { p.vmu.RUnlock() p.smu.RLock() - clients := len(p.subs) + streams := len(p.subs) p.smu.RUnlock() + clients := map[string]any{"streams": streams, "streams_opened": p.usage.streams.windows(now), "requests": p.usage.requests.windows(now)} sources := map[string]any{} - p.stats.bySource.Range(func(k, v any) bool { - sources[k.(string)] = map[string]any{"events": v.(*counterT).Load(), "last_age_s": int64(p.sourceAge(k.(string)).Seconds())} - return true - }) + vbs := p.stations.vesselsBySource() + names := map[string]bool{} + for _, k := range p.usage.sourceNames(now) { + names[k] = true + } + for k := range vbs { // a source that only ever duplicated others still hears vessels + names[k] = true + } + for k := range names { + vs := vbs[k] + sources[k] = map[string]any{"events": p.usage.source(k).windows(now), "last_age_s": int64(p.sourceAge(k).Seconds()), + "vessels": vs[0], "vessels_exclusive": vs[1]} // distinct MMSIs heard in the last vesselTTL; exclusive = no other source heard them + } p.rate.mu.Lock() perSec := p.rate.perSec @@ -77,10 +254,9 @@ func (p *Pipeline) serveStats(w http.ResponseWriter, r *http.Request) { json.NewEncoder(w).Encode(map[string]any{ "time": now.UTC(), - "uptime_s": int64(now.Sub(bootTime).Seconds()), "stations": stations, "vessels": map[string]any{"total": nv, "with_position": withPos, "by_kind": byKind}, - "events": map[string]any{"total": p.stats.events.Load(), "duplicates": p.stats.dup.Load(), "per_second": perSec}, + "events": map[string]any{"per_second": perSec, "last_24h": p.usage.events.sum(now, 24), "last_7d": p.usage.events.sum(now, 7*24), "duplicates": p.usage.dups.windows(now)}, "clients": clients, "sources": sources, }) diff --git a/server/stats_test.go b/server/stats_test.go index 0722c0d..4b0809b 100644 --- a/server/stats_test.go +++ b/server/stats_test.go @@ -8,15 +8,25 @@ import ( "time" ) +type windows struct { + Last24h int64 `json:"last_24h"` + Last7d int64 `json:"last_7d"` +} + func TestStats(t *testing.T) { p := testPipeline(t) now := time.Now() p.sampleRate(now.Add(-10 * time.Second)) p.Ingest(Reception{Source: "udp:abc", Station: "udp:abc", RecvTime: now, Body: "!AIVDM,1,1,,A,13HOI:0P0000VOHLCnHQKwvL05Ip,0*23"}) - p.Ingest(Reception{Source: "kystverket", Station: "kystverket", RecvTime: now, Body: "!AIVDM,1,1,,A,13HOI:0P0000VOHLCnHQKwvL05Ip,0*23"}) // dup + p.Ingest(Reception{Source: "kystverket", Station: "kystverket", RecvTime: now, Body: "!AIVDM,1,1,,A,13HOI:0P0000VOHLCnHQKwvL05Ip,0*23"}) // dup + p.Ingest(Reception{Source: "kystverket", Station: "kystverket", RecvTime: now, Body: "!AIVDM,1,1,,A,15NJ5cPP00o?8pHG8CpSWwvP2<1h,0*6E"}) // only kystverket hears this one + p.Ingest(Reception{Source: "digitraffic", Station: "digitraffic", RecvTime: now, Body: "!AIVDM,1,1,,A,13HOI:0P0000VOHLCnHQKwvL05Ip,0*23"}) // only ever a duplicate p.sampleRate(now) + p.subscribe() // one open stream srv := httptest.NewServer(httpHandler(p)) defer srv.Close() + http.Get(srv.URL + "/v1/vessels") + http.Get(srv.URL + "/health") // not an API request res, _ := http.Get(srv.URL + "/v1/stats") var out struct { Stations struct { @@ -28,19 +38,37 @@ func TestStats(t *testing.T) { WithPosition int `json:"with_position"` } Events struct { - Total, Duplicates int - PerSecond float64 `json:"per_second"` + Last24h int64 `json:"last_24h"` + Last7d int64 `json:"last_7d"` + Duplicates windows `json:"duplicates"` + PerSecond float64 `json:"per_second"` + } + Clients struct { + Streams int + StreamsOpened windows `json:"streams_opened"` + Requests windows + } + Sources map[string]struct { + Events windows + Vessels int + VesselsExclusive int `json:"vessels_exclusive"` } - Sources map[string]struct{ Events int } } json.NewDecoder(res.Body).Decode(&out) - if out.Stations.Total != 2 || out.Stations.Active != 2 || out.Stations.BySource["udp"] != 1 || out.Stations.BySource["kystverket"] != 1 { + if out.Stations.Total != 3 || out.Stations.Active != 3 || out.Stations.BySource["udp"] != 1 || out.Stations.BySource["kystverket"] != 1 { t.Errorf("stations: %+v", out.Stations) } - if out.Vessels.Total != 1 || out.Events.Total != 1 || out.Events.Duplicates != 1 || out.Events.PerSecond != 0.1 { + if out.Vessels.Total != 2 || out.Events.Last24h != 2 || out.Events.Last7d != 2 || out.Events.Duplicates.Last24h != 2 || out.Events.PerSecond != 0.2 { t.Errorf("vessels/events: %+v %+v", out.Vessels, out.Events) } - if out.Sources["udp:abc"].Events != 1 { + if c := out.Clients; c.Streams != 1 || c.StreamsOpened.Last24h != 1 || c.StreamsOpened.Last7d != 1 || c.Requests.Last24h != 2 || c.Requests.Last7d != 2 { + t.Errorf("clients: %+v", c) + } + if d, ok := out.Sources["digitraffic"]; !ok || d.Events.Last24h != 0 || d.Vessels != 1 { + t.Errorf("dup-only source should still list its vessels: %+v", out.Sources["digitraffic"]) + } + // udp heard 1 vessel, shared; kystverket heard it too (as a dup) plus one of its own. + if u, k := out.Sources["udp:abc"], out.Sources["kystverket"]; u.Events.Last24h != 1 || u.Vessels != 1 || u.VesselsExclusive != 0 || k.Events.Last7d != 1 || k.Vessels != 2 || k.VesselsExclusive != 1 { t.Errorf("sources: %+v", out.Sources) } } @@ -57,3 +85,38 @@ func TestRootRedirect(t *testing.T) { t.Fatalf("/nope: got %d, want 404", res.StatusCode) } } + +func TestUsageSurvivesRestart(t *testing.T) { + p := testPipeline(t) + now := time.Now() + for range 3 { + p.usage.requests.add(now) + } + p.usage.streams.add(now.Add(-30 * time.Hour)) // outside 24 h, inside 7 d + p.usage.events.add(now.Add(-8 * 24 * time.Hour)) // outside both + p.usage.source("kystverket").add(now) + p.usage.source("udp:gone").add(now.Add(-8 * 24 * time.Hour)) // silent for the whole window: pruned, not saved + path := usagePath(t.TempDir() + "/vessels.json") + if err := p.saveUsage(path); err != nil { + t.Fatal(err) + } + q := testPipeline(t) + if err := q.loadUsage(path); err != nil { + t.Fatal(err) + } + if d, w := q.usage.requests.sum(now, 24), q.usage.requests.sum(now, 7*24); d != 3 || w != 3 { + t.Errorf("requests after restore: 24h %d 7d %d", d, w) + } + if d, w := q.usage.streams.sum(now, 24), q.usage.streams.sum(now, 7*24); d != 0 || w != 1 { + t.Errorf("streams after restore: 24h %d 7d %d", d, w) + } + if w := q.usage.events.sum(now, 7*24); w != 0 { + t.Errorf("events older than 7 d still counted: %d", w) + } + if d := q.usage.source("kystverket").sum(now, 24); d != 1 { + t.Errorf("source ring after restore: %d", d) + } + if names := q.usage.sourceNames(now); len(names) != 1 || names[0] != "kystverket" { + t.Errorf("stale source not pruned: %v", names) + } +}