From 8a766ce9cc73b13fe0d6704708e97c6b638f47a3 Mon Sep 17 00:00:00 2001 From: GordonYuanyc Date: Thu, 1 Oct 2026 00:52:02 -0400 Subject: [PATCH 1/6] Add a live event-time barrier for Remote Write (#772) A pane is published only after a strictly later sample arrives in the same series, or after the wall-clock idle rule fires (window + 5 s). A series that stops sending leaves its last pane open, so every query over that window misses ("materialization population has unpublished input") and is forwarded. Remote Write has no way to say "all data up to T is written", and /api/v1/precompute/drain seals input for the process lifetime. POST /api/v1/precompute/watermark?time_ms=T lets the producer declare that every sample at or before T has been written. The receiver broadcasts a WorkerMessage::AdvanceWatermark to every worker; Remote Write acknowledges only after enqueueing, so the barrier queues behind every acknowledged write. Each worker closes, in every group, the stored buckets (panes or FullWindow windows) whose samples all lie at or before T, including groups that stopped, and replies after publishing. Right-closed PromQL groups file a sample at t as t - 1, so their barrier is T; other groups use T + 1. The scan is bounded by the latest open bucket, as in force_close_all. Input stays open; later samples for closed buckets follow the late-data policy. A failed worker answers the barrier with its error. flush_all and force_close_all are unchanged. Co-Authored-By: Claude Opus 5.5 --- .../drivers/ingest/prometheus_remote_write.rs | 36 ++ data_plane/src/drivers/query/servers/http.rs | 41 ++ .../src/precompute_engine/series_router.rs | 32 ++ data_plane/src/precompute_engine/worker.rs | 393 +++++++++++++++++- .../asapquery_compatibility_process_e2e.rs | 138 ++++++ 5 files changed, 639 insertions(+), 1 deletion(-) diff --git a/data_plane/src/drivers/ingest/prometheus_remote_write.rs b/data_plane/src/drivers/ingest/prometheus_remote_write.rs index 6f2c2ba2a..da73a51e7 100644 --- a/data_plane/src/drivers/ingest/prometheus_remote_write.rs +++ b/data_plane/src/drivers/ingest/prometheus_remote_write.rs @@ -208,6 +208,19 @@ impl PrometheusRemoteWriteReceiver { } } + /// Live completion barrier: the producer asserts that every sample at or + /// before `event_time_ms` has been written. Requests already acknowledged + /// are queued ahead of the barrier on every worker, so their completed + /// panes are published before this returns. Unlike `drain`, input stays + /// open. + pub async fn advance_watermark(&self, event_time_ms: i64) -> Result<(), String> { + self.inner + .ingest + .router + .advance_watermark(event_time_ms) + .await + } + /// Permanently seal this finite source before queuing worker barriers. pub async fn drain(&self) -> Result<(), String> { if self.inner.config.revisions.is_some() { @@ -1409,6 +1422,29 @@ mod tests { assert_eq!(drain.await.unwrap().unwrap_err(), "sink write failed"); } + // A live barrier queues behind acknowledged writes and leaves input open. + #[tokio::test] + async fn watermark_barrier_follows_acknowledged_writes_and_keeps_input_open() { + let (receiver, mut worker) = configured_receiver(); + receiver.accept(&one_sample(1.0)).unwrap(); + let handle = receiver.clone(); + let barrier = tokio::spawn(async move { handle.advance_watermark(100).await }); + assert!(matches!( + worker.recv().await.unwrap(), + WorkerMessage::BoundInput { .. } + )); + let WorkerMessage::AdvanceWatermark { + event_time_ms: 100, + reply, + } = worker.recv().await.unwrap() + else { + panic!("expected watermark barrier") + }; + receiver.accept(&one_sample(1.0)).unwrap(); + reply.send(Ok(())).unwrap(); + barrier.await.unwrap().unwrap(); + } + #[test] fn rejects_writes_without_an_active_physical_plan() { let (sender, _worker) = mpsc::channel(1); diff --git a/data_plane/src/drivers/query/servers/http.rs b/data_plane/src/drivers/query/servers/http.rs index c48391baa..e068a9b2f 100644 --- a/data_plane/src/drivers/query/servers/http.rs +++ b/data_plane/src/drivers/query/servers/http.rs @@ -445,6 +445,10 @@ impl HttpServer { ) .route("/api/v1/store/metrics", get(handle_store_metrics)) .route("/api/v1/precompute/drain", post(handle_precompute_drain)) + .route( + "/api/v1/precompute/watermark", + post(handle_precompute_watermark), + ) .route("/api/v1/physical-plan", post(handle_post_physical_plan)) .route( "/api/v1/physical-plan/discard", @@ -548,6 +552,10 @@ impl HttpServer { let app = Router::new() .route("/api/v1/precompute/drain", post(handle_precompute_drain)) + .route( + "/api/v1/precompute/watermark", + post(handle_precompute_watermark), + ) .route(query_endpoint, get(handle_instant_query)) .route(query_endpoint, post(handle_instant_query_post)) .route(range_query_endpoint, get(handle_range_query)) @@ -5358,6 +5366,39 @@ async fn handle_precompute_drain(State(state): State) -> Response { } } +/// Live completion barrier: `POST /api/v1/precompute/watermark?time_ms=T` +/// declares every Remote Write sample at or before `T` written. Input stays +/// open, unlike `/api/v1/precompute/drain`. +async fn handle_precompute_watermark( + State(state): State, + Query(params): Query>, +) -> Response { + let Some(receiver) = state.remote_write.as_ref() else { + return (StatusCode::NOT_FOUND, "Remote Write is disabled").into_response(); + }; + let Some(time_ms) = params.get("time_ms").and_then(|v| v.parse::().ok()) else { + return ( + StatusCode::BAD_REQUEST, + axum::Json(serde_json::json!({ + "status": "error", "error": "time_ms must be an integer number of milliseconds" + })), + ) + .into_response(); + }; + match receiver.advance_watermark(time_ms).await { + Ok(()) => ( + StatusCode::OK, + axum::Json(serde_json::json!({"status": "success", "time_ms": time_ms})), + ) + .into_response(), + Err(error) => ( + StatusCode::INTERNAL_SERVER_ERROR, + axum::Json(serde_json::json!({"status": "error", "error": error})), + ) + .into_response(), + } +} + /// Return list of metrics currently in the store. async fn handle_store_metrics(State(state): State) -> axum::response::Response { use axum::http::StatusCode; diff --git a/data_plane/src/precompute_engine/series_router.rs b/data_plane/src/precompute_engine/series_router.rs index 0334f5fb2..dfeea4f45 100644 --- a/data_plane/src/precompute_engine/series_router.rs +++ b/data_plane/src/precompute_engine/series_router.rs @@ -84,6 +84,12 @@ pub enum WorkerMessage { Flush, /// Finite-input barrier: acknowledge only after queued input and trailing panes reach the sink. Drain(tokio::sync::oneshot::Sender>), + /// Live-input barrier: every sample at or before `event_time_ms` has been + /// sent, so buckets complete by then close in every group. Input stays open. + AdvanceWatermark { + event_time_ms: i64, + reply: tokio::sync::oneshot::Sender>, + }, /// Execute downstream DAG work after all raw windows have been sealed. CompleteDag { plan: Arc, @@ -141,6 +147,10 @@ impl fmt::Debug for WorkerMessage { .finish(), Self::Flush => f.write_str("Flush"), Self::Drain(_) => f.write_str("Drain"), + Self::AdvanceWatermark { event_time_ms, .. } => f + .debug_struct("AdvanceWatermark") + .field("event_time_ms", event_time_ms) + .finish(), Self::CompleteDag { .. } => f.write_str("CompleteDag"), Self::Shutdown => f.write_str("Shutdown"), } @@ -275,6 +285,7 @@ impl SeriesRouter { } WorkerMessage::Flush | WorkerMessage::Drain(_) + | WorkerMessage::AdvanceWatermark { .. } | WorkerMessage::Shutdown | WorkerMessage::CompleteDag { .. } => 0, }; @@ -329,6 +340,27 @@ impl SeriesRouter { Ok(()) } + /// Send a live event-time barrier to every worker and wait for all of + /// them. Input already routed is queued ahead of the barrier. + pub async fn advance_watermark(&self, event_time_ms: i64) -> Result<(), String> { + let mut replies = Vec::new(); + for sender in &self.senders { + let (tx, rx) = tokio::sync::oneshot::channel(); + sender + .send(WorkerMessage::AdvanceWatermark { + event_time_ms, + reply: tx, + }) + .await + .map_err(|e| e.to_string())?; + replies.push(rx); + } + for reply in replies { + reply.await.map_err(|e| e.to_string())??; + } + Ok(()) + } + pub async fn shutdown(&self) -> Result<(), String> { let mut errors = Vec::new(); for (worker, sender) in self.senders.iter().enumerate() { diff --git a/data_plane/src/precompute_engine/worker.rs b/data_plane/src/precompute_engine/worker.rs index 5d71f6645..7fa05a1a8 100644 --- a/data_plane/src/precompute_engine/worker.rs +++ b/data_plane/src/precompute_engine/worker.rs @@ -303,7 +303,9 @@ impl Worker { // behind a successful drain. The engine must be restarted explicitly. if let Some(error) = &processing_error { match msg { - WorkerMessage::Drain(reply) | WorkerMessage::CompleteDag { reply, .. } => { + WorkerMessage::Drain(reply) + | WorkerMessage::CompleteDag { reply, .. } + | WorkerMessage::AdvanceWatermark { reply, .. } => { let _ = reply.send(Err(error.clone())); } WorkerMessage::Shutdown => break, @@ -448,6 +450,17 @@ impl Worker { } let _ = reply.send(processing_error.clone().map_or(Ok(()), Err)); } + WorkerMessage::AdvanceWatermark { + event_time_ms, + reply, + } => { + let result = self.close_through(event_time_ms).map_err(|e| e.to_string()); + if let Err(error) = &result { + processing_error = Some(error.clone()); + warn!("Worker {} watermark barrier error: {}", self.id, error); + } + let _ = reply.send(result); + } WorkerMessage::CompleteDag { plan, reply } => { let result = match &processing_error { Some(error) => Err(error.clone()), @@ -1234,6 +1247,93 @@ impl Worker { Ok(()) } + /// Live event-time barrier: the producer asserts that every sample with a + /// timestamp at or before `event_time_ms` has been written. Close, in + /// every group, each stored bucket whose samples all lie at or before + /// that time, including groups that stopped receiving input. Observed + /// event time is not changed; later input for a closed bucket follows the + /// configured late-data policy. A repeated or lower barrier is a no-op. + fn close_through( + &mut self, + event_time_ms: i64, + ) -> Result<(), Box> { + if self.pass_raw_samples { + return Ok(()); + } + + let mut emit_batch: Vec<(PrecomputedOutput, Box)> = Vec::new(); + + for state in self.group_states.values_mut() { + if state.max_event_time_ms == i64::MIN { + continue; // No samples received yet — no panes to close. + } + // Buckets are `[start, end)` of pane timestamps and close once the + // closure watermark reaches `end`. A PromQL right-closed group + // files a sample at `t` as `t - 1`, so its buckets ending at + // `event_time_ms` are complete; otherwise they end one later. + let right_closed = state + .config + .parameters + .get("promql_right_closed") + .and_then(serde_json::Value::as_bool) + .unwrap_or(false); + let barrier = if right_closed { + event_time_ms + } else { + event_time_ms.saturating_add(1) + }; + if barrier <= state.closure_watermark_ms { + continue; + } + + // Bound the scan by the latest open bucket, as force_close_all + // does: buckets after it hold no data, and the scan walks one + // slide at a time, so an unbounded barrier would not terminate. + let max_active = state.active_panes.keys().next_back().copied(); + let max_sketch = state.sketch_panes.keys().next_back().copied(); + if let Some(max_open) = max_active.max(max_sketch) { + let (_, max_open_end) = state.bucket_bounds(max_open); + let close_to = barrier.min(max_open_end); + let group_key = state.group_key.clone(); + for window_start in state.closed_buckets(state.closure_watermark_ms, close_to) { + let (_, window_end) = state.bucket_bounds(window_start); + let pane_starts = [window_start]; + let accumulators = [ + merge_panes_for_window(&mut state.active_panes, &pane_starts), + merge_sketch_panes_for_window(&mut state.sketch_panes, &pane_starts), + ]; + for accumulator in accumulators.into_iter().flatten() { + let output = precomputed_output_for_group( + window_start as u64, + window_end as u64, + build_group_key_label_values(&group_key), + PolicyFingerprint::from_config(&state.config), + &group_key, + &state.input_revisions, + state.series_id, + state.catalog_generation.as_ref(), + state.stored_output_reference.clone(), + ); + emit_batch.push((output, accumulator)); + } + } + } + state.closure_watermark_ms = barrier; + state.prune_pane_wall_clock(); + } + + if !emit_batch.is_empty() { + debug!( + "Worker {} watermark barrier emitting {} outputs", + self.id, + emit_batch.len() + ); + self.output_sink.emit_batch(emit_batch)?; + } + + Ok(()) + } + /// Force-close every window still open on shutdown. /// /// Unlike `flush_all` — which preserves observed event time and only applies @@ -3200,6 +3300,297 @@ mod tests { assert_eq!(emitted[1].0.end_timestamp, 10_000); } + fn drain_emitted_sums(sink: &CapturingOutputSink) -> Vec<(u64, u64, Vec, f64)> { + sink.drain() + .into_iter() + .map(|(output, accumulator)| { + ( + output.start_timestamp, + output.end_timestamp, + output.key.map(|key| key.labels).unwrap_or_default(), + accumulator + .as_any() + .downcast_ref::() + .unwrap() + .sum, + ) + }) + .collect() + } + + // The barrier closes exactly the stopped group's complete pane; repeats and lower barriers are no-ops. + #[test] + fn watermark_barrier_closes_complete_pane_of_stopped_group() { + let config = make_agg_config( + 1, + "cpu", + AggregationType::SingleSubpopulation, + "Sum", + 10, + 0, + vec![], + ); + let sink = Arc::new(CapturingOutputSink::new()); + let mut worker = make_worker( + HashMap::from([(1, config)]), + sink.clone(), + false, + 0, + LateDataPolicy::Drop, + ); + let pf = PolicyFingerprint(1); + // Group a keeps sending into [10s, 20s); group b stops inside [0, 10s). + for (sid, group, ts, value) in [ + (1, "a", 1_000, 1.0), + (2, "b", 2_000, 2.0), + (1, "a", 12_000, 3.0), + ] { + worker + .process_group_samples( + sid, + pf, + &test_group_key(group), + group_samples("cpu", vec![(ts, value)]), + ) + .unwrap(); + } + assert_eq!( + drain_emitted_sums(&sink), + vec![(0, 10_000, vec!["a".to_string()], 1.0)], + "only group a observed a later sample" + ); + worker.flush_all().unwrap(); + assert!( + sink.is_empty(), + "without wall-clock policy group b stays open" + ); + + // [0, 10s) holds samples up to 9_999 ms. + worker.close_through(9_998).unwrap(); + assert!(sink.is_empty()); + worker.close_through(9_999).unwrap(); + assert_eq!( + drain_emitted_sums(&sink), + vec![(0, 10_000, vec!["b".to_string()], 2.0)] + ); + worker.close_through(9_999).unwrap(); + worker.close_through(5_000).unwrap(); + assert!(sink.is_empty()); + + worker.force_close_all().unwrap(); + assert_eq!( + drain_emitted_sums(&sink), + vec![(10_000, 20_000, vec!["a".to_string()], 3.0)], + "the barrier left group a's later pane open" + ); + } + + // PromQL right-closed groups file a sample at t into the pane ending at t, so a barrier at t closes it. + #[test] + fn watermark_barrier_respects_promql_right_closed_membership() { + let mut config = make_agg_config( + 1, + "cpu", + AggregationType::SingleSubpopulation, + "Sum", + 10, + 0, + vec![], + ); + config + .parameters + .insert("promql_right_closed".into(), serde_json::json!(true)); + let sink = Arc::new(CapturingOutputSink::new()); + let mut worker = make_worker( + HashMap::from([(1, config)]), + sink.clone(), + false, + 0, + LateDataPolicy::Drop, + ); + worker + .process_group_samples( + 1, + PolicyFingerprint(1), + &test_group_key("a"), + group_samples("cpu", vec![(10_000, 4.0)]), + ) + .unwrap(); + worker.close_through(9_999).unwrap(); + assert!(sink.is_empty()); + worker.close_through(10_000).unwrap(); + assert_eq!( + drain_emitted_sums(&sink), + vec![(0, 10_000, vec!["a".to_string()], 4.0)] + ); + } + + // FullWindow buckets close one by one; an unbounded barrier terminates and later input is late. + #[test] + fn watermark_barrier_closes_full_windows_and_bounds_its_scan() { + let mut config = make_agg_config( + 1, + "cpu", + AggregationType::SingleSubpopulation, + "Sum", + 30, + 10, + vec![], + ); + config.window_layout = asap_types::WindowMaterializationLayout::FullWindow; + let sink = Arc::new(CapturingOutputSink::new()); + let mut worker = make_worker( + HashMap::from([(1, config)]), + sink.clone(), + false, + 0, + LateDataPolicy::Drop, + ); + let sample = |worker: &mut Worker, ts: i64| { + worker + .process_group_samples( + 1, + PolicyFingerprint(1), + &test_group_key("a"), + group_samples("cpu", vec![(ts, 42.0)]), + ) + .unwrap(); + }; + // 15s lies in [0, 30s) and [10s, 40s). + sample(&mut worker, 15_000); + worker.close_through(29_998).unwrap(); + assert!(sink.is_empty()); + worker.close_through(29_999).unwrap(); + assert_eq!( + drain_emitted_sums(&sink), + vec![(0, 30_000, vec!["a".to_string()], 42.0)] + ); + worker.close_through(i64::MAX).unwrap(); + assert_eq!( + drain_emitted_sums(&sink), + vec![(10_000, 40_000, vec!["a".to_string()], 42.0)] + ); + sample(&mut worker, 50_000); + worker.force_close_all().unwrap(); + assert!( + sink.is_empty(), + "input behind the barrier follows the drop policy" + ); + } + + // A barrier queued behind input publishes that input before acknowledging, and input stays open. + #[tokio::test] + async fn watermark_barrier_publishes_queued_input_before_acknowledging() { + let config = make_agg_config( + 1, + "cpu", + AggregationType::SingleSubpopulation, + "Sum", + 10, + 0, + vec![], + ); + let sink = Arc::new(CapturingOutputSink::new()); + let mut worker = make_worker( + HashMap::from([(1, config)]), + sink.clone(), + false, + 0, + LateDataPolicy::Drop, + ); + let (tx, rx) = tokio::sync::mpsc::channel(8); + worker.receiver = rx; + let task = tokio::spawn(worker.run()); + let samples = |ts: i64| WorkerMessage::GroupSamples { + sid: 1, + policy_fp: PolicyFingerprint(1), + group_key: test_group_key(""), + samples: group_samples("cpu", vec![(ts, 2.0)]), + ingest_received_at: std::time::Instant::now(), + }; + tx.send(samples(1_000)).await.unwrap(); + let (reply, ack) = tokio::sync::oneshot::channel(); + tx.send(WorkerMessage::AdvanceWatermark { + event_time_ms: 9_999, + reply, + }) + .await + .unwrap(); + ack.await.unwrap().unwrap(); + assert_eq!(drain_emitted_sums(&sink), vec![(0, 10_000, vec![], 2.0)]); + tx.send(samples(15_000)).await.unwrap(); + tx.send(WorkerMessage::Shutdown).await.unwrap(); + task.await.unwrap(); + assert_eq!( + drain_emitted_sums(&sink), + vec![(10_000, 20_000, vec![], 2.0)] + ); + } + + // A failed worker answers every queued barrier with its error instead of dropping the reply. + #[tokio::test] + async fn failed_worker_answers_queued_watermark_barriers_with_its_error() { + struct FailedSink; + impl OutputSink for FailedSink { + fn emit_batch( + &self, + _: Vec<(PrecomputedOutput, Box)>, + ) -> Result<(), Box> { + Err("deliberate sink failure".into()) + } + } + let config = make_agg_config( + 1, + "cpu", + AggregationType::SingleSubpopulation, + "Sum", + 10, + 0, + vec![], + ); + let mut worker = make_worker( + HashMap::from([(1, config)]), + Arc::new(CapturingOutputSink::new()), + false, + 0, + LateDataPolicy::Drop, + ); + worker.output_sink = Arc::new(FailedSink); + let (tx, rx) = tokio::sync::mpsc::channel(8); + worker.receiver = rx; + tx.send(WorkerMessage::GroupSamples { + sid: 1, + policy_fp: PolicyFingerprint(1), + group_key: test_group_key(""), + samples: group_samples("cpu", vec![(1_000, 2.0)]), + ingest_received_at: std::time::Instant::now(), + }) + .await + .unwrap(); + let mut acks = Vec::new(); + for event_time_ms in [9_999, 19_999] { + let (reply, ack) = tokio::sync::oneshot::channel(); + tx.send(WorkerMessage::AdvanceWatermark { + event_time_ms, + reply, + }) + .await + .unwrap(); + acks.push(ack); + } + // Both barriers are queued before the worker starts. + let task = tokio::spawn(worker.run()); + for ack in acks { + assert!(ack + .await + .unwrap() + .unwrap_err() + .contains("deliberate sink failure")); + } + drop(tx); + task.await.unwrap(); + } + #[test] fn test_flush_publishes_worker_watermark() { let config = make_agg_config( diff --git a/data_plane/tests/asapquery_compatibility_process_e2e.rs b/data_plane/tests/asapquery_compatibility_process_e2e.rs index b18ad599a..226a0c3ed 100644 --- a/data_plane/tests/asapquery_compatibility_process_e2e.rs +++ b/data_plane/tests/asapquery_compatibility_process_e2e.rs @@ -1915,5 +1915,143 @@ async fn collector_free_profile_serves_complete_matrix_and_falls_back_exactly() assert!(metrics.contains("asap_remote_write_rejected_requests_total 1")); } +// A series that stops sending keeps its last pane open, so a grouped query over +// that window falls back; the live watermark barrier closes it with input open. +#[tokio::test] +async fn watermark_barrier_serves_window_of_stopped_series() { + let fallback_app = Router::new() + .route("/-/healthy", get(|| async { "Prometheus is Healthy." })) + .route( + "/api/v1/query", + get(|| async { + Json(serde_json::json!({ + "status": "success", + "data": {"resultType": "vector", "result": []} + })) + }), + ); + let fallback_listener = tokio::net::TcpListener::bind("127.0.0.1:0") + .await + .expect("bind fallback"); + let fallback_address = fallback_listener.local_addr().expect("fallback address"); + tokio::spawn(async move { + axum::serve(fallback_listener, fallback_app) + .await + .expect("serve fallback") + }); + + let backend_port = unused_port(); + let output_dir = tempfile::tempdir().expect("backend output directory"); + let fixture: control_plane::physical::compiler::BackendLocalPlanningInput = + serde_json::from_str(include_str!( + "../../docs/examples/asapquery-compatibility-demo-snapshot.json" + )) + .unwrap(); + let snapshot = output_dir.path().join("snapshot.json"); + std::fs::write( + &snapshot, + serde_json::to_vec("e_snapshot_for_test(fixture)).unwrap(), + ) + .unwrap(); + let child = Command::new(env!("CARGO_BIN_EXE_data_plane")) + .args(["--profile", "asapquery", "--planning-snapshot"]) + .arg(&snapshot) + .arg("--prometheus-server") + .arg(format!("http://{fallback_address}")) + .args([ + "--forward-unsupported-queries", + "--precompute-allowed-lateness-ms", + "0", + "--precompute-flush-interval-ms", + "25", + "--http-port", + &backend_port.to_string(), + "--output-dir", + ]) + .arg(output_dir.path()) + .stdout(Stdio::null()) + .stderr(Stdio::inherit()) + .spawn() + .expect("start production backend"); + let mut child = ChildGuard(child); + let client = reqwest::Client::new(); + let backend = format!("http://127.0.0.1:{backend_port}"); + wait_until_ready(&client, &format!("{backend}/api/v1/health"), &mut child.0).await; + + let now_ms = std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .expect("system time") + .as_millis() as i64; + let base = now_ms - now_ms.rem_euclid(5_000) - 20_000; + // `api` continues past the first 5 s window; `worker` stops inside it. + let request = WriteRequest { + timeseries: vec![ + series_with_labels( + "asap_demo_gauge", + &[("job", "api")], + &[(base + 500, 1.0), (base + 2_900, 3.0), (base + 6_600, 6.0)], + ), + series_with_labels( + "asap_demo_gauge", + &[("job", "worker")], + &[(base + 700, 100.0), (base + 3_100, 100.0)], + ), + ], + }; + assert_eq!(remote_write(&client, &backend, &request).await, 204); + + let query = "topk(1, sum_over_time(asap_demo_gauge[5s]))"; + let eval = (base + 5_000) as f64 / 1_000.0; + let before: Value = client + .get(format!("{backend}/api/v1/query")) + .query(&[("query", query.to_string()), ("time", eval.to_string())]) + .send() + .await + .expect("instant query") + .json() + .await + .expect("instant JSON"); + assert!( + !is_warm(&before), + "the stopped series' pane is still open: {before}" + ); + + let missing = client + .post(format!("{backend}/api/v1/precompute/watermark")) + .send() + .await + .expect("watermark request"); + assert_eq!(missing.status().as_u16(), 400); + let barrier = client + .post(format!("{backend}/api/v1/precompute/watermark")) + .query(&[("time_ms", base + 5_000)]) + .send() + .await + .expect("watermark request"); + assert_eq!( + barrier.status().as_u16(), + 200, + "{}", + barrier.text().await.unwrap() + ); + + let backend_log = output_dir.path().join("query_engine.log"); + let warm = wait_for_warm_instant(&client, &backend, query, eval, &backend_log).await; + assert_eq!(first_value(&warm, "value"), Some(200.0), "{warm}"); + assert_eq!( + warm["data"]["result"][0]["metric"], + serde_json::json!({"job": "worker"}) + ); + // Input stays open after the barrier. + let later = WriteRequest { + timeseries: vec![series_with_labels( + "asap_demo_gauge", + &[("job", "api")], + &[(base + 7_000, 7.0)], + )], + }; + assert_eq!(remote_write(&client, &backend, &later).await, 204); +} + #[path = "support/univmon_erp_process.rs"] mod univmon_erp_process; From da99a171c8f10ad88bc1e2005e5ffea754850cab Mon Sep 17 00:00:00 2001 From: GordonYuanyc Date: Thu, 1 Oct 2026 00:52:02 -0400 Subject: [PATCH 2/6] docs: describe the Remote Write watermark barrier Document POST /api/v1/precompute/watermark for users of the asapquery profile, and correct the design note that said the profile has no public barrier endpoint. Co-Authored-By: Claude Opus 5.5 --- docs/design_docs/asapquery-compatibility-profile.md | 6 ++++-- docs/user_guide/asapquery-profile.md | 10 ++++++++++ 2 files changed, 14 insertions(+), 2 deletions(-) diff --git a/docs/design_docs/asapquery-compatibility-profile.md b/docs/design_docs/asapquery-compatibility-profile.md index e911510fd..ad3e50fb7 100644 --- a/docs/design_docs/asapquery-compatibility-profile.md +++ b/docs/design_docs/asapquery-compatibility-profile.md @@ -211,8 +211,10 @@ Remote Write carries neither a producer-partition roster nor an authoritative watermark. Finite `POST /api/v1/precompute/drain` closes the input generation; it is not continuous completion. Existing typed `SummaryWatermarkBarrier` and coordinator APIs need registered producer/partition identity and publication -ordering before they can establish continuous closure. There is no public HTTP -barrier endpoint in this profile. +ordering before they can establish continuous closure. The live +`POST /api/v1/precompute/watermark?time_ms=T` endpoint is a producer assertion +without that identity: each worker closes, in every group, the stored buckets +whose samples all lie at or before `T`, and input stays open. ### SummaryStore diff --git a/docs/user_guide/asapquery-profile.md b/docs/user_guide/asapquery-profile.md index 8a9a2a0d6..5e1493200 100644 --- a/docs/user_guide/asapquery-profile.md +++ b/docs/user_guide/asapquery-profile.md @@ -66,6 +66,16 @@ forwarded as lifecycle events. Native histograms and exemplars are rejected. For plan/SID inspection and finite-replay boundaries, use the [E2E physical-DAG walkthrough](../evaluation/e2e-physical-dag.md). +A summary window is published only after a later sample arrives in the same +series, or after it has been idle for the window length plus 5 s. A series that +stops sending therefore delays every query over its last window. A producer that +knows it has written every sample up to time `T` (milliseconds) can say so with +`POST /api/v1/precompute/watermark?time_ms=T`. The backend then publishes every +window that ends at or before `T`, in all series, before returning `200`; +earlier acknowledged writes are included. With several producers, send it only +after all of them have written through `T`. Input stays open; a later sample for +a published window is handled as late data. + The deduplication horizon must cover allowed lateness plus the deployment's maximum expected Prometheus retry interval. Configure those assumptions with `--precompute-allowed-lateness-ms`, From 92c49ec8498d99b7f203fccb5fc63d3ab2b8c39f Mon Sep 17 00:00:00 2001 From: GordonYuanyc Date: Thu, 1 Oct 2026 01:42:56 -0400 Subject: [PATCH 3/6] Bound the watermark barrier's closure advance by the latest open bucket close_through capped its scan at the end of a group's latest open bucket but still set the closure watermark to the full barrier. A barrier far in the future (for example microseconds or nanoseconds sent as time_ms) then made every later sample of every existing group late until restart: under ForwardToStore each became a standalone one-sample correction, and counter deltas were dropped. The closure watermark now advances only to min(barrier, latest open bucket end), and not at all for a group without open buckets. Published buckets are still detected as late. The worker barrier tests are merged into one left/right-closed boundary test, and the FullWindow test now checks that input after an i64::MAX barrier aggregates normally (it was dropped before this change). The async run-loop tests are removed; the process test covers that path. Co-Authored-By: Claude Opus 5.5 --- data_plane/src/precompute_engine/worker.rs | 298 +++++---------------- 1 file changed, 74 insertions(+), 224 deletions(-) diff --git a/data_plane/src/precompute_engine/worker.rs b/data_plane/src/precompute_engine/worker.rs index 7fa05a1a8..3a8ca93b6 100644 --- a/data_plane/src/precompute_engine/worker.rs +++ b/data_plane/src/precompute_engine/worker.rs @@ -1286,9 +1286,10 @@ impl Worker { continue; } - // Bound the scan by the latest open bucket, as force_close_all - // does: buckets after it hold no data, and the scan walks one - // slide at a time, so an unbounded barrier would not terminate. + // Advance no further than the end of the latest open bucket, as + // force_close_all does: later buckets hold no data, the scan walks + // one slide at a time, and a far-future barrier must not make + // input for never-opened buckets late. let max_active = state.active_panes.keys().next_back().copied(); let max_sketch = state.sketch_panes.keys().next_back().copied(); if let Some(max_open) = max_active.max(max_sketch) { @@ -1317,9 +1318,9 @@ impl Worker { emit_batch.push((output, accumulator)); } } + state.closure_watermark_ms = state.closure_watermark_ms.max(close_to); + state.prune_pane_wall_clock(); } - state.closure_watermark_ms = barrier; - state.prune_pane_wall_clock(); } if !emit_batch.is_empty() { @@ -3318,114 +3319,73 @@ mod tests { .collect() } - // The barrier closes exactly the stopped group's complete pane; repeats and lower barriers are no-ops. + // A stopped group's complete pane closes exactly at the barrier (T for PromQL + // right-closed membership, else T + 1); repeats and lower barriers are no-ops. #[test] fn watermark_barrier_closes_complete_pane_of_stopped_group() { - let config = make_agg_config( - 1, - "cpu", - AggregationType::SingleSubpopulation, - "Sum", - 10, - 0, - vec![], - ); - let sink = Arc::new(CapturingOutputSink::new()); - let mut worker = make_worker( - HashMap::from([(1, config)]), - sink.clone(), - false, - 0, - LateDataPolicy::Drop, - ); - let pf = PolicyFingerprint(1); - // Group a keeps sending into [10s, 20s); group b stops inside [0, 10s). - for (sid, group, ts, value) in [ - (1, "a", 1_000, 1.0), - (2, "b", 2_000, 2.0), - (1, "a", 12_000, 3.0), - ] { - worker - .process_group_samples( - sid, - pf, - &test_group_key(group), - group_samples("cpu", vec![(ts, value)]), - ) - .unwrap(); - } - assert_eq!( - drain_emitted_sums(&sink), - vec![(0, 10_000, vec!["a".to_string()], 1.0)], - "only group a observed a later sample" - ); - worker.flush_all().unwrap(); - assert!( - sink.is_empty(), - "without wall-clock policy group b stays open" - ); - - // [0, 10s) holds samples up to 9_999 ms. - worker.close_through(9_998).unwrap(); - assert!(sink.is_empty()); - worker.close_through(9_999).unwrap(); - assert_eq!( - drain_emitted_sums(&sink), - vec![(0, 10_000, vec!["b".to_string()], 2.0)] - ); - worker.close_through(9_999).unwrap(); - worker.close_through(5_000).unwrap(); - assert!(sink.is_empty()); - - worker.force_close_all().unwrap(); - assert_eq!( - drain_emitted_sums(&sink), - vec![(10_000, 20_000, vec!["a".to_string()], 3.0)], - "the barrier left group a's later pane open" - ); - } - - // PromQL right-closed groups file a sample at t into the pane ending at t, so a barrier at t closes it. - #[test] - fn watermark_barrier_respects_promql_right_closed_membership() { - let mut config = make_agg_config( - 1, - "cpu", - AggregationType::SingleSubpopulation, - "Sum", - 10, - 0, - vec![], - ); - config - .parameters - .insert("promql_right_closed".into(), serde_json::json!(true)); - let sink = Arc::new(CapturingOutputSink::new()); - let mut worker = make_worker( - HashMap::from([(1, config)]), - sink.clone(), - false, - 0, - LateDataPolicy::Drop, - ); - worker - .process_group_samples( + for (right_closed, stopped_ts, boundary) in [(false, 2_000, 9_999), (true, 10_000, 10_000)] + { + let mut config = make_agg_config( 1, - PolicyFingerprint(1), - &test_group_key("a"), - group_samples("cpu", vec![(10_000, 4.0)]), - ) - .unwrap(); - worker.close_through(9_999).unwrap(); - assert!(sink.is_empty()); - worker.close_through(10_000).unwrap(); - assert_eq!( - drain_emitted_sums(&sink), - vec![(0, 10_000, vec!["a".to_string()], 4.0)] - ); + "cpu", + AggregationType::SingleSubpopulation, + "Sum", + 10, + 0, + vec![], + ); + if right_closed { + config + .parameters + .insert("promql_right_closed".into(), serde_json::json!(true)); + } + let sink = Arc::new(CapturingOutputSink::new()); + let mut worker = make_worker( + HashMap::from([(1, config)]), + sink.clone(), + false, + 0, + LateDataPolicy::Drop, + ); + // Group a moves on to [10s, 20s); group b stops inside [0, 10s). + for (sid, group, ts, value) in [ + (1, "a", 1_000, 1.0), + (2, "b", stopped_ts, 2.0), + (1, "a", 12_000, 3.0), + ] { + worker + .process_group_samples( + sid, + PolicyFingerprint(1), + &test_group_key(group), + group_samples("cpu", vec![(ts, value)]), + ) + .unwrap(); + } + sink.drain(); + + worker.close_through(boundary - 1).unwrap(); + assert!(sink.is_empty(), "right_closed={right_closed}"); + worker.close_through(boundary).unwrap(); + assert_eq!( + drain_emitted_sums(&sink), + vec![(0, 10_000, vec!["b".to_string()], 2.0)], + "right_closed={right_closed}" + ); + worker.close_through(boundary).unwrap(); + worker.close_through(5_000).unwrap(); + assert!(sink.is_empty(), "right_closed={right_closed}"); + worker.force_close_all().unwrap(); + assert_eq!( + drain_emitted_sums(&sink), + vec![(10_000, 20_000, vec!["a".to_string()], 3.0)], + "the barrier left group a's later pane open" + ); + } } - // FullWindow buckets close one by one; an unbounded barrier terminates and later input is late. + // FullWindow buckets close too; an unbounded barrier terminates and does not + // make input for later, never-opened windows late. #[test] fn watermark_barrier_closes_full_windows_and_bounds_its_scan() { let mut config = make_agg_config( @@ -3458,8 +3418,6 @@ mod tests { }; // 15s lies in [0, 30s) and [10s, 40s). sample(&mut worker, 15_000); - worker.close_through(29_998).unwrap(); - assert!(sink.is_empty()); worker.close_through(29_999).unwrap(); assert_eq!( drain_emitted_sums(&sink), @@ -3470,125 +3428,17 @@ mod tests { drain_emitted_sums(&sink), vec![(10_000, 40_000, vec!["a".to_string()], 42.0)] ); + // 50s lies in [30s, 60s), [40s, 70s) and [50s, 80s), none opened before. sample(&mut worker, 50_000); + sample(&mut worker, 51_000); worker.force_close_all().unwrap(); - assert!( - sink.is_empty(), - "input behind the barrier follows the drop policy" - ); - } - - // A barrier queued behind input publishes that input before acknowledging, and input stays open. - #[tokio::test] - async fn watermark_barrier_publishes_queued_input_before_acknowledging() { - let config = make_agg_config( - 1, - "cpu", - AggregationType::SingleSubpopulation, - "Sum", - 10, - 0, - vec![], - ); - let sink = Arc::new(CapturingOutputSink::new()); - let mut worker = make_worker( - HashMap::from([(1, config)]), - sink.clone(), - false, - 0, - LateDataPolicy::Drop, - ); - let (tx, rx) = tokio::sync::mpsc::channel(8); - worker.receiver = rx; - let task = tokio::spawn(worker.run()); - let samples = |ts: i64| WorkerMessage::GroupSamples { - sid: 1, - policy_fp: PolicyFingerprint(1), - group_key: test_group_key(""), - samples: group_samples("cpu", vec![(ts, 2.0)]), - ingest_received_at: std::time::Instant::now(), - }; - tx.send(samples(1_000)).await.unwrap(); - let (reply, ack) = tokio::sync::oneshot::channel(); - tx.send(WorkerMessage::AdvanceWatermark { - event_time_ms: 9_999, - reply, - }) - .await - .unwrap(); - ack.await.unwrap().unwrap(); - assert_eq!(drain_emitted_sums(&sink), vec![(0, 10_000, vec![], 2.0)]); - tx.send(samples(15_000)).await.unwrap(); - tx.send(WorkerMessage::Shutdown).await.unwrap(); - task.await.unwrap(); assert_eq!( drain_emitted_sums(&sink), - vec![(10_000, 20_000, vec![], 2.0)] - ); - } - - // A failed worker answers every queued barrier with its error instead of dropping the reply. - #[tokio::test] - async fn failed_worker_answers_queued_watermark_barriers_with_its_error() { - struct FailedSink; - impl OutputSink for FailedSink { - fn emit_batch( - &self, - _: Vec<(PrecomputedOutput, Box)>, - ) -> Result<(), Box> { - Err("deliberate sink failure".into()) - } - } - let config = make_agg_config( - 1, - "cpu", - AggregationType::SingleSubpopulation, - "Sum", - 10, - 0, - vec![], + [30_000, 40_000, 50_000] + .map(|start| (start, start + 30_000, vec!["a".to_string()], 84.0)) + .to_vec(), + "input after the barrier aggregates normally" ); - let mut worker = make_worker( - HashMap::from([(1, config)]), - Arc::new(CapturingOutputSink::new()), - false, - 0, - LateDataPolicy::Drop, - ); - worker.output_sink = Arc::new(FailedSink); - let (tx, rx) = tokio::sync::mpsc::channel(8); - worker.receiver = rx; - tx.send(WorkerMessage::GroupSamples { - sid: 1, - policy_fp: PolicyFingerprint(1), - group_key: test_group_key(""), - samples: group_samples("cpu", vec![(1_000, 2.0)]), - ingest_received_at: std::time::Instant::now(), - }) - .await - .unwrap(); - let mut acks = Vec::new(); - for event_time_ms in [9_999, 19_999] { - let (reply, ack) = tokio::sync::oneshot::channel(); - tx.send(WorkerMessage::AdvanceWatermark { - event_time_ms, - reply, - }) - .await - .unwrap(); - acks.push(ack); - } - // Both barriers are queued before the worker starts. - let task = tokio::spawn(worker.run()); - for ack in acks { - assert!(ack - .await - .unwrap() - .unwrap_err() - .contains("deliberate sink failure")); - } - drop(tx); - task.await.unwrap(); } #[test] From 813a78d660fc9af604b8e7cf43c91aa4070ba0ab Mon Sep 17 00:00:00 2001 From: GordonYuanyc Date: Thu, 1 Oct 2026 01:42:56 -0400 Subject: [PATCH 4/6] test: trim the watermark barrier receiver and process tests Drop the receiver test (a one-line forward plus channel FIFO; the process test shows input stays open). In the process test, query once right after the barrier returns instead of polling, since panes are published before the 200, and drop the missing-time_ms check. Co-Authored-By: Claude Opus 5.5 --- .../drivers/ingest/prometheus_remote_write.rs | 23 ------------- .../asapquery_compatibility_process_e2e.rs | 33 +++++++++---------- 2 files changed, 16 insertions(+), 40 deletions(-) diff --git a/data_plane/src/drivers/ingest/prometheus_remote_write.rs b/data_plane/src/drivers/ingest/prometheus_remote_write.rs index da73a51e7..8be04952b 100644 --- a/data_plane/src/drivers/ingest/prometheus_remote_write.rs +++ b/data_plane/src/drivers/ingest/prometheus_remote_write.rs @@ -1422,29 +1422,6 @@ mod tests { assert_eq!(drain.await.unwrap().unwrap_err(), "sink write failed"); } - // A live barrier queues behind acknowledged writes and leaves input open. - #[tokio::test] - async fn watermark_barrier_follows_acknowledged_writes_and_keeps_input_open() { - let (receiver, mut worker) = configured_receiver(); - receiver.accept(&one_sample(1.0)).unwrap(); - let handle = receiver.clone(); - let barrier = tokio::spawn(async move { handle.advance_watermark(100).await }); - assert!(matches!( - worker.recv().await.unwrap(), - WorkerMessage::BoundInput { .. } - )); - let WorkerMessage::AdvanceWatermark { - event_time_ms: 100, - reply, - } = worker.recv().await.unwrap() - else { - panic!("expected watermark barrier") - }; - receiver.accept(&one_sample(1.0)).unwrap(); - reply.send(Ok(())).unwrap(); - barrier.await.unwrap().unwrap(); - } - #[test] fn rejects_writes_without_an_active_physical_plan() { let (sender, _worker) = mpsc::channel(1); diff --git a/data_plane/tests/asapquery_compatibility_process_e2e.rs b/data_plane/tests/asapquery_compatibility_process_e2e.rs index 226a0c3ed..2c3ab5a81 100644 --- a/data_plane/tests/asapquery_compatibility_process_e2e.rs +++ b/data_plane/tests/asapquery_compatibility_process_e2e.rs @@ -2002,26 +2002,24 @@ async fn watermark_barrier_serves_window_of_stopped_series() { let query = "topk(1, sum_over_time(asap_demo_gauge[5s]))"; let eval = (base + 5_000) as f64 / 1_000.0; - let before: Value = client - .get(format!("{backend}/api/v1/query")) - .query(&[("query", query.to_string()), ("time", eval.to_string())]) - .send() - .await - .expect("instant query") - .json() - .await - .expect("instant JSON"); + let (client_ref, backend_ref) = (&client, &backend); + let instant = move || async move { + client_ref + .get(format!("{backend_ref}/api/v1/query")) + .query(&[("query", query.to_string()), ("time", eval.to_string())]) + .send() + .await + .expect("instant query") + .json::() + .await + .expect("instant JSON") + }; + let before = instant().await; assert!( !is_warm(&before), "the stopped series' pane is still open: {before}" ); - let missing = client - .post(format!("{backend}/api/v1/precompute/watermark")) - .send() - .await - .expect("watermark request"); - assert_eq!(missing.status().as_u16(), 400); let barrier = client .post(format!("{backend}/api/v1/precompute/watermark")) .query(&[("time_ms", base + 5_000)]) @@ -2035,8 +2033,9 @@ async fn watermark_barrier_serves_window_of_stopped_series() { barrier.text().await.unwrap() ); - let backend_log = output_dir.path().join("query_engine.log"); - let warm = wait_for_warm_instant(&client, &backend, query, eval, &backend_log).await; + // Panes are published before the barrier returns, so one query suffices. + let warm = instant().await; + assert!(is_warm(&warm), "{warm}"); assert_eq!(first_value(&warm, "value"), Some(200.0), "{warm}"); assert_eq!( warm["data"]["result"][0]["metric"], From 70f6158b978ee33c8e1d5fd9ea4d8cb91384a9ea Mon Sep 17 00:00:00 2001 From: GordonYuanyc Date: Thu, 1 Oct 2026 01:42:56 -0400 Subject: [PATCH 5/6] docs: mention the absolute wall-clock deadline for open windows A window also closes at the latest window + 5 s after its first sample, not only after being idle that long. Co-Authored-By: Claude Opus 5.5 --- docs/user_guide/asapquery-profile.md | 9 +++++---- 1 file changed, 5 insertions(+), 4 deletions(-) diff --git a/docs/user_guide/asapquery-profile.md b/docs/user_guide/asapquery-profile.md index 5e1493200..16b79dc4a 100644 --- a/docs/user_guide/asapquery-profile.md +++ b/docs/user_guide/asapquery-profile.md @@ -66,10 +66,11 @@ forwarded as lifecycle events. Native histograms and exemplars are rejected. For plan/SID inspection and finite-replay boundaries, use the [E2E physical-DAG walkthrough](../evaluation/e2e-physical-dag.md). -A summary window is published only after a later sample arrives in the same -series, or after it has been idle for the window length plus 5 s. A series that -stops sending therefore delays every query over its last window. A producer that -knows it has written every sample up to time `T` (milliseconds) can say so with +A summary window is published when a later sample arrives in the same series, +or at the latest the window length plus 5 s after its first sample arrived. A +series that stops sending therefore delays every query over its last window. A +producer that knows it has written every sample up to time `T` (milliseconds) +can say so with `POST /api/v1/precompute/watermark?time_ms=T`. The backend then publishes every window that ends at or before `T`, in all series, before returning `200`; earlier acknowledged writes are included. With several producers, send it only From e60554db755bba32b9d02aff45ed0af78f44797c Mon Sep 17 00:00:00 2001 From: GordonYuanyc Date: Thu, 1 Oct 2026 01:45:10 -0400 Subject: [PATCH 6/6] test: shorten the watermark barrier process test assertions Co-Authored-By: Claude Opus 5.5 --- .../asapquery_compatibility_process_e2e.rs | 19 ++++--------------- 1 file changed, 4 insertions(+), 15 deletions(-) diff --git a/data_plane/tests/asapquery_compatibility_process_e2e.rs b/data_plane/tests/asapquery_compatibility_process_e2e.rs index 2c3ab5a81..a6ad4f2a1 100644 --- a/data_plane/tests/asapquery_compatibility_process_e2e.rs +++ b/data_plane/tests/asapquery_compatibility_process_e2e.rs @@ -2015,10 +2015,7 @@ async fn watermark_barrier_serves_window_of_stopped_series() { .expect("instant JSON") }; let before = instant().await; - assert!( - !is_warm(&before), - "the stopped series' pane is still open: {before}" - ); + assert!(!is_warm(&before), "stopped series' pane is open: {before}"); let barrier = client .post(format!("{backend}/api/v1/precompute/watermark")) @@ -2026,12 +2023,7 @@ async fn watermark_barrier_serves_window_of_stopped_series() { .send() .await .expect("watermark request"); - assert_eq!( - barrier.status().as_u16(), - 200, - "{}", - barrier.text().await.unwrap() - ); + assert_eq!(barrier.status().as_u16(), 200); // Panes are published before the barrier returns, so one query suffices. let warm = instant().await; @@ -2042,12 +2034,9 @@ async fn watermark_barrier_serves_window_of_stopped_series() { serde_json::json!({"job": "worker"}) ); // Input stays open after the barrier. + let later = series_with_labels("asap_demo_gauge", &[("job", "api")], &[(base + 7_000, 7.0)]); let later = WriteRequest { - timeseries: vec![series_with_labels( - "asap_demo_gauge", - &[("job", "api")], - &[(base + 7_000, 7.0)], - )], + timeseries: vec![later], }; assert_eq!(remote_write(&client, &backend, &later).await, 204); }