Skip to content
Open
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
134 changes: 117 additions & 17 deletions data_plane/src/precompute_engine/revisions.rs
Original file line number Diff line number Diff line change
Expand Up @@ -655,18 +655,54 @@ impl RevisionRuntime {
now_ms: u64,
ranges: &[(StoredOutputId, u64, u64)],
expected: &CatalogGeneration,
) -> Result<crate::storage_engines::sketch_db::index::SketchStore, RevisionError> {
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<StoredOutputId>,
now_ms: u64,
ranges: &[RevisionQueryRange],
expected: &CatalogGeneration,
) -> Result<crate::storage_engines::sketch_db::index::SketchStore, RevisionError> {
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
})
});
Expand Down Expand Up @@ -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() {
Expand Down Expand Up @@ -842,7 +907,7 @@ mod tests {
let eligible = |records: &BTreeMap<u64, Vec<RevisionRecord>>| {
outputs
.iter()
.all(|o| records_cover(&records[&o.0], 0, 100))
.all(|o| records_cover(records[&o.0].iter(), 0, 100))
};
assert_eq!(
store
Expand Down Expand Up @@ -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<u64>,
) -> Result<
Option<crate::storage_engines::sketch_db::index::SketchStore>,
crate::query_engines::EngineError,
Expand Down Expand Up @@ -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,
Expand All @@ -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<Item = &'a RevisionRecord> + 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<Item = &'a RevisionRecord>,
start: u64,
end: u64,
) -> bool {
let mut groups = BTreeMap::<&BTreeMap<String, String>, Vec<(u64, u64)>>::new();
for record in records {
let windows = groups.entry(&record.group).or_default();
Expand Down
2 changes: 2 additions & 0 deletions data_plane/src/query_engines/asap_query_engine/engine.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down Expand Up @@ -542,6 +543,7 @@ impl ASAPQueryEngine {
entry,
&[at],
physical.precompute_plan.summary_catalog.as_ref(),
self.max_stored_input_lag_ms,
)
})
.transpose()?
Expand Down
21 changes: 18 additions & 3 deletions data_plane/tests/support/lifecycle_placement_process.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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!(
Expand Down
Loading