diff --git a/data_plane/src/precompute_engine/raw_dag.rs b/data_plane/src/precompute_engine/raw_dag.rs index 5d94407a9..95c0fe0cb 100644 --- a/data_plane/src/precompute_engine/raw_dag.rs +++ b/data_plane/src/precompute_engine/raw_dag.rs @@ -223,11 +223,6 @@ impl RawDagProgram { self.updater().map(|_| ()) } - pub fn uses_counter_delta(&self) -> bool { - // Planner represents rate computation as an explicit upstream operator. - false - } - pub fn apply( &self, updater: &mut dyn AccumulatorUpdater, diff --git a/data_plane/src/precompute_engine/worker.rs b/data_plane/src/precompute_engine/worker.rs index ae9db38fe..0bb5ebaaa 100644 --- a/data_plane/src/precompute_engine/worker.rs +++ b/data_plane/src/precompute_engine/worker.rs @@ -546,15 +546,7 @@ impl Worker { let too_late = previous_event_time != i64::MIN && pane_timestamp(*ts) < watermark_for_event_time(previous_event_time, allowed_lateness_ms); - let value = if state.program.as_deref().map_or_else( - || { - matches!( - state.config.sample_update_rule(), - SampleUpdateRule::CounterDelta { .. } - ) - }, - |p| p.uses_counter_delta(), - ) { + let value = if legacy_counter_delta(state) { reset_aware_counter_delta(&mut state.counter_previous, series_key, *val, *ts) } else { Some(*val) @@ -605,15 +597,7 @@ impl Worker { // Never feed the raw counter value into a membership // heap; the authoritative ExactCounter branch remains // responsible for the visible result. - if state.program.as_deref().map_or_else( - || { - matches!( - state.config.sample_update_rule(), - SampleUpdateRule::CounterDelta { .. } - ) - }, - |p| p.uses_counter_delta(), - ) { + if legacy_counter_delta(state) { if let Some(input) = state.input_revisions.get_mut(&bucket_start) { Arc::make_mut(input).first_revision = 0; } @@ -1612,6 +1596,21 @@ pub(crate) fn apply_sample( /// Convert a cumulative counter sample into a non-negative, reset-aware /// increment. Only the immediately preceding sample per series is retained; /// pane rotation therefore cannot lose the boundary increment. +/// Whether a sample must be converted to a counter delta before it reaches the +/// accumulator. +/// +/// Only the legacy configured update rule does this. A selected Planner program +/// never does: rate is an explicit upstream operator in the DAG, so the sample +/// reaches the accumulator unchanged. Both call sites used to ask the program +/// and were always told `false`. +fn legacy_counter_delta(state: &GroupState) -> bool { + state.program.is_none() + && matches!( + state.config.sample_update_rule(), + SampleUpdateRule::CounterDelta { .. } + ) +} + pub(crate) fn reset_aware_counter_delta( previous: &mut HashMap, series_key: &str, diff --git a/data_plane/src/storage_engines/sketch_db/backfill/window_builder.rs b/data_plane/src/storage_engines/sketch_db/backfill/window_builder.rs index f6dc241ac..ae89f9fd6 100644 --- a/data_plane/src/storage_engines/sketch_db/backfill/window_builder.rs +++ b/data_plane/src/storage_engines/sketch_db/backfill/window_builder.rs @@ -167,21 +167,15 @@ pub fn build_dag_accumulator( samples: &[RawSample], ) -> Result, String> { let mut updater = program.updater()?; - let mut previous = std::collections::HashMap::new(); + // A selected program never converts samples to counter deltas; rate is an + // explicit upstream operator, so each sample is applied as it arrives. for sample in samples { - let value = if program.uses_counter_delta() { - crate::precompute_engine::worker::reset_aware_counter_delta( - &mut previous, - &sample.labels, - sample.value, - sample.timestamp_ms, - ) - } else { - Some(sample.value) - }; - if let Some(value) = value { - program.apply(&mut *updater, &sample.labels, value, sample.timestamp_ms)?; - } + program.apply( + &mut *updater, + &sample.labels, + sample.value, + sample.timestamp_ms, + )?; } Ok(updater.take_accumulator()) }