Skip to content
Merged
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
176 changes: 141 additions & 35 deletions data_plane/src/precompute_engine/worker.rs
Original file line number Diff line number Diff line change
Expand Up @@ -60,25 +60,27 @@ struct GroupState {
/// Per-group watermark: tracks the maximum timestamp seen across all
/// series in this group on this worker.
previous_watermark_ms: i64,
/// Wall-clock-time (ms since epoch) at which each currently-open pane
/// was first opened. Used by `flush_all`'s wall-clock fallback to
/// close panes that have been alive too long when event-time has
/// stagnated. Keyed by `pane_start_ms` and covers entries in BOTH
/// Wall-clock time (ms since epoch) at which each currently-open pane
/// was last touched by an input. Used by `flush_all`'s wall-clock fallback
/// to close panes that have been idle too long when event-time has
/// stagnated. Refreshing on every touch prevents an active bulk ingest
/// with a fixed event timestamp from being force-closed mid-stream.
/// Keyed by `pane_start_ms` and covers entries in BOTH
/// `active_panes` and `sketch_panes` (a single pane_start may have
/// either or both populated). Entries are GC'd by
/// `prune_pane_wall_clock_starts` after each window-close cycle.
pane_wall_clock_starts_ms: BTreeMap<i64, i64>,
/// `prune_pane_wall_clock_last_touches` after each window-close cycle.
pane_wall_clock_last_touches_ms: BTreeMap<i64, i64>,
}

