Compare commits

...
4 Commits
Author SHA1 Message Date
rouggy 7d33379fe1 feat(cluster): feed the locator store from PSK Reporter, filtered at the broker
The store shipped without its main source. Locators came only from this
station's own WSJT-X decodes, which is what the whole MQTT discussion was
about.

PSK Reporter now feeds it through a new OnGrid callback, fired before any
geographic filtering: what the store wants is "which square is this callsign
in", and that is true whoever happened to hear the report.

The subscription filters on the RECEIVER's square, a level the v2 topic
exposes. 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
that already existed here — the rest was received, TLS-decrypted, JSON-parsed
and discarded. One ring of squares is 0.2 to 1.2 a second.

By square rather than by DXCC, which was the obvious alternative: one country
measured 1.2 messages a second (OH) against 72.5 (K) on a single band, because
a DXCC can be a continent. By square the same measurement is 0.2 to 1.2, so the
load follows distance — what the feed is actually about — and is the same for
every operator.

Grid chasing subscribes with the "+" band wildcard, so one subscription per
square covers every band instead of one per band per square.

The store gains a source column (decode | mqtt), migrated in place on an
existing file.
2026-08-12 17:16:47 +02:00
rouggy 82a49150d2 feat(cluster): persist learnt locators behind a "Chase new grids" option
The callsign->grid map died with the process. Every restart began with an empty
Locator column that took an hour of listening to refill, and everything learnt
the day before was thrown away.

internal/gridcache is its own SQLite file in the data directory, not a table 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. Deleting the file costs a few days of listening and nothing
else.

Callsign is the primary key and the newest report wins — operators move, go
portable, go on expedition, and a stale locator is worse than none for grid
chasing because it reads as a square already worked. Entries not seen in two
years are pruned at open: that is the one way this cache can be actively wrong
rather than merely empty, and it is what bounds a store that would otherwise
only grow.

Only CHANGES are queued. The feeds repeat themselves, so writing every report
would put the whole stream in the batch instead of the news in it. Batches
flush on a timer and on shutdown — a restart is exactly when the cache is worth
the most.

Rotation is switched off while the store is attached: the cap would discard
callsigns the database still holds and the lookup would then miss what we know.
Age in the store is the bound instead.
2026-08-12 17:01:05 +02:00
rouggy 1ef9e553f7 perf(cluster): rotate the grid cache instead of emptying it, cap 100k
The cache dropped EVERY entry once it passed 20 000. That was survivable while
the only feed was this station's own WSJT-X decodes, which never reached the
ceiling — it is a cliff for anything larger, and every locator in the cluster
list would disappear at once, periodically, for no reason the operator could
see.

Two generations: when the current map fills it becomes the previous one and a
fresh map takes over; lookups consult both. A rotation therefore costs the
older half and nothing more. It needs no insertion order, no per-entry
timestamp and no bookkeeping on the write path — all of which "evict the
oldest thousand" would require, on a path that runs once per decode.

The cap is 100 000, measured at 82 bytes an entry: 8 MB a generation, 16 MB
for both. 20 000 was chosen when the ceiling was unreachable anyway.

