diff --git a/data_plane/src/drivers/ingest/prometheus_remote_write.rs b/data_plane/src/drivers/ingest/prometheus_remote_write.rs index 6f2c2ba2a..daec008c6 100644 --- a/data_plane/src/drivers/ingest/prometheus_remote_write.rs +++ b/data_plane/src/drivers/ingest/prometheus_remote_write.rs @@ -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::{ @@ -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, @@ -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(); diff --git a/data_plane/src/precompute_engine/output_sink.rs b/data_plane/src/precompute_engine/output_sink.rs index eab97b320..6ff80fcd5 100644 --- a/data_plane/src/precompute_engine/output_sink.rs +++ b/data_plane/src/precompute_engine/output_sink.rs @@ -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) @@ -603,6 +606,7 @@ mod tests { ); assert!(summary_store .query_exact_agg_range(meta.storage_handle, 1_000, 2_000) + .unwrap() .is_empty()); } diff --git a/data_plane/src/query_engines/asap_query_engine/summary_executor.rs b/data_plane/src/query_engines/asap_query_engine/summary_executor.rs index 0511e1aad..0aec76f4b 100644 --- a/data_plane/src/query_engines/asap_query_engine/summary_executor.rs +++ b/data_plane/src/query_engines/asap_query_engine/summary_executor.rs @@ -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 { diff --git a/data_plane/src/storage_engines/sketch_db/index/mod.rs b/data_plane/src/storage_engines/sketch_db/index/mod.rs index d831682f2..5ca1677ab 100644 --- a/data_plane/src/storage_engines/sketch_db/index/mod.rs +++ b/data_plane/src/storage_engines/sketch_db/index/mod.rs @@ -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>, + window_end: u64, + state: &Arc, +) -> 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. @@ -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, - BTreeMap>, - )> { + ) -> Result< + Vec<( + BTreeMap, + BTreeMap>, + )>, + 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 @@ -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(); @@ -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(); @@ -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 @@ -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" @@ -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::() + .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 @@ -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 \ @@ -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" @@ -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);