Skip to content
Merged
3 changes: 2 additions & 1 deletion cmd/ingestor/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -73,7 +73,8 @@ type Config struct {

// IATAWarnIntervalSec is how often a region dropped by
// ObserverIATAWhitelist is re-logged while it keeps arriving (#110).
// 0 or less means the default, 6 hours. See iata_drop_warn.go.
// 0 or less means the default, 6 hours; larger than 86400 (24 hours)
// is capped at 86400. See iata_drop_warn.go.
IATAWarnIntervalSec int `json:"iataWarnIntervalSec,omitempty"`

// iataDropWarn is the bounded per-region throttle for that warning.
Expand Down
10 changes: 9 additions & 1 deletion cmd/ingestor/iata_drop_warn.go
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,10 @@ const (
iataWarnMaxKeyLen = 32
// defaultIATAWarnIntervalSec is the default re-log interval (6h).
defaultIATAWarnIntervalSec = 6 * 60 * 60
// maxIATAWarnIntervalSec caps the configured interval (24h). A larger
// value buys nothing for a warning meant to stay visible, and a huge one
// overflows time.Duration into a negative interval that logs every drop.
maxIATAWarnIntervalSec = 24 * 60 * 60
)

// iataDropThrottle is the per-Config throttle state. The zero value is ready.
Expand All @@ -47,11 +51,15 @@ type iataDropThrottle struct {
sweeps int // full-table sweeps, for tests
}

// IATAWarnInterval returns how often a dropped region is re-logged.
// IATAWarnInterval returns how often a dropped region is re-logged:
// iataWarnIntervalSec capped at 24 hours, or the default for 0 or less.
func (c *Config) IATAWarnInterval() time.Duration {
if c == nil || c.IATAWarnIntervalSec <= 0 {
return defaultIATAWarnIntervalSec * time.Second
}
if c.IATAWarnIntervalSec > maxIATAWarnIntervalSec {
return maxIATAWarnIntervalSec * time.Second
}
return time.Duration(c.IATAWarnIntervalSec) * time.Second
}

Expand Down
117 changes: 117 additions & 0 deletions cmd/ingestor/iata_drop_warn_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ package main

import (
"fmt"
"math"
"os"
"path/filepath"
"strings"
Expand Down Expand Up @@ -32,6 +33,43 @@ func TestIATAWarnIntervalDefaultAndConfigured(t *testing.T) {
}
}

// A huge configured interval must not overflow time.Duration (int64
// nanoseconds) into a negative interval, which would log every drop: the
// flood the throttle exists to prevent (#150). It is capped at 24 hours.
func TestIATAWarnIntervalClampsHugeValues(t *testing.T) {
var huge int64 = 1e10 // seconds; 1e19 ns overflows int64
for _, sec := range []int{int(huge), math.MaxInt, 24*60*60 + 1, 365 * 24 * 60 * 60} {
if got := (&Config{IATAWarnIntervalSec: sec}).IATAWarnInterval(); got != 24*time.Hour {
t.Errorf("iataWarnIntervalSec=%d: interval %v, want the 24h cap", sec, got)
}
}
if got := (&Config{IATAWarnIntervalSec: 24 * 60 * 60}).IATAWarnInterval(); got != 24*time.Hour {
t.Errorf("exactly 24h: interval %v", got)
}
if got := (&Config{IATAWarnIntervalSec: 24*60*60 - 1}).IATAWarnInterval(); got != 24*time.Hour-time.Second {
t.Errorf("just under the cap: interval %v", got)
}

path := filepath.Join(t.TempDir(), "config.json")
if err := os.WriteFile(path, []byte(`{"observerIATAWhitelist":["ARN"],"iataWarnIntervalSec":10000000000}`), 0o600); err != nil {
t.Fatal(err)
}
cfg, err := LoadConfig(path)
if err != nil {
t.Fatal(err)
}
if got := cfg.IATAWarnInterval(); got != 24*time.Hour {
t.Fatalf("iataWarnIntervalSec=10000000000 from JSON: interval %v, want 24h", got)
}
// and the throttle built on it still suppresses a repeat
if w, _ := cfg.iataDropWarn.shouldWarn("GOT", t0IATA, cfg.IATAWarnInterval()); !w {
t.Fatal("first drop not logged")
}
if w, _ := cfg.iataDropWarn.shouldWarn("GOT", t0IATA.Add(time.Minute), cfg.IATAWarnInterval()); w {
t.Error("a repeat a minute later was logged again")
}
}

func TestIATAWarnIntervalFromJSON(t *testing.T) {
path := filepath.Join(t.TempDir(), "config.json")
if err := os.WriteFile(path, []byte(`{"observerIATAWhitelist":["ARN"],"iataWarnIntervalSec":600}`), 0o600); err != nil {
Expand Down Expand Up @@ -159,6 +197,85 @@ func TestIATADropThrottleSweepsAreAmortized(t *testing.T) {
}
}

// fillIATA tracks n new codes with the given prefix at ts; each must get its
// own slot.
func fillIATA(t *testing.T, th *iataDropThrottle, prefix string, n int, ts time.Time, iv time.Duration) {
t.Helper()
for i := 0; i < n; i++ {
if w, o := th.shouldWarn(fmt.Sprintf("%s%04d", prefix, i), ts, iv); !w || o {
t.Fatalf("%s%04d at %v: warn=%v overflow=%v, want an own slot", prefix, i, ts.Sub(t0IATA), w, o)
}
}
}

// oldest is a lower bound on the oldest tracked time (#150). With entries
// of different ages it must follow the oldest one still tracked: too high
// and an expired slot is not reclaimed (a new region falls into the shared
// overflow throttle), too low and a full table of fresh entries is swept
// on every drop.
func TestIATADropThrottleOldestTracksStaggeredEntries(t *testing.T) {
var th iataDropThrottle
iv := time.Hour
half := iataWarnMaxTracked / 2
fillIATA(t, &th, "A", half, t0IATA, iv)
fillIATA(t, &th, "B", iataWarnMaxTracked-half, t0IATA.Add(30*time.Minute), iv)

// full of fresh entries: overflow, and no sweep (nothing can be expired)
if w, o := th.shouldWarn("EARLY", t0IATA.Add(40*time.Minute), iv); !w || !o {
t.Fatalf("full table at +40m: warn=%v overflow=%v, want the overflow warning", w, o)
}
if th.sweeps != 0 {
t.Fatalf("%d sweeps before any entry could expire", th.sweeps)
}

// +1h: the A half has expired; a new region reclaims a slot
if w, o := th.shouldWarn("N1", t0IATA.Add(time.Hour), iv); !w || o {
t.Fatalf("+1h: warn=%v overflow=%v, want an own slot (A entries expired)", w, o)
}
if len(th.last) != iataWarnMaxTracked-half+1 || th.sweeps != 1 {
t.Fatalf("+1h: %d entries after %d sweeps, want %d after 1", len(th.last), th.sweeps, iataWarnMaxTracked-half+1)
}
// refill to the cap with entries at +1h
fillIATA(t, &th, "C", half-1, t0IATA.Add(time.Hour), iv)

// +1h30m: the B half (+30m) has expired now and must be reclaimed,
// although every entry swept at +1h is gone and the newest are fresh
if w, o := th.shouldWarn("N2", t0IATA.Add(90*time.Minute), iv); !w || o {
t.Fatalf("+1h30m: warn=%v overflow=%v, want an own slot (B entries expired)", w, o)
}
if _, ok := th.last["B0000"]; ok {
t.Error("expired B entries were not reclaimed at +1h30m")
}
if len(th.last) != half+1 || th.sweeps != 2 {
t.Fatalf("+1h30m: %d entries after %d sweeps, want %d after 2", len(th.last), th.sweeps, half+1)
}
}

// A drop can carry an earlier time than the current oldest entry (the
// caller takes the time before the lock). oldest must drop to it, or that
// entry outlives its interval in a full table.
func TestIATADropThrottleOldestFollowsAnEarlierDrop(t *testing.T) {
var th iataDropThrottle
iv := time.Hour
fillIATA(t, &th, "A", iataWarnMaxTracked, t0IATA, iv)
// +1h: everything expired; sweep, keep only LATE (taken at +1h)
if w, o := th.shouldWarn("LATE", t0IATA.Add(time.Hour), iv); !w || o {
t.Fatalf("+1h: warn=%v overflow=%v", w, o)
}
// a drop stamped +50m arrives after it
if w, o := th.shouldWarn("EARLIER", t0IATA.Add(50*time.Minute), iv); !w || o {
t.Fatalf("out-of-order drop: warn=%v overflow=%v", w, o)
}
fillIATA(t, &th, "F", iataWarnMaxTracked-2, t0IATA.Add(time.Hour), iv)
// +1h50m: EARLIER has expired and must free its slot
if w, o := th.shouldWarn("NEW", t0IATA.Add(110*time.Minute), iv); !w || o {
t.Fatalf("+1h50m: warn=%v overflow=%v, want EARLIER's expired slot", w, o)
}
if _, ok := th.last["EARLIER"]; ok {
t.Error("EARLIER outlived its interval in a full table")
}
}

func TestIATADropThrottleConcurrent(t *testing.T) {
var th iataDropThrottle
iv := time.Hour
Expand Down
2 changes: 1 addition & 1 deletion config.example.json
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,7 @@
"observerIATAWhitelist": [],
"_comment_observerIATAWhitelist": "Global IATA region whitelist. When non-empty, only observers whose IATA code (from MQTT topic) matches are processed. Case-insensitive. Empty = allow all. Unlike per-source iataFilter, this applies across all MQTT sources. A dropped region is logged by the ingestor as one '[region-filter] dropping region \"XYZ\"' line, repeated at most every iataWarnIntervalSec while it keeps arriving.",
"iataWarnIntervalSec": 21600,
"_comment_iataWarnIntervalSec": "Optional (ingestor). Seconds between repeated '[region-filter]' warnings for a region dropped by observerIATAWhitelist. Default 21600 (6h); 0 or omitted = default. At most 512 regions are tracked individually (region codes come from the MQTT topic, so this is bounded); beyond that, drops share one overflow warning with the same interval.",
"_comment_iataWarnIntervalSec": "Optional (ingestor). Seconds between repeated '[region-filter]' warnings for a region dropped by observerIATAWhitelist. Default 21600 (6h); 0 or omitted = default; maximum 86400 (24h), larger values are capped. At most 512 regions are tracked individually (region codes come from the MQTT topic, so this is bounded); beyond that, drops share one overflow warning with the same interval.",
"retention": {
"nodeDays": 7,
"observerDays": 14,
Expand Down
7 changes: 6 additions & 1 deletion public/live.js
Original file line number Diff line number Diff line change
Expand Up @@ -1182,6 +1182,9 @@
localStorage.setItem('live-matrix-mode', matrixMode);
applyMatrixTheme(matrixMode);
syncHeatToggleToMatrix(matrixMode);
// Matrix ON hid the heat layer; OFF brings it back as Heat is set.
// During init applyLiveControlEffects() (re)builds it after loadNodes().
if (!matrixMode && heatEnabled) showHeatMap();
} },
{ id: 'liveMatrixRainToggle', restore: () => matrixRain, onChange: (v) => {
matrixRain = v;
Expand All @@ -1202,14 +1205,16 @@
}

// Matrix mode owns the heat map: while it is on, the heat layer is hidden
// and its toggle unchecked and disabled. Safe before the map exists.
// and its toggle unchecked and disabled. When it is off, the toggle shows
// the Heat setting again (#150). Safe before the map exists.
function syncHeatToggleToMatrix(on) {
const ht = document.getElementById('liveHeatToggle');
if (on) {
hideHeatMap();
if (ht) { ht.checked = false; ht.disabled = true; }
} else if (ht) {
ht.disabled = false; // recover from stale state
ht.checked = heatEnabled;
}
}

Expand Down
17 changes: 12 additions & 5 deletions public/rx-coverage.js
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,11 @@
// coverage and leaderboard responses, the settle timer, move handlers)
// captures it and does nothing once a newer mount or a destroy happened.
var generation = 0;
// Request sequence per data stream (#150): a coverage or leaderboard
// response renders only if no newer request of its kind has started since,
// so a slow response for an earlier days (or rx, or viewport) cannot
// overwrite newer data.
var coverageSeq = 0, boardSeq = 0;

// Initial viewport (#124), first valid one wins: explicit URL lat/lon/zoom,
// this page's own saved view, /api/config/map, then this offline fallback
Expand All @@ -22,6 +27,8 @@
var MIN_ZOOM = 1, MAX_ZOOM = 19;

function isLive(gen) { return !destroyed && gen === generation; }
function isLatestCoverage(gen, seq) { return isLive(gen) && seq === coverageSeq; }
function isLatestBoard(gen, seq) { return isLive(gen) && seq === boardSeq; }

// validView returns {lat, lon, zoom} when all three are present, numeric and
// in range; anything invalid, partial or out of range gives null.
Expand Down Expand Up @@ -124,12 +131,12 @@

function drawCoverage() {
if (!map || destroyed) return;
var gen = generation;
var gen = generation, seq = ++coverageSeq;
var b = map.getBounds();
var bbox = [b.getSouth(), b.getWest(), b.getNorth(), b.getEast()].join(',');
var url = '/api/rx-coverage?bbox=' + bbox + '&z=' + map.getZoom() + '&days=' + days + (selectedRx ? '&rx=' + encodeURIComponent(selectedRx) : '');
fetch(url).then(function (r) { return r.json(); }).then(function (fc) {
if (!isLive(gen) || !covLayer) return;
if (!isLatestCoverage(gen, seq) || !covLayer) return;
covLayer.clearLayers();
(fc.features || []).forEach(function (f) {
var ring = (f.geometry.coordinates[0] || []).map(function (c) { return [c[1], c[0]]; });
Expand Down Expand Up @@ -251,11 +258,11 @@
}

function loadBoard() {
var gen = generation;
var gen = generation, seq = ++boardSeq;
fetch('/api/rx-leaderboard?days=' + days + '&limit=25').then(function (r) { return r.json(); })
.then(function (d) { if (!isLive(gen)) return; boardCache = d.observers || []; renderBoard(); })
.then(function (d) { if (!isLatestBoard(gen, seq)) return; boardCache = d.observers || []; renderBoard(); })
.catch(function (e) {
if (!isLive(gen)) return;
if (!isLatestBoard(gen, seq)) return;
console.warn('rx-coverage: leaderboard fetch failed', e);
var el = document.getElementById('rxBoard');
if (el) el.innerHTML = '<div class="muted" style="color:var(--text-muted);font-size:13px">Could not load mobile observers.</div>';
Expand Down
4 changes: 3 additions & 1 deletion public/style.css
Original file line number Diff line number Diff line change
Expand Up @@ -5101,7 +5101,9 @@ td[data-filter-field] { cursor: context-menu; }
justify-content: space-between;
padding: 14px 16px;
border-bottom: 1px solid var(--border);
background: var(--surface-2, var(--surface));
/* nav-drawer.css paints the title and close button with --nav-text(-muted),
so the header sits on the nav colours in every theme (#150). */
background: var(--nav-bg2);
}
.nav-drawer-title {
font-weight: 700;
Expand Down
54 changes: 54 additions & 0 deletions test-issue-124-rx-coverage-viewport.js
Original file line number Diff line number Diff line change
Expand Up @@ -281,6 +281,60 @@ function assertView(env, want, tag) {
assert.strictEqual(newLayer.cleared + newLayer.added, 0, 'the old coverage response was drawn on the new layer (cleared ' + newLayer.cleared + ', added ' + newLayer.added + ')');
});

await test('13. #150: switching days quickly, the older days response cannot overwrite the newer data', async () => {
const env = makeEnv({ storage: { 'rx-coverage-view': JSON.stringify({ lat: 56.1, lng: 9.9, zoom: 10 }) } });
const layers = [];
env.sandbox.L.layerGroup = () => {
const l = { cleared: 0, polys: [], addTo() { return l; }, clearLayers() { l.cleared++; l.polys = []; } };
layers.push(l);
return l;
};
env.sandbox.L.polygon = (ring) => {
const pg = { addTo(l) { l.polys.push(ring); return pg; }, bindTooltip() { return pg; } };
return pg;
};
// the days bar: record its click handler
const bar = env.sandbox.document.getElementById('rxDays');
bar.addEventListener = (ev, fn) => { if (ev === 'click') bar.onclick = fn; };
const pickDays = (d) => bar.onclick({ target: { closest: () => ({ dataset: { days: String(d) } }) } });
const feature = (lat) => ({ features: [{ properties: {}, geometry: { coordinates: [[[10, lat], [11, lat], [10, lat + 1]]] } }] });
const observers = (name) => ({ observers: [{ pubkey: name.toLowerCase(), name, score: 1, cells: 1, nodes: 1, receptions: 1 }] });

await mount(env, null);
env.runTimers(); // first coverage request (days=7)
assert(env.pending.some((p) => /rx-leaderboard\?days=7&/.test(p.url)), 'no days=7 leaderboard request');
assert(env.pending.some((p) => /rx-coverage\?bbox=.*&days=7$/.test(p.url)), 'no days=7 coverage request');

pickDays(30);
assert(env.pending.some((p) => /rx-leaderboard\?days=30&/.test(p.url)), 'no days=30 leaderboard request');
assert(env.pending.some((p) => /rx-coverage\?bbox=.*&days=30$/.test(p.url)), 'no days=30 coverage request');

// the newer (days=30) responses arrive first
const board = env.sandbox.document.getElementById('rxBoard');
assert(env.respond(/rx-leaderboard\?days=30&/, observers('NEWER')));
assert(env.respond(/rx-coverage\?bbox=.*&days=30$/, feature(60)));
await flush();
const layer = layers[0];
assert(/NEWER/.test(board.innerHTML), 'the days=30 leaderboard did not render');
assert(layer.polys.length === 1 && layer.polys[0][0][0] === 60, 'the days=30 coverage did not render: ' + JSON.stringify(layer.polys));

// then the slow days=7 responses are released
assert(env.respond(/rx-leaderboard\?days=7&/, observers('OLDER')));
assert(env.respond(/rx-coverage\?bbox=.*&days=7$/, feature(50)));
await flush();
assert(/NEWER/.test(board.innerHTML) && !/OLDER/.test(board.innerHTML), 'the older days=7 leaderboard replaced the days=30 one');
assert(layer.polys.length === 1 && layer.polys[0][0][0] === 60, 'the older days=7 coverage replaced the days=30 one: ' + JSON.stringify(layer.polys));

// a late failure of an older leaderboard request does not replace it either
pickDays(14);
pickDays(1);
assert(env.respond(/rx-leaderboard\?days=1&/, observers('NEWEST')));
await flush();
assert(env.failFetch(/rx-leaderboard\?days=14&/));
await flush();
assert(/NEWEST/.test(board.innerHTML), 'a failed older leaderboard request replaced the newest one: ' + board.innerHTML);
});

await test('11. coverage filtering and leaderboard requests are unchanged', async () => {
const env = makeEnv({ hash: '#/rx-coverage?days=14' });
await mount(env, { center: [55.68, 12.57], zoom: 9 });
Expand Down
Loading
Loading