Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 13 additions & 0 deletions data_plane/src/drivers/ingest/prometheus_remote_write.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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() {
Expand Down
41 changes: 41 additions & 0 deletions data_plane/src/drivers/query/servers/http.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down Expand Up @@ -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))
Expand Down Expand Up @@ -5358,6 +5366,39 @@ async fn handle_precompute_drain(State(state): State<AppState>) -> 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<AppState>,
Query(params): Query<HashMap<String, String>>,
) -> 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::<i64>().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<AppState>) -> axum::response::Response {
use axum::http::StatusCode;
Expand Down
32 changes: 32 additions & 0 deletions data_plane/src/precompute_engine/series_router.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<Result<(), String>>),
/// 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<Result<(), String>>,
},
/// Execute downstream DAG work after all raw windows have been sealed.
CompleteDag {
plan: Arc<crate::storage_engines::types::RuntimePhysicalPlan>,
Expand Down Expand Up @@ -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"),
}
Expand Down Expand Up @@ -275,6 +285,7 @@ impl SeriesRouter {
}
WorkerMessage::Flush
| WorkerMessage::Drain(_)
| WorkerMessage::AdvanceWatermark { .. }
| WorkerMessage::Shutdown
| WorkerMessage::CompleteDag { .. } => 0,
};
Expand Down Expand Up @@ -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() {
Expand Down
Loading
Loading