diff --git a/data_plane/src/precompute_engine/revisions.rs b/data_plane/src/precompute_engine/revisions.rs index 6403ec199..24bee38f4 100644 --- a/data_plane/src/precompute_engine/revisions.rs +++ b/data_plane/src/precompute_engine/revisions.rs @@ -655,18 +655,54 @@ impl RevisionRuntime { now_ms: u64, ranges: &[(StoredOutputId, u64, u64)], expected: &CatalogGeneration, + ) -> Result { + let ranges: Vec<_> = ranges + .iter() + .map(|&(output, start, end)| RevisionQueryRange { + output, + start, + end, + max_lag: 0, + }) + .collect(); + self.query_view_ranges(required, now_ms, &ranges, expected) + } + + fn query_view_ranges( + &self, + required: &BTreeSet, + now_ms: u64, + ranges: &[RevisionQueryRange], + expected: &CatalogGeneration, ) -> Result { let (plan, store) = self.installed()?; if plan.precompute_plan.summary_catalog.as_ref() != Some(expected) { return Err("query snapshot generation differs from the selected QueryPlan".into()); } let mut pinned = store.pin_matching(required, now_ms, |outputs| { - ranges - .iter() - .all(|(output, start, end)| records_cover(&outputs[&output.0], *start, *end)) + ranges.iter().all(|range| { + range + .covered_window(outputs[&range.output.0].iter()) + .is_some() + }) })?; + // Keep only the admitted windows from the single pinned revision. + // An incomplete newer window must not shadow a complete older one. + let selected: Vec<_> = ranges + .iter() + .map(|range| { + let records = pinned + .records + .iter() + .filter(|record| record.reference.stored_output_id == range.output); + let (start, end) = range + .covered_window(records) + .expect("pinned revision covered every requested range"); + (range.output, start, end) + }) + .collect(); pinned.records.retain(|r| { - ranges.iter().any(|(output, start, end)| { + selected.iter().any(|(output, start, end)| { r.reference.stored_output_id == *output && r.start_ms >= *start && r.end_ms <= *end }) }); @@ -722,6 +758,35 @@ mod tests { BTreeSet::from([StoredOutputId(1), StoredOutputId(2)]) } + // Off-grid mixed reads pin the newest complete window within the bound; + // exact reads and tighter bounds reject the same lagged records. + #[test] + fn revision_windows_honor_stored_input_lag_and_group_completeness() { + let mut first = record(1, 1); + first.start_ms = 100; + first.end_ms = 200; + let mut second = first.clone(); + second.group.insert("job".into(), "second".into()); + let mut partial = first.clone(); + partial.start_ms = 110; + partial.end_ms = 210; + let records = [first, second, partial]; + let mut range = RevisionQueryRange { + output: StoredOutputId(1), + start: 115, + end: 215, + max_lag: 15, + }; + assert_eq!(range.covered_window(records.iter()), Some((100, 200))); + range.max_lag = 14; + assert_eq!(range.covered_window(records.iter()), None); + range.max_lag = 0; + assert_eq!(range.covered_window(records.iter()), None); + range.start = 100; + range.end = 200; + assert_eq!(range.covered_window(records.iter()), Some((100, 200))); + } + /// A durable partial r2 cannot force a two-branch query to mix r1 and r2. #[test] fn partial_revision_restart_and_pinned_read_are_consistent() { @@ -842,7 +907,7 @@ mod tests { let eligible = |records: &BTreeMap>| { outputs .iter() - .all(|o| records_cover(&records[&o.0], 0, 100)) + .all(|o| records_cover(records[&o.0].iter(), 0, 100)) }; assert_eq!( store @@ -974,6 +1039,7 @@ pub(crate) fn pin_query( entry: &asap_types::query_plan::QueryPlanEntry, times: &[u64], generation: Option<&CatalogGeneration>, + max_stored_input_lag_ms: Option, ) -> Result< Option, crate::query_engines::EngineError, @@ -1018,21 +1084,24 @@ pub(crate) fn pin_query( .materialization_bindings() .into_iter() .flat_map(|binding| { - times.iter().map(move |time| { - ( - binding.materialization, - time.saturating_sub( - binding - .readout_lookback_ms - .unwrap_or(entry.instant.lookback_ms), - ), - *time, - ) + times.iter().map(move |time| RevisionQueryRange { + output: binding.materialization, + start: time.saturating_sub( + binding + .readout_lookback_ms + .unwrap_or(entry.instant.lookback_ms), + ), + end: *time, + max_lag: if entry.mixes_raw_and_stored_inputs() { + max_stored_input_lag_ms.unwrap_or_else(|| binding.slide_ms()) + } else { + 0 + }, }) }) .collect(); runtime - .query_view( + .query_view_ranges( &required, now, &ranges, @@ -1056,7 +1125,38 @@ pub(crate) fn pin_query( }) } -fn records_cover(records: &[RevisionRecord], start: u64, end: u64) -> bool { +/// A mixed query may shift its stored window back within its lag bound; +/// ordinary stored-only queries require the exact requested interval. +struct RevisionQueryRange { + output: StoredOutputId, + start: u64, + end: u64, + max_lag: u64, +} + +impl RevisionQueryRange { + fn covered_window<'a>( + &self, + records: impl Iterator + Clone, + ) -> Option<(u64, u64)> { + if self.max_lag == 0 { + return records_cover(records, self.start, self.end).then_some((self.start, self.end)); + } + let ends: BTreeSet<_> = records.clone().map(|record| record.end_ms).collect(); + ends.into_iter().rev().find_map(|end| { + let lag = self.end.checked_sub(end)?; + let start = self.start.checked_sub(lag)?; + (lag <= self.max_lag && records_cover(records.clone(), start, end)) + .then_some((start, end)) + }) + } +} + +fn records_cover<'a>( + records: impl Iterator, + start: u64, + end: u64, +) -> bool { let mut groups = BTreeMap::<&BTreeMap, Vec<(u64, u64)>>::new(); for record in records { let windows = groups.entry(&record.group).or_default(); diff --git a/data_plane/src/query_engines/asap_query_engine/engine.rs b/data_plane/src/query_engines/asap_query_engine/engine.rs index 8bb3bfec0..fac582e18 100644 --- a/data_plane/src/query_engines/asap_query_engine/engine.rs +++ b/data_plane/src/query_engines/asap_query_engine/engine.rs @@ -250,6 +250,7 @@ impl ASAPQueryEngine { entry, times, physical.precompute_plan.summary_catalog.as_ref(), + self.max_stored_input_lag_ms, )? else { return Ok(None); @@ -542,6 +543,7 @@ impl ASAPQueryEngine { entry, &[at], physical.precompute_plan.summary_catalog.as_ref(), + self.max_stored_input_lag_ms, ) }) .transpose()? diff --git a/data_plane/tests/support/lifecycle_placement_process.rs b/data_plane/tests/support/lifecycle_placement_process.rs index 9414f6f9d..5291ff9af 100644 --- a/data_plane/tests/support/lifecycle_placement_process.rs +++ b/data_plane/tests/support/lifecycle_placement_process.rs @@ -209,9 +209,20 @@ async fn expensive_summary_store_rebuilds_state_from_raw_series() { // the default bound of one 10 s slide. #[tokio::test] async fn mixed_placement_combines_raw_series_with_stored_state_within_the_lag_bound() { + mixed_placement_at_offset(0).await; +} + +// Revision pinning admits the latest complete stored window for an off-grid +// evaluation, while raw data is still fetched at the requested timestamp. +#[tokio::test] +async fn mixed_placement_pins_stored_revision_between_window_boundaries() { + mixed_placement_at_offset(3_000).await; +} + +async fn mixed_placement_at_offset(offset_ms: i64) { const MIXED: &str = "sum(rate(a[1m])) + sum(rate(b[10m]))"; let origin = origin_ms(); - let at_ms = origin + 650_000; + let at_ms = origin + 650_000 + offset_ms; // Raw `a` rises 1/s. Stored `b` rises 2/s, then 4/s over the last 5 s // before t_q, so its rate tells which stored window was read. Its samples // sit mid-second, off every window boundary. @@ -260,11 +271,15 @@ async fn mixed_placement_combines_raw_series_with_stored_state_within_the_lag_bo panic!("no mixed answer: {last}\nbackend log:\n{log}") }); assert!(lag <= 10_000, "lag {lag} beyond one slide"); - assert_eq!(lag % 10_000, 0, "stored windows end on the 10 s grid"); + assert_eq!( + lag % 10_000, + offset_ms as u64, + "stored windows end on the 10 s grid" + ); // Exact reference: rate(a) over (t_q - 1m, t_q] plus rate(b) over the // stored window (t_s - 10m, t_s], t_s = t_q - lag. Its samples span 599 s // and extrapolate half a second to each window edge. - let end = 650 - lag as i64 / 1000; + let end = 650 + offset_ms / 1000 - lag as i64 / 1000; let expected = 1.0 + (counter(end - 1) - counter(end - 600)) / 599.0; let value = first_value(&last, "value").unwrap(); assert!(