diff --git a/internal/monitoring/monitor.go b/internal/monitoring/monitor.go index e50f2f007..9f231a364 100644 --- a/internal/monitoring/monitor.go +++ b/internal/monitoring/monitor.go @@ -5029,7 +5029,7 @@ func (m *Monitor) syncUnifiedAgentMetrics(store ResourceStoreInterface) { } } if len(storeWrites) > 0 { - m.metricsStore.WriteBatchSync(storeWrites) + m.metricsStore.WriteBatchBounded(storeWrites) } } @@ -5144,7 +5144,7 @@ func (m *Monitor) syncUnifiedVMMetrics(store ResourceStoreInterface) { } } if len(storeWrites) > 0 { - m.metricsStore.WriteBatchSync(storeWrites) + m.metricsStore.WriteBatchBounded(storeWrites) } } @@ -5242,7 +5242,7 @@ func (m *Monitor) syncUnifiedStorageMetrics(store ResourceStoreInterface) { } } if len(storeWrites) > 0 { - m.metricsStore.WriteBatchSync(storeWrites) + m.metricsStore.WriteBatchBounded(storeWrites) } } @@ -5433,7 +5433,7 @@ func (m *Monitor) syncUnifiedAppContainerMetrics(store ResourceStoreInterface) { } } if len(storeWrites) > 0 { - m.metricsStore.WriteBatchSync(storeWrites) + m.metricsStore.WriteBatchBounded(storeWrites) } } diff --git a/pkg/metrics/store.go b/pkg/metrics/store.go index 6f3dc713f..57ae21f7e 100644 --- a/pkg/metrics/store.go +++ b/pkg/metrics/store.go @@ -753,13 +753,11 @@ func (s *Store) WriteWithTier(resourceType, resourceID, metricType string, value s.enqueueWrite(writeRequest{metrics: toWrite}) } -// WriteBatchSync bypasses the in-memory sample buffer but still serializes the -// batch through the ingestion worker. Platform pollers call this concurrently; -// letting every caller open its own SQLite transaction exhausts the shared -// reader pool while those transactions wait on the single WAL writer lock. -func (s *Store) WriteBatchSync(metrics []WriteMetric) { +// prepareWriteBatch validates and normalizes a caller batch, logging and +// dropping invalid entries. Shared by the synchronous and bounded batch paths. +func (s *Store) prepareWriteBatch(metrics []WriteMetric) []bufferedMetric { if len(metrics) == 0 { - return + return nil } batch := make([]bufferedMetric, 0, len(metrics)) @@ -796,13 +794,36 @@ func (s *Store) WriteBatchSync(metrics []WriteMetric) { Int("dropped", droppedInvalid). Msg("Dropped invalid metrics writes from batch") } + return batch +} + +// WriteBatchSync bypasses the in-memory sample buffer but still serializes the +// batch through the ingestion worker, waiting for the commit however long it +// takes. Callers rely on read-your-writes: mock seeding reads store coverage +// straight back, and the write-path invariant tests count committed rows. The +// live monitoring pipeline must NOT use this — it calls WriteBatchBounded so a +// slow metrics disk can never stall polling (#1437). +func (s *Store) WriteBatchSync(metrics []WriteMetric) { + batch := s.prepareWriteBatch(metrics) if len(batch) == 0 { return } - s.enqueueAndWait(writeRequest{metrics: batch}) } +// WriteBatchBounded is the monitoring-pipeline variant of WriteBatchSync: it +// hands the batch to the ingestion worker but never blocks the caller past +// syncWriteWaitTimeout. Within the budget it behaves like WriteBatchSync; past +// it the batch stays queued (or, if the queue cannot even accept it, is +// dropped with a warning) and the caller moves on. +func (s *Store) WriteBatchBounded(metrics []WriteMetric) { + batch := s.prepareWriteBatch(metrics) + if len(batch) == 0 { + return + } + s.boundedEnqueueAndWait(writeRequest{metrics: batch}) +} + func (s *Store) enqueueMaintenance(run func()) { if run == nil { return @@ -917,26 +938,46 @@ func (s *Store) enqueueWrite(req writeRequest) { } } -// syncWriteWaitTimeout bounds how long a synchronous batch write blocks on the -// ingestion worker. The monitoring pipeline calls WriteBatchSync inline (state -// broadcast, agent ingest, poll publish), so an unbounded wait lets a slow -// metrics disk starve polling entirely: on #1437's instance SQLite commits ran -// for seconds to minutes and the monitor froze after its first cycle while -// history writes queued behind retention maintenance. Within the budget the -// call keeps read-your-writes; past it the caller moves on and history lands -// whenever the worker catches up. +// syncWriteWaitTimeout bounds how long a WriteBatchBounded call blocks on the +// ingestion worker. The monitoring pipeline writes inline (state broadcast, +// agent ingest, poll publish), so an unbounded wait lets a slow metrics disk +// starve polling entirely: on #1437's instance SQLite commits ran for seconds +// to minutes and the monitor froze after its first cycle while history writes +// queued behind retention maintenance. Within the budget the call keeps +// read-your-writes; past it the caller moves on and history lands whenever +// the worker catches up. const syncWriteWaitTimeout = 2 * time.Second // enqueueAndWait hands the batch to the ingestion worker and waits for its -// commit, but never longer than syncWriteWaitTimeout in total. If the queue -// cannot even accept the batch within the budget the batch is dropped, matching -// enqueueWrite's saturation behavior. If only the commit is outstanding the -// batch stays queued and is not lost. +// commit with no deadline. WriteBatchSync callers depend on read-your-writes +// regardless of disk speed. Only store shutdown releases the wait early. func (s *Store) enqueueAndWait(req writeRequest) { if req.done == nil { req.done = make(chan struct{}) } + select { + case s.writeCh <- req: + case <-s.stopCh: + return + } + + select { + case <-req.done: + case <-s.stopCh: + } +} + +// boundedEnqueueAndWait hands the batch to the ingestion worker and waits for +// its commit, but never longer than syncWriteWaitTimeout in total. If the +// queue cannot even accept the batch within the budget the batch is dropped, +// matching enqueueWrite's saturation behavior. If only the commit is +// outstanding the batch stays queued and is not lost. +func (s *Store) boundedEnqueueAndWait(req writeRequest) { + if req.done == nil { + req.done = make(chan struct{}) + } + timer := time.NewTimer(syncWriteWaitTimeout) defer timer.Stop() @@ -951,7 +992,7 @@ func (s *Store) enqueueAndWait(req writeRequest) { Int("batch_size", len(req.metrics)). Int("write_queue_depth", len(s.writeCh)). Int("write_queue_capacity", cap(s.writeCh)). - Msg("Metrics write queue saturated, dropping synchronous batch to keep monitoring live") + Msg("Metrics write queue saturated, dropping bounded batch to keep monitoring live") return } diff --git a/pkg/metrics/store_write_backpressure_test.go b/pkg/metrics/store_write_backpressure_test.go index f387e6b4c..561fbbad3 100644 --- a/pkg/metrics/store_write_backpressure_test.go +++ b/pkg/metrics/store_write_backpressure_test.go @@ -5,10 +5,12 @@ import ( "time" ) -// Regression coverage for #1437: WriteBatchSync sits on the monitoring +// Regression coverage for #1437: WriteBatchBounded sits on the monitoring // pipeline (state broadcast, agent ingest, poll publish). When the ingestion // worker cannot keep up with the disk, the call must return within its wait -// budget instead of stalling the monitor until polling stops. +// budget instead of stalling the monitor until polling stops. WriteBatchSync +// keeps its full commit wait — mock seeding and the write-path invariant +// tests read their writes straight back. // newUnservicedStore builds a bare store whose ingestion worker never runs, // modelling a writer wedged behind slow SQLite maintenance. @@ -30,57 +32,91 @@ func backpressureProbeMetric() WriteMetric { } } -func requireReturnsWithinWaitBudget(t *testing.T, s *Store, label string) { +func requireBoundedReturn(t *testing.T, s *Store, label string) { t.Helper() done := make(chan struct{}) go func() { - s.WriteBatchSync([]WriteMetric{backpressureProbeMetric()}) + s.WriteBatchBounded([]WriteMetric{backpressureProbeMetric()}) close(done) }() select { case <-done: case <-time.After(syncWriteWaitTimeout + 3*time.Second): - t.Fatalf("WriteBatchSync stalled past its wait budget (%s)", label) + t.Fatalf("WriteBatchBounded stalled past its wait budget (%s)", label) } } -func TestWriteBatchSyncReturnsWhenQueueSaturated(t *testing.T) { +func TestWriteBatchBoundedReturnsWhenQueueSaturated(t *testing.T) { s := newUnservicedStore(1) s.writeCh <- writeRequest{metrics: []bufferedMetric{{}}} - requireReturnsWithinWaitBudget(t, s, "saturated queue") + requireBoundedReturn(t, s, "saturated queue") if got := len(s.writeCh); got != 1 { t.Fatalf("saturated queue depth = %d, want the pre-existing batch only", got) } } -func TestWriteBatchSyncReturnsWhenCommitLags(t *testing.T) { +func TestWriteBatchBoundedReturnsWhenCommitLags(t *testing.T) { s := newUnservicedStore(4) - requireReturnsWithinWaitBudget(t, s, "lagging commit") + requireBoundedReturn(t, s, "lagging commit") if got := len(s.writeCh); got != 1 { t.Fatalf("queued batches = %d, want 1 (batch must stay queued, not be dropped)", got) } } -func TestWriteBatchSyncReturnsOnClosedStore(t *testing.T) { +func TestWriteBatchBoundedReturnsOnClosedStore(t *testing.T) { s := newUnservicedStore(1) s.writeCh <- writeRequest{metrics: []bufferedMetric{{}}} close(s.stopCh) done := make(chan struct{}) go func() { - s.WriteBatchSync([]WriteMetric{backpressureProbeMetric()}) + s.WriteBatchBounded([]WriteMetric{backpressureProbeMetric()}) close(done) }() select { case <-done: case <-time.After(time.Second): - t.Fatal("WriteBatchSync did not observe store shutdown") + t.Fatal("WriteBatchBounded did not observe store shutdown") + } +} + +// WriteBatchSync must keep waiting for the commit well past the bounded +// budget — an early return would break read-your-writes for mock seeding and +// the write-path invariant tests on slow disks, which is exactly what a +// bounded WriteBatchSync did to CI (run 31494553977). +func TestWriteBatchSyncWaitsForCommitBeyondBoundedBudget(t *testing.T) { + s := newUnservicedStore(4) + + returned := make(chan struct{}) + go func() { + s.WriteBatchSync([]WriteMetric{backpressureProbeMetric()}) + close(returned) + }() + + select { + case <-returned: + t.Fatal("WriteBatchSync returned before its batch committed") + case <-time.After(syncWriteWaitTimeout + time.Second): + } + + // Simulate the worker committing the queued batch: closing done must + // release the waiting caller. + req := <-s.writeCh + if req.done == nil { + t.Fatal("queued request carries no done channel") + } + close(req.done) + + select { + case <-returned: + case <-time.After(time.Second): + t.Fatal("WriteBatchSync did not return after its commit completed") } }