Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
7d33379fe1 | ||
|
|
82a49150d2 | ||
|
|
1ef9e553f7 | ||
|
|
39fab68bd4 |
@@ -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
@@ -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.
|
||||
|
||||
@@ -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": "",
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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. */}
|
||||
|
||||
@@ -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.',
|
||||
|
||||
Vendored
+6
@@ -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>;
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
@@ -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
@@ -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.
|
||||
//
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user