From 257ef35878d60a7bdf785ec6fa43a46a7e283fba Mon Sep 17 00:00:00 2001 From: GordonYuanyc Date: Thu, 1 Oct 2026 01:01:31 -0400 Subject: [PATCH 1/2] fix(store): merge late ForwardToStore corrections in exact-agg reads With LateDataPolicy::ForwardToStore, input for a pane that the idle rule or absolute deadline already closed is emitted as a separate correction state and appended beside the published state for the same window. The sketch read path merges every frame for a window, but query_exact_agg_range kept only the last state per window end, so each correction replaced the published sum instead of adding to it (Remote Write repro: 15 published, then 40 more late, read 40 instead of 55). Merge in-memory exact states that share a window, as the precompute design doc specifies for ForwardToStore. States that cannot merge return an error and the bound read fails closed instead of returning a partial answer. Co-Authored-By: Claude Opus 5.5 --- .../drivers/ingest/prometheus_remote_write.rs | 189 ++++++++++++++++++ .../src/precompute_engine/output_sink.rs | 8 +- .../asap_query_engine/summary_executor.rs | 3 + .../storage_engines/sketch_db/index/mod.rs | 90 +++++++-- 4 files changed, 271 insertions(+), 19 deletions(-) diff --git a/data_plane/src/drivers/ingest/prometheus_remote_write.rs b/data_plane/src/drivers/ingest/prometheus_remote_write.rs index 6f2c2ba2a..cf702e67f 100644 --- a/data_plane/src/drivers/ingest/prometheus_remote_write.rs +++ b/data_plane/src/drivers/ingest/prometheus_remote_write.rs @@ -1700,6 +1700,195 @@ mod tests { slow_task.await.unwrap(); } + /// One Remote Write worker whose wall clock the test controls, so the idle + /// rule closes pane [0, 60s) of `requests_total{job="a"}` on demand. + struct IdleCloseHarness { + receiver: PrometheusRemoteWriteReceiver, + queued: mpsc::Receiver, + worker: mpsc::Sender, + clock: Arc, + groups: Arc, + } + + impl IdleCloseHarness { + fn start() -> Self { + use crate::precompute_engine::{ + config::LateDataPolicy, + output_sink::SketchStoreSink, + worker::{Worker, WorkerRuntimeConfig}, + }; + use std::sync::atomic::{AtomicI64, AtomicUsize}; + let (receiver, 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 groups = Arc::new(AtomicUsize::new(0)); + let mut runner = Worker::new( + 0, + rx, + sink, + ingest.hot_reload_config.clone(), + 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: 5_000, + wall_clock_max_open_grace_period_ms: i64::MAX, + }, + groups.clone(), + Arc::new(AtomicI64::new(i64::MIN)), + ); + let clock = Arc::new(AtomicI64::new(1_000_000)); + let now = clock.clone(); + runner.set_now_ms_fn(Box::new(move || now.load(Ordering::Relaxed))); + tokio::spawn(runner.run()); + Self { + receiver, + queued, + worker, + clock, + groups, + } + } + + async fn write(&mut self, samples: &[(i64, f64)]) { + let request = WriteRequest { + timeseries: vec![TimeSeries { + labels: vec![ + Label { + name: "__name__".into(), + value: "requests_total".into(), + }, + Label { + name: "job".into(), + value: "a".into(), + }, + ], + samples: samples + .iter() + .map(|&(timestamp, value)| Sample { timestamp, value }) + .collect(), + exemplars: vec![], + histograms: vec![], + }], + }; + self.receiver.accept(&compressed(request)).unwrap(); + while let Ok(message) = self.queued.try_recv() { + self.worker.send(message).await.unwrap(); + } + } + + /// Waits for every message sent so far. Only used once no pane is + /// open, so the drain cannot close one itself. + async fn barrier(&self) { + let (done, result) = tokio::sync::oneshot::channel(); + self.worker.send(WorkerMessage::Drain(done)).await.unwrap(); + result.await.unwrap().unwrap(); + } + + fn read_sum( + &self, + ) -> Result< + f64, + crate::query_engines::asap_query_engine::summary_executor::SummaryExecutorError, + > { + let ingest = &self.receiver.inner.ingest; + 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 context = + crate::query_engines::asap_query_engine::summary_executor::QueryExecutionContext { + index: &ingest.summary_store, + t0_ms: 0, + t1_ms: 60_000, + is_cumulative: true, + allowed_materializations: None, + }; + let groups = context.read_bound_materialization(&binding)?; + assert_eq!(groups.len(), 1); + Ok(groups[0].1.exact_value(&None).unwrap()) + } + + /// Once the worker has touched the pane, advances the clock past + /// window + idle grace and waits until the idle rule publishes it. + async fn idle_close(&self) { + while self.groups.load(Ordering::Relaxed) == 0 { + tokio::task::yield_now().await; + } + self.clock.fetch_add(65_000, Ordering::Relaxed); + self.worker.send(WorkerMessage::Flush).await.unwrap(); + for _ in 0..500 { + if self.read_sum().is_ok() { + return; + } + tokio::time::sleep(std::time::Duration::from_millis(10)).await; + } + panic!("idle rule did not publish the pane"); + } + } + + /// Samples arriving after the idle rule published their pane are merged + /// into that pane's warm read instead of replacing it. + #[tokio::test] + async fn late_samples_after_idle_close_merge_into_published_pane() { + let mut stack = IdleCloseHarness::start(); + stack.write(&[(1_000, 1.0), (2_000, 2.0)]).await; + stack.idle_close().await; + assert_eq!(stack.read_sum().unwrap(), 3.0); + stack.write(&[(3_000, 4.0)]).await; + stack.barrier().await; + assert_eq!(stack.read_sum().unwrap(), 7.0); + stack.write(&[(4_000, 8.0)]).await; + stack.barrier().await; + assert_eq!(stack.read_sum().unwrap(), 15.0); + } + + /// Identical and partial retries of samples already counted in an + /// idle-closed pane or its correction do not change the warm read. + #[tokio::test] + async fn retried_samples_after_idle_close_are_counted_once() { + let mut stack = IdleCloseHarness::start(); + let first = [(1_000, 1.0), (2_000, 2.0)]; + stack.write(&first).await; + stack.idle_close().await; + stack.write(&[(3_000, 4.0)]).await; + stack.write(&first).await; + stack.write(&first[..1]).await; + stack.write(&[(3_000, 4.0)]).await; + stack.barrier().await; + assert_eq!(stack.read_sum().unwrap(), 7.0); + stack.write(&[(2_000, 2.0), (5_000, 16.0)]).await; + stack.barrier().await; + assert_eq!(stack.read_sum().unwrap(), 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..2aec869e1 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,9 @@ impl QueryExecutionContext<'_> { let Some((labels, windows)) = self .index .query_exact_agg_range(sid, self.t0_ms, self.t1_ms) + .map_err(|_| { + SummaryExecutorError::Unsupported("stored exact states do not merge") + })? .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..1736adb76 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(); @@ -2279,7 +2302,7 @@ impl SketchStore { // of the range. In-memory wins on a window-end collision. 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 +6158,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 +6300,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 +6381,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 +6448,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 +6690,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); From 3f44147e8961b535630303fe59ecb549d56f9d2a Mon Sep 17 00:00:00 2001 From: GordonYuanyc Date: Thu, 1 Oct 2026 01:31:37 -0400 Subject: [PATCH 2/2] test: fold late-correction receiver tests into one; review nits Replace the two clock-driven Remote Write tests with one test that closes the pane through the worker drain, then checks the late merge and exactly-once retries. Share the bound-read setup with the existing queued-population test. Surface exact merge failures as Decode errors and note that the disk tier keeps one record per window. Co-Authored-By: Claude Opus 5.5 --- .../drivers/ingest/prometheus_remote_write.rs | 295 ++++++------------ .../asap_query_engine/summary_executor.rs | 4 +- .../storage_engines/sketch_db/index/mod.rs | 3 +- 3 files changed, 95 insertions(+), 207 deletions(-) diff --git a/data_plane/src/drivers/ingest/prometheus_remote_write.rs b/data_plane/src/drivers/ingest/prometheus_remote_write.rs index cf702e67f..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,193 +1696,86 @@ mod tests { slow_task.await.unwrap(); } - /// One Remote Write worker whose wall clock the test controls, so the idle - /// rule closes pane [0, 60s) of `requests_total{job="a"}` on demand. - struct IdleCloseHarness { - receiver: PrometheusRemoteWriteReceiver, - queued: mpsc::Receiver, - worker: mpsc::Sender, - clock: Arc, - groups: Arc, - } - - impl IdleCloseHarness { - fn start() -> Self { - use crate::precompute_engine::{ - config::LateDataPolicy, - output_sink::SketchStoreSink, - worker::{Worker, WorkerRuntimeConfig}, + /// 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![], }; - use std::sync::atomic::{AtomicI64, AtomicUsize}; - let (receiver, 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 groups = Arc::new(AtomicUsize::new(0)); - let mut runner = Worker::new( - 0, - rx, - sink, - ingest.hot_reload_config.clone(), - 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: 5_000, - wall_clock_max_open_grace_period_ms: i64::MAX, - }, - groups.clone(), - Arc::new(AtomicI64::new(i64::MIN)), - ); - let clock = Arc::new(AtomicI64::new(1_000_000)); - let now = clock.clone(); - runner.set_now_ms_fn(Box::new(move || now.load(Ordering::Relaxed))); - tokio::spawn(runner.run()); - Self { - receiver, - queued, - worker, - clock, - groups, - } - } - - async fn write(&mut self, samples: &[(i64, f64)]) { let request = WriteRequest { - timeseries: vec![TimeSeries { - labels: vec![ - Label { - name: "__name__".into(), - value: "requests_total".into(), - }, - Label { - name: "job".into(), - value: "a".into(), - }, - ], - samples: samples - .iter() - .map(|&(timestamp, value)| Sample { timestamp, value }) - .collect(), - exemplars: vec![], - histograms: vec![], - }], + timeseries: vec![series], }; - self.receiver.accept(&compressed(request)).unwrap(); - while let Ok(message) = self.queued.try_recv() { - self.worker.send(message).await.unwrap(); + receiver.accept(&compressed(request)).unwrap(); + while let Ok(message) = queued.try_recv() { + worker.send(message).await.unwrap(); } - } - - /// Waits for every message sent so far. Only used once no pane is - /// open, so the drain cannot close one itself. - async fn barrier(&self) { let (done, result) = tokio::sync::oneshot::channel(); - self.worker.send(WorkerMessage::Drain(done)).await.unwrap(); + worker.send(WorkerMessage::Drain(done)).await.unwrap(); result.await.unwrap().unwrap(); - } - - fn read_sum( - &self, - ) -> Result< - f64, - crate::query_engines::asap_query_engine::summary_executor::SummaryExecutorError, - > { - let ingest = &self.receiver.inner.ingest; - 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 context = - crate::query_engines::asap_query_engine::summary_executor::QueryExecutionContext { - index: &ingest.summary_store, - t0_ms: 0, - t1_ms: 60_000, - is_cumulative: true, - allowed_materializations: None, - }; - let groups = context.read_bound_materialization(&binding)?; - assert_eq!(groups.len(), 1); - Ok(groups[0].1.exact_value(&None).unwrap()) - } - - /// Once the worker has touched the pane, advances the clock past - /// window + idle grace and waits until the idle rule publishes it. - async fn idle_close(&self) { - while self.groups.load(Ordering::Relaxed) == 0 { - tokio::task::yield_now().await; - } - self.clock.fetch_add(65_000, Ordering::Relaxed); - self.worker.send(WorkerMessage::Flush).await.unwrap(); - for _ in 0..500 { - if self.read_sum().is_ok() { - return; - } - tokio::time::sleep(std::time::Duration::from_millis(10)).await; - } - panic!("idle rule did not publish the pane"); - } - } - - /// Samples arriving after the idle rule published their pane are merged - /// into that pane's warm read instead of replacing it. - #[tokio::test] - async fn late_samples_after_idle_close_merge_into_published_pane() { - let mut stack = IdleCloseHarness::start(); - stack.write(&[(1_000, 1.0), (2_000, 2.0)]).await; - stack.idle_close().await; - assert_eq!(stack.read_sum().unwrap(), 3.0); - stack.write(&[(3_000, 4.0)]).await; - stack.barrier().await; - assert_eq!(stack.read_sum().unwrap(), 7.0); - stack.write(&[(4_000, 8.0)]).await; - stack.barrier().await; - assert_eq!(stack.read_sum().unwrap(), 15.0); - } - - /// Identical and partial retries of samples already counted in an - /// idle-closed pane or its correction do not change the warm read. - #[tokio::test] - async fn retried_samples_after_idle_close_are_counted_once() { - let mut stack = IdleCloseHarness::start(); - let first = [(1_000, 1.0), (2_000, 2.0)]; - stack.write(&first).await; - stack.idle_close().await; - stack.write(&[(3_000, 4.0)]).await; - stack.write(&first).await; - stack.write(&first[..1]).await; - stack.write(&[(3_000, 4.0)]).await; - stack.barrier().await; - assert_eq!(stack.read_sum().unwrap(), 7.0); - stack.write(&[(2_000, 2.0), (5_000, 16.0)]).await; - stack.barrier().await; - assert_eq!(stack.read_sum().unwrap(), 23.0); + }; + 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] 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 2aec869e1..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,9 +489,7 @@ impl QueryExecutionContext<'_> { let Some((labels, windows)) = self .index .query_exact_agg_range(sid, self.t0_ms, self.t1_ms) - .map_err(|_| { - SummaryExecutorError::Unsupported("stored exact states do not merge") - })? + .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 1736adb76..5ca1677ab 100644 --- a/data_plane/src/storage_engines/sketch_db/index/mod.rs +++ b/data_plane/src/storage_engines/sketch_db/index/mod.rs @@ -2299,7 +2299,8 @@ 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); Ok(by_label_map.into_iter().collect())