diff --git a/docs/release-control/v6/internal/subsystems/performance-and-scalability.md b/docs/release-control/v6/internal/subsystems/performance-and-scalability.md index ef912446f..c40cb46b1 100644 --- a/docs/release-control/v6/internal/subsystems/performance-and-scalability.md +++ b/docs/release-control/v6/internal/subsystems/performance-and-scalability.md @@ -3168,3 +3168,14 @@ The existing 500-node mixed-endpoint workload and latency budgets remain intact. `pkg/metrics/store_additional_test.go` holds every history connection to prove availability without relying on favourable scheduling, and covers committed-data visibility, write rejection, clear and pool shutdown. + +### Metrics-store shutdown never closes the ingestion channel + +The metrics store's ingestion worker must not close its write channel on +shutdown. A concurrent writer that passed the `stopping` check before `Close`, +and the `WriteBatchSync`/`WriteBatchBounded` paths that do not consult +`stopping` at all, may still send to that channel. Closing it races those sends +and panics the process during shutdown. The worker instead drains already-queued +requests without closing the channel; a late write lands in the buffered channel +and is discarded with the store. `pkg/metrics/store_additional_test.go` pins the +post-shutdown enqueue and the concurrent-close race. diff --git a/pkg/metrics/store.go b/pkg/metrics/store.go index b478da9c2..e998c8629 100644 --- a/pkg/metrics/store.go +++ b/pkg/metrics/store.go @@ -1842,12 +1842,21 @@ func (s *Store) backgroundWorker() { if batch := s.drainBuffer(); len(batch) > 0 { remaining = append(remaining, writeRequest{metrics: batch}) } - close(s.writeCh) - for req := range s.writeCh { - remaining = append(remaining, req) + // Do NOT close writeCh. Concurrent writers (WriteWithTier after + // its stopping check, and WriteBatchSync/WriteBatchBounded which + // never check stopping) may still be sending; closing the channel + // races those sends and panics the process on shutdown. Drain what + // is already queued instead, and let any late write land in the + // buffered channel to be discarded with the store. + for { + select { + case req := <-s.writeCh: + remaining = append(remaining, req) + default: + s.processWriteRequests(remaining) + return + } } - s.processWriteRequests(remaining) - return case req, ok := <-s.writeCh: if !ok { diff --git a/pkg/metrics/store_additional_test.go b/pkg/metrics/store_additional_test.go index 2f8c2d7b7..83eca705b 100644 --- a/pkg/metrics/store_additional_test.go +++ b/pkg/metrics/store_additional_test.go @@ -1455,3 +1455,90 @@ func TestStoreStatsReaderObservesCommittedDataAndCloses(t *testing.T) { t.Fatal("main pool remained open after store shutdown") } } + +// Regression coverage for the shutdown send-on-closed-channel review of +// pkg/metrics/store.go. The ingestion worker used to close writeCh on +// <-stopCh. A writer that passed the stopping check just before Close, or a +// WriteBatchSync/WriteBatchBounded caller (which never checks stopping), could +// then send on the closed channel and panic the process during shutdown. + +func newShutdownRaceStore(t *testing.T) *Store { + t.Helper() + dir := t.TempDir() + cfg := DefaultConfig(dir) + cfg.DBPath = filepath.Join(dir, "metrics-shutdown-race.db") + cfg.FlushInterval = time.Hour + store, err := NewStore(cfg) + if err != nil { + t.Fatalf("NewStore returned error: %v", err) + } + return store +} + +func shutdownRaceMetric() bufferedMetric { + return bufferedMetric{ + resourceType: "vm", + resourceID: "shutdown-race", + metricType: "cpu", + value: 1, + timestamp: time.Unix(1_700_000_000, 0), + tier: TierRaw, + } +} + +// A send that reaches the ingestion channel after the worker has finished +// shutdown must not panic. Before the fix the worker closed writeCh, so this +// deterministic post-Close enqueue hit a closed channel. +func TestEnqueueWriteAfterShutdownDoesNotPanic(t *testing.T) { + store := newShutdownRaceStore(t) + if err := store.Close(); err != nil { + t.Fatalf("Close returned error: %v", err) + } + + defer func() { + if r := recover(); r != nil { + t.Fatalf("enqueueWrite panicked after shutdown: %v", r) + } + }() + store.enqueueWrite(writeRequest{metrics: []bufferedMetric{shutdownRaceMetric()}}) +} + +// Concurrent writers racing a Close exercise the real window: WriteBatchSync +// and WriteBatchBounded do not consult stopping, and WriteWithTier checks it +// before releasing bufferMu. None of them may panic. +func TestConcurrentWritesDuringShutdownDoNotPanic(t *testing.T) { + store := newShutdownRaceStore(t) + + var wg sync.WaitGroup + stop := make(chan struct{}) + + for i := 0; i < 8; i++ { + wg.Add(1) + go func() { + defer wg.Done() + for { + select { + case <-stop: + return + default: + } + store.WriteWithTier("vm", "shutdown-race", "cpu", 1, time.Unix(1_700_000_000, 0), TierRaw) + store.WriteBatchBounded([]WriteMetric{{ + ResourceType: "vm", + ResourceID: "shutdown-race", + MetricType: "cpu", + Value: 1, + Timestamp: time.Unix(1_700_000_000, 0), + Tier: TierRaw, + }}) + } + }() + } + + time.Sleep(50 * time.Millisecond) + if err := store.Close(); err != nil { + t.Fatalf("Close returned error: %v", err) + } + close(stop) + wg.Wait() +}