Skip to content

Commit 49b381b

Browse files
authored
Merge pull request #376 from ProjectASAP/fix/shutdown-force-close
fix(precompute): force-close open windows on worker shutdown
2 parents 4292b07 + 045dec2 commit 49b381b

1 file changed

Lines changed: 245 additions & 0 deletions

File tree

‎data_plane/src/precompute_engine/worker.rs‎

Lines changed: 245 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -296,6 +296,15 @@ impl Worker {
296296
if let Err(e) = self.flush_all() {
297297
warn!("Worker {} final flush error: {}", self.id, e);
298298
}
299+
// Force-close any windows still open after the final flush.
300+
// `flush_all` only advances the watermark by +1ms (plus the
301+
// wall-clock fallback, whose grace may not have elapsed for a
302+
// one-shot batch), so the trailing window can remain open and
303+
// its data would never reach the store. No more samples will
304+
// arrive after shutdown, so close every remaining pane.
305+
if let Err(e) = self.force_close_all() {
306+
warn!("Worker {} shutdown force-close error: {}", self.id, e);
307+
}
299308
break;
300309
}
301310
}
@@ -865,6 +874,106 @@ impl Worker {
865874
Ok(())
866875
}
867876

877+
/// Force-close every window still open on shutdown.
878+
///
879+
/// Unlike `flush_all` — which only advances the watermark by `+1ms` (plus
880+
/// the wall-clock fallback, gated on grace having elapsed) — this emits the
881+
/// window for every remaining pane unconditionally, because no further
882+
/// samples will arrive once the engine is shutting down. Without it, a
883+
/// one-shot batch whose records all fall in a single window (so event-time
884+
/// never advances past the window end) would leave that window open forever
885+
/// and never write it to the store. Covers both `active_panes` (sample
886+
/// aggregation) and `sketch_panes` (OTLP-delivered sketches).
887+
///
888+
/// To advance past the open windows we use a *finite* bound derived from
889+
/// the largest open pane (`max_pane + window_size_ms`) rather than
890+
/// `i64::MAX`: `WindowManager::closed_windows` enumerates window starts up
891+
/// to `current_wm` one slide at a time, so passing `i64::MAX` would loop
892+
/// ~`i64::MAX / slide` times and overflow. `max_pane + window_size_ms` is
893+
/// the smallest watermark that closes the latest open window.
894+
///
895+
/// Idempotent: closed panes are drained from both pane maps and their
896+
/// wall-clock bookkeeping is pruned, so a second call emits nothing.
897+
fn force_close_all(&mut self) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
898+
if self.pass_raw_samples {
899+
return Ok(());
900+
}
901+
902+
let mut emit_batch: Vec<(PrecomputedOutput, Box<dyn AggregateCore>)> = Vec::new();
903+
904+
for (&sid, state) in &mut self.group_states {
905+
let _ = sid; // sid is the bucket key; group_key/policy_fp live on `state`
906+
if state.previous_watermark_ms == i64::MIN {
907+
continue; // never received data — nothing to close
908+
}
909+
910+
// The latest window start equals the largest open pane start across
911+
// both pane maps; closing `[start, start + size)` needs
912+
// `wm >= start + size`.
913+
let max_active = state.active_panes.keys().next_back().copied();
914+
let max_sketch = state.sketch_panes.keys().next_back().copied();
915+
let max_pane = match (max_active, max_sketch) {
916+
(Some(a), Some(b)) => a.max(b),
917+
(Some(a), None) => a,
918+
(None, Some(b)) => b,
919+
(None, None) => continue, // no open panes
920+
};
921+
let force_wm = max_pane.saturating_add(state.window_manager.window_size_ms());
922+
923+
let group_key = state.group_key.clone();
924+
let closed = state
925+
.window_manager
926+
.closed_windows(state.previous_watermark_ms, force_wm);
927+
928+
for window_start in &closed {
929+
let (_, window_end) = state.window_manager.window_bounds(*window_start);
930+
let pane_starts = state.window_manager.panes_for_window(*window_start);
931+
932+
if let Some(accumulator) =
933+
merge_panes_for_window(&mut state.active_panes, &pane_starts)
934+
{
935+
let key = build_group_key_label_values(&group_key);
936+
let output = PrecomputedOutput::new(
937+
*window_start as u64,
938+
window_end as u64,
939+
Some(key),
940+
PolicyFingerprint::from_config(&state.config),
941+
);
942+
emit_batch.push((output, accumulator));
943+
}
944+
945+
if let Some(accumulator) =
946+
merge_sketch_panes_for_window(&mut state.sketch_panes, &pane_starts)
947+
{
948+
let key = build_group_key_label_values(&group_key);
949+
let output = PrecomputedOutput::new(
950+
*window_start as u64,
951+
window_end as u64,
952+
Some(key),
953+
PolicyFingerprint::from_config(&state.config),
954+
);
955+
emit_batch.push((output, accumulator));
956+
}
957+
}
958+
959+
if force_wm > state.previous_watermark_ms {
960+
state.previous_watermark_ms = force_wm;
961+
}
962+
state.prune_pane_wall_clock_starts();
963+
}
964+
965+
if !emit_batch.is_empty() {
966+
debug!(
967+
"Worker {} shutdown force-close emitting {} outputs",
968+
self.id,
969+
emit_batch.len()
970+
);
971+
self.output_sink.emit_batch(emit_batch)?;
972+
}
973+
974+
Ok(())
975+
}
976+
868977
/// Compute the global watermark as min(all worker watermarks), ignoring
869978
/// workers that haven't started yet (still at i64::MIN).
870979
fn compute_global_watermark(&self) -> i64 {
@@ -2749,4 +2858,140 @@ aggregations:
27492858
"grace=0 must disable the fallback — event-time-only semantics"
27502859
);
27512860
}
2861+
2862+
// -----------------------------------------------------------------------
2863+
// Test: shutdown force-close emits the trailing window
2864+
//
2865+
// The immediate-shutdown batch case: every record falls in one window and
2866+
// no later timestamp ever advances the watermark, so flush_all (with the
2867+
// wall-clock fallback disabled, grace=0) leaves the window open. On
2868+
// shutdown, force_close_all must close and emit it so the data reaches the
2869+
// store instead of being lost. Covers both the sample (`active_panes`) and
2870+
// sketch (`sketch_panes`) paths.
2871+
// -----------------------------------------------------------------------
2872+
2873+
#[test]
2874+
fn shutdown_force_close_emits_trailing_sample_window() {
2875+
// 10s tumbling window; make_worker uses grace=0, isolating the
2876+
// force-close from the wall-clock fallback.
2877+
let config = make_agg_config(
2878+
1,
2879+
"cpu",
2880+
AggregationType::SingleSubpopulation,
2881+
"Sum",
2882+
10,
2883+
0,
2884+
vec![],
2885+
);
2886+
let mut agg_configs = HashMap::new();
2887+
agg_configs.insert(1, config);
2888+
let sink = Arc::new(CapturingOutputSink::new());
2889+
let mut worker = make_worker(agg_configs, sink.clone(), false, 0, LateDataPolicy::Drop);
2890+
2891+
// All samples land in window [0, 10_000); the watermark freezes below
2892+
// the window end because no later timestamp ever arrives.
2893+
let pf = PolicyFingerprint(1);
2894+
for i in 0..5 {
2895+
worker
2896+
.process_group_samples(1, pf, "", group_samples("cpu", vec![(1_000 + i * 100, 1.0)]))
2897+
.unwrap();
2898+
}
2899+
2900+
worker.flush_all().unwrap();
2901+
assert_eq!(
2902+
sink.len(),
2903+
0,
2904+
"trailing window must remain open after the final flush"
2905+
);
2906+
2907+
worker.force_close_all().unwrap();
2908+
let captured = sink.drain();
2909+
assert_eq!(
2910+
captured.len(),
2911+
1,
2912+
"shutdown force-close must emit the trailing window"
2913+
);
2914+
let (output, acc) = &captured[0];
2915+
assert!(!output.policy_fp.is_unset());
2916+
assert_eq!(output.start_timestamp, 0);
2917+
assert_eq!(output.end_timestamp, 10_000);
2918+
let sum_acc = acc
2919+
.as_any()
2920+
.downcast_ref::<SumAccumulator>()
2921+
.expect("should be SumAccumulator");
2922+
assert!(
2923+
(sum_acc.sum - 5.0).abs() < 1e-10,
2924+
"5 samples of 1.0 → sum 5, got {}",
2925+
sum_acc.sum
2926+
);
2927+
2928+
// Idempotent: panes are drained, so a second force-close emits nothing.
2929+
worker.force_close_all().unwrap();
2930+
assert_eq!(
2931+
sink.len(),
2932+
0,
2933+
"force-close must be idempotent once panes are drained"
2934+
);
2935+
}
2936+
2937+
#[test]
2938+
fn shutdown_force_close_emits_trailing_sketch_window() {
2939+
// 30s tumbling window; grace=0 isolates the force-close.
2940+
let cfg = make_agg_config(
2941+
1,
2942+
"http_requests_total_latency_ms_quantile",
2943+
AggregationType::DDSketch,
2944+
"",
2945+
30,
2946+
0,
2947+
vec!["zone"],
2948+
);
2949+
let agg_configs = HashMap::from([(1, cfg)]);
2950+
let sink = Arc::new(CapturingOutputSink::new());
2951+
let mut worker = make_worker_with_grace(agg_configs, sink.clone(), 0);
2952+
2953+
// 10 sketches, all stamped at frozen event-time 0 → window [0, 30_000).
2954+
let pf = PolicyFingerprint(1);
2955+
for i in 0..10 {
2956+
let s = make_ddsketch(0.01, &[1.0 + i as f64]);
2957+
worker
2958+
.process_accumulator_input(51, pf, "us-east", 0, Box::new(s))
2959+
.unwrap();
2960+
}
2961+
2962+
worker.flush_all().unwrap();
2963+
assert_eq!(
2964+
sink.len(),
2965+
0,
2966+
"trailing sketch window must remain open after flush (grace=0, event-time frozen)"
2967+
);
2968+
2969+
worker.force_close_all().unwrap();
2970+
let captured = sink.drain();
2971+
assert_eq!(
2972+
captured.len(),
2973+
1,
2974+
"shutdown force-close must emit the trailing sketch window"
2975+
);
2976+
let (output, acc) = &captured[0];
2977+
assert_eq!(output.start_timestamp, 0);
2978+
assert_eq!(output.end_timestamp, 30_000);
2979+
assert_eq!(acc.type_name(), "DDSketchAccumulator");
2980+
let dd = acc
2981+
.as_any()
2982+
.downcast_ref::<DDSketchAccumulator>()
2983+
.expect("must downcast to DDSketchAccumulator");
2984+
assert_eq!(
2985+
dd.inner.total_count(),
2986+
10,
2987+
"all 10 frozen-time sketches must merge into the single emitted output"
2988+
);
2989+
2990+
worker.force_close_all().unwrap();
2991+
assert_eq!(
2992+
sink.len(),
2993+
0,
2994+
"force-close must be idempotent once panes are drained"
2995+
);
2996+
}
27522997
}

0 commit comments

Comments
 (0)