Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,8 @@ func (p *asapEdgeProcessor) flushLoop() {
return
case <-subC:
p.flushSubWindow(context.Background())
case <-p.wakeCh:
p.flushSubWindow(context.Background())
case <-t.C:
p.flushAll(context.Background())
}
Expand All @@ -67,6 +69,11 @@ func (p *asapEdgeProcessor) flushLoop() {
// each tick must cover every active series; threshold-gated emits
// are small).
p.flushSubWindow(context.Background())
case <-p.wakeCh:
// Out-of-cycle wake (e.g. an insert-time GOS threshold crossing).
// Same handler as subC: emits whatever's currently divergent/dirty
// across every shard, just triggered early instead of by the timer.
p.flushSubWindow(context.Background())
case <-t.C:
shardIdx := tick % n
p.flushShardWarmCold(context.Background(), shardIdx)
Expand All @@ -75,6 +82,20 @@ func (p *asapEdgeProcessor) flushLoop() {
}
}

// wakeSubWindow requests an out-of-cycle sub-window flush — e.g. an
// insert-time GOS threshold crossing that shouldn't wait for the next
// SubWindowInterval tick (or, for a GOS-only family with no sub-window
// ticker configured at all, that would otherwise have no flush path short of
// window close). Non-blocking: if a wake is already pending, this is a
// no-op — the pending flushSubWindow call will pick up every series'
// current dirty state, including whatever just crossed threshold.
func (p *asapEdgeProcessor) wakeSubWindow() {
select {
case p.wakeCh <- struct{}{}:
default:
}
}

// subWindowEnabled reports whether the threshold-driven sub-window producer is
// active: a positive interval shorter than the window (validated at config).
func (p *asapEdgeProcessor) subWindowEnabled() bool {
Expand All @@ -83,10 +104,17 @@ func (p *asapEdgeProcessor) subWindowEnabled() bool {

// flushSubWindow fires a divergence-gated sub-window delta emit for every
// shard's sketch aggregators and forwards the result, WITHOUT rotating windows.
//
// No p.subWindowEnabled() guard here: this now runs from two triggers (the
// legacy subC ticker, gated at the call site by subWindowEnabled(), and the
// wakeCh out-of-cycle signal, which is NOT gated by it — a GOS-driven family
// must be able to wake a flush even with SubWindowInterval unset). Per-series
// gating happens inside sa.emitSubWindow / s.subWindowEnabled(); that gate
// still requires SubWindowInterval>0 today, so a GOS-only family with no
// sub-window interval configured won't yet see any effect from a wake — that
// per-series gate is next in line to be decoupled from SubWindowInterval as
// GOS families land (tracked starting with the CountSketch conversion).
func (p *asapEdgeProcessor) flushSubWindow(ctx context.Context) {
if !p.subWindowEnabled() {
return
}
nowMs := uint64(time.Now().UnixMilli())
out := pmetric.NewMetrics()
for _, sh := range p.shards {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -83,6 +83,14 @@ type asapEdgeProcessor struct {
doneCh chan struct{}
flushStarted bool

// wakeCh requests an out-of-cycle sub-window flush (see wakeSubWindow in
// flush.go). Buffered 1 and drained non-blocking: a pending wake already
// covers any crossing that arrives before the flush loop gets to it, so
// callers never block on send. Always present, independent of whether
// SubWindowInterval/subC is configured — a family with insert-time GOS
// detection wakes the loop regardless of the legacy sub-window ticker.
wakeCh chan struct{}

// ctrlChan is the optional control-plane poll channel (nil when the
// ControlChannel config block is unset). When set, Start() spawns a poll
// loop that applies received config updates to the live Precompute
Expand All @@ -109,6 +117,7 @@ func newProcessor(cfg *Config, set processor.Settings, next consumer.Metrics) (*
coldExtLabels: cfg.Cold.ExternalLabels,
stopCh: make(chan struct{}),
doneCh: make(chan struct{}),
wakeCh: make(chan struct{}, 1),
}
for i := range cfg.Metrics {
m := &cfg.Metrics[i]
Expand Down Expand Up @@ -165,6 +174,9 @@ func newProcessor(cfg *Config, set processor.Settings, next consumer.Metrics) (*
edgeID: cfg.EdgeID,
subWindowInterval: cfg.SubWindowInterval,
subWindowEpsilon: cfg.SubWindowEpsilon,
gosDeltaEpsilon: fam.GosDeltaEpsilon,
gosSites: fam.GosSites,
gosAnisotropic: fam.GosAnisotropic,
}
if sa, ok := newSketchAggregator(name, fam, opts, p.logger); ok {
sa.procDropCount = &p.sketchDropCount
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,81 @@
// Copyright The OpenTelemetry Authors
// SPDX-License-Identifier: Apache-2.0

package asapedgeprocessor

import (
"context"
"testing"
"time"

"go.opentelemetry.io/collector/component"
"go.opentelemetry.io/collector/pdata/pcommon"
"go.opentelemetry.io/collector/pdata/pmetric"
"go.opentelemetry.io/collector/processor"
"go.uber.org/zap"
)

// TestWakeSubWindowFlushesBeforeTicker proves wakeSubWindow triggers a
// sub-window emit immediately, independent of the SubWindowInterval ticker's
// own cadence — the mechanism the (not-yet-landed) per-family GOS insert-time
// checks will call. SubWindowInterval is set far longer than the test would
// ever run, so any output observed can only have come from the wake, not the
// ticker firing on its own.
func TestWakeSubWindowFlushesBeforeTicker(t *testing.T) {
tru := true
cap := &capMetrics{}
cfg := &Config{
ShardCount: 1,
WindowDuration: time.Hour,
SubWindowInterval: 30 * time.Minute,
DropOriginal: true,
Metrics: []MetricFamily{
{Metric: "reqs", Family: FamilySum, AggregateBy: []string{"zone"}, DeltaTransmission: &tru},
},
Cold: ColdConfig{Enabled: false},
}
if err := cfg.Validate(); err != nil {
t.Fatal(err)
}
set := processor.Settings{TelemetrySettings: component.TelemetrySettings{Logger: zap.NewNop()}}
p, err := newProcessor(cfg, set, cap)
if err != nil {
t.Fatal(err)
}
if err := p.Start(context.Background(), nil); err != nil {
t.Fatal(err)
}

md := pmetric.NewMetrics()
sm := md.ResourceMetrics().AppendEmpty().ScopeMetrics().AppendEmpty()
m := sm.Metrics().AppendEmpty()
m.SetName("reqs")
s := m.SetEmptySum()
s.SetAggregationTemporality(pmetric.AggregationTemporalityDelta)
dp := s.DataPoints().AppendEmpty()
dp.Attributes().PutStr("zone", "z0")
dp.SetDoubleValue(7)
dp.SetTimestamp(pcommon.NewTimestampFromTime(time.Now()))
if err := p.ConsumeMetrics(context.Background(), md); err != nil {
t.Fatal(err)
}

// Rapid double-wake must not block (buffered-1, non-blocking send).
p.wakeSubWindow()
p.wakeSubWindow()

deadline := time.Now().Add(2 * time.Second)
for len(cap.got) == 0 && time.Now().Before(deadline) {
time.Sleep(5 * time.Millisecond)
}

ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second)
defer cancel()
if err := p.Shutdown(ctx); err != nil {
t.Fatalf("shutdown: %v", err)
}

if len(cap.got) == 0 {
t.Fatal("wakeSubWindow did not produce a flush within 2s; SubWindowInterval (30m) could not have fired on its own")
}
}