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
5 changes: 0 additions & 5 deletions data_plane/src/precompute_engine/raw_dag.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
35 changes: 17 additions & 18 deletions data_plane/src/precompute_engine/worker.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -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;
}
Expand Down Expand Up @@ -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<String, (i64, f64)>,
series_key: &str,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -167,21 +167,15 @@ pub fn build_dag_accumulator(
samples: &[RawSample],
) -> Result<Box<dyn AggregateCore>, 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())
}
Loading