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
122 changes: 100 additions & 22 deletions data_plane/src/drivers/ingest/prometheus_remote_write.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1580,6 +1580,23 @@ mod tests {
assert_eq!(receiver.stats().duplicates.load(Ordering::Relaxed), 1);
}

/// Bound read of `configured_receiver`'s pane [0, 60s), grouped by `job`.
fn job_pane_binding(ingest: &IngestState) -> asap_types::query_plan::MaterializationBinding {
let snapshot = ingest.hot_reload_config.snapshot();
let policy = *snapshot.materializations_by_output.keys().next().unwrap();
let catalog = ingest.summary_store.summary_catalog_snapshot().unwrap();
asap_types::query_plan::MaterializationBinding {
full_window_slide_ms: None,
materialization: policy,
stored_output_reference: catalog.output_reference(policy).unwrap(),
output_grouping: asap_types::query_plan::PhysicalGrouping::Reduce(vec!["job".into()]),
item_labels: vec![],
window_ms: 60_000,
pane_origin_ms: Some(0),
readout_lookback_ms: None,
}
}

#[tokio::test]
async fn queued_population_blocks_partial_warm_read_until_both_workers_publish() {
use crate::precompute_engine::{
Expand Down Expand Up @@ -1651,28 +1668,7 @@ mod tests {
let (done, result) = tokio::sync::oneshot::channel();
fast.send(WorkerMessage::Drain(done)).await.unwrap();
result.await.unwrap().unwrap();
let policy = *ingest
.hot_reload_config
.snapshot()
.materializations_by_output
.keys()
.next()
.unwrap();
let binding = asap_types::query_plan::MaterializationBinding {
full_window_slide_ms: None,
materialization: policy,
stored_output_reference: ingest
.summary_store
.summary_catalog_snapshot()
.unwrap()
.output_reference(policy)
.unwrap(),
output_grouping: asap_types::query_plan::PhysicalGrouping::Reduce(vec!["job".into()]),
item_labels: vec![],
window_ms: 60_000,
pane_origin_ms: Some(0),
readout_lookback_ms: None,
};
let binding = job_pane_binding(ingest);
let context = QueryExecutionContext {
index: &ingest.summary_store,
t0_ms: 0,
Expand Down Expand Up @@ -1700,6 +1696,88 @@ mod tests {
slow_task.await.unwrap();
}

/// A late sample for an already closed and published pane is merged into
/// that pane's read, and retried samples are counted once.
#[tokio::test]
async fn late_sample_merges_into_closed_pane_and_retries_count_once() {
use crate::precompute_engine::{
config::LateDataPolicy,
output_sink::SketchStoreSink,
worker::{Worker, WorkerRuntimeConfig},
};
use crate::query_engines::asap_query_engine::summary_executor::QueryExecutionContext;
use std::sync::atomic::AtomicI64;
let (receiver, mut queued) = configured_receiver();
let ingest = &receiver.inner.ingest;
let sink = Arc::new(SketchStoreSink::new(
ingest.summary_store.clone(),
ingest.hot_reload_config.clone(),
ingest.series_resolver.clone(),
));
let (worker, rx) = mpsc::channel(8);
let config = WorkerRuntimeConfig {
max_buffer_per_series: 100,
allowed_lateness_ms: 0,
pass_raw_samples: false,
raw_mode_aggregation_id: 0,
late_data_policy: LateDataPolicy::ForwardToStore,
wall_clock_idle_grace_period_ms: i64::MAX,
wall_clock_max_open_grace_period_ms: i64::MAX,
};
let (groups, watermark) = (Arc::default(), Arc::new(AtomicI64::new(i64::MIN)));
let plan = ingest.hot_reload_config.clone();
tokio::spawn(Worker::new(0, rx, sink, plan, config, groups, watermark).run());
// Each write is drained, so the first write's pane is closed and published.
let mut write = async |samples: &[(i64, f64)]| {
let series = TimeSeries {
labels: [("__name__", "requests_total"), ("job", "a")]
.map(|(name, value)| Label {
name: name.into(),
value: value.into(),
})
.into(),
samples: samples
.iter()
.map(|&(timestamp, value)| Sample { timestamp, value })
.collect(),
exemplars: vec![],
histograms: vec![],
};
let request = WriteRequest {
timeseries: vec![series],
};
receiver.accept(&compressed(request)).unwrap();
while let Ok(message) = queued.try_recv() {
worker.send(message).await.unwrap();
}
let (done, result) = tokio::sync::oneshot::channel();
worker.send(WorkerMessage::Drain(done)).await.unwrap();
result.await.unwrap().unwrap();
};
let binding = job_pane_binding(ingest);
let context = QueryExecutionContext {
index: &ingest.summary_store,
t0_ms: 0,
t1_ms: 60_000,
is_cumulative: true,
allowed_materializations: None,
};
let read = || {
context.read_bound_materialization(&binding).unwrap()[0]
.1
.exact_value(&None)
};
let (first, late) = ([(1_000, 1.0), (2_000, 2.0)], [(3_000, 4.0)]);
write(&first).await;
write(&late).await;
assert_eq!(read(), Some(7.0));
write(&first).await;
write(&late).await;
assert_eq!(read(), Some(7.0));
write(&[(2_000, 2.0), (5_000, 16.0)]).await;
assert_eq!(read(), Some(23.0));
}

#[tokio::test]
async fn valid_request_routes_canonical_sample_to_installed_plan() {
let (receiver, mut worker) = configured_receiver();
Expand Down
8 changes: 6 additions & 2 deletions data_plane/src/precompute_engine/output_sink.rs
Original file line number Diff line number Diff line change
Expand Up @@ -550,8 +550,11 @@ mod tests {
)])
.is_err());
assert_ne!(old_sid, new_sid);
assert!(store.query_exact_agg_range(old_sid, 1000, 2000).is_empty());
let values = store.query_exact_agg_range(new_sid, 1000, 2000);
assert!(store
.query_exact_agg_range(old_sid, 1000, 2000)
.unwrap()
.is_empty());
let values = store.query_exact_agg_range(new_sid, 1000, 2000).unwrap();
assert_eq!(
values[0].1.values().next().unwrap().aux_stats().sum,
Some(11.0)
Expand Down Expand Up @@ -603,6 +606,7 @@ mod tests {
);
assert!(summary_store
.query_exact_agg_range(meta.storage_handle, 1_000, 2_000)
.unwrap()
.is_empty());
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -489,6 +489,7 @@ impl QueryExecutionContext<'_> {
let Some((labels, windows)) = self
.index
.query_exact_agg_range(sid, self.t0_ms, self.t1_ms)
.map_err(SummaryExecutorError::Decode)?
.into_iter()
.next()
else {
Expand Down
93 changes: 75 additions & 18 deletions data_plane/src/storage_engines/sketch_db/index/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -50,6 +50,28 @@ fn now_ms() -> u64 {
.unwrap_or(0)
}

/// Add one stored exact state to its window, merging with any state already
/// read for that window.
fn merge_exact_window(
windows: &mut BTreeMap<i64, Arc<dyn crate::storage_engines::types::AggregateCore>>,
window_end: u64,
state: &Arc<dyn crate::storage_engines::types::AggregateCore>,
) -> Result<(), String> {
match windows.entry(window_end as i64) {
std::collections::btree_map::Entry::Vacant(entry) => {
entry.insert(Arc::clone(state));
}
std::collections::btree_map::Entry::Occupied(mut entry) => {
let merged = entry
.get()
.merge_with(state.as_ref())
.map_err(|error| error.to_string())?;
entry.insert(Arc::from(merged));
}
}
Ok(())
}

/// Map a [`SketchEncoding`] to the on-disk encoding tag stored per part
/// entry, so the disk read-back path can reconstruct the Full-vs-Delta
/// distinction the delta-stitching carry-in relies on.
Expand Down Expand Up @@ -2210,15 +2232,22 @@ impl SketchStore {
/// in `[start, end]` (or is sketch-backed). Defensive — caller
/// is responsible for confirming the sid's `agg_kind` is
/// `AggKind::ExactAgg { .. }` before calling.
///
/// In-memory states for one window are merged: a `ForwardToStore`
/// correction is appended beside the window's earlier state. States
/// that cannot merge are an error rather than a partial result.
pub fn query_exact_agg_range(
&self,
sid: u64,
start_unix_ms: u64,
end_unix_ms: u64,
) -> Vec<(
BTreeMap<String, String>,
BTreeMap<i64, Arc<dyn crate::storage_engines::types::AggregateCore>>,
)> {
) -> Result<
Vec<(
BTreeMap<String, String>,
BTreeMap<i64, Arc<dyn crate::storage_engines::types::AggregateCore>>,
)>,
String,
> {
// Key by the resolved label MAP (not `LabelValuesId`) so the
// in-memory tier and the durable disk tier — which carry
// independent intern spaces — union by label identity. Mirrors
Expand All @@ -2244,10 +2273,7 @@ impl SketchStore {
.range_query_into(start_unix_ms, end_unix_ms, &mut buf);
for (win, label_id, payload) in &buf {
if let Some(p) = payload.as_exact_agg_arc() {
by_label_id
.entry(*label_id)
.or_default()
.insert(win.1 as i64, Arc::clone(p));
merge_exact_window(by_label_id.entry(*label_id).or_default(), win.1, p)?;
}
}
buf.clear();
Expand All @@ -2256,10 +2282,7 @@ impl SketchStore {
sealed.range_query_into(start_unix_ms, end_unix_ms, &mut buf);
for (win, label_id, payload) in &buf {
if let Some(p) = payload.as_exact_agg_arc() {
by_label_id
.entry(*label_id)
.or_default()
.insert(win.1 as i64, Arc::clone(p));
merge_exact_window(by_label_id.entry(*label_id).or_default(), win.1, p)?;
}
}
buf.clear();
Expand All @@ -2276,10 +2299,11 @@ impl SketchStore {
}

// Union the durable disk tier for the flushed-then-evicted portion
// of the range. In-memory wins on a window-end collision.
// of the range. In-memory wins on a window-end collision; the disk
// tier keeps one record per window and does not merge.
self.union_disk_exact_agg_into(sid, start_unix_ms, end_unix_ms, &mut by_label_map);

by_label_map.into_iter().collect()
Ok(by_label_map.into_iter().collect())
}

/// Union the durable disk tier's exact-aggregation entries into
Expand Down Expand Up @@ -6135,7 +6159,7 @@ mod tests {
));

// The Sum exact-agg range query must resolve from disk.
let series = idx2.query_exact_agg_range(8100, 0, 150_000);
let series = idx2.query_exact_agg_range(8100, 0, 150_000).unwrap();
assert!(
!series.is_empty(),
"recovered exact-agg query returned No result after fresh reopen"
Expand Down Expand Up @@ -6277,6 +6301,37 @@ mod tests {
drop(p2);
}

/// A late correction appended for an already-published exact window is
/// read merged with that window's earlier state, not in place of it.
#[test]
fn exact_agg_range_merges_late_correction_for_same_window() {
use crate::storage_engines::types::AggregationType;
let idx = SketchStore::new();
let mut m = meta(8002);
m.agg_kind = AggKind::ExactAgg {
agg_type: AggregationType::Sum,
parameters_canonical: String::new(),
spatial_filter_canonical: String::new(),
};
idx.register(m);
for value in [15.0, 40.0] {
assert!(idx.append_precompute(
8002,
BTreeMap::new(),
(0, 10_000),
Box::new(asap_physical_operators::summary_kernels::SumAccumulator::with_sum(value)),
));
}
let series = idx.query_exact_agg_range(8002, 0, 10_000).unwrap();
assert_eq!(series.len(), 1);
let sum = series[0].1[&10_000]
.as_any()
.downcast_ref::<asap_physical_operators::summary_kernels::SumAccumulator>()
.unwrap()
.sum;
assert_eq!(sum, 55.0);
}

/// BUG #2: after flush+evict, an exact-agg (`sum by (zone)` shape)
/// range query must still resolve from disk. On origin/main
/// `query_exact_agg_range` reads ONLY the in-memory current+sealed
Expand Down Expand Up @@ -6327,7 +6382,7 @@ mod tests {
"exact-agg windows never fully evicted"
);
// Query the EVICTED portion [0, 150_000) — must come back from disk.
let series = idx.query_exact_agg_range(8001, 0, 150_000);
let series = idx.query_exact_agg_range(8001, 0, 150_000).unwrap();
assert!(
!series.is_empty(),
"LIVE BUG #2: exact-agg query returned No result after flush+evict \
Expand Down Expand Up @@ -6394,7 +6449,7 @@ mod tests {
"exact-agg windows never fully evicted"
);
// Query the EVICTED portion [0, 150_000) — must come back from disk.
let series = idx.query_exact_agg_range(8001, 0, 150_000);
let series = idx.query_exact_agg_range(8001, 0, 150_000).unwrap();
assert!(
!series.is_empty(),
"exact-agg query returned no result after flush and eviction"
Expand Down Expand Up @@ -6636,7 +6691,9 @@ mod tests {
.start_persistence(durable_cfg(temp.path().to_path_buf()))
.unwrap();
for (i, kind) in kinds.iter().enumerate() {
let series = store.query_exact_agg_range(9000 + i as u64, 0, 30001);
let series = store
.query_exact_agg_range(9000 + i as u64, 0, 30001)
.unwrap();
assert_eq!(series.len(), 1, "{kind:?}");
let state = &series[0].1[&30000];
assert_eq!(state.get_accumulator_type(), *kind);
Expand Down
Loading