Both maps now go through rememberDecodeGrid / lookupDecodeGrid. Rotation swaps
the map headers, so the read had to be guarded; routing every access through
one accessor is what makes that checkable rather than remembered.
2026-08-12 16:48:23 +02:00
rouggy 39fab68bd4 chore: open 0.24.8 2026-08-12 14:46:05 +02:00
15 changed files with 1017 additions and 31 deletions
+181 -18
View File
@@ -41,6 +41,7 @@ import (
"hamlog/internal/email"
"hamlog/internal/extsvc"
"hamlog/internal/geo"
"hamlog/internal/gridcache"
"hamlog/internal/integrations/udp"
"hamlog/internal/lookup"
"hamlog/internal/lotwusers"
@@ -235,6 +236,7 @@ const (
keyUltrabeamStep = "ultrabeam.step_khz" // re-tune hysteresis: 25 | 50 | 100 kHz
keyMotorTrackMode = "motor.track_mode" // "always" | "step" | "band"
keyMotorBandFreqs = "motor.band_freqs" // per-band tune frequency: "40m=7100,20m=14150"
keyChaseNewGrids = "cluster.chase_grids" // "1" → persist learnt locators across restarts
keyMotorType = "ultrabeam.type" // "ultrabeam" | "steppir" (default ultrabeam)
keyMotorTransport = "ultrabeam.transport" // "tcp" | "serial" (default tcp)
keyMotorCOM = "ultrabeam.com" // serial device name (COM3, /dev/ttyUSB0)
@@ -611,8 +613,23 @@ type App struct {
//
// In memory only, and bounded: it is a session-local view of who is on the
// air now, not a database.
decodeGrids map[string]string
//
// Bounded in TWO generations. The cap used to drop the whole map, which was
// survivable while the only source was this station's own decodes — it never
// reached the ceiling. It is a cliff for any larger feed: every locator in the
// list would vanish at once, periodically. Keeping the previous generation
// means a rotation costs the older half and nothing more, and it needs no
// insertion order, no per-entry timestamp and no bookkeeping on the write
// path — which "evict the oldest thousand" would all require.
decodeGrids map[string]string // current generation, written to
decodeGridsOld map[string]string // previous generation, still readable
decodeGridsMu sync.RWMutex
// gridStore persists the map across restarts when grid chasing is on. With
// it, rotation is switched OFF: rotating would drop callsigns the database
// still holds, and a lookup would then miss something we know. The store
// bounds itself by age instead, so memory follows how many distinct stations
// have actually been heard in two years rather than a made-up ceiling.
gridStore *gridcache.Store
// pskr is the PSK Reporter MQTT feed, up only while the opening watch is on.
// It is the source that makes VHF detection work at all: the cluster and RBN
// carry a handful of 6 m spots where PSK Reporter carries hundreds.
@@ -1375,8 +1392,11 @@ 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.
// 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()
@@ -1625,6 +1645,14 @@ func (a *App) shutdown(ctx context.Context) {
if a.qsoRec != nil {
a.qsoRec.Stop()
}
// Before the databases: Close flushes what the last minute learnt, and a
// restart is exactly when the grid cache is worth the most.
if a.gridStore != nil {
if err := a.gridStore.Close(); err != nil {
applog.Printf("gridcache: close: %v", err)
}
a.gridStore = nil
}
if a.logDb != nil && a.logDb != a.db {
_ = a.logDb.Close() // shared MySQL logbook (separate from the local config DB)
}
@@ -11937,15 +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.decodeGridsMu.Lock()
if a.decodeGrids == nil {
a.decodeGrids = make(map[string]string, 512)
}
if len(a.decodeGrids) > 20000 {
a.decodeGrids = make(map[string]string, 512) // bound a long session
}
a.decodeGrids[strings.ToUpper(ev.DecodeCall)] = ev.DecodeGrid
a.decodeGridsMu.Unlock()
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
@@ -17014,6 +17034,151 @@ func (a *App) clusterStatusMaps() *clusterStatusCache {
// was ambiguous and the frontend couldn't infer) we degrade gracefully
// to band-only — saying "worked" rather than wrongly flagging "new-slot"
// just because we don't know the mode.
// decodeGridsCap is how many callsigns one generation holds before it rotates.
//
// 100 000 measured at 82 bytes an entry — 8 MB a generation, 16 MB for both,
// which is nothing next to what it buys: a locator on a spot instead of an empty
// column. The old 20 000 was chosen when the only feed was this station's own
// decodes and the ceiling was never reached anyway.
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
}
// startGridCache brings the locator store up or down to match the setting.
//
// Switching it ON seeds the in-memory map from the database, which is the whole
// point: the cluster's locator column is populated in the first second instead
// of after an hour of listening. Switching it OFF closes the file and leaves the
// map alone — what has already been learnt this session stays usable, it simply
// stops being remembered. Nothing is deleted: turning the option back on picks
// up where it left off.
func (a *App) startGridCache() {
if a.gridStore != nil {
if err := a.gridStore.Close(); err != nil {
applog.Printf("gridcache: close: %v", err)
}
a.gridStore = nil
}
if !a.GetChaseNewGrids() {
return
}
path := filepath.Join(a.dataDir, "grids.db")
st, err := gridcache.Open(path, applog.Printf)
if err != nil {
applog.Printf("gridcache: disabled — %v", err)
return
}
seed, err := st.LoadAll()
if err != nil {
applog.Printf("gridcache: cannot read %s (%v) — starting empty", path, err)
seed = map[string]string{}
}
a.decodeGridsMu.Lock()
if a.decodeGrids == nil {
a.decodeGrids = make(map[string]string, len(seed)+512)
}
// Seed UNDER what this session already heard: a locator decoded a minute ago
// is newer than one stored days back, and the newest report is the one that
// counts when a station has moved.
for call, grid := range seed {
if _, live := a.decodeGrids[call]; !live {
a.decodeGrids[call] = grid
}
}
n := len(a.decodeGrids)
a.decodeGridsMu.Unlock()
a.gridStore = st
st.Start(a.ctx)
applog.Printf("gridcache: grid chasing on — %d locators loaded from %s (%d known in total)", len(seed), path, n)
}
// rememberDecodeGrid records the grid a station announced.
//
// Rotation, not eviction: when the current generation fills it becomes the
// 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, source string) {
call = strings.ToUpper(strings.TrimSpace(call))
grid = strings.TrimSpace(grid)
if call == "" || grid == "" {
return
}
store := a.gridStore
a.decodeGridsMu.Lock()
if a.decodeGrids == nil {
a.decodeGrids = make(map[string]string, 512)
}
// Rotate only while nothing is persisting the map. With a store the cap would
// discard callsigns the database still holds, so the lookup would miss what
// we know; age in the store is the bound instead.
if store == nil && len(a.decodeGrids) >= decodeGridsCap {
a.decodeGridsOld = a.decodeGrids
a.decodeGrids = make(map[string]string, 512)
}
// Only a CHANGE is worth writing. The feeds repeat themselves — the same
// station is reported by dozens of receivers a minute — so queueing every
// report would put the whole stream in the batch instead of the news in it.
changed := a.decodeGrids[call] != grid
if changed && a.decodeGridsOld != nil && a.decodeGridsOld[call] == grid {
// Known already, just in the older generation: promote it without calling
// it news.
changed = false
}
a.decodeGrids[call] = grid
a.decodeGridsMu.Unlock()
if changed && store != nil {
store.Put(call, grid, source)
}
}
// lookupDecodeGrid returns the last grid heard for a callsign, "" if unknown.
// The previous generation is consulted second, so a station that has gone quiet
// keeps its locator across one rotation instead of losing it at the cliff.
func (a *App) lookupDecodeGrid(call string) string {
call = strings.ToUpper(call)
a.decodeGridsMu.RLock()
defer a.decodeGridsMu.RUnlock()
if g := a.decodeGrids[call]; g != "" {
return g
}
return a.decodeGridsOld[call]
}
func (a *App) ClusterSpotStatuses(spots []SpotQuery) []SpotStatus {
out := make([]SpotStatus, len(spots))
if a.qso == nil {
@@ -17102,12 +17267,10 @@ func (a *App) ClusterSpotStatuses(spots []SpotQuery) []SpotStatus {
// is part of the key, so the "group digital modes" option decides whether a
// grid worked on FT8 still counts as new on FT4 — one rule, no branch here.
{
// The length check has to be INSIDE the lock: the decode goroutine
// replaces this map wholesale when it grows too large, so reading len()
// unguarded is a race on the map header, not a cheap fast path.
a.decodeGridsMu.RLock()
g := a.decodeGrids[strings.ToUpper(q.Call)]
a.decodeGridsMu.RUnlock()
// Read through the accessor: the decode goroutine swaps these maps on
// rotation, so touching them unguarded is a race on the map header,
// not a cheap fast path.
g := a.lookupDecodeGrid(q.Call)
if g != "" {
out[i].Grid = g
cm := out[i].Mode
+38 -4
View File
@@ -20,6 +20,8 @@ import (
"hamlog/internal/applog"
"hamlog/internal/bandopen"
"hamlog/internal/geo"
"hamlog/internal/gridcache"
"hamlog/internal/cluster"
"hamlog/internal/pskr"
)
@@ -133,17 +135,45 @@ 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,
Bands: bands,
RxGrids: rxGrids,
OpLat: a.opLat, OpLon: a.opLon,
Geo: func(grid string) (int, int, bool) {
lat, lon, ok := gridToLatLon(grid)
@@ -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.
+14
View File
@@ -1,4 +1,18 @@
[
{
"version": "0.24.8",
"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 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 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."
]
},
{
"version": "0.24.7",
"date": "",
+160
View File
@@ -0,0 +1,160 @@
package main
import (
"fmt"
"path/filepath"
"sync"
"testing"
"hamlog/internal/gridcache"
)
// The grid cache used to drop EVERYTHING at its ceiling. That was survivable
// while the only feed was this station's own decodes, which never reached it.
// Fed by anything larger it is a cliff: every locator in the cluster list
// disappears at once, periodically, and the operator sees the column empty
// itself for no reason.
//
// Rotation keeps the previous generation, so a full cache costs the older half
// and nothing more.
func TestDecodeGridRotationKeepsThePreviousGeneration(t *testing.T) {
a := &App{}
// Fill exactly one generation.
for i := 0; i < decodeGridsCap; i++ {
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)
}
if a.decodeGridsOld != nil {
t.Fatal("rotated early — the cap is the ceiling of ONE generation")
}
// One more entry rotates.
a.rememberDecodeGrid("NEWCALL", "IO91", gridcache.SourceDecode)
if a.decodeGridsOld == nil {
t.Fatal("did not rotate at the cap")
}
if got := a.lookupDecodeGrid("NEWCALL"); got != "IO91" {
t.Errorf("the entry that caused the rotation was lost: %q", got)
}
// The whole previous generation is still readable — this is the point.
if got := a.lookupDecodeGrid("CALL000000"); got != "JN36" {
t.Errorf("a locator from the previous generation was dropped: %q — that is the cliff again", got)
}
if got := a.lookupDecodeGrid("CALL099999"); got != "JN36" {
t.Errorf("previous generation incomplete: %q", got)
}
// A second rotation is what finally retires the oldest half.
for i := 0; i < decodeGridsCap; i++ {
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)
}
if got := a.lookupDecodeGrid("NEWCALL"); got != "IO91" {
t.Errorf("an entry one generation old was retired too early: %q", got)
}
}
// Callsigns are normalised on the way in AND on the way out, or a spot for
// "f4bpo" would miss a grid learnt as "F4BPO".
func TestDecodeGridCaseAndBlanks(t *testing.T) {
a := &App{}
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", gridcache.SourceDecode)
a.rememberDecodeGrid("K1ABC", "", gridcache.SourceDecode)
if got := a.lookupDecodeGrid("K1ABC"); got != "" {
t.Errorf("stored an empty grid: %q", got)
}
}
// The write path runs on the decode goroutine while the cluster status builder
// reads. Rotation swaps the map headers, so an unguarded read is a data race
// rather than a stale value.
//
// This build is CGO-free, so -race is not available here; the test exercises the
// interleaving and would fault on a concurrent map access even without it.
func TestDecodeGridConcurrentAccess(t *testing.T) {
a := &App{}
var wg sync.WaitGroup
wg.Add(2)
go func() {
defer wg.Done()
for i := 0; i < 20000; i++ {
a.rememberDecodeGrid(fmt.Sprintf("W%05d", i), "FN31", gridcache.SourceDecode)
}
}()
go func() {
defer wg.Done()
for i := 0; i < 20000; i++ {
_ = a.lookupDecodeGrid(fmt.Sprintf("W%05d", i))
}
}()
wg.Wait()
}
// Only a CHANGE may reach the write batch.
//
// The feeds repeat themselves — the same station is reported by dozens of
// receivers a minute — so queueing every report would put the whole stream in
// the batch instead of the news in it, and turn a cache into a write amplifier.
func TestOnlyChangesAreQueued(t *testing.T) {
st, err := gridcache.Open(filepath.Join(t.TempDir(), "grids.db"), nil)
if err != nil {
t.Fatal(err)
}
defer st.Close()
a := &App{gridStore: st}
a.rememberDecodeGrid("F4BPO", "JN36", gridcache.SourceDecode)
if n := st.Pending(); n != 1 {
t.Fatalf("a new locator queued %d writes, want 1", n)
}
if err := st.Flush(); err != nil {
t.Fatal(err)
}
// The same report, a hundred times over, is not news.
for i := 0; i < 100; i++ {
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", gridcache.SourceDecode)
if n := st.Pending(); n != 1 {
t.Errorf("a changed locator queued %d writes, want 1", n)
}
if got := a.lookupDecodeGrid("F4BPO"); got != "KP30" {
t.Errorf("the map kept the old locator: %q", got)
}
}
// Rotation must be OFF while a store is attached: it would discard callsigns the
// database still holds, and the lookup would then miss something we know.
func TestNoRotationWhilePersisting(t *testing.T) {
st, err := gridcache.Open(filepath.Join(t.TempDir(), "grids.db"), nil)
if err != nil {
t.Fatal(err)
}
defer st.Close()
a := &App{gridStore: st}
for i := 0; i < decodeGridsCap+10; i++ {
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")
}
if got := a.lookupDecodeGrid("CALL000000"); got != "JN36" {
t.Errorf("the first entry was dropped: %q", got)
}
}
+22 -1
View File
@@ -52,7 +52,7 @@ import {
GetADIFMonitor, SaveADIFMonitor, PickADIFMonitorFile,
GetRelayAuto, SaveRelayAuto, GetStationDevices,
GetAwardDefs, GetTrackedAwards, SaveTrackedAwards,
GetBandOpenSettings, SaveBandOpenSettings, GetPSKReporterStatus,
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';
@@ -1552,6 +1552,8 @@ export function SettingsModal({ onClose, onSaved, initialSection, onMainPaneChan
// has side effects there — adding the RBN nodes, bringing the PSK Reporter
// feed up or down — so the write has to go where those live.
const [bandOpen, setBandOpen] = useState<any>({ enabled: false, bands: [], available: [] });
const [chaseGrids, setChaseGrids] = useState(false);
const [gridStat, setGridStat] = useState<any>(null);
const [pskrStatus, setPskrStatus] = useState<any>(null);
const saveBandOpen = async (next: any) => {
setBandOpen(next);
@@ -1560,11 +1562,13 @@ export function SettingsModal({ onClose, onSaved, initialSection, onMainPaneChan
useEffect(() => {
(async () => {
try { setBandOpen(await GetBandOpenSettings()); } catch { /* defaults stand */ }
try { setChaseGrids(await GetChaseNewGrids()); } catch { /* defaults stand */ }
})();
// Poll the feed while the panel is open: a live count is the only thing that
// 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);
}, []);
@@ -4227,6 +4231,23 @@ export function SettingsModal({ onClose, onSaved, initialSection, onMainPaneChan
things set up once. A preferences dialog you reopen every ten minutes
is a filter in the wrong place. */}
{/* Grid chasing. Here rather than in the filter panel because it is set
up once: it decides whether locators learnt from decodes are KEPT
across restarts, not what the list shows right now. */}
<div className="border-t border-border/60 pt-3 space-y-2">
<label className="flex items-start gap-2 text-sm cursor-pointer">
<Checkbox checked={chaseGrids} className="mt-0.5"
onCheckedChange={(c) => { setChaseGrids(!!c); SetChaseNewGrids(!!c).catch(() => {}); }} />
<span>{t('clu.chaseGrids')} <span className="text-xs text-muted-foreground">{t('clu.chaseGridsHint')}</span></span>
</label>
{chaseGrids && (
<p className="pl-6 text-xs text-muted-foreground">
{t('clu.chaseGridsStat', { n: gridStat?.known ?? 0, p: gridStat?.pending ?? 0 })}
{pskrStatus?.running ? ` · ${t('bo.feedUp', { n: pskrStatus.received ?? 0 })}` : ''}
</p>
)}
</div>
{/* Band-opening watch. It lives HERE, with the cluster nodes, because
switching it on adds two of them the operator should see that
happen where it happens rather than find nodes they did not add. */}
+1 -1
View File
@@ -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.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-<date>.adi — so your contacts are protected even though they live on the server. Rotation keeps the last N of each.',
+6
View File
@@ -404,6 +404,8 @@ export function GetCatalogCodes():Promise<Array<string>>;
export function GetChangelog():Promise<Array<main.ChangelogEntry>>;
export function GetChaseNewGrids():Promise<boolean>;
export function GetChatHistory(arg1:number):Promise<Array<main.ChatMessage>>;
export function GetClublogCtyInfo():Promise<main.ClublogCtyInfo>;
@@ -440,6 +442,8 @@ export function GetFlexBandPower():Promise<Record<string, main.FlexBandPower>>;
export function GetFlexState():Promise<cat.FlexTXState>;
export function GetGridCacheStatus():Promise<main.GridCacheStatus>;
export function GetIcomState():Promise<cat.IcomTXState>;
export function GetListsSettings():Promise<main.ListsSettings>;
@@ -970,6 +974,8 @@ export function SetCIVTrace(arg1:boolean):Promise<void>;
export function SetCWDecoderPitch(arg1:number):Promise<void>;
export function SetChaseNewGrids(arg1:boolean):Promise<void>;
export function SetClublogCtyEnabled(arg1:boolean):Promise<void>;
export function SetClublogMostWantedEnabled(arg1:boolean):Promise<void>;
+12
View File
@@ -750,6 +750,10 @@ export function GetChangelog() {
return window['go']['main']['App']['GetChangelog']();
}
export function GetChaseNewGrids() {
return window['go']['main']['App']['GetChaseNewGrids']();
}
export function GetChatHistory(arg1) {
return window['go']['main']['App']['GetChatHistory'](arg1);
}
@@ -822,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']();
}
@@ -1882,6 +1890,10 @@ export function SetCWDecoderPitch(arg1) {
return window['go']['main']['App']['SetCWDecoderPitch'](arg1);
}
export function SetChaseNewGrids(arg1) {
return window['go']['main']['App']['SetChaseNewGrids'](arg1);
}
export function SetClublogCtyEnabled(arg1) {
return window['go']['main']['App']['SetClublogCtyEnabled'](arg1);
}
+16
View File
@@ -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;
+51
View File
@@ -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
}
+49
View File
@@ -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)
}
}
}
+222
View File
@@ -0,0 +1,222 @@
// Package gridcache is the long-term callsign→grid store behind grid chasing.
//
// 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.
//
// 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 (
"context"
"database/sql"
"fmt"
"strings"
"sync"
"time"
_ "modernc.org/sqlite"
)
// Retention bounds a store that would otherwise only ever grow. Two years is
// chosen to be far longer than any propagation interest and short enough that a
// reassigned callsign eventually stops carrying its previous holder's square —
// the one way this cache can be actively wrong rather than merely empty.
const Retention = 2 * 365 * 24 * time.Hour
// FlushEvery is the batch interval. Long enough that a burst of reports costs
// 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]Entry // call → what to write
stop chan struct{}
stopOnce sync.Once
wg sync.WaitGroup
logf func(string, ...any)
}
// Open creates or opens the store at path and prunes what has aged out.
func Open(path string, logf func(string, ...any)) (*Store, error) {
if logf == nil {
logf = func(string, ...any) {}
}
// WAL so a flush never blocks a read, and a busy timeout because the flush
// goroutine and the startup load can overlap on a slow disk.
db, err := sql.Open("sqlite", path+"?_pragma=journal_mode(WAL)&_pragma=busy_timeout(5000)")
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, " +
"source TEXT NOT NULL DEFAULT '')"); err != nil {
db.Close()
return nil, fmt.Errorf("gridcache: schema: %w", err)
}
// 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 {
s.logf("gridcache: pruned %d locators not heard in %d days", n, int(Retention.Hours()/24))
}
return s, nil
}
func (s *Store) prune() (int64, error) {
cut := time.Now().Add(-Retention).Unix()
res, err := s.db.Exec(`DELETE FROM grids WHERE updated_at < ?`, cut)
if err != nil {
return 0, err
}
return res.RowsAffected()
}
// LoadAll returns every stored locator, for seeding the in-memory map at
// startup. One query and one map build — the point of the whole package is that
// nothing afterwards has to ask the database anything.
func (s *Store) LoadAll() (map[string]string, error) {
rows, err := s.db.Query(`SELECT call, grid FROM grids`)
if err != nil {
return nil, fmt.Errorf("gridcache: load: %w", err)
}
defer rows.Close()
out := make(map[string]string, 4096)
for rows.Next() {
var call, grid string
if err := rows.Scan(&call, &grid); err != nil {
return nil, err
}
out[call] = grid
}
return out, rows.Err()
}
// 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, source string) {
call = strings.ToUpper(strings.TrimSpace(call))
grid = strings.TrimSpace(grid)
if call == "" || grid == "" {
return
}
s.mu.Lock()
s.dirty[call] = Entry{Grid: grid, Source: source}
s.mu.Unlock()
}
// Start runs the flush loop until Close.
func (s *Store) Start(ctx context.Context) {
s.wg.Add(1)
go func() {
defer s.wg.Done()
t := time.NewTicker(FlushEvery)
defer t.Stop()
for {
select {
case <-s.stop:
return
case <-ctx.Done():
return
case <-t.C:
if err := s.Flush(); err != nil {
s.logf("gridcache: flush failed: %v", err)
}
}
}
}()
}
// Flush writes the pending batch in one transaction. Safe to call with nothing
// pending, which is most of the time on a quiet band.
func (s *Store) Flush() error {
s.mu.Lock()
if len(s.dirty) == 0 {
s.mu.Unlock()
return nil
}
batch := s.dirty
s.dirty = map[string]Entry{}
s.mu.Unlock()
tx, err := s.db.Begin()
if err != nil {
s.requeue(batch)
return err
}
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)
return err
}
defer st.Close()
now := time.Now().Unix()
for call, e := range batch {
if _, err := st.Exec(call, e.Grid, now, e.Source); err != nil {
tx.Rollback()
s.requeue(batch)
return err
}
}
if err := tx.Commit(); err != nil {
s.requeue(batch)
return err
}
return nil
}
// 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]Entry) {
s.mu.Lock()
defer s.mu.Unlock()
for call, e := range batch {
if _, newer := s.dirty[call]; !newer {
s.dirty[call] = e
}
}
}
// Pending reports how many locators are waiting to be written.
func (s *Store) Pending() int {
s.mu.Lock()
defer s.mu.Unlock()
return len(s.dirty)
}
// Close stops the loop and writes whatever is pending. A restart is the moment
// the cache is most valuable, so losing the last minute of learning to a clean
// shutdown would be a poor trade.
func (s *Store) Close() error {
s.stopOnce.Do(func() { close(s.stop) })
s.wg.Wait()
err := s.Flush()
if cerr := s.db.Close(); err == nil {
err = cerr
}
return err
}
+153
View File
@@ -0,0 +1,153 @@
package gridcache
import (
"path/filepath"
"testing"
"time"
)
func open(t *testing.T) (*Store, string) {
t.Helper()
path := filepath.Join(t.TempDir(), "grids.db")
s, err := Open(path, nil)
if err != nil {
t.Fatalf("open: %v", err)
}
return s, path
}
// The whole reason the store exists: what was learnt is still there after a
// restart, so the cluster's locator column is full in the first second instead
// of after an hour of listening.
func TestSurvivesRestart(t *testing.T) {
s, path := open(t)
s.Put("F4BPO", "JN36", SourceDecode)
s.Put("OH5CX", "KP30", SourceDecode)
if err := s.Flush(); err != nil {
t.Fatalf("flush: %v", err)
}
if err := s.Close(); err != nil {
t.Fatalf("close: %v", err)
}
again, err := Open(path, nil)
if err != nil {
t.Fatalf("reopen: %v", err)
}
defer again.Close()
got, err := again.LoadAll()
if err != nil {
t.Fatalf("load: %v", err)
}
if got["F4BPO"] != "JN36" || got["OH5CX"] != "KP30" {
t.Errorf("locators lost across a restart: %v", got)
}
}
// A callsign has ONE grid and the newest report wins. Operators move, go
// portable, go on expedition — and a stale locator is worse than none for grid
// chasing, because it reads as a square already worked.
func TestNewestReportWins(t *testing.T) {
s, _ := open(t)
defer s.Close()
s.Put("F4BPO", "JN36", SourceDecode)
if err := s.Flush(); err != nil {
t.Fatal(err)
}
s.Put("F4BPO", "KP30", SourceDecode) // moved
if err := s.Flush(); err != nil {
t.Fatal(err)
}
got, err := s.LoadAll()
if err != nil {
t.Fatal(err)
}
if got["F4BPO"] != "KP30" {
t.Errorf("grid = %q, want the newer KP30", got["F4BPO"])
}
if len(got) != 1 {
t.Errorf("a callsign must hold one row, got %d: %v", len(got), got)
}
}
// Nothing reaches the disk until a flush, and a flush with nothing pending is
// not an error — that is most minutes on a quiet band.
func TestBatching(t *testing.T) {
s, _ := open(t)
defer s.Close()
for _, c := range []string{"A1AA", "B2BB", "C3CC"} {
s.Put(c, "JN36", SourceDecode)
}
if n := s.Pending(); n != 3 {
t.Errorf("pending = %d, want 3 queued and unwritten", n)
}
if got, _ := s.LoadAll(); len(got) != 0 {
t.Errorf("wrote before the flush: %v", got)
}
if err := s.Flush(); err != nil {
t.Fatal(err)
}
if n := s.Pending(); n != 0 {
t.Errorf("pending = %d after a flush, want 0", n)
}
if got, _ := s.LoadAll(); len(got) != 3 {
t.Errorf("flush wrote %d rows, want 3", len(got))
}
if err := s.Flush(); err != nil {
t.Errorf("empty flush must be a no-op, got %v", err)
}
}
// Age is what bounds a store that would otherwise only grow. A callsign not
// heard in two years is likely reassigned, and carrying its previous holder's
// 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", SourceDecode)
s.Put("STALE", "IO91", SourceDecode)
if err := s.Flush(); err != nil {
t.Fatal(err)
}
// Backdate one row past the retention window.
old := time.Now().Add(-Retention - 24*time.Hour).Unix()
if _, err := s.db.Exec(`UPDATE grids SET updated_at = ? WHERE call = 'STALE'`, old); err != nil {
t.Fatal(err)
}
s.Close()
again, err := Open(path, nil)
if err != nil {
t.Fatal(err)
}
defer again.Close()
got, _ := again.LoadAll()
if _, ok := got["STALE"]; ok {
t.Error("an entry past the retention window survived — the store is unbounded")
}
if got["FRESH"] != "JN36" {
t.Error("pruning took a live entry with it")
}
}
// Close has to write what the last minute learnt: a restart is exactly when the
// cache is worth the most, so losing it to a clean shutdown would be a poor
// trade.
func TestCloseFlushes(t *testing.T) {
s, path := open(t)
s.Put("LATE", "JN36", SourceDecode)
if err := s.Close(); err != nil {
t.Fatalf("close: %v", err)
}
again, err := Open(path, nil)
if err != nil {
t.Fatal(err)
}
defer again.Close()
got, _ := again.LoadAll()
if got["LATE"] != "JN36" {
t.Error("what was pending at shutdown was dropped")
}
}
+43 -4
View File
@@ -81,6 +81,15 @@ 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)
// 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)
}
@@ -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/<band>/<mode>/<tx call>/<rx call>/<tx grid>/<rx grid>/<tx dxcc>/<rx dxcc>
//
// 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.
//
+46
View File
@@ -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)
}
}