impl GroupState {
/// Drop wall-clock-start entries whose pane no longer exists in
/// Drop wall-clock-last-touch entries whose pane no longer exists in
/// either pane map. Called after window-close cycles in
/// `process_group_samples`, `process_accumulator_input`, and
/// `flush_all` so the bookkeeping doesn't leak as panes turn over.
fn prune_pane_wall_clock_starts(&mut self) {
fn prune_pane_wall_clock_last_touches(&mut self) {
let active = &self.active_panes;
let sketch = &self.sketch_panes;
self.pane_wall_clock_starts_ms
self.pane_wall_clock_last_touches_ms
.retain(|ps, _| active.contains_key(ps) || sketch.contains_key(ps));
}
}
Expand Down Expand Up @@ -348,7 +350,7 @@ impl Worker {
active_panes: BTreeMap::new(),
sketch_panes: BTreeMap::new(),
previous_watermark_ms: i64::MIN,
pane_wall_clock_starts_ms: BTreeMap::new(),
pane_wall_clock_last_touches_ms: BTreeMap::new(),
};
self.group_states.insert(sid, gs);
self.group_count
Expand Down Expand Up @@ -453,13 +455,12 @@ impl Worker {
}

// Normal path: route sample to its single pane accumulator.
// Record the pane's wall-clock birth time the first time
// we touch it (so the wall-clock fallback in `flush_all`
// can age it out even if event-time freezes).
// Refresh the pane's wall-clock last-touch time so the fallback
// only closes an idle pane, not a long-running bulk ingest whose
// records share one event timestamp.
state
.pane_wall_clock_starts_ms
.entry(pane_start)
.or_insert(now_ms);
.pane_wall_clock_last_touches_ms
.insert(pane_start, now_ms);
let updater = state
.active_panes
.entry(pane_start)
Expand Down Expand Up @@ -489,7 +490,7 @@ impl Worker {
}

state.previous_watermark_ms = current_wm;
state.prune_pane_wall_clock_starts();
state.prune_pane_wall_clock_last_touches();

// Emit to output sink
if !emit_batch.is_empty() {
Expand Down Expand Up @@ -591,14 +592,11 @@ impl Worker {
return Ok(());
}

// Record the pane's wall-clock birth time the first time we
// touch it (so the wall-clock fallback in `flush_all` can age
// it out even if event-time freezes — e.g. agents stamping
// every sketch with the same `time_unix_nano`).
// Refresh the pane's wall-clock last-touch time so an active sketch
// stream with a fixed event timestamp is not force-closed mid-ingest.
state
.pane_wall_clock_starts_ms
.entry(pane_start)
.or_insert(now_ms);
.pane_wall_clock_last_touches_ms
.insert(pane_start, now_ms);

// Merge into the sketch pane covering this timestamp.
match state.sketch_panes.remove(&pane_start) {
Expand Down Expand Up @@ -649,7 +647,7 @@ impl Worker {
}

state.previous_watermark_ms = current_wm;
state.prune_pane_wall_clock_starts();
state.prune_pane_wall_clock_last_touches();

if !emit_batch.is_empty() {
debug!(
Expand Down Expand Up @@ -801,14 +799,14 @@ impl Worker {
// forever — the 30s window never closes and ASAP-tier
// queries come back empty even though sketches keep
// arriving (sweep blocker #2). Force `effective_wm` past
// `pane_start + window_size_ms` for any pane older than
// `pane_start + window_size_ms` for any pane untouched for
// `window_size + grace` of WALL-CLOCK time. Set
// `wall_clock_grace_period_ms <= 0` to opt out and keep
// strict event-time semantics.
if grace_ms > 0 {
let window_size_ms = state.window_manager.window_size_ms();
for (&pane_start, &pane_birth_ms) in &state.pane_wall_clock_starts_ms {
if now_ms.saturating_sub(pane_birth_ms) >= window_size_ms + grace_ms {
for (&pane_start, &pane_last_touch_ms) in &state.pane_wall_clock_last_touches_ms {
if now_ms.saturating_sub(pane_last_touch_ms) >= window_size_ms + grace_ms {
let force_to = pane_start.saturating_add(window_size_ms);
if force_to > effective_wm {
effective_wm = force_to;
Expand Down Expand Up @@ -859,7 +857,7 @@ impl Worker {
state.previous_watermark_ms = effective_wm;
}

state.prune_pane_wall_clock_starts();
state.prune_pane_wall_clock_last_touches();
}

if !emit_batch.is_empty() {
Expand Down Expand Up @@ -959,7 +957,7 @@ impl Worker {
if force_wm > state.previous_watermark_ms {
state.previous_watermark_ms = force_wm;
}
state.prune_pane_wall_clock_starts();
state.prune_pane_wall_clock_last_touches();
}

if !emit_batch.is_empty() {
Expand Down Expand Up @@ -2664,8 +2662,8 @@ aggregations:
// per_key store, and ASAP-tier queries come back empty even though
// `worker_process_accumulator` keeps logging.
//
// The fix tracks each pane's wall-clock birth time and force-closes its
// window once `now - pane_birth >= window_size + grace`. This test
// The fix tracks each pane's wall-clock last-touch time and force-closes
// its window once `now - last_touch >= window_size + grace`. This test
// injects a fake clock so it runs in milliseconds instead of needing
// `std::thread::sleep(35s)`.
// -----------------------------------------------------------------------
Expand All @@ -2688,7 +2686,11 @@ aggregations:
make_hot_reload(agg_configs),
WorkerRuntimeConfig {
max_buffer_per_series: 10_000,
allowed_lateness_ms: 0,
// `flush_all` advances a stagnant event-time watermark by
// 1ms. Permit that one boundary nudge so a subsequent
// same-timestamp touch exercises wall-clock idleness rather
// than the unrelated late-data path.
allowed_lateness_ms: 1,
pass_raw_samples: false,
raw_mode_aggregation_id: 0,
late_data_policy: LateDataPolicy::Drop,
Expand Down Expand Up @@ -2753,7 +2755,7 @@ aggregations:
assert_eq!(
sink.len(),
0,
"flush at t_wall=pane_birth must not close the window — fallback fires only after grace"
"flush at t_wall=last_touch must not close the window — fallback fires only after grace"
);

// Advance fake wall-clock by exactly window_size + grace = 35s.
Expand Down Expand Up @@ -2800,7 +2802,7 @@ aggregations:
// Calling flush_all again with no new data must NOT re-emit
// the same window — once a window has been closed via the
// fallback, its pane is drained from `sketch_panes` and from
// `pane_wall_clock_starts_ms`, so there's nothing left to
// `pane_wall_clock_last_touches_ms`, so there's nothing left to
// close. This pins the monotonicity invariant: emitted
// windows have monotonically non-decreasing close times and
// each window is emitted at most once.
Expand All @@ -2812,6 +2814,110 @@ aggregations:
);
}

#[test]
fn wall_clock_fallback_does_not_close_active_raw_ingest() {
let cfg = make_agg_config(
7,
"netflow_bytes",
AggregationType::SingleSubpopulation,
"Sum",
1,
0,
vec![],
);
let sink = Arc::new(CapturingOutputSink::new());
let mut worker = make_worker_with_grace(HashMap::from([(7, cfg)]), sink.clone(), 5_000);
let wall_clock = Arc::new(AtomicI64::new(1_000_000));
let wc_clone = wall_clock.clone();
worker.set_now_ms_fn(Box::new(move || wc_clone.load(Ordering::Relaxed)));

let pf = PolicyFingerprint(7);
let mut expected_sum = 0.0;
for i in 0..7 {
wall_clock.store(1_000_000 + i * 1_000, Ordering::Relaxed);
let value = 1.0 + i as f64;
expected_sum += value;
worker
.process_group_samples(7, pf, "", group_samples("netflow_bytes", vec![(0, value)]))
.unwrap();
}

// The pane is older than the old creation-time deadline, but was
// touched only 500ms ago. Active ingest must keep it open.
wall_clock.store(1_006_500, Ordering::Relaxed);
worker.flush_all().unwrap();
assert_eq!(sink.len(), 0, "active raw pane was force-closed");

wall_clock.store(1_007_000, Ordering::Relaxed);
expected_sum += 8.0;
worker
.process_group_samples(7, pf, "", group_samples("netflow_bytes", vec![(0, 8.0)]))
.unwrap();

wall_clock.store(1_013_001, Ordering::Relaxed);
worker.flush_all().unwrap();
let captured = sink.drain();
assert_eq!(captured.len(), 1, "idle raw pane must close once");
let sum = captured[0]
.1
.as_any()
.downcast_ref::<SumAccumulator>()
.expect("must emit SumAccumulator");
assert!((sum.sum - expected_sum).abs() < 1e-9);
}

#[test]
fn wall_clock_fallback_does_not_close_active_sketch_ingest() {
let cfg = make_agg_config(
8,
"latency",
AggregationType::DDSketch,
"",
1,
0,
vec!["zone"],
);
let sink = Arc::new(CapturingOutputSink::new());
let mut worker = make_worker_with_grace(HashMap::from([(8, cfg)]), sink.clone(), 5_000);
let wall_clock = Arc::new(AtomicI64::new(2_000_000));
let wc_clone = wall_clock.clone();
worker.set_now_ms_fn(Box::new(move || wc_clone.load(Ordering::Relaxed)));

let pf = PolicyFingerprint(8);
for i in 0..7 {
wall_clock.store(2_000_000 + i * 1_000, Ordering::Relaxed);
worker
.process_accumulator_input(
80,
pf,
"us-east",
0,
Box::new(make_ddsketch(0.01, &[1.0 + i as f64])),
)
.unwrap();
}

wall_clock.store(2_006_500, Ordering::Relaxed);
worker.flush_all().unwrap();
assert_eq!(sink.len(), 0, "active sketch pane was force-closed");

wall_clock.store(2_007_000, Ordering::Relaxed);
worker
.process_accumulator_input(80, pf, "us-east", 0, Box::new(make_ddsketch(0.01, &[8.0])))
.unwrap();

wall_clock.store(2_013_001, Ordering::Relaxed);
worker.flush_all().unwrap();
let captured = sink.drain();
assert_eq!(captured.len(), 1, "idle sketch pane must close once");
let dd = captured[0]
.1
.as_any()
.downcast_ref::<DDSketchAccumulator>()
.expect("must emit DDSketchAccumulator");
assert_eq!(dd.inner.total_count(), 8);
}

/// Pin the wall-clock-fallback opt-out: setting
/// `wall_clock_grace_period_ms = 0` reverts to event-time-only
/// semantics, matching pre-fix behaviour. This keeps
Expand Down
Loading