diff --git a/cmd/server/analytics_after_startup_load_116_test.go b/cmd/server/analytics_after_startup_load_116_test.go new file mode 100644 index 000000000..71a4dd301 --- /dev/null +++ b/cmd/server/analytics_after_startup_load_116_test.go @@ -0,0 +1,501 @@ +package main + +import ( + "net/http" + "net/http/httptest" + "path/filepath" + "reflect" + "sort" + "sync" + "sync/atomic" + "testing" + "time" + + "github.com/gorilla/mux" +) + +// Issue #116: HTTP binds after the first chunk and the analytics +// recomputers start then, but LoadComplete() flips at the end of the hot +// window, before loadBackgroundChunks fills the retention window. A +// one-shot "startup load terminated" signal (StartupLoadDone) now fires on +// every RunStartupLoad exit path, every default-shape recomputer runs once +// after it, and the #1659 warm-up gate only opens on a pass that STARTED +// after it. backgroundLoadDone/Failed keep their health meaning. + +func isClosed(ch <-chan struct{}) bool { + select { + case <-ch: + return true + default: + return false + } +} + +// hotAndBackgroundDB seeds hot rows (last 30 min) and older rows (1–5 days) +// so RunStartupLoad with HotStartupHours=1 loads the older ones in the +// background fill. +func hotAndBackgroundDB(t *testing.T, hot, older int) string { + t.Helper() + dbPath := filepath.Join(t.TempDir(), "test.db") + now := time.Now().UTC() + seedTestDBRows(t, dbPath, hot+older, 2, func(i int) (string, int64) { + ago := 30 * time.Minute + if i >= hot { + ago = time.Duration(24+(i-hot)%96) * time.Hour + } + ts := now.Add(-ago) + return ts.Format(time.RFC3339), ts.Unix() + }) + return dbPath +} + +func openStartupStore(t *testing.T, dbPath string, retention, hot float64) *PacketStore { + t.Helper() + db, err := OpenDB(dbPath) + if err != nil { + t.Fatalf("OpenDB: %v", err) + } + t.Cleanup(func() { db.conn.Close() }) + return NewPacketStore(db, &PacketStoreConfig{RetentionHours: retention, HotStartupHours: hot}) +} + +func TestStartupLoadDoneFiresOnEveryPath_116(t *testing.T) { + type want struct{ done, failed bool } + cases := []struct { + name string + setup func(t *testing.T) *PacketStore + err bool + want want // backgroundLoadDone/Failed must keep their meaning + }{ + {"hot window + background fill", func(t *testing.T) *PacketStore { + return openStartupStore(t, hotAndBackgroundDB(t, 20, 30), 168, 1) + }, false, want{true, false}}, + {"hot window disabled (hotStartupHours=0)", func(t *testing.T) *PacketStore { + return openStartupStore(t, hotAndBackgroundDB(t, 5, 5), 168, 0) + }, false, want{true, false}}, + {"retention disabled", func(t *testing.T) *PacketStore { + return openStartupStore(t, hotAndBackgroundDB(t, 5, 5), 0, 1) + }, false, want{true, false}}, + {"empty database", func(t *testing.T) *PacketStore { + return openStartupStore(t, hotAndBackgroundDB(t, 0, 0), 168, 1) + }, false, want{true, false}}, + {"LoadChunked error", func(t *testing.T) *PacketStore { + s := openStartupStore(t, hotAndBackgroundDB(t, 5, 5), 168, 1) + s.db.conn.Close() + return s + }, true, want{true, true}}, + {"background fill failure", func(t *testing.T) *PacketStore { + s := openStartupStore(t, hotAndBackgroundDB(t, 5, 30), 168, 1) + s.bgLoaderEntryHook = func() { s.db.conn.Close() } // every background chunk fails + return s + }, false, want{false, true}}, + } + for _, c := range cases { + t.Run(c.name, func(t *testing.T) { + s := c.setup(t) + if isClosed(s.StartupLoadDone()) { + t.Fatal("StartupLoadDone closed before RunStartupLoad") + } + err := s.RunStartupLoad(7) + if (err != nil) != c.err { + t.Fatalf("RunStartupLoad err = %v, want error %v", err, c.err) + } + if !isClosed(s.StartupLoadDone()) { + t.Fatal("StartupLoadDone not closed after RunStartupLoad returned") + } + if got := (want{s.backgroundLoadDone.Load(), s.backgroundLoadFailed.Load()}); got != c.want { + t.Fatalf("health semantics changed: done/failed = %+v, want %+v", got, c.want) + } + s.signalStartupLoadDone() // a second signal is a no-op, not a panic + }) + } +} + +// The signal is separate from LoadComplete(): it stays open while the +// background fill runs, after the hot window has already reported complete. +func TestStartupLoadDoneWaitsForBackgroundFill_116(t *testing.T) { + s := openStartupStore(t, hotAndBackgroundDB(t, 20, 40), 168, 1) + var sawComplete, sawDone atomic.Bool + s.bgLoaderEntryHook = func() { + sawComplete.Store(s.LoadComplete()) + sawDone.Store(isClosed(s.StartupLoadDone())) + } + if err := s.RunStartupLoad(7); err != nil { + t.Fatal(err) + } + if !sawComplete.Load() { + t.Fatal("fixture: LoadComplete() was false when the background fill started") + } + if sawDone.Load() { + t.Fatal("StartupLoadDone was closed while the background fill was still to run") + } +} + +func TestStartupLoadDoneDropsDependentCaches_116(t *testing.T) { + db := setupTestDB(t) + defer db.Close() + s := NewPacketStore(db, nil) + s.hashSizeInfoMu.Lock() + s.hashSizeInfoCache = map[string]*hashSizeNodeInfo{"x": {}} + s.hashSizeInfoAt = time.Now() + s.hashSizeInfoMu.Unlock() + s.clockSkew.mu.Lock() + s.clockSkew.lastComputed = time.Now() + s.clockSkew.mu.Unlock() + s.cacheMu.Lock() + s.rfCache["AAR|"] = &cachedResult{expiresAt: time.Now().Add(time.Hour)} + s.cacheMu.Unlock() + + s.signalStartupLoadDone() + + s.hashSizeInfoMu.Lock() + hs := s.hashSizeInfoCache + s.hashSizeInfoMu.Unlock() + s.clockSkew.mu.Lock() + last := s.clockSkew.lastComputed + s.clockSkew.mu.Unlock() + s.cacheMu.Lock() + rf := len(s.rfCache) + s.cacheMu.Unlock() + if hs != nil { + t.Error("the hash-size info cache (15 s TTL) survived the startup-load signal") + } + if !last.IsZero() { + t.Error("the clock-skew recompute throttle survived the startup-load signal") + } + if rf != 0 { + t.Error("region-keyed analytics TTL caches computed on the partial store survived the signal") + } +} + +// RecomputeNow runs a pass that starts after the call, on the recomputer's +// own goroutine (never concurrently with a periodic pass), and restarts the +// ticker so no periodic pass follows right behind it. +func TestRecomputeNowRunsOnLoopAndResetsTicker_116(t *testing.T) { + var mu sync.Mutex + var starts []time.Time + var running, maxRunning atomic.Int32 + rc := newAnalyticsRecomputer("t", 400*time.Millisecond, func() interface{} { + if n := running.Add(1); n > maxRunning.Load() { + maxRunning.Store(n) + } + mu.Lock() + starts = append(starts, time.Now()) + mu.Unlock() + time.Sleep(5 * time.Millisecond) + running.Add(-1) + return 1 + }) + rc.Start() + defer rc.Stop() + time.Sleep(300 * time.Millisecond) + before := time.Now() + rc.RecomputeNow() + mu.Lock() + n := len(starts) + last := starts[n-1] + mu.Unlock() + if n != 2 || last.Before(before) { + t.Fatalf("RecomputeNow: %d passes, last started %v before the call", n, before.Sub(last)) + } + // Without the reset the periodic tick at ~400 ms would run now. + time.Sleep(250 * time.Millisecond) + mu.Lock() + n = len(starts) + mu.Unlock() + if n != 2 { + t.Fatalf("a periodic pass ran %d time(s) right after RecomputeNow: the ticker was not reset", n-2) + } + time.Sleep(400 * time.Millisecond) + mu.Lock() + n = len(starts) + mu.Unlock() + if n < 3 { + t.Fatal("the periodic loop stopped after RecomputeNow") + } + if maxRunning.Load() != 1 { + t.Fatalf("passes overlapped (%d at once)", maxRunning.Load()) + } + rc.Stop() + done := make(chan struct{}) + go func() { rc.RecomputeNow(); close(done) }() + select { + case <-done: + case <-time.After(2 * time.Second): + t.Fatal("RecomputeNow blocked on a stopped recomputer") + } +} + +// A pass that started before the startup load finished must not open the +// warm-up gate, even if the load finishes while it runs. +func TestWarmupGateIgnoresPassStartedBeforeLoad_116(t *testing.T) { + db := setupTestDB(t) + defer db.Close() + s := NewPacketStore(db, nil) + release := make(chan struct{}) + entered := make(chan struct{}, 4) + var calls atomic.Int32 + rc := newAnalyticsRecomputer("gated", time.Hour, func() interface{} { + if calls.Add(1) == 2 { + entered <- struct{}{} + <-release + } + return 1 + }) + loaded := s.StartupLoadDone() + rc.setWarmupReadyGate_1659(func() bool { return isClosed(loaded) }) + rc.Start() // pass 1: before the load, does not count + if !rc.FirstPassDoneAt_1659().IsZero() { + t.Fatal("a pass before the load opened the gate") + } + go rc.RecomputeNow() // pass 2 starts before the load ... + <-entered + s.signalStartupLoadDone() // ... and the load finishes while it runs + close(release) + time.Sleep(50 * time.Millisecond) + if !rc.FirstPassDoneAt_1659().IsZero() { + t.Fatal("a pass that STARTED before the load finished opened the warm-up gate") + } + rc.RecomputeNow() // pass 3 starts after: opens it + if rc.FirstPassDoneAt_1659().IsZero() { + t.Fatal("a pass started after the load did not open the gate") + } + rc.Stop() +} + +func recomputerByName(s *PacketStore) map[string]*analyticsRecomputer { + s.analyticsRecomputerMu.RLock() + defer s.analyticsRecomputerMu.RUnlock() + out := map[string]*analyticsRecomputer{} + for _, rc := range []*analyticsRecomputer{s.recompTopology, s.recompRF, s.recompDistance, s.recompChannels, + s.recompHashCollisions, s.recompHashSizes, s.recompRoles, s.recompObserversClockSkew, s.recompNodesClockSkew} { + out[rc.name] = rc + } + return out +} + +// Every default-shape recomputer runs exactly once more, promptly, after +// the signal; the gated three first and roles after nodes-clock-skew (roles +// reads that snapshot). +func TestEveryRecomputerRunsOnceAfterLoad_116(t *testing.T) { + db := setupTestDB(t) + defer db.Close() + s := NewPacketStore(db, nil) + stop := s.StartAnalyticsRecomputers(time.Hour) + defer stop() + rcs := recomputerByName(s) + before := map[string]int64{} + for n, rc := range rcs { + before[n] = rc.ComputeRuns() + } + time.Sleep(50 * time.Millisecond) + for n, rc := range rcs { + if rc.ComputeRuns() != before[n] { + t.Fatalf("%s recomputed before the load finished", n) + } + } + signalled := time.Now() + s.signalStartupLoadDone() + deadline := time.Now().Add(10 * time.Second) + for n, rc := range rcs { + for rc.ComputeRuns() < before[n]+1 { + if time.Now().After(deadline) { + t.Fatalf("%s did not recompute after the startup load finished (runs %d)", n, rc.ComputeRuns()) + } + time.Sleep(5 * time.Millisecond) + } + } + time.Sleep(100 * time.Millisecond) + type started struct { + name string + at time.Time + } + var order []started + for n, rc := range rcs { + if got := rc.ComputeRuns(); got != before[n]+1 { + t.Errorf("%s ran %d extra passes after the load, want 1", n, got-before[n]) + } + if rc.LastStartAt().Before(signalled) { + t.Errorf("%s's post-load pass started before the signal", n) + } + order = append(order, started{n, rc.LastStartAt()}) + } + sort.Slice(order, func(i, j int) bool { return order[i].at.Before(order[j].at) }) + pos := map[string]int{} + for i, o := range order { + pos[o.name] = i + } + for _, g := range []string{"rf", "topology", "channels"} { + if pos[g] > 2 { + t.Errorf("gated recomputer %s ran at position %d; the gated three go first", g, pos[g]) + } + } + if pos["roles"] < pos["nodes-clock-skew"] { + t.Error("roles recomputed before nodes-clock-skew, whose snapshot it reads") + } +} + +func analyticsRouter(t *testing.T, s *PacketStore) *mux.Router { + t.Helper() + srv := NewServer(s.db, &Config{Port: 3000}, NewHub()) + srv.store = s + r := mux.NewRouter() + srv.RegisterRoutes(r) + return r +} + +func get(r http.Handler, path string) *httptest.ResponseRecorder { + w := httptest.NewRecorder() + r.ServeHTTP(w, httptest.NewRequest("GET", path, nil)) + return w +} + +// End to end the way main.go runs it: recomputers start at the first chunk, +// the hot window completes (LoadComplete()==true), the background fill is +// held. RF stays 503 (on master it opened on the hot-window snapshot); +// ungated endpoints keep answering 200; after the fill RF serves a snapshot +// of the full store. +func TestGatedRFWaitsForBackgroundFill_116(t *testing.T) { + s := openStartupStore(t, hotAndBackgroundDB(t, 20, 60), 168, 1) + hold := make(chan struct{}) + inBg := make(chan struct{}) + s.bgLoaderEntryHook = func() { close(inBg); <-hold } + loadErr := make(chan error, 1) + go func() { loadErr <- s.RunStartupLoad(7) }() + <-s.FirstChunkReady() + <-inBg + if !s.LoadComplete() { + t.Fatal("fixture: hot window not complete when the background fill started") + } + stop := s.StartAnalyticsRecomputers(time.Hour) + defer stop() + r := analyticsRouter(t, s) + if w := get(r, "/api/analytics/rf"); w.Code != http.StatusServiceUnavailable { + t.Fatalf("RF during the background fill: %d, want 503 (the hot-window snapshot must not open the gate)", w.Code) + } + for _, p := range []string{"/api/analytics/hash-sizes", "/api/analytics/hash-collisions", "/api/analytics/roles"} { + if w := get(r, p); w.Code != http.StatusOK { + t.Errorf("ungated %s during the load: %d, want 200 (no new 503s)", p, w.Code) + } + } + close(hold) + if err := <-loadErr; err != nil { + t.Fatal(err) + } + deadline := time.Now().Add(10 * time.Second) + for get(r, "/api/analytics/rf").Code != http.StatusOK { + if time.Now().After(deadline) { + t.Fatal("RF never opened after the startup load") + } + time.Sleep(10 * time.Millisecond) + } + got := s.recompRF.Load().(map[string]interface{})["totalTransmissions"] + want := s.computeAnalyticsRF("", "", TimeWindow{})["totalTransmissions"] + if !reflect.DeepEqual(got, want) { + t.Fatalf("RF snapshot totalTransmissions=%v, full store gives %v", got, want) + } +} + +// When the force timeout opened the gate on partial data, the post-load +// recompute still replaces that snapshot promptly. +func TestForcedOpenSnapshotReplacedAfterLoad_116(t *testing.T) { + old := warmupForceTimeout + warmupForceTimeout = time.Millisecond + defer func() { warmupForceTimeout = old }() + db := setupTestDB(t) + defer db.Close() + s := NewPacketStore(db, nil) + stop := s.StartAnalyticsRecomputers(time.Hour) + defer stop() + time.Sleep(5 * time.Millisecond) + if s.recompRF.IsWarmingUp_1659() { + t.Fatal("fixture: the force timeout did not open the gate") + } + runs := s.recompRF.ComputeRuns() + s.signalStartupLoadDone() + deadline := time.Now().Add(5 * time.Second) + for s.recompRF.ComputeRuns() == runs { + if time.Now().After(deadline) { + t.Fatal("the forced-open RF snapshot was not replaced after the load") + } + time.Sleep(5 * time.Millisecond) + } + if s.recompRF.FirstPassDoneAt_1659().IsZero() { + t.Fatal("the post-load pass did not mark the first real pass") + } +} + +// The lazy distance build refreshes the distance snapshot before the +// index reports built, so the handler never goes from 202 to an older +// snapshot. +func TestDistanceSnapshotRefreshedBeforeReady_116(t *testing.T) { + db := setupTestDB(t) + defer db.Close() + s := NewPacketStore(db, nil) + if err := s.Load(); err != nil { + t.Fatal(err) + } + stop := s.StartAnalyticsRecomputers(time.Hour) + defer stop() + runs := s.recompDistance.ComputeRuns() + var runsAtReady int64 = -1 + var snapAtReady interface{} + s.TriggerDistanceIndexBuild() + deadline := time.Now().Add(5 * time.Second) + for runsAtReady < 0 { + if s.DistanceIndexBuilt() { + runsAtReady = s.recompDistance.ComputeRuns() + snapAtReady = s.recompDistance.Load() + } + if time.Now().After(deadline) { + t.Fatal("distance index never built") + } + time.Sleep(time.Millisecond) + } + if runsAtReady <= runs { + t.Fatal("the distance index reported built before the distance snapshot was refreshed") + } + if !reflect.DeepEqual(snapAtReady, s.computeAnalyticsDistance("", "")) { + t.Fatal("the distance snapshot at ready does not match the built index") + } +} + +// Region- and area-keyed distance results bypass the snapshot and live in +// the distance TTL cache. One computed from the index as it was before a +// build must not be served once that build reports the index built, so +// the build drops the cache before it flips DistanceIndexBuilt. +func TestDistanceBuildDropsRegionAreaCache_116(t *testing.T) { + store := newDistBuildStore(t) + store.config = &Config{Areas: map[string]AreaEntry{ + "BAY": {Label: "Bay", Polygon: [][2]float64{{36.0, -124.0}, {39.0, -124.0}, {39.0, -120.0}, {36.0, -120.0}}}, + }} + gate := installDistBuildGate(store) + store.TriggerDistanceIndexBuild() + distBuildWaitClosed(t, "the build to reach distanceBuildHook", gate.entered) + + // The build is held before it reads the dataset, so these come from + // the index as it was before it (nothing indexed yet). + keys := [][2]string{{"SJC", ""}, {"", "BAY"}, {"SJC", "BAY"}} + old := make(map[[2]string]map[string]interface{}, len(keys)) + for _, k := range keys { + old[k] = store.GetAnalyticsDistance(k[0], k[1]) + if again := store.GetAnalyticsDistance(k[0], k[1]); reflect.ValueOf(again).Pointer() != reflect.ValueOf(old[k]).Pointer() { + t.Fatalf("fixture: %v was not served from the TTL cache", k) + } + } + close(gate.release) + distBuildWaitCurrent(t, store) + + for _, k := range keys { + got := store.GetAnalyticsDistance(k[0], k[1]) + want := store.computeAnalyticsDistance(k[0], k[1]) + if reflect.DeepEqual(old[k], want) { + t.Fatalf("fixture: %v gives the same result before and after the build", k) + } + if reflect.ValueOf(got).Pointer() == reflect.ValueOf(old[k]).Pointer() { + t.Errorf("%v: served the result cached from the pre-build index after the index reported built", k) + } else if !reflect.DeepEqual(got, want) { + t.Errorf("%v: result after ready does not match the built index", k) + } + } +} diff --git a/cmd/server/analytics_recomputer.go b/cmd/server/analytics_recomputer.go index 88ac631e4..d4acc3a04 100644 --- a/cmd/server/analytics_recomputer.go +++ b/cmd/server/analytics_recomputer.go @@ -10,6 +10,9 @@ package main import ( + "fmt" + "log" + "strings" "sync" "sync/atomic" "time" @@ -37,6 +40,9 @@ type analyticsRecomputer struct { cache atomic.Value // holds interface{} — the latest snapshot stop chan struct{} done chan struct{} + // recomputeReq carries RecomputeNow requests to the loop goroutine, + // which closes the inner channel when that pass is done (#116). + recomputeReq chan chan struct{} startOnce sync.Once stopOnce sync.Once @@ -44,6 +50,7 @@ type analyticsRecomputer struct { // Stats (atomic). computeRuns atomic.Int64 lastComputeNs atomic.Int64 // duration of last compute in nanoseconds + lastStartNs atomic.Int64 // wall clock when the last compute started (#116) // Issue #1659 (PR #1688 r1) — warmup gate state, inlined here so // hot-path readers (IsWarmingUp_1659) do lock-free atomic loads @@ -66,6 +73,8 @@ func newAnalyticsRecomputer(name string, interval time.Duration, compute func() compute: compute, stop: make(chan struct{}), done: make(chan struct{}), + + recomputeReq: make(chan chan struct{}), } } @@ -96,6 +105,13 @@ func (r *analyticsRecomputer) loop() { select { case <-t.C: r.runOnce() + case ack := <-r.recomputeReq: + // #116: on-demand pass, on this goroutine so it never runs + // concurrently with a periodic one; restart the ticker so no + // periodic pass follows right behind it. + r.runOnce() + t.Reset(r.interval) + close(ack) case <-r.stop: return } @@ -114,7 +130,12 @@ func (r *analyticsRecomputer) runOnce() { // reach markFirstPassDone otherwise). _ = recover() }() + // #116: sample the warm-up readiness gate BEFORE computing, so a + // pass that started on a partially loaded store never ends the + // warm-up, even if the load finishes while it runs. + ready := r.warmupReadyGateOpen_1659() t0 := time.Now() + r.lastStartNs.Store(t0.UnixNano()) result := r.compute() r.lastComputeNs.Store(int64(time.Since(t0))) r.computeRuns.Add(1) @@ -128,9 +149,64 @@ func (r *analyticsRecomputer) runOnce() { // PR #1688 r1: called on EVERY successful pass (even nil // result) so a compute that returns nil but doesn't panic // still lifts the gate — banner-stuck-forever fix (munger #2). - // The markFirstPassDone helper is idempotent and additionally - // consults the chunked-loader readiness gate (munger #5). - r.markFirstPassDone_1659() + // The markFirstPassDone helper is idempotent; it only runs for a + // pass that started with the readiness gate open (munger #5, #116). + if ready { + r.markFirstPassDone_1659() + } +} + +// RecomputeNow has the loop goroutine run a pass that starts after this +// call and waits for it to finish; the periodic ticker restarts from +// that pass (#116). Returns early once the recomputer is stopped; blocks +// until Start if it has not started yet. +func (r *analyticsRecomputer) RecomputeNow() { + ack := make(chan struct{}) + select { + case r.recomputeReq <- ack: + case <-r.stop: + return + } + select { + case <-ack: + case <-r.stop: + } +} + +// LastStartAt returns when the most recent compute started (zero before +// the first one). +func (r *analyticsRecomputer) LastStartAt() time.Time { + ns := r.lastStartNs.Load() + if ns == 0 { + return time.Time{} + } + return time.Unix(0, ns) +} + +// recomputeWhenLoaded waits for loaded to close, then recomputes each +// recomputer once, one at a time, in slice order (#116). Sequential so +// the post-load passes do not all hold the store read lock at once, and +// so a recomputer that reads another's snapshot can follow it. Returns +// when done or when stop closes. +func recomputeWhenLoaded(loaded, stop <-chan struct{}, rcs []*analyticsRecomputer) { + select { + case <-loaded: + case <-stop: + return + } + t0 := time.Now() + parts := make([]string, 0, len(rcs)) + for _, rc := range rcs { + select { + case <-stop: + return + default: + } + rc.RecomputeNow() + parts = append(parts, fmt.Sprintf("%s=%s", rc.name, rc.LastComputeDuration().Round(time.Millisecond))) + } + log.Printf("[analytics-recompute] startup load done: recomputed %d snapshots in %s (%s)", + len(rcs), time.Since(t0).Round(time.Millisecond), strings.Join(parts, " ")) } // Load returns the most recently computed snapshot, or nil if Start @@ -260,34 +336,60 @@ func (s *PacketStore) StartAnalyticsRecomputers(defaultInterval time.Duration, o "nodes-clock-skew", pickInterval(ov.NodesClockSkew, defaultInterval), func() interface{} { return s.computeFleetClockSkew() }, ) + // Start and post-load recompute order (#116): the three warm-up-gated + // recomputers first, and roles after nodes-clock-skew because + // computeAnalyticsRoles reads that snapshot (GetFleetClockSkew). all := []*analyticsRecomputer{ - s.recompTopology, s.recompRF, s.recompDistance, - s.recompChannels, s.recompHashCollisions, s.recompHashSizes, - s.recompRoles, + s.recompRF, s.recompTopology, s.recompChannels, + s.recompDistance, s.recompHashCollisions, s.recompHashSizes, s.recompObserversClockSkew, s.recompNodesClockSkew, + s.recompRoles, } s.analyticsRecomputerMu.Unlock() - // Issue #1659 (PR #1688 r1, munger #5): wire the chunked-loader - // readiness gate on the three warmup-gated recomputers (RF, - // Topology, Channels). markFirstPassDone_1659 will refuse to - // flip first-pass-done until s.LoadComplete() reports true — - // i.e. the cold-load has populated all observations. Otherwise - // the FIRST recomputer pass runs against the post-restart in-RAM - // slice and the gate opens on partial data (the original #1659 - // bug class). - loadCompleteGate := s.LoadComplete - s.recompRF.setWarmupReadyGate_1659(loadCompleteGate) - s.recompTopology.setWarmupReadyGate_1659(loadCompleteGate) - s.recompChannels.setWarmupReadyGate_1659(loadCompleteGate) + // Issue #1659 (PR #1688 r1, munger #5): wire the loader readiness + // gate on the three warmup-gated recomputers (RF, Topology, + // Channels). #116: only a pass that STARTS after the whole startup + // load terminated (hot window AND background fill) ends the + // warm-up. LoadComplete() is not enough: it flips at the end of the + // hot window, so the gate used to open on a snapshot that missed the + // background fill. + loaded := s.StartupLoadDone() + loadedGate := func() bool { + select { + case <-loaded: + return true + default: + return false + } + } + s.recompRF.setWarmupReadyGate_1659(loadedGate) + s.recompTopology.setWarmupReadyGate_1659(loadedGate) + s.recompChannels.setWarmupReadyGate_1659(loadedGate) for _, rc := range all { rc.Start() } + // #116: main.go starts the recomputers at the first load chunk, so + // the initial computes above saw part of the data. Recompute every + // one as soon as the startup load terminates instead of a full + // interval later. + stopPostLoad := make(chan struct{}) + postLoadDone := make(chan struct{}) + go func() { + defer close(postLoadDone) + recomputeWhenLoaded(loaded, stopPostLoad, all) + }() + + var stopOnce sync.Once return func() { - for _, rc := range all { - rc.Stop() - } + stopOnce.Do(func() { + close(stopPostLoad) + for _, rc := range all { + rc.Stop() + } + <-postLoadDone + }) } } diff --git a/cmd/server/analytics_warmup_1659.go b/cmd/server/analytics_warmup_1659.go index 0709d5fc8..623c14913 100644 --- a/cmd/server/analytics_warmup_1659.go +++ b/cmd/server/analytics_warmup_1659.go @@ -19,14 +19,17 @@ // // Two correctness invariants: // -// 1. (#1688 munger #5) Only mark first-pass-done when BOTH: +// 1. (#1688 munger #5, #116) Only mark first-pass-done when BOTH: // a. a recomputer pass has completed, AND -// b. the chunked loader has finished (s.LoadComplete()). +// b. that pass STARTED after the whole startup load terminated, +// hot window and background fill (store.StartupLoadDone()). // The gate's `readyGate` callback is wired by -// StartAnalyticsRecomputers to `store.LoadComplete`. Passes that -// complete while loadComplete is still false leave the gate in -// the warming-up state; the NEXT pass after loadComplete flips -// true is the one that opens the gate. +// StartAnalyticsRecomputers to that signal and sampled before +// each compute. Passes that start earlier leave the gate in the +// warming-up state; StartAnalyticsRecomputers recomputes as soon +// as the signal fires, and that pass opens the gate. (The gate +// used to be store.LoadComplete, which flips at the end of the +// hot window, before the background fill.) // // 2. (#1688 munger #2 + kent-beck #2) The gate MUST lift in bounded // time. If compute() panics on every pass, hangs indefinitely, @@ -40,7 +43,8 @@ // default) elapsed since the recomputer was constructed // forces IsWarmingUp_1659() to false — degraded mode // (serve whatever cache exists, possibly empty) is -// strictly better than a permanent 503. +// strictly better than a permanent 503. A snapshot served +// that way is replaced by the post-load recompute (#116). // // Concurrency (#1688 munger #3): // @@ -105,26 +109,27 @@ func (r *analyticsRecomputer) loadWarmupReadyGate_1659() func() bool { } // markFirstPassDone_1659 is called from analyticsRecomputer.runOnce() -// after every compute attempt (success OR nil result; panics are -// caught upstream and never reach here). -// -// The gate flip is conditional on the readyGate (when set) reporting -// true — this implements the munger #5 fix: first-pass-done must -// require BOTH a recomputer pass complete AND the chunked loader to -// have finished populating the in-RAM observation set. +// after a compute attempt (success OR nil result; panics are caught +// upstream and never reach here) that STARTED with the readyGate open +// (see warmupReadyGateOpen_1659): first-pass-done requires a pass that +// began after the startup load terminated (munger #5, #116). // // Idempotent: only the FIRST successful flip wins; subsequent calls // observe a non-zero firstPassDoneNs and return immediately. func (r *analyticsRecomputer) markFirstPassDone_1659() { - if r.firstPassDoneNs.Load() != 0 { - return - } - if gate := r.loadWarmupReadyGate_1659(); gate != nil && !gate() { - return - } r.firstPassDoneNs.CompareAndSwap(0, time.Now().UnixNano()) } +// warmupReadyGateOpen_1659 reports whether the readyGate (when set) is +// open. runOnce samples it BEFORE computing and calls +// markFirstPassDone_1659 only when it was open (#116): a pass that began +// on a partially loaded store must not end the warm-up just because the +// load finished while it ran. +func (r *analyticsRecomputer) warmupReadyGateOpen_1659() bool { + gate := r.loadWarmupReadyGate_1659() + return gate == nil || gate() +} + // FirstPassDoneAt_1659 reports the time the first full compute pass // completed (subject to the readyGate). Returns zero time if no // qualifying pass has completed yet. diff --git a/cmd/server/analytics_warmup_1659_test.go b/cmd/server/analytics_warmup_1659_test.go index e78c71d98..3fae696b3 100644 --- a/cmd/server/analytics_warmup_1659_test.go +++ b/cmd/server/analytics_warmup_1659_test.go @@ -82,11 +82,11 @@ func TestAnalyticsRF_AfterFirstPassReturns200(t *testing.T) { db := setupTestDB(t) defer db.Close() store := NewPacketStore(db, nil) - // #1688 r1: the warmup gate now ALSO requires LoadComplete() to be - // true before first-pass-done flips (munger #5). Tests that don't - // exercise the chunked loader must flip it manually to model a - // production server that has finished cold-loading. - store.loadComplete.Store(true) + // #1688 r1 / #116: the warmup gate only opens on a pass that starts + // after the startup load terminated (munger #5). Tests that don't run + // RunStartupLoad signal it manually to model a production server that + // has finished loading (hot window and background fill). + store.signalStartupLoadDone() stop := store.StartAnalyticsRecomputers(50 * time.Millisecond) defer stop() diff --git a/cmd/server/chunked_load.go b/cmd/server/chunked_load.go index 300d58f47..314481508 100644 --- a/cmd/server/chunked_load.go +++ b/cmd/server/chunked_load.go @@ -96,9 +96,42 @@ func (s *PacketStore) OnChunkLoaded(fn func(rowsThisChunk, totalRows int)) { func (s *PacketStore) chunkedLoadInit() { s.chunkInitOnce.Do(func() { s.firstChunkReady = make(chan struct{}) + s.startupLoadDone = make(chan struct{}) }) } +// StartupLoadDone returns a channel closed once RunStartupLoad has +// returned (#116): LoadChunked AND the background fill loader have +// terminated, whether they succeeded or not, so nothing more is loaded +// from SQLite at start-up. LoadComplete() is not a substitute: it flips at +// the end of the hot window, before the background fill starts. Nor is +// backgroundLoadDone, which stays false on a failed fill (health +// semantics, #1690) -- a failure must still wake the recomputers. +func (s *PacketStore) StartupLoadDone() <-chan struct{} { + s.chunkedLoadInit() + return s.startupLoadDone +} + +// signalStartupLoadDone closes StartupLoadDone exactly once. Before +// closing it drops what was computed from the partial store and is read +// by the post-load recomputes or later requests: the hash-size info cache +// (15 s TTL), the clock-skew recompute throttle (30 s) and the +// region/window analytics TTL caches. +func (s *PacketStore) signalStartupLoadDone() { + s.chunkedLoadInit() + if !s.startupLoadSignaled.CompareAndSwap(false, true) { + return + } + s.hashSizeInfoMu.Lock() + s.hashSizeInfoCache = nil + s.hashSizeInfoMu.Unlock() + if s.clockSkew != nil { + s.clockSkew.Invalidate() + } + s.invalidateCachesFor(cacheInvalidation{eviction: true}) + close(s.startupLoadDone) +} + func (s *PacketStore) signalFirstChunk() { if s.firstChunkSignaled.CompareAndSwap(false, true) { close(s.firstChunkReady) @@ -162,6 +195,7 @@ func (s *PacketStore) fireChunkCallbacks(rowsThisChunk, totalRows int) { // - hotStartupHours > 0 success: terminal state is whatever // loadBackgroundChunks set (done=true on full coverage, // failed=true on partial / chunk errors — see #1690). +// - on every path: StartupLoadDone() is closed (#116). // // Issue #1809 root cause: previously main.go spawned loadBackgroundChunks // at FirstChunkReady while LoadChunked was still merging the remainder @@ -172,6 +206,8 @@ func (s *PacketStore) fireChunkCallbacks(rowsThisChunk, totalRows int) { // parallelism while ensuring oldestLoaded has a valid floor when the // bg loader starts. func (s *PacketStore) RunStartupLoad(chunkSize int) error { + // #116: runs last (deferred first), on every return path. + defer s.signalStartupLoadDone() // #89: once this returns, on every path (including a failed or partial // background fill), nothing more is loaded at start-up; parked // route_mask changes for transmissions still missing can be dropped. diff --git a/cmd/server/clock_skew.go b/cmd/server/clock_skew.go index f8a6408fc..1eb5fbb4c 100644 --- a/cmd/server/clock_skew.go +++ b/cmd/server/clock_skew.go @@ -220,6 +220,14 @@ func NewClockSkewEngine() *ClockSkewEngine { } } +// Invalidate makes the next Recompute run even if the last one is younger +// than computeInterval (#116: after the startup load). +func (e *ClockSkewEngine) Invalidate() { + e.mu.Lock() + e.lastComputed = time.Time{} + e.mu.Unlock() +} + // Recompute recalculates all clock skew data from the packet store. // Called periodically or on demand. Holds store RLock externally. // Uses read-copy-update: heavy computation runs outside the write lock, diff --git a/cmd/server/store.go b/cmd/server/store.go index 2bbf04171..8383d866d 100644 --- a/cmd/server/store.go +++ b/cmd/server/store.go @@ -521,10 +521,16 @@ type PacketStore struct { chunkInitOnce sync.Once firstChunkReady chan struct{} firstChunkSignaled atomic.Bool - loadComplete atomic.Bool - loadProgressRows atomic.Int64 - chunkCBMu sync.Mutex - chunkCallbacks []func(rowsThisChunk, totalRows int) + // startupLoadDone is closed once RunStartupLoad returns, on every + // path: hot window AND background fill have terminated, successfully + // or not (#116). Separate from backgroundLoadDone, which also means + // "coverage reached" for health reporting. See StartupLoadDone. + startupLoadDone chan struct{} + startupLoadSignaled atomic.Bool + loadComplete atomic.Bool + loadProgressRows atomic.Int64 + chunkCBMu sync.Mutex + chunkCallbacks []func(rowsThisChunk, totalRows int) // Eviction config and stats retentionHours float64 // 0 = unlimited @@ -4868,6 +4874,20 @@ func (s *PacketStore) runDistanceIndexBuild() { obs := s.totalObs s.mu.Unlock() + // #116: the distance snapshot (and any region/area result) was + // computed from the index as it was before this build. Refresh + // before reporting the index built, so the handler never goes + // from 202 to serving that older snapshot for up to an interval. + s.cacheMu.Lock() + s.distCache = make(map[string]*cachedResult) + s.cacheMu.Unlock() + s.analyticsRecomputerMu.RLock() + rc := s.recompDistance + s.analyticsRecomputerMu.RUnlock() + if rc != nil { + rc.RecomputeNow() + } + s.distLazyMu.Lock() s.distLazyBuilt = true s.distLazyBuiltGen = gen