Split bounded pipeline metrics writes from synchronous batch writes

6b79aa997 bounded WriteBatchSync itself, which broke its read-your-writes
contract on slow disks: CI's metrics write-amplification and 500-node
load tests count committed rows after writing, and mock seeding reads
store coverage straight back, so the 2-second early return failed both
(runs 31475700902, 31494553977). Fast local disks masked it.

WriteBatchSync returns to a full commit wait. The monitoring pipeline's
four sync sites move to WriteBatchBounded, which carries the bounded
enqueue-plus-wait semantics, so the #1437 slow-disk stall fix stays
exactly where the hazard is. Both paths share prepareWriteBatch
validation, and a new regression test pins WriteBatchSync waiting past
the bounded budget.

Refs #1437

Contract-Neutral: behavioral fix: split bounded pipeline writes from synchronous batch writes, restores read-your-writes (#1437 follow-up), no public contract delta
This commit is contained in:
rcourtman 2026-08-11 14:42:42 +01:00
parent b6da28da0a
commit 2bc4ed7254
3 changed files with 113 additions and 36 deletions

View file

@ -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)
}
}

View file

@ -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
}

View file

@ -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")
}
}