diff --git a/data_plane/src/drivers/ingest/prometheus_remote_write.rs b/data_plane/src/drivers/ingest/prometheus_remote_write.rs index 6f2c2ba2a..8be04952b 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() { 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..3a8ca93b6 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,94 @@ 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; + } + + // 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) { + 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 = state.closure_watermark_ms.max(close_to); + 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 +3301,146 @@ 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() + } + + // 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() { + for (right_closed, stopped_ts, boundary) in [(false, 2_000, 9_999), (true, 10_000, 10_000)] + { + let mut config = make_agg_config( + 1, + "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 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( + 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_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)] + ); + // 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_eq!( + drain_emitted_sums(&sink), + [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" + ); + } + #[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..a6ad4f2a1 100644 --- a/data_plane/tests/asapquery_compatibility_process_e2e.rs +++ b/data_plane/tests/asapquery_compatibility_process_e2e.rs @@ -1915,5 +1915,131 @@ 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 (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), "stopped series' pane is open: {before}"); + + 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); + + // 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"], + 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![later], + }; + assert_eq!(remote_write(&client, &backend, &later).await, 204); +} + #[path = "support/univmon_erp_process.rs"] mod univmon_erp_process; 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..16b79dc4a 100644 --- a/docs/user_guide/asapquery-profile.md +++ b/docs/user_guide/asapquery-profile.md @@ -66,6 +66,17 @@ 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 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 +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`,