mirror of
https://github.com/rcourtman/Pulse.git
synced 2026-10-03 04:38:48 +00:00
fix(metrics): stop closing writeCh on shutdown
The ingestion worker closed writeCh when it observed stopCh. A concurrent WriteWithTier that passed its stopping check, or a WriteBatchSync/WriteBatchBounded caller (which never checks stopping), could then send on the closed channel and panic the process during shutdown. Reproduced deterministically: enqueueWrite after Close panics with send on closed channel, and 8 concurrent monitoring writers racing Close panic in boundedEnqueueAndWait (store.go:1034). Never close writeCh. Drain already-queued requests non-blockingly, then process them. Late writes land in the buffered channel and are discarded with the store rather than crashing it. Record the invariant in the performance-and-scalability contract and pin it in the accepted store proof file. Change-source: pulse-maintainer
This commit is contained in:
parent
54f0619e57
commit
5bc6f919bb
3 changed files with 112 additions and 5 deletions
|
|
@ -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.
|
||||
|
|
|
|||
|
|
@ -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 {
|
||||
|
|
|
|||
|
|
@ -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()
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue