From 1e8d619448d29421f5ad63c601539a3eecfbc918 Mon Sep 17 00:00:00 2001 From: zz_y Date: Fri, 4 Sep 2026 05:42:08 -0600 Subject: [PATCH] fix(precompute): enforce envelope series cap --- asap-precompute-go/fixes_review_test.go | 31 +++++++++++++++++ asap-precompute-go/window.go | 44 ++++++++++++++++++------- 2 files changed, 63 insertions(+), 12 deletions(-) diff --git a/asap-precompute-go/fixes_review_test.go b/asap-precompute-go/fixes_review_test.go index cd2718e6..72361705 100644 --- a/asap-precompute-go/fixes_review_test.go +++ b/asap-precompute-go/fixes_review_test.go @@ -146,3 +146,34 @@ func TestObserveEnvelope_DeltaDoubleCountGuard(t *testing.T) { t.Fatalf("same-range re-delivery: want %q, got %q", "BBBCCC", string(fs.state)) } } + +func TestObserveEnvelope_EvictOldestHonorsSeriesCap(t *testing.T) { + cfg := &PrecomputeConfig{AggID: 1, SketchType: SketchTypeDDSketch, Mode: Tumbling, + Window: WindowSpec{Size: 10 * time.Second}, MaxSeries: 1, OnOverflow: OnOverflowEvictOldest} + p := New(cfg, newFakeFactory(), &fakeObserver{}).(*precompute) + mk := func(label string, timestamp uint64) *SketchEnvelope { + return &SketchEnvelope{SchemaVersion: 1, SketchType: SketchTypeDDSketch, AggID: 1, + Labels: []KeyValue{{Key: "k", Value: label}}, WindowStartMs: 0, + WindowEndMs: timestamp, Encoding: EncodingProtoDelta, Payload: []byte(label)} + } + if err := p.ObserveEnvelope(mk("a", 1_000)); err != nil { + t.Fatal(err) + } + keyA := cfg.SeriesKeyForEntry(nil, []KeyValue{{Key: "k", Value: "a"}}) + p.snapshotCache.CacheInbound(keyA, []byte("cached")) + if err := p.ObserveEnvelope(mk("b", 2_000)); err != nil { + t.Fatal(err) + } + if got := p.window.activeSeriesCount(); got != 1 { + t.Fatalf("active series=%d, want 1", got) + } + if _, exists := p.window.series[keyA]; exists { + t.Fatal("oldest envelope series was not evicted") + } + if got := p.snapshotCache.LenInbound(); got != 1 { + t.Fatalf("inbound cache=%d, want only new series", got) + } + if got := p.Stats().Snapshot().ActiveSeries; got != 1 { + t.Fatalf("active-series stat=%d", got) + } +} diff --git a/asap-precompute-go/window.go b/asap-precompute-go/window.go index 0665d468..d3d54a86 100644 --- a/asap-precompute-go/window.go +++ b/asap-precompute-go/window.go @@ -356,18 +356,8 @@ func (w *windowState) admitSeriesLocked( // blocking the runtime hot path. return nil, ErrSeriesCapExceeded case OnOverflowEvictOldest: - var ( - oldestKey string - oldestMs uint64 = ^uint64(0) - ) - for k, e := range w.series { - if e.LastSeenMs < oldestMs { - oldestMs = e.LastSeenMs - oldestKey = k - } - } + oldestKey, _ := w.evictOldestLocked() if oldestKey != "" { - delete(w.series, oldestKey) if stats != nil { stats.ActiveSeries.Add(-1) } @@ -418,6 +408,22 @@ func (w *windowState) admitSeriesLocked( return entry, nil } +func (w *windowState) evictOldestLocked() (string, *seriesEntry) { + var oldestKey string + oldestMs := ^uint64(0) + for key, entry := range w.series { + if entry.LastSeenMs < oldestMs { + oldestMs = entry.LastSeenMs + oldestKey = key + } + } + entry := w.series[oldestKey] + if oldestKey != "" { + delete(w.series, oldestKey) + } + return oldestKey, entry +} + // recordLocked feeds one observation into a series' sketch and advances its // bookkeeping. Caller holds w.mu. Shared by observe and observeKeyed. func (w *windowState) recordLocked(entry *seriesEntry, obs *Observation, observer SketchObserver) error { @@ -483,8 +489,22 @@ func (w *windowState) observeEnvelope( entry, ok := w.series[key] if !ok { if cfg.MaxSeries > 0 && uint64(len(w.series)) >= cfg.MaxSeries { - if cfg.OnOverflow == OnOverflowDrop || cfg.OnOverflow == OnOverflowBlock { + switch cfg.OnOverflow { + case OnOverflowDrop, OnOverflowBlock: return ErrSeriesCapExceeded + case OnOverflowEvictOldest: + evictedKey, evicted := w.evictOldestLocked() + if evictedKey != "" { + if snapshotCache != nil { + snapshotCache.Delete(evictedKey) + } + if evicted != nil && evicted.Sketch != nil { + evicted.Sketch.Reset() + } + if stats != nil { + stats.ActiveSeries.Add(-1) + } + } } } sketch := sketchFactory()