diff --git a/app.go b/app.go index 191fbe2..f9e8839 100644 --- a/app.go +++ b/app.go @@ -1392,13 +1392,12 @@ func (a *App) startup(ctx context.Context) { go a.sendTelemetryHeartbeat() go a.liveStatusLoop() // multi-op: heartbeat current activity to shared MySQL go a.chatLoop() // multi-op: poll the shared chat + heartbeat presence - // PSK Reporter, when the opening watch is on. After the operator's grid is - // known: without it there is no distance to measure and the feed stays down. - a.startBandOpenFeed() - // Locator store. After settings are scoped, so the option is read from the - // right profile, and before the cluster starts serving statuses so the first - // spots already carry their locators. + // Locator store BEFORE the feed: the feed asks whether it exists to decide + // which bands to subscribe to. a.startGridCache() + // PSK Reporter. After the operator's grid is known: without it there is no + // distance to measure and no receiver squares to filter on, so it stays down. + a.startBandOpenFeed() // One-time tidy-up of a field nothing used to record. Background, once. a.backfillDistancesOnce() @@ -11966,7 +11965,7 @@ func (a *App) consumeUDPEvents() { // Remember the grid before anything else: a CQ is the one message that // carries it, and the station may never send another. if ev.DecodeGrid != "" { - a.rememberDecodeGrid(ev.DecodeCall, ev.DecodeGrid) + a.rememberDecodeGrid(ev.DecodeCall, ev.DecodeGrid, gridcache.SourceDecode) } // A WSJT-X decode (heard station). Render it on the FlexRadio // panadapter when the option is on; green + SNR comment, auto-expiring @@ -17046,10 +17045,33 @@ const decodeGridsCap = 100000 // GetChaseNewGrids reports whether learnt locators are kept across restarts. func (a *App) GetChaseNewGrids() bool { return a.settingOr(keyChaseNewGrids, "") == "1" } +// GridCacheStatus is the live count under the option. A store that is on but +// has learnt nothing looks exactly like a broken one until a number moves. +type GridCacheStatus struct { + Enabled bool `json:"enabled"` + Known int `json:"known"` // locators in memory, stored + learnt this session + Pending int `json:"pending"` // waiting for the next batch write +} + +// GetGridCacheStatus reports what the locator store holds. +func (a *App) GetGridCacheStatus() GridCacheStatus { + out := GridCacheStatus{Enabled: a.gridStore != nil} + a.decodeGridsMu.RLock() + out.Known = len(a.decodeGrids) + len(a.decodeGridsOld) + a.decodeGridsMu.RUnlock() + if a.gridStore != nil { + out.Pending = a.gridStore.Pending() + } + return out +} + // SetChaseNewGrids turns grid chasing on or off and applies it immediately. func (a *App) SetChaseNewGrids(on bool) error { a.setSetting(keyChaseNewGrids, boolStr(on)) a.startGridCache() + // Resubscribe: with grid chasing the feed covers every band, without it only + // the opening bands — and with neither consumer it comes down entirely. + a.startBandOpenFeed() return nil } @@ -17108,7 +17130,7 @@ func (a *App) startGridCache() { // previous one and a fresh map takes over. Nothing is scanned, nothing is // timestamped, and the write stays a single map assignment — which matters // because this runs once per decode. -func (a *App) rememberDecodeGrid(call, grid string) { +func (a *App) rememberDecodeGrid(call, grid, source string) { call = strings.ToUpper(strings.TrimSpace(call)) grid = strings.TrimSpace(grid) if call == "" || grid == "" { @@ -17140,7 +17162,7 @@ func (a *App) rememberDecodeGrid(call, grid string) { a.decodeGridsMu.Unlock() if changed && store != nil { - store.Put(call, grid) + store.Put(call, grid, source) } } diff --git a/bandopen_sources.go b/bandopen_sources.go index 8dbdd3f..930542c 100644 --- a/bandopen_sources.go +++ b/bandopen_sources.go @@ -20,6 +20,8 @@ import ( "hamlog/internal/applog" "hamlog/internal/bandopen" + "hamlog/internal/geo" + "hamlog/internal/gridcache" "hamlog/internal/cluster" "hamlog/internal/pskr" ) @@ -133,18 +135,46 @@ func (a *App) startBandOpenFeed() { // on a timer fed by spots this path no longer looks at — so they would // hang there until the app restarted. a.clearBandOpenings() + } + chaseGrids := a.gridStore != nil + if !s.Enabled && !chaseGrids { return } // Every spot is measured from the operator's position. Without one there is // nothing to measure, and a detector fed unmeasurable spots reports nothing // while looking like it is working. if !a.opSet { - applog.Printf("bandopen: no station grid set — the opening watch needs one to measure a path") + applog.Printf("pskr: no station grid set — the feed needs one to measure a path") return } + + // Grid chasing wants every band; the opening watch wants its four. "+" is the + // MQTT single-level wildcard, so one subscription per square covers the lot. + bands := s.Bands + if chaseGrids { + bands = []string{"+"} + } + // Filter at the BROKER on the receiver's square rather than receiving the + // world and discarding it here. Measured on the live feed: the four opening + // bands unfiltered are 83 messages a second, of which roughly one in a + // hundred survived the NearKm test below. One ring of squares — about the + // same 300 km — is 0.2 to 1.2 a second, and the same for every operator, + // where filtering by DXCC ranged from 1.2 (OH) to 72.5 (K). + rxGrids := geo.NeighbourGrids(a.opLat, a.opLon, 1) + + var onGrid func(call, grid string) + if chaseGrids { + onGrid = func(call, grid string) { a.rememberDecodeGrid(call, grid, gridcache.SourceMQTT) } + } + var onSpot func(pskr.Spot) + if s.Enabled { + onSpot = a.feedBandOpen + } + a.pskr = pskr.New(pskr.Config{ - Bands: s.Bands, - OpLat: a.opLat, OpLon: a.opLon, + Bands: bands, + RxGrids: rxGrids, + OpLat: a.opLat, OpLon: a.opLon, Geo: func(grid string) (int, int, bool) { lat, lon, ok := gridToLatLon(grid) if !ok { @@ -156,12 +186,16 @@ func (a *App) startBandOpenFeed() { b := int(initialBearingDeg(a.opLat, a.opLon, lat, lon) + 0.5) return d, b, true }, - OnSpot: a.feedBandOpen, + OnSpot: onSpot, + OnGrid: onGrid, Logf: applog.Printf, }) if err := a.pskr.Start(); err != nil { - applog.Printf("bandopen: PSK Reporter feed did not start: %v", err) + applog.Printf("pskr: feed did not start: %v", err) + return } + applog.Printf("pskr: feed up — bands %v, %d receiver squares (openings=%v, grids=%v)", + bands, len(rxGrids), s.Enabled, chaseGrids) } // feedBandOpen hands one PSK Reporter decode to the detector. diff --git a/changelog.json b/changelog.json index 1b2e9c8..37b0946 100644 --- a/changelog.json +++ b/changelog.json @@ -4,11 +4,13 @@ "date": "", "en": [ "Cluster: the grid cache now holds 100,000 callsigns and rotates instead of emptying itself, so locators stop vanishing from the list.", - "New option \"Chase new grids\": locators learnt from decodes are kept in their own database, so the cluster shows them from the first second." + "New option \"Chase new grids\": locators learnt from your decodes and from PSK Reporter are kept in their own database, with their source.", + "Band openings: the PSK Reporter feed is now filtered at the broker, which cuts it from about 83 messages a second to under two." ], "fr": [ "Cluster : le cache de locators garde 100 000 indicatifs et tourne au lieu de se vider, les locators ne disparaissent donc plus de la liste.", - "Nouvelle option « Chasse aux nouveaux carrés » : les locators appris des décodes sont conservés dans leur propre base, le cluster les affiche donc dès la première seconde." + "Nouvelle option « Chasse aux nouveaux carrés » : les locators appris de tes décodes et de PSK Reporter sont conservés dans leur propre base, avec leur source.", + "Ouvertures de bande : le flux PSK Reporter est désormais filtré chez le broker, ce qui le fait passer d environ 83 messages par seconde à moins de deux." ] }, { diff --git a/decodegrids_test.go b/decodegrids_test.go index 9518f0c..b875c76 100644 --- a/decodegrids_test.go +++ b/decodegrids_test.go @@ -22,7 +22,7 @@ func TestDecodeGridRotationKeepsThePreviousGeneration(t *testing.T) { // Fill exactly one generation. for i := 0; i < decodeGridsCap; i++ { - a.rememberDecodeGrid(fmt.Sprintf("CALL%06d", i), "JN36") + a.rememberDecodeGrid(fmt.Sprintf("CALL%06d", i), "JN36", gridcache.SourceDecode) } if got := a.lookupDecodeGrid("CALL000000"); got != "JN36" { t.Fatalf("first entry lost before any rotation: %q", got) @@ -32,7 +32,7 @@ func TestDecodeGridRotationKeepsThePreviousGeneration(t *testing.T) { } // One more entry rotates. - a.rememberDecodeGrid("NEWCALL", "IO91") + a.rememberDecodeGrid("NEWCALL", "IO91", gridcache.SourceDecode) if a.decodeGridsOld == nil { t.Fatal("did not rotate at the cap") } @@ -49,7 +49,7 @@ func TestDecodeGridRotationKeepsThePreviousGeneration(t *testing.T) { // A second rotation is what finally retires the oldest half. for i := 0; i < decodeGridsCap; i++ { - a.rememberDecodeGrid(fmt.Sprintf("SECOND%06d", i), "KP20") + a.rememberDecodeGrid(fmt.Sprintf("SECOND%06d", i), "KP20", gridcache.SourceDecode) } if got := a.lookupDecodeGrid("CALL000000"); got != "" { t.Errorf("the cache is unbounded: %q survived two rotations", got) @@ -63,12 +63,12 @@ func TestDecodeGridRotationKeepsThePreviousGeneration(t *testing.T) { // "f4bpo" would miss a grid learnt as "F4BPO". func TestDecodeGridCaseAndBlanks(t *testing.T) { a := &App{} - a.rememberDecodeGrid(" f4bpo ", "JN36") + a.rememberDecodeGrid(" f4bpo ", "JN36", gridcache.SourceDecode) if got := a.lookupDecodeGrid("F4BPO"); got != "JN36" { t.Errorf("lookup of the upper-case form failed: %q", got) } - a.rememberDecodeGrid("", "JN36") - a.rememberDecodeGrid("K1ABC", "") + a.rememberDecodeGrid("", "JN36", gridcache.SourceDecode) + a.rememberDecodeGrid("K1ABC", "", gridcache.SourceDecode) if got := a.lookupDecodeGrid("K1ABC"); got != "" { t.Errorf("stored an empty grid: %q", got) } @@ -87,7 +87,7 @@ func TestDecodeGridConcurrentAccess(t *testing.T) { go func() { defer wg.Done() for i := 0; i < 20000; i++ { - a.rememberDecodeGrid(fmt.Sprintf("W%05d", i), "FN31") + a.rememberDecodeGrid(fmt.Sprintf("W%05d", i), "FN31", gridcache.SourceDecode) } }() go func() { @@ -112,7 +112,7 @@ func TestOnlyChangesAreQueued(t *testing.T) { defer st.Close() a := &App{gridStore: st} - a.rememberDecodeGrid("F4BPO", "JN36") + a.rememberDecodeGrid("F4BPO", "JN36", gridcache.SourceDecode) if n := st.Pending(); n != 1 { t.Fatalf("a new locator queued %d writes, want 1", n) } @@ -122,14 +122,14 @@ func TestOnlyChangesAreQueued(t *testing.T) { // The same report, a hundred times over, is not news. for i := 0; i < 100; i++ { - a.rememberDecodeGrid("F4BPO", "JN36") + a.rememberDecodeGrid("F4BPO", "JN36", gridcache.SourceDecode) } if n := st.Pending(); n != 0 { t.Errorf("unchanged reports queued %d writes — the batch would carry the whole feed", n) } // A station that moved is. - a.rememberDecodeGrid("F4BPO", "KP30") + a.rememberDecodeGrid("F4BPO", "KP30", gridcache.SourceDecode) if n := st.Pending(); n != 1 { t.Errorf("a changed locator queued %d writes, want 1", n) } @@ -149,7 +149,7 @@ func TestNoRotationWhilePersisting(t *testing.T) { a := &App{gridStore: st} for i := 0; i < decodeGridsCap+10; i++ { - a.rememberDecodeGrid(fmt.Sprintf("CALL%06d", i), "JN36") + a.rememberDecodeGrid(fmt.Sprintf("CALL%06d", i), "JN36", gridcache.SourceDecode) } if a.decodeGridsOld != nil { t.Error("rotated while a store was attached — locators the database holds would go missing") diff --git a/frontend/src/components/SettingsModal.tsx b/frontend/src/components/SettingsModal.tsx index 380f573..603e019 100644 --- a/frontend/src/components/SettingsModal.tsx +++ b/frontend/src/components/SettingsModal.tsx @@ -52,7 +52,7 @@ import { GetADIFMonitor, SaveADIFMonitor, PickADIFMonitorFile, GetRelayAuto, SaveRelayAuto, GetStationDevices, GetAwardDefs, GetTrackedAwards, SaveTrackedAwards, - GetBandOpenSettings, SaveBandOpenSettings, GetPSKReporterStatus, GetChaseNewGrids, SetChaseNewGrids, + GetBandOpenSettings, SaveBandOpenSettings, GetPSKReporterStatus, GetChaseNewGrids, SetChaseNewGrids, GetGridCacheStatus, } from '../../wailsjs/go/main/App'; import type { profile as profileModels } from '../../wailsjs/go/models'; import type { LookupSettingsForm, StationSettingsForm, ListsSettingsForm, ModePresetForm } from '@/types'; @@ -1553,6 +1553,7 @@ export function SettingsModal({ onClose, onSaved, initialSection, onMainPaneChan // feed up or down — so the write has to go where those live. const [bandOpen, setBandOpen] = useState({ enabled: false, bands: [], available: [] }); const [chaseGrids, setChaseGrids] = useState(false); + const [gridStat, setGridStat] = useState(null); const [pskrStatus, setPskrStatus] = useState(null); const saveBandOpen = async (next: any) => { setBandOpen(next); @@ -1567,6 +1568,7 @@ export function SettingsModal({ onClose, onSaved, initialSection, onMainPaneChan // distinguishes "connected" from "connected and receiving nothing". const t = window.setInterval(async () => { try { setPskrStatus(await GetPSKReporterStatus()); } catch { /* ignore */ } + try { setGridStat(await GetGridCacheStatus()); } catch { /* ignore */ } }, 3000); return () => window.clearInterval(t); }, []); @@ -4238,6 +4240,12 @@ export function SettingsModal({ onClose, onSaved, initialSection, onMainPaneChan onCheckedChange={(c) => { setChaseGrids(!!c); SetChaseNewGrids(!!c).catch(() => {}); }} /> {t('clu.chaseGrids')} {t('clu.chaseGridsHint')} + {chaseGrids && ( +

+ {t('clu.chaseGridsStat', { n: gridStat?.known ?? 0, p: gridStat?.pending ?? 0 })} + {pskrStatus?.running ? ` · ${t('bo.feedUp', { n: pskrStatus.received ?? 0 })}` : ''} +

+ )} {/* Band-opening watch. It lives HERE, with the cluster nodes, because diff --git a/frontend/src/lib/i18n.tsx b/frontend/src/lib/i18n.tsx index 5d0a202..3d55cde 100644 --- a/frontend/src/lib/i18n.tsx +++ b/frontend/src/lib/i18n.tsx @@ -270,7 +270,7 @@ const en: Dict = { 'clu.muteWorkedHint': '(they stay in the list, just quiet — leaves the colour for what is left to do)', 'clu.slotHighlight': 'Colour the stations not worked on this band and mode', 'clu.slotHighlightHint': '(by callsign, whatever the entity status says)', - 'bo.open': 'open', 'bo.liveTip': '{band} is open — {n} stations, ~{km} km, {sector}{season}. Click for the band map.', 'bo.enable': 'Watch for band openings', 'bo.enableHint': '(10, 12, 6, 4 and 2 m. Switching this on adds the two RBN nodes and subscribes to the PSK Reporter feed — the detection needs far more ears than a cluster can give it.)', 'bo.feedUp': 'PSK Reporter feed up — {n} decodes seen', 'bo.feedDown': 'PSK Reporter feed down — needs your station grid, and a moment to connect', 'clu.chaseGrids': 'Chase new grids', 'clu.chaseGridsHint': '(keeps the locators learnt from WSJT-X decodes in their own database, so the cluster shows them from the first second instead of after an hour of listening)', 'clu.workedSameSlot': 'Already worked only on the same slot', + 'bo.open': 'open', 'bo.liveTip': '{band} is open — {n} stations, ~{km} km, {sector}{season}. Click for the band map.', 'bo.enable': 'Watch for band openings', 'bo.enableHint': '(10, 12, 6, 4 and 2 m. Switching this on adds the two RBN nodes and subscribes to the PSK Reporter feed — the detection needs far more ears than a cluster can give it.)', 'bo.feedUp': 'PSK Reporter feed up — {n} decodes seen', 'bo.feedDown': 'PSK Reporter feed down — needs your station grid, and a moment to connect', 'clu.chaseGrids': 'Chase new grids', 'clu.chaseGridsHint': '(learns locators from your own WSJT-X decodes AND from PSK Reporter, and keeps them in their own database so the cluster shows them from the first second)', 'clu.chaseGridsStat': '{n} locators known — {p} waiting to be written', 'clu.workedSameSlot': 'Already worked only on the same slot', 'clu.workedSameSlotHint': '— a spot shows "worked" only if you worked that call on the SAME band and mode, not just anywhere. Combines with digital-mode grouping (Settings → General): with it on, a call worked on 20m FT8 also counts as worked for a 20m FT4 spot; with it off, FT8 and FT4 are separate slots.', // Backup panel 'bk.hintMysql': 'On close (once/day) OpsLog snapshots the local SQLite (config) AND exports the shared MySQL log to ADIF — opslog-log-.adi — so your contacts are protected even though they live on the server. Rotation keeps the last N of each.', diff --git a/frontend/wailsjs/go/main/App.d.ts b/frontend/wailsjs/go/main/App.d.ts index a076cc8..f4ed933 100644 --- a/frontend/wailsjs/go/main/App.d.ts +++ b/frontend/wailsjs/go/main/App.d.ts @@ -442,6 +442,8 @@ export function GetFlexBandPower():Promise>; export function GetFlexState():Promise; +export function GetGridCacheStatus():Promise; + export function GetIcomState():Promise; export function GetListsSettings():Promise; diff --git a/frontend/wailsjs/go/main/App.js b/frontend/wailsjs/go/main/App.js index 10b9c87..80269b6 100644 --- a/frontend/wailsjs/go/main/App.js +++ b/frontend/wailsjs/go/main/App.js @@ -826,6 +826,10 @@ export function GetFlexState() { return window['go']['main']['App']['GetFlexState'](); } +export function GetGridCacheStatus() { + return window['go']['main']['App']['GetGridCacheStatus'](); +} + export function GetIcomState() { return window['go']['main']['App']['GetIcomState'](); } diff --git a/frontend/wailsjs/go/models.ts b/frontend/wailsjs/go/models.ts index aaa9223..09eab0d 100644 --- a/frontend/wailsjs/go/models.ts +++ b/frontend/wailsjs/go/models.ts @@ -2419,6 +2419,22 @@ export namespace main { this.body = source["body"]; } } + export class GridCacheStatus { + enabled: boolean; + known: number; + pending: number; + + static createFrom(source: any = {}) { + return new GridCacheStatus(source); + } + + constructor(source: any = {}) { + if ('string' === typeof source) source = JSON.parse(source); + this.enabled = source["enabled"]; + this.known = source["known"]; + this.pending = source["pending"]; + } + } export class ModePreset { name: string; default_rst_sent?: string; diff --git a/internal/geo/geo.go b/internal/geo/geo.go index 56c3dd5..fc3a479 100644 --- a/internal/geo/geo.go +++ b/internal/geo/geo.go @@ -67,3 +67,54 @@ func DistanceBetweenGrids(a, b string) (km float64, ok bool) { } return HaversineKm(lat1, lon1, lat2, lon2), true } + +// LatLonToGrid returns the 4-character Maidenhead square for a position. +func LatLonToGrid(lat, lon float64) string { + lon = math.Mod(lon+180, 360) + if lon < 0 { + lon += 360 + } + lat = lat + 90 + if lat < 0 { + lat = 0 + } else if lat > 180 { + lat = 180 + } + return string([]byte{ + byte('A' + int(lon/20)), + byte('A' + int(lat/10)), + byte('0' + int(math.Mod(lon, 20)/2)), + byte('0' + int(math.Mod(lat, 10)/1)), + }) +} + +// NeighbourGrids returns the square holding (lat, lon) and the ring of squares +// around it — 9 squares for ring 1, 25 for ring 2. +// +// Used to filter the PSK Reporter feed at the BROKER rather than in OpsLog. A +// square is about 111 km tall and 150 km wide at mid latitudes, so one ring is +// roughly the 300 km "around here" the feed already meant, and the traffic that +// used to be received and discarded is never sent. +func NeighbourGrids(lat, lon float64, ring int) []string { + if ring < 0 { + ring = 0 + } + seen := map[string]bool{} + out := []string{} + for dLat := -ring; dLat <= ring; dLat++ { + for dLon := -ring; dLon <= ring; dLon++ { + // One square step: 1° of latitude, 2° of longitude. + la := lat + float64(dLat) + lo := lon + float64(dLon)*2 + if la > 90 || la < -90 { + continue // past a pole there is no square, not a wrapped one + } + g := LatLonToGrid(la, lo) + if !seen[g] { + seen[g] = true + out = append(out, g) + } + } + } + return out +} diff --git a/internal/geo/neighbour_test.go b/internal/geo/neighbour_test.go new file mode 100644 index 0000000..69430a7 --- /dev/null +++ b/internal/geo/neighbour_test.go @@ -0,0 +1,49 @@ +package geo + +import "testing" + +// The squares the PSK Reporter feed is filtered on. A wrong ring means either +// receiving the world again or hearing nothing. +func TestNeighbourGrids(t *testing.T) { + // JN36 is around 46.5N 5.5E. + lat, lon, ok := GridToLatLon("JN36") + if !ok { + t.Fatal("JN36 did not resolve") + } + if got := LatLonToGrid(lat, lon); got != "JN36" { + t.Fatalf("round trip gave %q, want JN36", got) + } + + ring := NeighbourGrids(lat, lon, 1) + if len(ring) != 9 { + t.Errorf("ring 1 has %d squares, want 9: %v", len(ring), ring) + } + found := false + for _, g := range ring { + if g == "JN36" { + found = true + } + if len(g) != 4 { + t.Errorf("not a 4-character square: %q", g) + } + } + if !found { + t.Errorf("the operator's own square is missing from %v", ring) + } + if n := len(NeighbourGrids(lat, lon, 0)); n != 1 { + t.Errorf("ring 0 has %d squares, want just the operator's", n) + } + if n := len(NeighbourGrids(lat, lon, 2)); n != 25 { + t.Errorf("ring 2 has %d squares, want 25", n) + } +} + +// Near a pole a step north has nowhere to go; it must be dropped, not wrapped +// onto a square on the far side of the world. +func TestNeighbourGridsNearThePole(t *testing.T) { + for _, g := range NeighbourGrids(89.5, 25, 1) { + if len(g) != 4 { + t.Errorf("bad square near the pole: %q", g) + } + } +} diff --git a/internal/gridcache/gridcache.go b/internal/gridcache/gridcache.go index ddc90d7..b96646d 100644 --- a/internal/gridcache/gridcache.go +++ b/internal/gridcache/gridcache.go @@ -1,27 +1,11 @@ // Package gridcache is the long-term callsign→grid store behind grid chasing. // -// A DX-cluster line never carries the DX's locator, so a spot can only show one -// if OpsLog learnt it elsewhere: a CQ decoded over the WSJT-X UDP link, or a -// PSK Reporter report. Learning it is easy; the problem is that the knowledge -// died with the process. Every restart began with an empty column that took an -// hour of listening to fill, and everything learnt yesterday was thrown away. +// Its own SQLite file, not a table in the settings database: that one sits +// wherever the operator put it, often a synchronised folder, and this rewrites +// itself every minute. Deleting the file costs a few days of listening. // -// So the map is persisted. Three things follow from what it is: -// -// - It is a CACHE, not user data. Deleting the file costs a few days of -// listening and nothing else, which is why it lives in its own file rather -// than in the settings database — that one sits wherever the operator chose -// to put it, often a synchronised folder, and a store that rewrites itself -// every minute has no business there. -// -// - A callsign has ONE grid. The key is unique and the newest report wins: -// operators move, go portable, go on expedition. A stale locator is worse -// than none for grid chasing, because it reads as a square already worked. -// -// - Writes are batched. The feeds repeat themselves — the same station is -// reported by dozens of receivers a minute — so the store accumulates -// changes in memory and flushes them on a timer. Nothing on the ingest path -// touches the disk. +// A callsign has one grid and the newest report wins — a stale locator reads as +// a square already worked, which is worse than none. package gridcache import ( @@ -45,11 +29,23 @@ const Retention = 2 * 365 * 24 * time.Hour // one transaction, short enough that a crash loses a minute of learning. const FlushEvery = 60 * time.Second +// Sources a locator can come from. +const ( + SourceDecode = "decode" // a CQ this station's own receiver decoded + SourceMQTT = "mqtt" // a PSK Reporter report +) + +// Entry is one locator and where it came from. +type Entry struct { + Grid string + Source string +} + type Store struct { db *sql.DB mu sync.Mutex - dirty map[string]string // call → grid, waiting to be written + dirty map[string]Entry // call → what to write stop chan struct{} stopOnce sync.Once @@ -69,15 +65,16 @@ func Open(path string, logf func(string, ...any)) (*Store, error) { if err != nil { return nil, fmt.Errorf("gridcache: open %s: %w", path, err) } - if _, err := db.Exec(`CREATE TABLE IF NOT EXISTS grids ( - call TEXT PRIMARY KEY, - grid TEXT NOT NULL, - updated_at INTEGER NOT NULL - )`); err != nil { + if _, err := db.Exec("CREATE TABLE IF NOT EXISTS grids (" + + "call TEXT PRIMARY KEY, grid TEXT NOT NULL, updated_at INTEGER NOT NULL, " + + "source TEXT NOT NULL DEFAULT '')"); err != nil { db.Close() return nil, fmt.Errorf("gridcache: schema: %w", err) } - s := &Store{db: db, dirty: map[string]string{}, stop: make(chan struct{}), logf: logf} + // A file written before the column existed keeps its rows; this fails + // harmlessly when the column is already there. + db.Exec("ALTER TABLE grids ADD COLUMN source TEXT NOT NULL DEFAULT ''") + s := &Store{db: db, dirty: map[string]Entry{}, stop: make(chan struct{}), logf: logf} if n, err := s.prune(); err != nil { s.logf("gridcache: prune failed: %v", err) } else if n > 0 { @@ -118,14 +115,14 @@ func (s *Store) LoadAll() (map[string]string, error) { // Put queues a locator for writing. Callers pass only what CHANGED — an // unchanged report is the common case by a wide margin and must not reach here, // or the batch would carry the whole feed instead of the news in it. -func (s *Store) Put(call, grid string) { +func (s *Store) Put(call, grid, source string) { call = strings.ToUpper(strings.TrimSpace(call)) grid = strings.TrimSpace(grid) if call == "" || grid == "" { return } s.mu.Lock() - s.dirty[call] = grid + s.dirty[call] = Entry{Grid: grid, Source: source} s.mu.Unlock() } @@ -160,7 +157,7 @@ func (s *Store) Flush() error { return nil } batch := s.dirty - s.dirty = map[string]string{} + s.dirty = map[string]Entry{} s.mu.Unlock() tx, err := s.db.Begin() @@ -168,8 +165,9 @@ func (s *Store) Flush() error { s.requeue(batch) return err } - st, err := tx.Prepare(`INSERT INTO grids (call, grid, updated_at) VALUES (?, ?, ?) - ON CONFLICT(call) DO UPDATE SET grid = excluded.grid, updated_at = excluded.updated_at`) + st, err := tx.Prepare(`INSERT INTO grids (call, grid, updated_at, source) VALUES (?, ?, ?, ?) + ON CONFLICT(call) DO UPDATE SET grid = excluded.grid, updated_at = excluded.updated_at, + source = excluded.source`) if err != nil { tx.Rollback() s.requeue(batch) @@ -177,8 +175,8 @@ func (s *Store) Flush() error { } defer st.Close() now := time.Now().Unix() - for call, grid := range batch { - if _, err := st.Exec(call, grid, now); err != nil { + for call, e := range batch { + if _, err := st.Exec(call, e.Grid, now, e.Source); err != nil { tx.Rollback() s.requeue(batch) return err @@ -193,17 +191,17 @@ func (s *Store) Flush() error { // requeue puts a failed batch back, without overwriting anything learnt while it // was in flight — the newer value is the right one. -func (s *Store) requeue(batch map[string]string) { +func (s *Store) requeue(batch map[string]Entry) { s.mu.Lock() defer s.mu.Unlock() - for call, grid := range batch { + for call, e := range batch { if _, newer := s.dirty[call]; !newer { - s.dirty[call] = grid + s.dirty[call] = e } } } -// Pending reports how many locators are waiting to be written (for diagnostics). +// Pending reports how many locators are waiting to be written. func (s *Store) Pending() int { s.mu.Lock() defer s.mu.Unlock() diff --git a/internal/gridcache/gridcache_test.go b/internal/gridcache/gridcache_test.go index 1a000f3..c71d9f6 100644 --- a/internal/gridcache/gridcache_test.go +++ b/internal/gridcache/gridcache_test.go @@ -21,8 +21,8 @@ func open(t *testing.T) (*Store, string) { // of after an hour of listening. func TestSurvivesRestart(t *testing.T) { s, path := open(t) - s.Put("F4BPO", "JN36") - s.Put("OH5CX", "KP30") + s.Put("F4BPO", "JN36", SourceDecode) + s.Put("OH5CX", "KP30", SourceDecode) if err := s.Flush(); err != nil { t.Fatalf("flush: %v", err) } @@ -51,11 +51,11 @@ func TestNewestReportWins(t *testing.T) { s, _ := open(t) defer s.Close() - s.Put("F4BPO", "JN36") + s.Put("F4BPO", "JN36", SourceDecode) if err := s.Flush(); err != nil { t.Fatal(err) } - s.Put("F4BPO", "KP30") // moved + s.Put("F4BPO", "KP30", SourceDecode) // moved if err := s.Flush(); err != nil { t.Fatal(err) } @@ -79,7 +79,7 @@ func TestBatching(t *testing.T) { defer s.Close() for _, c := range []string{"A1AA", "B2BB", "C3CC"} { - s.Put(c, "JN36") + s.Put(c, "JN36", SourceDecode) } if n := s.Pending(); n != 3 { t.Errorf("pending = %d, want 3 queued and unwritten", n) @@ -106,8 +106,8 @@ func TestBatching(t *testing.T) { // square is the one way this cache can be actively wrong rather than empty. func TestPruneOnOpen(t *testing.T) { s, path := open(t) - s.Put("FRESH", "JN36") - s.Put("STALE", "IO91") + s.Put("FRESH", "JN36", SourceDecode) + s.Put("STALE", "IO91", SourceDecode) if err := s.Flush(); err != nil { t.Fatal(err) } @@ -137,7 +137,7 @@ func TestPruneOnOpen(t *testing.T) { // trade. func TestCloseFlushes(t *testing.T) { s, path := open(t) - s.Put("LATE", "JN36") + s.Put("LATE", "JN36", SourceDecode) if err := s.Close(); err != nil { t.Fatalf("close: %v", err) } diff --git a/internal/pskr/pskr.go b/internal/pskr/pskr.go index 6bba928..e7b0f69 100644 --- a/internal/pskr/pskr.go +++ b/internal/pskr/pskr.go @@ -81,7 +81,16 @@ type Config struct { // OnSpot receives every accepted decode. Called from the MQTT goroutine, so // it must not block: the broker's buffer is what pays for it if it does. OnSpot func(Spot) - Logf func(string, ...any) + // OnGrid receives the transmitter of EVERY message, before any geographic + // filtering, for the callsign-to-locator store. Same goroutine as OnSpot and + // the same rule: do not block. + OnGrid func(call, grid string) + // RxGrids filters at the BROKER: only reports collected by a receiver in one + // of these squares are sent at all. Empty keeps the old behaviour, which was + // to receive the world and discard it here — measured at 83 messages a second + // for the four opening bands, of which about one in a hundred survived. + RxGrids []string + Logf func(string, ...any) } // Watcher owns the MQTT connection and its subscriptions. @@ -115,6 +124,32 @@ func New(cfg Config) *Watcher { return &Watcher{cfg: cfg} } +// topics builds the subscription list. +// +// The v2 topic is +// +// pskr/filter/v2//////// +// +// so the receiver's square is a level the broker can filter on, and a band of +// "+" means every band. Filtering by RECEIVER square rather than by receiver +// DXCC is deliberate: measured on 20 m, one country ranged from 1.2 messages a +// second (OH) to 72.5 (K), because a DXCC can be a continent. By square the +// same measurement is 0.2 to 1.2 — the load follows distance, which is what the +// feed is actually about, and it is the same for every operator. +func (w *Watcher) topics() []string { + out := []string{} + for _, b := range w.cfg.Bands { + if len(w.cfg.RxGrids) == 0 { + out = append(out, "pskr/filter/v2/"+b+"/#") + continue + } + for _, g := range w.cfg.RxGrids { + out = append(out, "pskr/filter/v2/"+b+"/+/+/+/+/"+strings.ToUpper(g)+"/+/+") + } + } + return out +} + // Start connects and subscribes. Safe to call when already running. func (w *Watcher) Start() error { w.mu.Lock() @@ -143,10 +178,7 @@ func (w *Watcher) Start() error { opts.OnConnect = func(c mqtt.Client) { w.cfg.Logf("pskr: connected to %s", w.cfg.Broker) - for _, b := range w.cfg.Bands { - // Every mode, every pair of stations, on this band. That firehose IS - // the point: the detector's job is to find the shape in it. - topic := "pskr/filter/v2/" + b + "/#" + for _, topic := range w.topics() { if tok := c.Subscribe(topic, 0, w.handle); tok.Wait() && tok.Error() != nil { w.cfg.Logf("pskr: subscribe %s failed: %v", topic, tok.Error()) continue @@ -206,6 +238,13 @@ func (w *Watcher) handle(_ mqtt.Client, m mqtt.Message) { return } + // The locator store takes every transmitter, before any of the geography + // below. What it wants is "which square is this callsign in", and that is + // true whoever happened to hear the report. + if w.cfg.OnGrid != nil { + w.cfg.OnGrid(call, grid[:4]) + } + // THE RECEIVER HAS TO BE NEAR THE OPERATOR. This is the whole difference // between a useful feed and a world map. // diff --git a/internal/pskr/topics_test.go b/internal/pskr/topics_test.go new file mode 100644 index 0000000..baaae60 --- /dev/null +++ b/internal/pskr/topics_test.go @@ -0,0 +1,46 @@ +package pskr + +import ( + "strings" + "testing" +) + +// The receiver square is a level the BROKER can filter on, which is the whole +// point: measured on the live feed, the four opening bands unfiltered are 83 +// messages a second of which about one in a hundred survives the NearKm test. +// One ring of squares is under two a second, and the same for every operator — +// where filtering by DXCC ranged from 1.2 (OH) to 72.5 (K) on one band. +func TestTopicsFilterOnTheReceiverSquare(t *testing.T) { + w := New(Config{Bands: []string{"6m"}, RxGrids: []string{"jn36", "JN37"}}) + got := w.topics() + if len(got) != 2 { + t.Fatalf("want one subscription per band × square, got %v", got) + } + // Level order: band/mode/txcall/rxcall/txgrid/RXGRID/txdxcc/rxdxcc + want := "pskr/filter/v2/6m/+/+/+/+/JN36/+/+" + if got[0] != want { + t.Errorf("topic = %q, want %q", got[0], want) + } + if !strings.Contains(got[1], "/JN37/") { + t.Errorf("square not upper-cased into the topic: %q", got[1]) + } +} + +// With no squares the old behaviour stands: receive the band and decide here. +func TestTopicsWithoutSquares(t *testing.T) { + w := New(Config{Bands: []string{"10m", "2m"}}) + got := w.topics() + if len(got) != 2 || got[0] != "pskr/filter/v2/10m/#" { + t.Errorf("unfiltered topics = %v", got) + } +} + +// Grid chasing wants every band. "+" is the MQTT single-level wildcard, so one +// subscription per square covers the lot instead of one per band per square. +func TestTopicsAllBands(t *testing.T) { + w := New(Config{Bands: []string{"+"}, RxGrids: []string{"JN36"}}) + got := w.topics() + if len(got) != 1 || got[0] != "pskr/filter/v2/+/+/+/+/+/JN36/+/+" { + t.Errorf("all-band topic = %v", got) + } +}