diff --git a/cmd/server/distance_lazy_build_149_test.go b/cmd/server/distance_lazy_build_149_test.go new file mode 100644 index 000000000..0edd4a6f6 --- /dev/null +++ b/cmd/server/distance_lazy_build_149_test.go @@ -0,0 +1,783 @@ +package main + +import ( + "context" + "net/http/httptest" + "os" + "os/exec" + "runtime" + "strings" + "sync" + "sync/atomic" + "testing" + "time" + + "github.com/gorilla/mux" +) + +// Issue #149: two defects in the gate state of the lazy distance-index build +// (#1011). +// +// 1. distLazyOnce (a sync.Once) was reassigned to a zero value by the +// background-load completion while Do could be running on it. Do holds +// the Once's internal mutex for the whole build, so its deferred unlock +// hit the zeroed mutex: "fatal error: sync: unlock of unlocked mutex", +// which cannot be recovered — the server exits. +// 2. TriggerDistanceIndexBuild took s.mu.RLock while holding distLazyMu, and +// the background-load completion took distLazyMu while holding s.mu.Lock. +// Opposite lock order: both block forever. +// +// On a broken gate these scenarios crash or hang the process, so they run in +// a child process (runDistBuildChild). Every wait inside the child has a +// deadline and dumps all goroutines when it expires; the parent only reports +// what the child printed. + +const distBuildChildEnv = "CORESCOPE_DIST_BUILD_149_CHILD" + +// distBuildDeadline bounds each wait. The scenarios finish in milliseconds +// when the gate is correct; the deadline only turns a hang into a failure +// and leaves room for a loaded -race CI runner. +const distBuildDeadline = 10 * time.Second + +// runDistBuildChild runs scenario in a child process that re-executes only +// the calling test, and fails the test if the child fails, crashes or hangs. +// In the child it runs scenario directly. +func runDistBuildChild(t *testing.T, scenario func(t *testing.T)) { + t.Helper() + if os.Getenv(distBuildChildEnv) == t.Name() { + scenario(t) + return + } + ctx, cancel := context.WithTimeout(context.Background(), 2*time.Minute) + defer cancel() + cmd := exec.CommandContext(ctx, os.Args[0], + "-test.run=^"+t.Name()+"$", "-test.count=1", "-test.v", "-test.timeout=90s") + cmd.Env = append(os.Environ(), distBuildChildEnv+"="+t.Name()) + out, err := cmd.CombinedOutput() + if err != nil { + t.Fatalf("child process failed (%v); child output (tail):\n%s", err, lastLines(string(out), 150)) + } +} + +func lastLines(s string, n int) string { + lines := strings.Split(strings.TrimRight(s, "\n"), "\n") + if len(lines) > n { + lines = append([]string{"..."}, lines[len(lines)-n:]...) + } + return strings.Join(lines, "\n") +} + +func distBuildGoroutines() string { + buf := make([]byte, 1<<16) + for { + n := runtime.Stack(buf, true) + if n < len(buf) { + return string(buf[:n]) + } + buf = make([]byte, 2*len(buf)) + } +} + +// distBuildParked reports whether a goroutine with frame on its stack is +// blocked with the given runtime wait reason ("sync.Mutex.Lock", +// "sync.RWMutex.Lock", "sync.RWMutex.RLock"). +// +// Maintenance: the match depends on the runtime.Stack text format (the wait +// reason in the goroutine header) and the method names passed as frame. +// Re-verify it when changing the Go version or renaming those methods. +func distBuildParked(waitReason, frame string) bool { + return distBuildParkedCount(waitReason, frame) > 0 +} + +// distBuildParkedCount counts the goroutines distBuildParked matches; a +// waitReason of "sync." matches any mutex wait. +func distBuildParkedCount(waitReason, frame string) int { + n := 0 + for _, g := range strings.Split(distBuildGoroutines(), "\n\n") { + header, _, _ := strings.Cut(g, "\n") + if strings.Contains(header, "["+waitReason) && strings.Contains(g, frame) { + n++ + } + } + return n +} + +const ( + triggerFrame = ".(*PacketStore).TriggerDistanceIndexBuild(" + loaderFrame = ".(*PacketStore).loadBackgroundChunks(" + buildFrame = ".(*PacketStore).runDistanceIndexBuild(" + indexBuildsCreator = "github.com/corescope/server.(*PacketStore).startBackgroundIndexBuilds" +) + +// distBuildWait yields until cond holds and fails with a goroutine dump when +// the deadline expires first. +func distBuildWait(t *testing.T, what string, cond func() bool) { + t.Helper() + deadline := time.Now().Add(distBuildDeadline) + for !cond() { + if time.Now().After(deadline) { + t.Fatalf("timed out after %v waiting for %s; goroutines:\n%s", distBuildDeadline, what, distBuildGoroutines()) + } + runtime.Gosched() + } +} + +func distBuildWaitClosed(t *testing.T, what string, ch <-chan struct{}) { + t.Helper() + select { + case <-ch: + case <-time.After(distBuildDeadline): + t.Fatalf("timed out after %v waiting for %s; goroutines:\n%s", distBuildDeadline, what, distBuildGoroutines()) + } +} + +// distBuildWaitCurrent waits until no build runs and the index is current. +func distBuildWaitCurrent(t *testing.T, store *PacketStore) { + t.Helper() + distBuildWait(t, "a finished build with a current distance index", func() bool { + return !store.DistanceIndexBuilding() && store.DistanceIndexBuilt() + }) +} + +// distBuildSettle waits, best effort, until the goroutine count is back at +// baseline, so a build goroutine that already did its bookkeeping has also +// returned. Before #149 the build goroutine died in sync.Once.Do's deferred +// unlock right after that bookkeeping; waiting for it to return makes the +// fatal error happen before the child exits instead of racing it. +func distBuildSettle(baseline int) { + deadline := time.Now().Add(time.Second) + for runtime.NumGoroutine() > baseline && time.Now().Before(deadline) { + time.Sleep(time.Millisecond) + } +} + +// newDistBuildStore returns a loaded store whose background loader goes +// straight to its completion step: retention is set, and oldestLoaded is +// already past the retention cutoff, so the chunk loop exits on its first +// check and the only s.mu.Lock the loader takes is the completion's. Load()'s +// own index-build goroutines (#1008) have finished and returned, so nothing +// else queues for s.mu or counts as a leftover goroutine. +func newDistBuildStore(t *testing.T) *PacketStore { + t.Helper() + db := setupRichTestDB(t) + t.Cleanup(func() { db.Close() }) + store := NewPacketStore(db, &PacketStoreConfig{RetentionHours: 48}) + if err := store.Load(); err != nil { + t.Fatalf("Load(): %v", err) + } + if !store.WaitIndexesReady(distBuildDeadline) { + t.Fatal("Load()'s background index builds did not finish") + } + distBuildWait(t, "Load()'s index-build goroutines to return", func() bool { + return !strings.Contains(distBuildGoroutines(), "created by "+indexBuildsCreator) + }) + store.mu.Lock() + store.oldestLoaded = time.Now().UTC().Add(-96 * time.Hour).Format(time.RFC3339) + store.mu.Unlock() + return store +} + +// distBuildGate is installed as distanceBuildHook: it counts builds and holds +// the first one at its start (before it reads the dataset) until release is +// closed. Later builds only count. +type distBuildGate struct { + builds atomic.Int32 + entered chan struct{} + release chan struct{} +} + +func installDistBuildGate(store *PacketStore) *distBuildGate { + g := &distBuildGate{entered: make(chan struct{}), release: make(chan struct{})} + store.distanceBuildHook = func() { + if g.builds.Add(1) == 1 { + close(g.entered) + <-g.release + } + } + return g +} + +// distBuildLoadWhileQueued starts a build, holds it before it reads the +// dataset, runs the background-load completion to the end, then lets the +// build finish. +func distBuildLoadWhileQueued(t *testing.T) (*PacketStore, *distBuildGate) { + t.Helper() + store := newDistBuildStore(t) + gate := installDistBuildGate(store) + baseline := runtime.NumGoroutine() + + store.TriggerDistanceIndexBuild() + distBuildWaitClosed(t, "the build to reach distanceBuildHook", gate.entered) + // Before #149 the build was inside distLazyOnce.Do here, holding the + // Once's mutex, and this completion overwrote distLazyOnce. + store.loadBackgroundChunks() + close(gate.release) + + distBuildWaitCurrent(t, store) + distBuildSettle(baseline) + return store, gate +} + +// Crash reproduction: before #149 the child died with "fatal error: sync: +// unlock of unlocked mutex" from sync.(*Once).doSlow. +func TestDistanceBuild149_LoadCompletionDuringBuildDoesNotCrash(t *testing.T) { + runDistBuildChild(t, func(t *testing.T) { + distBuildLoadWhileQueued(t) + }) +} + +// The load completed before the queued build read the dataset, so that build +// already saw the full dataset: no second build. +func TestDistanceBuild149_LoadBeforeBuildSnapshotNoExtraBuild(t *testing.T) { + runDistBuildChild(t, func(t *testing.T) { + _, gate := distBuildLoadWhileQueued(t) + if n := gate.builds.Load(); n != 1 { + t.Fatalf("builds = %d, want 1: the build read the dataset after the load completed, so a rebuild is redundant", n) + } + }) +} + +// The load completed after the build read the dataset: the index that build +// produces is stale, so exactly one more build must follow, and the gate must +// not report the stale index as built. +func TestDistanceBuild149_LoadAfterBuildSnapshotRebuildsOnce(t *testing.T) { + runDistBuildChild(t, func(t *testing.T) { + store := newDistBuildStore(t) + gate := installDistBuildGate(store) + baseline := runtime.NumGoroutine() + + store.TriggerDistanceIndexBuild() + distBuildWaitClosed(t, "the build to reach distanceBuildHook", gate.entered) + + // Park the build between reading the dataset and its bookkeeping: + // the bookkeeping needs distLazyMu. The build has read the dataset + // once distHops is set and s.mu is free again. + store.distLazyMu.Lock() + close(gate.release) + distBuildWait(t, "the build to read the dataset", func() bool { + store.mu.RLock() + defer store.mu.RUnlock() + return store.distHops != nil + }) + + // Complete the background load now, after the build's snapshot. If + // the completion itself waits for distLazyMu, let both go and let + // them race for it: the outcome must be the same either way. + loaderDone := make(chan struct{}) + go func() { + defer close(loaderDone) + store.loadBackgroundChunks() + }() + distBuildWait(t, "the load completion to finish or wait for distLazyMu", func() bool { + select { + case <-loaderDone: + return true + default: + return distBuildParked("sync.Mutex.Lock", loaderFrame) + } + }) + store.distLazyMu.Unlock() + distBuildWaitClosed(t, "the load completion to finish", loaderDone) + + distBuildWaitCurrent(t, store) + distBuildSettle(baseline) + if n := gate.builds.Load(); n != 2 { + t.Fatalf("builds = %d, want 2: the first build read the dataset before the load completed, so exactly one rebuild must follow", n) + } + }) +} + +// The build is already waiting for s.mu when the load completion takes it and +// bumps the generation. The build must read the generation inside its own s.mu +// section: a value read before it predates the bump, and would force a +// redundant rebuild. +func TestDistanceBuild149_LoadHoldsMuWhileBuildQueuedNoExtraBuild(t *testing.T) { + runDistBuildChild(t, func(t *testing.T) { + store := newDistBuildStore(t) + gate := installDistBuildGate(store) + baseline := runtime.NumGoroutine() + + store.TriggerDistanceIndexBuild() + distBuildWaitClosed(t, "the build to reach distanceBuildHook", gate.entered) + + // Queue the load completion as a writer behind a held read lock, then + // queue the build behind the load completion. + store.mu.RLock() + loaderDone := make(chan struct{}) + go func() { + defer close(loaderDone) + store.loadBackgroundChunks() + }() + distBuildWait(t, "the load completion to queue for s.mu.Lock", func() bool { + return distBuildParked("sync.RWMutex.Lock", loaderFrame) + }) + close(gate.release) + distBuildWait(t, "the build to queue behind the load completion", func() bool { + return distBuildParkedCount("sync.", buildFrame) == 1 + }) + store.mu.RUnlock() + distBuildWaitClosed(t, "the load completion to finish", loaderDone) + + distBuildWaitCurrent(t, store) + distBuildSettle(baseline) + if n := gate.builds.Load(); n != 1 { + t.Fatalf("builds = %d, want 1: the build took s.mu after the load completed", n) + } + }) +} + +// A build whose dataset went stale during its first pass runs a second pass. +// When that pass starts, distLazyBuilding is still true, so triggers during it +// start no build of their own, and the index built by the first pass does not +// count as built. +func TestDistanceBuild149_TriggersDuringStaleRepassStartNoBuild(t *testing.T) { + runDistBuildChild(t, func(t *testing.T) { + store := newDistBuildStore(t) + var builds atomic.Int32 + entered1, release1 := make(chan struct{}), make(chan struct{}) + entered2, release2 := make(chan struct{}), make(chan struct{}) + store.distanceBuildHook = func() { + switch builds.Add(1) { + case 1: + close(entered1) + <-release1 + case 2: + close(entered2) + <-release2 + } + } + baseline := runtime.NumGoroutine() + + store.TriggerDistanceIndexBuild() + distBuildWaitClosed(t, "pass 1 to reach distanceBuildHook", entered1) + // Park pass 1 before its bookkeeping (it needs distLazyMu), once it + // has read the dataset, and complete the load behind its snapshot. + store.distLazyMu.Lock() + close(release1) + distBuildWait(t, "pass 1 to read the dataset", func() bool { + store.mu.RLock() + defer store.mu.RUnlock() + return store.distHops != nil + }) + // If the completion itself takes distLazyMu (allowed: s.mu → + // distLazyMu), let both go; it bumps the generation first either way. + loaderDone := make(chan struct{}) + go func() { + defer close(loaderDone) + store.loadBackgroundChunks() + }() + distBuildWait(t, "the load completion to finish or wait for distLazyMu", func() bool { + select { + case <-loaderDone: + return true + default: + return distBuildParked("sync.Mutex.Lock", loaderFrame) + } + }) + store.distLazyMu.Unlock() + distBuildWaitClosed(t, "the load completion to finish", loaderDone) + distBuildWaitClosed(t, "the stale second pass to reach distanceBuildHook", entered2) + + building, built := store.DistanceIndexBuilding(), store.DistanceIndexBuilt() + for i := 0; i < 8; i++ { + store.TriggerDistanceIndexBuild() + } + close(release2) + distBuildWaitCurrent(t, store) + distBuildSettle(baseline) + if !building { + t.Error("DistanceIndexBuilding() = false during the stale second pass") + } + if built { + t.Error("DistanceIndexBuilt() = true during the stale second pass: the first pass's index predates the load") + } + if n := builds.Load(); n != 2 { + t.Fatalf("builds = %d, want 2 (pass 1 and one second pass; triggers during the second pass must not start another)", n) + } + }) +} + +// Deadlock reproduction: a trigger on the debounce path holds distLazyMu and +// waits for s.mu.RLock while the load completion holds s.mu.Lock and waits for +// distLazyMu. Before #149 both blocked forever and the child timed out with +// both goroutines in the dump. +func TestDistanceBuild149_TriggerAndLoadCompletionDoNotDeadlock(t *testing.T) { + runDistBuildChild(t, func(t *testing.T) { + store := newDistBuildStore(t) + // A completed, current build puts the next trigger on the debounce + // path, which reads totalObs under s.mu. + store.TriggerDistanceIndexBuild() + distBuildWaitCurrent(t, store) + + // Hold a read lock so the load completion queues as a writer. + store.mu.RLock() + loaderDone := make(chan struct{}) + go func() { + defer close(loaderDone) + store.loadBackgroundChunks() + }() + distBuildWait(t, "the load completion to queue for s.mu.Lock", func() bool { + return distBuildParked("sync.RWMutex.Lock", loaderFrame) + }) + + // With a writer queued, RLock blocks: the trigger parks in RLock. + triggerDone := make(chan struct{}) + go func() { + defer close(triggerDone) + store.TriggerDistanceIndexBuild() + }() + distBuildWait(t, "the trigger to wait for s.mu.RLock", func() bool { + return distBuildParked("sync.RWMutex.RLock", triggerFrame) + }) + + // Let the load completion take s.mu. If the trigger holds distLazyMu + // while it waits, and the completion needs distLazyMu, neither can + // continue. + store.mu.RUnlock() + distBuildWaitClosed(t, "the load completion to finish (lock-order deadlock)", loaderDone) + distBuildWaitClosed(t, "the trigger to finish (lock-order deadlock)", triggerDone) + distBuildWait(t, "no build in flight", func() bool { return !store.DistanceIndexBuilding() }) + }) +} + +// distLazyBuilding must be set before the build goroutine is started, or +// triggers that run before that goroutine is scheduled start builds of their +// own. With one P the goroutine cannot run until the test goroutine blocks or +// yields, which makes that window deterministic. +func TestDistanceBuild149_BuildingIsSetBeforeTheBuildGoroutineRuns(t *testing.T) { + runDistBuildChild(t, func(t *testing.T) { + store := newDistBuildStore(t) + gate := installDistBuildGate(store) + baseline := runtime.NumGoroutine() + + defer runtime.GOMAXPROCS(runtime.GOMAXPROCS(1)) + runtime.Gosched() // start a fresh time slice: no preemption below + store.TriggerDistanceIndexBuild() + building := store.DistanceIndexBuilding() + for i := 0; i < 20; i++ { + store.TriggerDistanceIndexBuild() + } + if !building { + t.Error("DistanceIndexBuilding() = false right after the first trigger: distLazyBuilding is not set before the build goroutine starts") + } + + distBuildWaitClosed(t, "the build to reach distanceBuildHook", gate.entered) + close(gate.release) + distBuildWaitCurrent(t, store) + distBuildSettle(baseline) + if n := gate.builds.Load(); n != 1 { + t.Fatalf("builds = %d after 21 back-to-back triggers, want 1", n) + } + }) +} + +// Triggers on the debounce path read totalObs after releasing distLazyMu and +// re-check the gate before they start a build. When a rebuild is due and many +// of them get past the first check together, exactly one starts it. +func TestDistanceBuild149_DebouncedRebuildStartsOnce(t *testing.T) { + runDistBuildChild(t, func(t *testing.T) { + store := newDistBuildStore(t) + var builds atomic.Int32 + release := make(chan struct{}) + store.distanceBuildHook = func() { + if builds.Add(1) > 1 { + <-release // hold every rebuild until all triggers returned + } + } + baseline := runtime.NumGoroutine() + store.TriggerDistanceIndexBuild() + distBuildWaitCurrent(t, store) + store.distLazyMu.Lock() + store.distLazyLastBuilt = time.Now().Add(-6 * time.Minute) // rebuild due + store.distLazyMu.Unlock() + + // Park every trigger behind a held write lock, then let them all go. + const N = 16 + store.mu.Lock() + var wg sync.WaitGroup + wg.Add(N) + for i := 0; i < N; i++ { + go func() { + defer wg.Done() + store.TriggerDistanceIndexBuild() + }() + } + distBuildWait(t, "every trigger to wait for a lock", func() bool { + return distBuildParkedCount("sync.", triggerFrame) == N + }) + store.mu.Unlock() + wg.Wait() + close(release) + + distBuildWaitCurrent(t, store) + distBuildSettle(baseline) + if n := builds.Load(); n != 2 { + t.Fatalf("builds = %d, want 2 (the first build and one debounced rebuild for %d triggers)", n, N) + } + }) +} + +// A debounced trigger re-checks the gate in its second section: when a build +// finished while it read totalObs, its own rebuild is redundant. +func TestDistanceBuild149_DebouncedTriggerSkipsWhenABuildFinishedMeanwhile(t *testing.T) { + runDistBuildChild(t, func(t *testing.T) { + store := newDistBuildStore(t) + var builds atomic.Int32 + store.distanceBuildHook = func() { builds.Add(1) } + baseline := runtime.NumGoroutine() + store.TriggerDistanceIndexBuild() + distBuildWaitCurrent(t, store) + store.distLazyMu.Lock() + store.distLazyLastBuilt = store.distLazyLastBuilt.Add(-6 * time.Minute) // rebuild due + store.distLazyMu.Unlock() + + // Park the trigger in its s.mu read, after its first section ... + store.mu.Lock() + done := make(chan struct{}) + go func() { + defer close(done) + store.TriggerDistanceIndexBuild() + }() + distBuildWait(t, "the trigger to wait for s.mu.RLock", func() bool { + return distBuildParked("sync.RWMutex.RLock", triggerFrame) + }) + // ... then at its second section (s.mu → distLazyMu is allowed). + store.distLazyMu.Lock() + store.mu.Unlock() + distBuildWait(t, "the trigger to wait for distLazyMu", func() bool { + return distBuildParked("sync.Mutex.Lock", triggerFrame) + }) + // A build completes in between: this is what runDistanceIndexBuild records. + store.distLazyLastBuilt = time.Now() + store.distLazyMu.Unlock() + distBuildWaitClosed(t, "the trigger to return", done) + + distBuildWaitCurrent(t, store) + distBuildSettle(baseline) + if n := builds.Load(); n != 1 { + t.Fatalf("builds = %d, want 1: a build finished between the trigger's sections, so its rebuild is redundant", n) + } + }) +} + +// Many concurrent triggers during the first build start exactly one build, and +// triggers after it are debounced. +func TestDistanceBuild149_ConcurrentTriggersStartOneBuild(t *testing.T) { + store := newDistBuildStore(t) + gate := installDistBuildGate(store) + + const N = 64 + start := make(chan struct{}) + var wg sync.WaitGroup + wg.Add(N) + for i := 0; i < N; i++ { + go func() { + defer wg.Done() + <-start + store.TriggerDistanceIndexBuild() + }() + } + close(start) + wg.Wait() + distBuildWaitClosed(t, "the build to reach distanceBuildHook", gate.entered) + if !store.DistanceIndexBuilding() { + t.Fatal("DistanceIndexBuilding() = false while the build is held open") + } + close(gate.release) + distBuildWaitCurrent(t, store) + + wg.Add(N) + for i := 0; i < N; i++ { + go func() { + defer wg.Done() + store.TriggerDistanceIndexBuild() + }() + } + wg.Wait() + distBuildWaitCurrent(t, store) + if n := gate.builds.Load(); n != 1 { + t.Fatalf("builds = %d, want 1 (%d concurrent triggers, then %d debounced ones)", n, N, N) + } +} + +// The HTTP contract is unchanged: 202 + Retry-After while the index is not +// built (and while the build runs), 200 once it is, and 202 again with one +// rebuild after the background load completes. +func TestDistanceBuild149_HandlerContractAndLoadInvalidation(t *testing.T) { + store := newDistBuildStore(t) + gate := installDistBuildGate(store) + srv := NewServer(store.db, &Config{Port: 3000}, NewHub()) + srv.store = store + r := mux.NewRouter() + srv.RegisterRoutes(r) + get := func() *httptest.ResponseRecorder { + w := httptest.NewRecorder() + r.ServeHTTP(w, httptest.NewRequest("GET", "/api/analytics/distance", nil)) + return w + } + want202 := func(when string) { + t.Helper() + w := get() + if w.Code != 202 { + t.Fatalf("%s: status %d, want 202 (body=%s)", when, w.Code, w.Body.String()) + } + if ra := w.Header().Get("Retry-After"); ra != "5" { + t.Fatalf("%s: Retry-After = %q, want \"5\"", when, ra) + } + if !strings.Contains(w.Body.String(), `"status":"building"`) { + t.Fatalf("%s: body %s, want status building", when, w.Body.String()) + } + } + want200 := func(when string) { + t.Helper() + if w := get(); w.Code != 200 { + t.Fatalf("%s: status %d, want 200 (body=%s)", when, w.Code, w.Body.String()) + } + } + + want202("first request") + distBuildWaitClosed(t, "the build to reach distanceBuildHook", gate.entered) + want202("request during the build") + close(gate.release) + distBuildWaitCurrent(t, store) + want200("after the build") + if n := gate.builds.Load(); n != 1 { + t.Fatalf("builds = %d after the first build, want 1", n) + } + + store.loadBackgroundChunks() + if store.DistanceIndexBuilt() { + t.Fatal("DistanceIndexBuilt() = true after the background load completed; the index predates the load") + } + want202("first request after the background load") + distBuildWaitCurrent(t, store) + want200("after the rebuild") + if n := gate.builds.Load(); n != 2 { + t.Fatalf("builds = %d, want 2 (one rebuild after the background load)", n) + } +} + +// The debounce policy is unchanged: once a build completed, a trigger only +// rebuilds when Δobs since that build reached 5 % or 5 minutes have passed. +// While a debounced rebuild runs, the index reports not built, so the +// handler answers 202 as for any other build. +func TestDistanceBuild149_DebounceUnchanged(t *testing.T) { + store := newDistBuildStore(t) + var builds atomic.Int32 + var hold sync.Mutex // held by the test to keep a rebuild in flight + store.distanceBuildHook = func() { + builds.Add(1) + hold.Lock() + hold.Unlock() + } + store.TriggerDistanceIndexBuild() + distBuildWaitCurrent(t, store) + lastObsIs := func(when string, want int) { + t.Helper() + store.distLazyMu.Lock() + got := store.distLazyLastObs + store.distLazyMu.Unlock() + if got != want { + t.Errorf("%s: distLazyLastObs = %d, want totalObs %d: a build records the observation count it read", when, got, want) + } + } + store.mu.RLock() + loadedObs := store.totalObs + store.mu.RUnlock() + lastObsIs("after the first build", loadedObs) + + cases := []struct { + name string + lastObs int + curObs int + age time.Duration + rebuilds bool + }{ + {name: "just built, no new observations", lastObs: 1000, curObs: 1000}, + {name: "Δobs 4 %", lastObs: 1000, curObs: 1040, age: time.Minute}, + {name: "Δobs 6 %", lastObs: 1000, curObs: 1060, age: time.Minute, rebuilds: true}, + {name: "4 minutes old", lastObs: 1000, curObs: 1000, age: 4 * time.Minute}, + {name: "6 minutes old", lastObs: 1000, curObs: 1000, age: 6 * time.Minute, rebuilds: true}, + {name: "no Δobs baseline, just built", lastObs: 0, curObs: 1000}, + } + for _, c := range cases { + store.distLazyMu.Lock() + store.distLazyLastObs = c.lastObs + store.distLazyLastBuilt = time.Now().Add(-c.age) + store.distLazyMu.Unlock() + store.mu.Lock() + store.totalObs = c.curObs + store.mu.Unlock() + + before := builds.Load() + hold.Lock() + store.TriggerDistanceIndexBuild() + built := store.DistanceIndexBuilt() + hold.Unlock() + if built == c.rebuilds { + t.Errorf("%s: DistanceIndexBuilt() = %v right after the trigger, want %v", c.name, built, !c.rebuilds) + } + distBuildWaitCurrent(t, store) + want := before + if c.rebuilds { + want++ + } + if got := builds.Load(); got != want { + t.Errorf("%s: builds %d → %d, want %d", c.name, before, got, want) + } + if c.rebuilds { + lastObsIs(c.name, c.curObs) + } + } +} + +// The debounce boundaries are master's: its suppression condition was +// elapsed < 5 min && Δobs < 5 %, so exactly 5 minutes or exactly 5 % rebuilds. +func TestDistanceBuild149_RebuildDueBoundaries(t *testing.T) { + cases := []struct { + elapsed time.Duration + lastObs, curObs int + want bool + }{ + {5 * time.Minute, 1000, 1000, true}, + {5*time.Minute - time.Nanosecond, 1000, 1000, false}, + {0, 1000, 1050, true}, + {0, 1000, 1049, false}, + {0, 1000, 900, false}, // fewer observations (eviction) + {0, 0, 1000, false}, // no Δobs baseline + } + for _, c := range cases { + if got := distanceRebuildDue(c.elapsed, c.lastObs, c.curObs); got != c.want { + t.Errorf("distanceRebuildDue(%v, %d, %d) = %v, want %v", c.elapsed, c.lastObs, c.curObs, got, c.want) + } + } +} + +// BenchmarkTriggerDistanceIndexBuild covers the two trigger paths a burst of +// requests hits repeatedly: a build already in flight (early return under +// distLazyMu) and a current index inside the debounce window (reads totalObs +// under s.mu.RLock, suppressed). Neither starts a build, so a bare store is +// enough. +func BenchmarkTriggerDistanceIndexBuild(b *testing.B) { + b.Run("building", func(b *testing.B) { + s := &PacketStore{} + s.distLazyBuilding = true + b.ReportAllocs() + for i := 0; i < b.N; i++ { + s.TriggerDistanceIndexBuild() + } + }) + b.Run("debounced", func(b *testing.B) { + s := &PacketStore{totalObs: 1000} + s.distLazyBuilt = true + s.distLazyLastBuilt = time.Now() + s.distLazyLastObs = 1000 + b.ReportAllocs() + for i := 0; i < b.N; i++ { + s.TriggerDistanceIndexBuild() + } + if s.DistanceIndexBuilding() { + b.Fatal("the debounced path started a build") + } + }) +} diff --git a/cmd/server/distance_lock_order_guard_test.go b/cmd/server/distance_lock_order_guard_test.go new file mode 100644 index 000000000..1d1a6d061 --- /dev/null +++ b/cmd/server/distance_lock_order_guard_test.go @@ -0,0 +1,552 @@ +package main + +import ( + "fmt" + "go/ast" + "go/parser" + "go/token" + "path/filepath" + "sort" + "strings" + "testing" +) + +// Issue #149 lock-order rule: s.mu (Lock or RLock) is never acquired while +// distLazyMu is held. The other nesting, distLazyMu taken while s.mu is held, +// stays allowed; a distLazyMu → s.mu path deadlocks against any path that +// uses it (the pre-#149 background-load completion did). +// +// The guard parses this package's non-test sources and walks every function +// that uses distLazyMu statement by statement, tracking whether distLazyMu +// may be held (the union over branches). While it may be held, it rejects a +// .mu.Lock() / .mu.RLock() call and a call x.m() on a plain identifier x +// where m is a PacketStore method that takes s.mu, directly or through other +// PacketStore methods on the same receiver (a method of a field, such as +// s.distDataGen.Load(), is not one). The same applies to a call deferred +// after defer distLazyMu.Unlock(), which runs before that unlock. Calls in +// go statements and function literals are not inherited: they run on another +// goroutine or later. Function literals are checked as functions of their +// own. Not seen: s.mu taken inside a function literal that a distLazyMu +// holder calls, by a package-level function, or through a method value. +func TestDistLazyMuNeverHeldWhileTakingStoreMu(t *testing.T) { + names, err := filepath.Glob("*.go") + if err != nil { + t.Fatal(err) + } + fset := token.NewFileSet() + var files []*ast.File + for _, name := range names { + if strings.HasSuffix(name, "_test.go") { + continue + } + f, err := parser.ParseFile(fset, name, nil, parser.SkipObjectResolution) + if err != nil { + t.Fatalf("parse %s: %v", name, err) + } + files = append(files, f) + } + + takers := storeMuTakers(files) + for _, m := range []string{"Load", "loadBackgroundChunks"} { + if !takers[m] { + t.Fatalf("storeMuTakers misses %s, which takes s.mu: the call-graph scan is broken", m) + } + } + + var checked []string + var violations []string + for _, f := range files { + ast.Inspect(f, func(n ast.Node) bool { + var name string + var body *ast.BlockStmt + switch fn := n.(type) { + case *ast.FuncDecl: + name, body = fn.Name.Name, fn.Body + case *ast.FuncLit: + name, body = "func literal", fn.Body + } + if body == nil || !usesField(body, "distLazyMu") { + return true + } + checked = append(checked, name) + violations = append(violations, lockOrderViolations(fset, body, takers)...) + return true + }) + } + + sort.Strings(checked) + if !containsString(checked, "TriggerDistanceIndexBuild") { + t.Fatalf("TriggerDistanceIndexBuild was not checked (checked: %v): the guard no longer sees the distance-build gate", checked) + } + if len(violations) > 0 { + t.Fatalf("s.mu is taken while distLazyMu may be held (lock-order rule, #149):\n %s", + strings.Join(violations, "\n ")) + } +} + +// The walker itself: a positive and a negative control on synthetic code. +func TestLockOrderWalkerControls(t *testing.T) { + const src = `package p +func (s *PacketStore) direct() { + s.distLazyMu.Lock() + if s.distLazyBuilt { + s.mu.RLock() + s.mu.RUnlock() + } + s.distLazyMu.Unlock() +} +func (s *PacketStore) viaMethod() { + s.distLazyMu.Lock() + defer s.distLazyMu.Unlock() + _ = s.reads() +} +func (s *PacketStore) reads() int { + s.mu.RLock() + defer s.mu.RUnlock() + return s.totalObs +} +func (s *PacketStore) okEarlyReturn() { + s.distLazyMu.Lock() + if s.distLazyBuilding { + s.distLazyMu.Unlock() + return + } + s.distLazyMu.Unlock() + s.mu.RLock() + s.mu.RUnlock() +} +func (s *PacketStore) okGoroutine() { + s.distLazyMu.Lock() + go s.reads() + go func() { s.mu.Lock(); s.mu.Unlock() }() + s.distLazyMu.Unlock() +} +func (s *PacketStore) okOtherOrder() { + s.mu.Lock() + s.distLazyMu.Lock() + s.distLazyMu.Unlock() + s.mu.Unlock() +} +func (s *PacketStore) Load() { + s.mu.Lock() + s.mu.Unlock() +} +func (s *PacketStore) viaNamesake() { + s.distLazyMu.Lock() + s.Load() + s.distLazyMu.Unlock() +} +func (s *PacketStore) viaDeferAfterUnlockDefer() { + s.distLazyMu.Lock() + defer s.distLazyMu.Unlock() + defer s.reads() +} +func (s *PacketStore) okDeferBeforeUnlockDefer() { + defer s.reads() + s.distLazyMu.Lock() + defer s.distLazyMu.Unlock() +} +func (s *PacketStore) okDeferThenExplicitUnlock() { + s.distLazyMu.Lock() + defer s.reads() + s.distLazyMu.Unlock() +} +func (s *PacketStore) okFieldMethod() { + s.distLazyMu.Lock() + _ = s.distDataGen.Load() + s.distLazyMu.Unlock() +} +func (s *PacketStore) loopHeldAcrossIterations() { + for i := 0; i < 2; i++ { + s.mu.RLock() + s.mu.RUnlock() + s.distLazyMu.Lock() + } +} +` + fset := token.NewFileSet() + f, err := parser.ParseFile(fset, "synthetic.go", src, parser.SkipObjectResolution) + if err != nil { + t.Fatal(err) + } + takers := storeMuTakers([]*ast.File{f}) + got := map[string]int{} + for _, d := range f.Decls { + fn := d.(*ast.FuncDecl) + got[fn.Name.Name] = len(lockOrderViolations(fset, fn.Body, takers)) + } + want := map[string]int{ + "direct": 1, "viaMethod": 1, "reads": 0, "okEarlyReturn": 0, + "okGoroutine": 0, "okOtherOrder": 0, "loopHeldAcrossIterations": 1, + "Load": 0, "viaNamesake": 1, "okFieldMethod": 0, + "viaDeferAfterUnlockDefer": 1, "okDeferBeforeUnlockDefer": 0, "okDeferThenExplicitUnlock": 0, + } + for name, n := range want { + if got[name] != n { + t.Errorf("%s: %d violation(s), want %d", name, got[name], n) + } + } +} + +func containsString(list []string, s string) bool { + for _, v := range list { + if v == s { + return true + } + } + return false +} + +// usesField reports whether n selects a field or method named field, outside +// nested function literals. +func usesField(n ast.Node, field string) bool { + found := false + inspectSameGoroutine(n, func(n ast.Node) { + if sel, ok := n.(*ast.SelectorExpr); ok && sel.Sel.Name == field { + found = true + } + }) + return found +} + +// inspectSameGoroutine visits n in source order, skipping function literals +// and the calls of go statements (their arguments are still visited). +func inspectSameGoroutine(n ast.Node, visit func(ast.Node)) { + ast.Inspect(n, func(n ast.Node) bool { + switch n := n.(type) { + case *ast.FuncLit: + return false + case *ast.GoStmt: + for _, a := range n.Call.Args { + inspectSameGoroutine(a, visit) + } + return false + case nil: + return false + } + visit(n) + return true + }) +} + +// callOn matches ..(...) and returns true for it. +func callOn(call *ast.CallExpr, field string, methods ...string) bool { + sel, ok := call.Fun.(*ast.SelectorExpr) + if !ok { + return false + } + inner, ok := sel.X.(*ast.SelectorExpr) + if !ok || inner.Sel.Name != field { + return false + } + for _, m := range methods { + if sel.Sel.Name == m { + return true + } + } + return false +} + +// storeMuTakers returns the PacketStore methods that take s.mu (Lock or +// RLock) on their own goroutine, directly or through PacketStore methods +// called on the same receiver. +func storeMuTakers(files []*ast.File) map[string]bool { + takers := map[string]bool{} + calls := map[string][]string{} + for _, f := range files { + for _, d := range f.Decls { + fn, ok := d.(*ast.FuncDecl) + if !ok || fn.Body == nil || fn.Recv == nil || len(fn.Recv.List) != 1 { + continue + } + star, ok := fn.Recv.List[0].Type.(*ast.StarExpr) + if !ok { + continue + } + if id, ok := star.X.(*ast.Ident); !ok || id.Name != "PacketStore" || len(fn.Recv.List[0].Names) != 1 { + continue + } + recv := fn.Recv.List[0].Names[0].Name + name := fn.Name.Name + inspectSameGoroutine(fn.Body, func(n ast.Node) { + call, ok := n.(*ast.CallExpr) + if !ok { + return + } + sel, ok := call.Fun.(*ast.SelectorExpr) + if !ok { + return + } + if callOn(call, "mu", "Lock", "RLock") { + if inner := sel.X.(*ast.SelectorExpr); isIdent(inner.X, recv) { + takers[name] = true + } + } + if isIdent(sel.X, recv) { + calls[name] = append(calls[name], sel.Sel.Name) + } + }) + } + } + for changed := true; changed; { + changed = false + for m, callees := range calls { + if takers[m] { + continue + } + for _, c := range callees { + if takers[c] { + takers[m] = true + changed = true + break + } + } + } + } + return takers +} + +func isIdent(e ast.Expr, name string) bool { + id, ok := e.(*ast.Ident) + return ok && id.Name == name +} + +// lockOrderViolations walks body and reports each place where s.mu is taken +// while distLazyMu may be held. +func lockOrderViolations(fset *token.FileSet, body *ast.BlockStmt, takers map[string]bool) []string { + w := &lockOrderWalker{fset: fset, takers: takers, seen: map[token.Pos]bool{}} + w.block(body.List, false) + return w.violations +} + +type lockOrderWalker struct { + fset *token.FileSet + takers map[string]bool + seen map[token.Pos]bool + violations []string + // branchHeld collects the state at break/continue/goto, which leave + // the enclosing statement list; loops and switches merge it. + branchHeld bool + // deferredUnlock: defer distLazyMu.Unlock() was registered, so a call + // deferred later runs before that unlock (LIFO), with the lock held. + deferredUnlock bool +} + +// block walks stmts with distLazyMu possibly held on entry. It returns +// whether the lock may be held at the end and whether every path ends early +// (return, panic, break, continue, goto). +func (w *lockOrderWalker) block(stmts []ast.Stmt, held bool) (bool, bool) { + for _, s := range stmts { + var ends bool + held, ends = w.stmt(s, held) + if ends { + return held, true + } + } + return held, false +} + +func (w *lockOrderWalker) stmt(s ast.Stmt, held bool) (bool, bool) { + switch s := s.(type) { + case *ast.BlockStmt: + return w.block(s.List, held) + case *ast.LabeledStmt: + return w.stmt(s.Stmt, held) + case *ast.IfStmt: + if s.Init != nil { + held, _ = w.stmt(s.Init, held) + } + held = w.expr(s.Cond, held) + thenHeld, thenEnds := w.block(s.Body.List, held) + elseHeld, elseEnds := held, false + if s.Else != nil { + elseHeld, elseEnds = w.stmt(s.Else, held) + } + return mergeLockState([]bool{thenHeld, elseHeld}, []bool{thenEnds, elseEnds}) + case *ast.ForStmt: + if s.Init != nil { + held, _ = w.stmt(s.Init, held) + } + return w.loop(held, func(h bool) (bool, bool) { + h = w.expr(s.Cond, h) + h, ends := w.block(s.Body.List, h) + if s.Post != nil && !ends { + h, _ = w.stmt(s.Post, h) + } + return h, ends + }), false + case *ast.RangeStmt: + held = w.expr(s.X, held) + return w.loop(held, func(h bool) (bool, bool) { return w.block(s.Body.List, h) }), false + case *ast.SwitchStmt: + if s.Init != nil { + held, _ = w.stmt(s.Init, held) + } + held = w.expr(s.Tag, held) + return w.clauses(held, s.Body) + case *ast.TypeSwitchStmt: + if s.Init != nil { + held, _ = w.stmt(s.Init, held) + } + held, _ = w.stmt(s.Assign, held) + return w.clauses(held, s.Body) + case *ast.SelectStmt: + return w.clauses(held, s.Body) + case *ast.ReturnStmt: + for _, r := range s.Results { + held = w.expr(r, held) + } + return held, true + case *ast.BranchStmt: + w.branchHeld = w.branchHeld || held + return held, true + case *ast.GoStmt: + // The goroutine does not hold distLazyMu; only its arguments are + // evaluated here. + for _, a := range s.Call.Args { + held = w.expr(a, held) + } + return held, false + case *ast.DeferStmt: + // A deferred call runs at return. defer distLazyMu.Unlock() keeps + // the lock held until then, which is what "held" already says. + for _, a := range s.Call.Args { + held = w.expr(a, held) + } + if callOn(s.Call, "distLazyMu", "Unlock") { + w.deferredUnlock = true + } else if held && w.deferredUnlock { + w.check(s.Call) + } + return held, false + case *ast.ExprStmt: + held = w.expr(s.X, held) + if call, ok := s.X.(*ast.CallExpr); ok && isIdent(call.Fun, "panic") { + return held, true + } + return held, false + default: + return w.expr(s, held), false + } +} + +// loop runs body once with the entry state and, when the body can end with +// distLazyMu held, again with it held, as the next iteration would. +func (w *lockOrderWalker) loop(held bool, body func(bool) (bool, bool)) bool { + saved := w.branchHeld + w.branchHeld = false + out, _ := body(held) + if (out || w.branchHeld) && !held { + out2, _ := body(true) + out = out || out2 + } + out = held || out || w.branchHeld + w.branchHeld = saved + return out +} + +func (w *lockOrderWalker) clauses(held bool, body *ast.BlockStmt) (bool, bool) { + saved := w.branchHeld + w.branchHeld = false + var outs, ends []bool + hasDefault := false + for _, c := range body.List { + h := held + var list []ast.Stmt + switch c := c.(type) { + case *ast.CaseClause: + for _, e := range c.List { + h = w.expr(e, h) + } + hasDefault = hasDefault || c.List == nil + list = c.Body + case *ast.CommClause: + if c.Comm != nil { + h, _ = w.stmt(c.Comm, h) + } + hasDefault = hasDefault || c.Comm == nil + list = c.Body + } + h, e := w.block(list, h) + outs, ends = append(outs, h), append(ends, e) + } + if !hasDefault { + outs, ends = append(outs, held), append(ends, false) + } + out, allEnd := mergeLockState(outs, ends) + // A break leaves the switch/select and continues after it. + if w.branchHeld { + out, allEnd = true, false + } + w.branchHeld = saved + return out, allEnd +} + +// mergeLockState joins branches: held if any branch that falls through may +// hold the lock; ends only if every branch ends early. +func mergeLockState(held, ends []bool) (bool, bool) { + out, allEnd := false, true + for i := range held { + if !ends[i] { + allEnd = false + out = out || held[i] + } + } + return out, allEnd +} + +// expr applies the lock effects of the calls in n, in source order. +func (w *lockOrderWalker) expr(n ast.Node, held bool) bool { + if n == nil { + return held + } + inspectSameGoroutine(n, func(n ast.Node) { + call, ok := n.(*ast.CallExpr) + if !ok { + return + } + switch { + case callOn(call, "distLazyMu", "Lock", "TryLock"): + held = true + case callOn(call, "distLazyMu", "Unlock"): + held = false + case held: + w.check(call) + } + }) + return held +} + +// check reports call if it takes s.mu; the caller knows distLazyMu may be +// held when it runs. +func (w *lockOrderWalker) check(call *ast.CallExpr) { + if callOn(call, "mu", "Lock", "RLock") { + w.report(call, "takes "+exprString(call.Fun)) + return + } + sel, ok := call.Fun.(*ast.SelectorExpr) + if !ok || !w.takers[sel.Sel.Name] { + return + } + if _, onIdent := sel.X.(*ast.Ident); onIdent { + w.report(call, "calls "+sel.Sel.Name+", which takes s.mu") + } +} + +func (w *lockOrderWalker) report(call *ast.CallExpr, what string) { + if w.seen[call.Pos()] { + return + } + w.seen[call.Pos()] = true + w.violations = append(w.violations, fmt.Sprintf("%s: %s with distLazyMu held", w.fset.Position(call.Pos()), what)) +} + +func exprString(e ast.Expr) string { + switch e := e.(type) { + case *ast.SelectorExpr: + return exprString(e.X) + "." + e.Sel.Name + case *ast.Ident: + return e.Name + } + return "?" +} diff --git a/cmd/server/store.go b/cmd/server/store.go index 4130a1824..ba70232eb 100644 --- a/cmd/server/store.go +++ b/cmd/server/store.go @@ -319,24 +319,35 @@ type PacketStore struct { // overwritten. Incremented only under s.mu.RLock, read under s.mu.Lock. distSnapReaders atomic.Int64 - // Lazy-build gate for the distance index (#1011). distLazyBuilt is - // true once the first /api/analytics/distance request has completed - // its build (or a debounced rebuild). distLazyBuilding is true - // while a build is currently running — concurrent requests in this - // window receive 202 + Retry-After rather than racing N parallel - // O(n²) computations. distLazyOnce serialises the first build; - // reset by the background loader (Load() chunked merge) and by the - // debounced-rebuild policy so subsequent rebuilds can re-fire. + // Lazy-build gate for the distance index (#1011), guarded by + // distLazyMu. distLazyBuilding is true from the trigger that starts a + // build until that build has read the current dataset — concurrent + // requests in this window receive 202 + Retry-After rather than racing + // N parallel O(n²) computations. distLazyBuilt is true once a build has + // completed; distLazyBuiltGen is the distDataGen that build read, so the + // index is current only while the two match (#149). + // + // Lock order (#149): s.mu is never acquired (Lock or RLock) while + // distLazyMu is held. The other nesting, distLazyMu taken while s.mu is + // held, is allowed; a path that does it, as the background-load + // completion used to, deadlocks against any distLazyMu → s.mu path. distLazyMu sync.Mutex - distLazyOnce sync.Once distLazyBuilt bool distLazyBuilding bool + distLazyBuiltGen uint64 distLazyLastBuilt time.Time distLazyLastObs int // totalObs at last build, for Δobs debounce - // distanceBuildHook, if non-nil, runs at the start of the lazy build - // goroutine (after distLazyBuilding is set, before any lock is held). Tests - // use it to hold the build open so concurrent requests deterministically - // observe the "building" window; nil (and zero overhead) in production. + // distDataGen counts dataset changes that the incremental distance + // maintenance does not cover — today only the background chunk load's + // completion. Written under s.mu.Lock, so a build reads a value that + // matches its snapshot; atomic, so the gate can compare it under + // distLazyMu without taking s.mu. + distDataGen atomic.Uint64 + // distanceBuildHook, if non-nil, runs at the start of each build pass + // (after distLazyBuilding is set, before any lock is held). Tests use it + // to hold the build open so concurrent requests deterministically + // observe the "building" window, and to count builds; nil (and zero + // overhead) in production. distanceBuildHook func() // Cached GetNodeHashSizeInfo result — recomputed at most once every 15s @@ -1671,13 +1682,13 @@ func (s *PacketStore) loadBackgroundChunks() { s.buildPathHopIndex() // Distance index is now lazy (#1011) — built on first // /api/analytics/distance request, not at background-load - // completion. If a previous request already triggered the build, - // invalidate the gate so the next request rebuilds against the - // fuller dataset. - s.distLazyMu.Lock() - s.distLazyBuilt = false - s.distLazyOnce = sync.Once{} - s.distLazyMu.Unlock() + // completion. Bumping the generation makes an index built from the + // smaller dataset stale, so the next request rebuilds it, and a build + // in flight that already read it builds once more (#149). Taking + // distLazyMu here would be allowed (s.mu → distLazyMu; only the reverse + // is forbidden, see the gate fields); the atomic generation just makes + // it unnecessary. + s.distDataGen.Add(1) s.mu.Unlock() // #1008 review m3: flip the ready flags after the synchronous // rebuild for symmetry with startBackgroundIndexBuilds. Safe @@ -4629,12 +4640,21 @@ func (s *PacketStore) compactDistIndex(remove map[*StoreTx]bool) { } // DistanceIndexBuilt reports whether the distance analytics index has -// been constructed. Used by tests and /api/perf to verify the lazy -// build invariant from issue #1011 (eager Load() build removed). +// been built from the current dataset: false before the first build, and +// after the background load completed until a build has read the fuller +// dataset. Used by the /api/analytics/distance handler and by tests to +// verify the lazy build invariant from issue #1011 (eager Load() build +// removed). func (s *PacketStore) DistanceIndexBuilt() bool { s.distLazyMu.Lock() defer s.distLazyMu.Unlock() - return s.distLazyBuilt + return s.distIndexCurrentLocked() +} + +// distIndexCurrentLocked reports whether a completed build read the +// current dataset. Caller holds distLazyMu. +func (s *PacketStore) distIndexCurrentLocked() bool { + return s.distLazyBuilt && s.distLazyBuiltGen == s.distDataGen.Load() } // DistanceIndexBuilding reports whether a lazy distance-index build @@ -4649,63 +4669,101 @@ func (s *PacketStore) DistanceIndexBuilding() bool { // TriggerDistanceIndexBuild kicks off a lazy build of the distance // index in a background goroutine if one is not already running and -// the debounce policy permits. Returns immediately. Idempotent: -// concurrent callers see only one build, gated by sync.Once. +// the debounce policy permits. Returns immediately. Concurrent callers +// start at most one build: distLazyBuilding is set under distLazyMu +// before the goroutine starts. +// +// A missing or stale index (see distDataGen) always builds. Debounce +// policy for a current index (#1011 triage Fix path): rebuild only if +// Δobs ≥ 5% since the last build or 5 minutes have passed. // -// Debounce policy (#1011 triage Fix path): rebuild if Δobs > 5% since -// the last build OR at most once per 5 minutes — whichever is more -// restrictive. The first-ever build always runs. +// totalObs lives under s.mu, which must not be taken while distLazyMu is +// held (#149), so the debounce path reads it between two distLazyMu +// sections and re-checks the gate before it starts a build. func (s *PacketStore) TriggerDistanceIndexBuild() { s.distLazyMu.Lock() if s.distLazyBuilding { s.distLazyMu.Unlock() return } - // Debounce: if a build has already completed, suppress re-trigger - // unless Δobs > 5% or >5min has elapsed. - if s.distLazyBuilt { - s.mu.RLock() - curObs := s.totalObs - s.mu.RUnlock() - elapsed := time.Since(s.distLazyLastBuilt) - deltaPct := 0.0 - if s.distLazyLastObs > 0 { - deltaPct = float64(curObs-s.distLazyLastObs) / float64(s.distLazyLastObs) - } - if elapsed < 5*time.Minute && deltaPct < 0.05 { - s.distLazyMu.Unlock() - return - } - // Reset the gate so a new build can fire. - s.distLazyOnce = sync.Once{} - s.distLazyBuilt = false + if !s.distIndexCurrentLocked() { + s.startDistanceBuildLocked() + s.distLazyMu.Unlock() + return } + lastBuilt, lastObs := s.distLazyLastBuilt, s.distLazyLastObs s.distLazyMu.Unlock() - // Fire-and-forget; sync.Once collapses concurrent goroutines into - // a single build. The Once is reset above (under the mutex) before - // each rebuild cycle. - go s.distLazyOnce.Do(func() { - s.distLazyMu.Lock() - s.distLazyBuilding = true - s.distLazyMu.Unlock() + s.mu.RLock() + curObs := s.totalObs + s.mu.RUnlock() + if !distanceRebuildDue(time.Since(lastBuilt), lastObs, curObs) { + return + } + + s.distLazyMu.Lock() + // While distLazyMu was released another trigger may have started a + // build, or a build may have finished; either makes this one redundant. + if !s.distLazyBuilding && s.distLazyLastBuilt.Equal(lastBuilt) { + s.startDistanceBuildLocked() + } + s.distLazyMu.Unlock() +} + +// distanceRebuildDue applies the debounce policy to a current index that +// was built elapsed ago from lastObs observations, now that there are +// curObs. +func distanceRebuildDue(elapsed time.Duration, lastObs, curObs int) bool { + if elapsed >= 5*time.Minute { + return true + } + return lastObs > 0 && float64(curObs-lastObs)/float64(lastObs) >= 0.05 +} +// startDistanceBuildLocked marks a build in flight and starts it. Caller +// holds distLazyMu. Clearing distLazyBuilt keeps the handler answering +// 202 for the whole build, a debounced rebuild included (#1011). +func (s *PacketStore) startDistanceBuildLocked() { + s.distLazyBuilding = true + s.distLazyBuilt = false + go s.runDistanceIndexBuild() +} + +// runDistanceIndexBuild builds the distance index and records the +// dataset generation it read. If the background load completed after +// that read, the new index is already stale and it builds once more; +// distLazyBuilding stays true until a pass has read the current dataset. +// s.mu and distLazyMu are never held together. +func (s *PacketStore) runDistanceIndexBuild() { + for { if s.distanceBuildHook != nil { s.distanceBuildHook() // test seam: hold the build window open } s.mu.Lock() s.buildDistanceIndex() - obsAtBuild := s.totalObs + // Keep this read inside the s.mu section of the build: the load + // completion bumps the generation under s.mu, so here it matches the + // data just read. Read after the Unlock, a load completing in between + // would count as seen by an index built without it (a lost + // invalidation); read before the Lock, a load the build did see + // would force a redundant rebuild. + gen := s.distDataGen.Load() + obs := s.totalObs s.mu.Unlock() s.distLazyMu.Lock() - s.distLazyBuilding = false s.distLazyBuilt = true + s.distLazyBuiltGen = gen s.distLazyLastBuilt = time.Now() - s.distLazyLastObs = obsAtBuild + s.distLazyLastObs = obs + stale := gen != s.distDataGen.Load() + s.distLazyBuilding = stale s.distLazyMu.Unlock() - }) + if !stale { + return + } + } } // buildDistanceIndex precomputes haversine distances for all packets. diff --git a/cmd/server/sync_once_reset_guard_test.go b/cmd/server/sync_once_reset_guard_test.go new file mode 100644 index 000000000..208d9850f --- /dev/null +++ b/cmd/server/sync_once_reset_guard_test.go @@ -0,0 +1,220 @@ +package main + +import ( + "fmt" + "go/ast" + "go/parser" + "go/token" + "go/types" + "path/filepath" + "strings" + "testing" +) + +// A sync.Once must never be reassigned. Do holds the Once's internal mutex +// for the whole callback; overwriting the struct zeroes that mutex under a +// running Do, whose deferred unlock then dies with "fatal error: sync: unlock +// of unlocked mutex" (#149). When something has to be able to run again, use +// explicit state under a mutex instead. +// +// The guard parses this package's non-test sources. A type "holds a Once" +// when it is sync.Once, or a struct or array holding one by value. The guard +// rejects any assignment to a field or variable declared with such a type, +// and any assignment of a composite literal of such a type. Not covered: +// copies through pointers (*p = *q), which need type information; go vet's +// copylocks sees those, but go test's vet subset does not run it. +func TestNoSyncOnceIsReassigned(t *testing.T) { + names, err := filepath.Glob("*.go") + if err != nil { + t.Fatal(err) + } + fset := token.NewFileSet() + var files []*ast.File + for _, name := range names { + if strings.HasSuffix(name, "_test.go") { + continue + } + f, err := parser.ParseFile(fset, name, nil, parser.SkipObjectResolution) + if err != nil { + t.Fatalf("parse %s: %v", name, err) + } + files = append(files, f) + } + if len(files) == 0 { + t.Fatal("no production sources parsed: the guard scaffold is broken") + } + holdsOnce := onceHolderCheck(files) + for _, typ := range []string{"StoreTx", "PacketStore"} { + if !holdsOnce(ast.NewIdent(typ)) { + t.Fatalf("%s holds a sync.Once but the type scan misses it: the guard is broken", typ) + } + } + if found := syncOnceReassignments(fset, files); len(found) > 0 { + t.Fatalf("sync.Once reassigned (a running Do dies on its deferred unlock, #149):\n %s", + strings.Join(found, "\n ")) + } +} + +func TestSyncOnceGuardControls(t *testing.T) { + const src = `package p +import "sync" +type gate struct { + once sync.Once + built bool +} +type T struct { + once sync.Once + ptr *sync.Once + n int + g gate + gates [2]gate + gp *gate +} +var global sync.Once +func (t *T) resetField() { t.once = sync.Once{} } +func (t *T) resetGlobal() { global = sync.Once{} } +func (t *T) resetDeref() { *t.ptr = sync.Once{} } +func (t *T) copyInto(o T) { t.once = o.once } +func (t *T) resetHolder() { t.g = gate{} } +func (t *T) copyHolder(o gate) { t.g = o } +func (t *T) resetHolderArray() { t.gates = [2]gate{} } +func (t *T) resetHolderDeref() { *t.gp = gate{} } +func (t *T) resetOuter() { *t = T{} } +func (t *T) newPointer() { t.ptr = new(sync.Once) } +func (t *T) newHolderPointer() { t.gp = &gate{} } +func (t *T) otherField() { t.n = 1; t.g.built = true } +func (t *T) localOnce() { var o sync.Once; o.Do(func() {}) } +func (t *T) definedOnce() { o := sync.Once{}; o.Do(func() {}) } +` + fset := token.NewFileSet() + f, err := parser.ParseFile(fset, "synthetic.go", src, parser.SkipObjectResolution) + if err != nil { + t.Fatal(err) + } + found := syncOnceReassignments(fset, []*ast.File{f}) + got := map[string]bool{} + for _, d := range f.Decls { + fn, ok := d.(*ast.FuncDecl) + if !ok { + continue + } + start, end := fset.Position(fn.Pos()).Line, fset.Position(fn.End()).Line + for _, v := range found { + var line int + fmt.Sscanf(v[strings.Index(v, ":")+1:], "%d", &line) + if line >= start && line <= end { + got[fn.Name.Name] = true + } + } + } + for name, want := range map[string]bool{ + "resetField": true, "resetGlobal": true, "resetDeref": true, "copyInto": true, + "resetHolder": true, "copyHolder": true, "resetHolderArray": true, + "resetHolderDeref": true, "resetOuter": true, + "newPointer": false, "newHolderPointer": false, "otherField": false, + "localOnce": false, "definedOnce": false, + } { + if got[name] != want { + t.Errorf("%s: flagged = %v, want %v (found: %v)", name, got[name], want, found) + } + } +} + +// syncOnceReassignments returns "file:line: ..." for each assignment in files +// that overwrites a value holding a sync.Once. +func syncOnceReassignments(fset *token.FileSet, files []*ast.File) []string { + holdsOnce := onceHolderCheck(files) + onceNames := map[string]bool{} + for _, f := range files { + ast.Inspect(f, func(n ast.Node) bool { + switch n := n.(type) { + case *ast.Field: + if holdsOnce(n.Type) { + for _, id := range n.Names { + onceNames[id.Name] = true + } + } + case *ast.ValueSpec: + if n.Type != nil && holdsOnce(n.Type) { + for _, id := range n.Names { + onceNames[id.Name] = true + } + } + } + return true + }) + } + var found []string + for _, f := range files { + ast.Inspect(f, func(n ast.Node) bool { + as, ok := n.(*ast.AssignStmt) + if !ok || as.Tok == token.DEFINE { + return true + } + reassigns := false + for _, lhs := range as.Lhs { + if sel, ok := lhs.(*ast.SelectorExpr); ok && onceNames[sel.Sel.Name] { + reassigns = true + } + if id, ok := lhs.(*ast.Ident); ok && onceNames[id.Name] { + reassigns = true + } + } + for _, rhs := range as.Rhs { + if lit, ok := rhs.(*ast.CompositeLit); ok && lit.Type != nil && holdsOnce(lit.Type) { + reassigns = true + } + } + if reassigns { + found = append(found, fmt.Sprintf("%s: %s", fset.Position(as.Pos()), types.ExprString(as.Lhs[0]))) + } + return true + }) + } + return found +} + +// onceHolderCheck returns a predicate reporting whether a type expression +// holds a sync.Once by value: sync.Once itself, or a package type, struct or +// array that holds one. +func onceHolderCheck(files []*ast.File) func(ast.Expr) bool { + typeDecls := map[string]ast.Expr{} + for _, f := range files { + ast.Inspect(f, func(n ast.Node) bool { + if ts, ok := n.(*ast.TypeSpec); ok { + typeDecls[ts.Name.Name] = ts.Type + } + return true + }) + } + holders := map[string]bool{} // package types that hold a Once by value + var holdsOnce func(e ast.Expr) bool + holdsOnce = func(e ast.Expr) bool { + switch e := e.(type) { + case *ast.SelectorExpr: + return isIdent(e.X, "sync") && e.Sel.Name == "Once" + case *ast.Ident: + return holders[e.Name] + case *ast.ParenExpr: + return holdsOnce(e.X) + case *ast.ArrayType: + return e.Len != nil && holdsOnce(e.Elt) // an array, not a slice + case *ast.StructType: + for _, f := range e.Fields.List { + if holdsOnce(f.Type) { + return true + } + } + } + return false + } + for changed := true; changed; { + changed = false + for name, typ := range typeDecls { + if !holders[name] && holdsOnce(typ) { + holders[name], changed = true, true + } + } + } + return holdsOnce +}