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.
154 lines
4.0 KiB
Go
154 lines
4.0 KiB
Go
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")
|
|
}
|
|
}
|