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
70 changes: 70 additions & 0 deletions crates/asap_types/src/query_plan.rs
Original file line number Diff line number Diff line change
Expand Up @@ -756,6 +756,25 @@ impl MaterializationBinding {
}
}
}

/// Interval between the ends of successive complete stored windows.
pub fn slide_ms(&self) -> u64 {
self.full_window_slide_ms.unwrap_or(self.window_ms)
}

/// End of the newest window on this output's grid that ends at or before
/// `at_ms`; None when the grid is unknown or that window starts before 0.
pub fn latest_window_end_at_or_before(&self, at_ms: u64) -> Option<u64> {
let origin = i128::from(self.pane_origin_ms?);
let (window, slide) = (i128::from(self.window_ms), i128::from(self.slide_ms()));
if window == 0 || slide == 0 {
return None;
}
// Pane and full-window ends both lie on origin + window + n * slide.
let at = i128::from(at_ms);
let end = at - (at - origin - window).rem_euclid(slide);
u64::try_from(end).ok().filter(|end| *end >= self.window_ms)
}
}

#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
Expand Down Expand Up @@ -1096,3 +1115,54 @@ mod retired_plan_tests {
}
}
}

#[cfg(test)]
mod window_grid_tests {
use super::*;

fn binding(window_ms: u64, slide: Option<u64>, origin: Option<i64>) -> MaterializationBinding {
MaterializationBinding {
stored_output_reference: serde_json::from_value(serde_json::json!({
"stored_output_id": 1,
"definition_id": format!("sds-v1:{}", "a".repeat(64)),
}))
.unwrap(),
full_window_slide_ms: slide,
materialization: StoredOutputId(1),
output_grouping: PhysicalGrouping::Reduce(vec![]),
item_labels: vec![],
window_ms,
pane_origin_ms: origin,
readout_lookback_ms: Some(window_ms),
}
}

// The newest complete window end follows the pane or full-window grid at any
// origin, and every returned window is one the grid can store.
#[test]
fn latest_window_end_follows_the_stored_grid() {
let panes = binding(60_000, None, Some(0));
assert_eq!(panes.slide_ms(), 60_000);
assert_eq!(panes.latest_window_end_at_or_before(120_000), Some(120_000));
assert_eq!(panes.latest_window_end_at_or_before(179_999), Some(120_000));
assert_eq!(panes.latest_window_end_at_or_before(59_999), None);
let sliding = binding(300_000, Some(60_000), Some(5_000));
assert_eq!(sliding.slide_ms(), 60_000);
let end = sliding.latest_window_end_at_or_before(1_000_000).unwrap();
assert_eq!(end, 965_000);
assert!(sliding.covers_range(end - 300_000, end));
assert_eq!(
binding(60_000, None, None).latest_window_end_at_or_before(120_000),
None
);
let shifted = binding(60_000, None, Some(-30_000));
assert_eq!(
shifted.latest_window_end_at_or_before(100_000),
Some(90_000)
);
assert_eq!(
binding(60_000, None, Some(500_000)).latest_window_end_at_or_before(100_000),
Some(80_000)
);
}
}
58 changes: 41 additions & 17 deletions crates/asap_types/src/query_plan/native.rs
Original file line number Diff line number Diff line change
Expand Up @@ -116,6 +116,32 @@ impl QueryPlanEntry {
}
}

/// A query-time raw input: a named range selector read from the raw-series
/// endpoint at the evaluation time.
fn is_query_time_raw(&self, input: &QueryNodeId) -> bool {
matches!(
self.nodes.get(input),
Some(QueryPlanNode::Logical {
operator: query_time::QueryTimeOperator::Scan {
metric: Some(_),
range_ms: Some(_),
..
},
..
})
)
}

/// True when the physical inputs include both query-time raw inputs and
/// stored inputs.
pub fn mixes_raw_and_stored_inputs(&self) -> bool {
self.physical_vector_binding()
.is_some_and(|(inputs, _, _)| {
inputs.iter().any(|input| self.is_query_time_raw(input))
&& !inputs.iter().all(|input| self.is_query_time_raw(input))
})
}

/// Whether the root physical result drops `__name__` from series identities.
pub fn drops_metric_name(&self) -> bool {
matches!(
Expand All @@ -142,24 +168,22 @@ impl QueryPlanEntry {
{
return Err(invalid("invalid physical vector source mapping or budget"));
}
let raw = |input: &QueryNodeId| {
matches!(
self.nodes.get(input),
Some(QueryPlanNode::Logical {
operator: query_time::QueryTimeOperator::Scan {
metric: Some(_),
range_ms: Some(_),
..
},
..
})
)
};
// Raw samples are read from the external endpoint at query time; like
// exact cuts, they share no snapshot with installed summary state.
if inputs.iter().any(raw) && !inputs.iter().all(raw) {
let raw = |input: &QueryNodeId| self.is_query_time_raw(input);
// Only a stored native batch has a lag the executor can check against
// the evaluation time of the raw inputs. Its admitted families (Sum and
// heap sketches, checked below) read out independently of the
// evaluation time, so a lagged batch is read as of its own window.
if inputs.iter().any(raw)
&& inputs.iter().any(|input| {
!raw(input)
&& !matches!(
self.nodes.get(input),
Some(QueryPlanNode::ReadMaterialization { .. })
)
})
{
return Err(invalid(
"query-time raw inputs cannot be mixed with installed state",
"query-time raw inputs mix only with stored native batches",
));
}
for input in inputs {
Expand Down
39 changes: 39 additions & 0 deletions data_plane/src/drivers/query/servers/http.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1488,6 +1488,23 @@ async fn process_via_router(
}
}

/// Remove the engine's stored-input lag note from the client warnings.
fn extract_stored_input_lag(value: &mut serde_json::Value) -> Option<u64> {
let warnings = value.get_mut("warnings")?.as_array_mut()?;
let mut lag = None;
warnings.retain(|warning| {
let Some(text) = warning
.as_str()
.and_then(|text| text.strip_prefix("asap_stored_input_lag_ms:"))
else {
return true;
};
lag = text.parse().ok();
false
});
lag
}

fn extract_logical_provenance(
value: &mut serde_json::Value,
) -> Option<Result<(u64, u64, u64, u64, u64, u64, u64), ()>> {
Expand Down Expand Up @@ -1622,6 +1639,12 @@ async fn annotate_data_source(response: Response, data_source_id: &'static str)
return Response::from_parts(parts, axum::body::Body::from(bytes));
}
if data_source_id == "asap_query" {
if let Some(lag) = extract_stored_input_lag(&mut value) {
parts.headers.insert(
"x-asap-stored-input-lag-ms",
axum::http::HeaderValue::from(lag),
);
}
if let Some(provenance) = extract_logical_provenance(&mut value) {
let (route, detail) = match provenance {
Ok((raw, summary, memo, remote, _legacy_indexes, rpcs, branches)) => {
Expand Down Expand Up @@ -6477,6 +6500,22 @@ mod logical_provenance_tests {
assert_eq!(value["warnings"], serde_json::json!(["partial data"]));
}

// The stored-input lag of a mixed program becomes response metadata, not a client warning.
#[tokio::test]
async fn stored_input_lag_is_a_response_header() {
let response = Json(serde_json::json!({"status":"success", "warnings":[
"asap_stored_input_lag_ms:40000", "asap_logical_stats:raw=0,summary=2,memo_hits=0,remote=0,remote_rpcs=0,remote_branches=0"
], "data":{"resultType":"vector", "result":[]}})).into_response();
let response = annotate_data_source(response, "asap_query").await;
assert_eq!(response.headers()["x-asap-stored-input-lag-ms"], "40000");
assert_eq!(response.headers()["x-asap-execution"], "warm");
let bytes = axum::body::to_bytes(response.into_body(), 4096)
.await
.unwrap();
let value: serde_json::Value = serde_json::from_slice(&bytes).unwrap();
assert!(value["warnings"].as_array().unwrap().is_empty());
}

#[test]
fn contradictory_provenance_is_not_warm() {
// Any observed local raw branch invalidates a deployed plan.
Expand Down
7 changes: 7 additions & 0 deletions data_plane/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -159,6 +159,12 @@ struct Args {
#[arg(long)]
disable_query_forwarding: bool,

/// Largest staleness of stored state combined with raw series read at the
/// query time, applied to every stored output; beyond it the query uses
/// the exact fallback. Unset: one slide interval of each stored output.
#[arg(long)]
max_stored_input_lag_ms: Option<u64>,

/// Database path (currently unused, kept for compatibility)
#[arg(long, default_value = "sketchdb.db")]
db_path: String,
Expand Down Expand Up @@ -762,6 +768,7 @@ async fn main() -> Result<()> {
.with_sketch_index(summary_store.clone())
.with_active_physical_plan(active_physical_plan.clone())
.with_query_forwarding_policy(query_forwarding_policy)
.with_max_stored_input_lag_ms(args.max_stored_input_lag_ms)
.with_exact_subquery_endpoint(args.prometheus_server.clone())
.with_metricsql_exact_subquery_endpoint(args.victoriametrics_url.clone());

Expand Down
66 changes: 63 additions & 3 deletions data_plane/src/query_engines/asap_query_engine/engine.rs
Original file line number Diff line number Diff line change
Expand Up @@ -78,6 +78,36 @@ mod readiness_coverage_tests {
}
}

#[cfg(test)]
mod stored_input_lag_tests {
// Only an execution that bound stored state beside raw inputs reports its lag.
#[test]
fn logical_execution_reports_stored_input_lag_when_mixed() {
use crate::query_engines::query_result::QueryResult;
let lag_warnings = |lag| {
let mut result = QueryResult::vector(vec![], 0);
let stats = super::super::logical_dag::ExecutionStats {
stored_input_lag_ms: lag,
..Default::default()
};
super::annotate_logical_execution(&mut result, &stats);
let QueryResult::Vector(result) = result else {
unreachable!()
};
result
.warnings
.into_iter()
.filter(|warning| warning.starts_with("asap_stored_input_lag_ms:"))
.collect::<Vec<_>>()
};
assert_eq!(
lag_warnings(Some(40_000)),
["asap_stored_input_lag_ms:40000"]
);
assert!(lag_warnings(None).is_empty());
}
}

#[cfg(test)]
mod forwarding_policy_tests {
use super::ASAPQueryEngine;
Expand Down Expand Up @@ -198,6 +228,9 @@ pub struct ASAPQueryEngine {
query_forwarding_policy: crate::query_engines::QueryForwardingPolicy,
exact_subquery_client: reqwest::Client,
execution_limits: asap_physical_operators::dag::Limits,
/// Deployment override of the staleness bound for stored inputs bound
/// beside query-time raw inputs; None uses each stored output's slide.
max_stored_input_lag_ms: Option<u64>,
}

impl ASAPQueryEngine {
Expand Down Expand Up @@ -278,6 +311,7 @@ impl ASAPQueryEngine {
Self {
prometheus_scrape_interval,
execution_limits: Default::default(),
max_stored_input_lag_ms: None,
summary_store: None,
active_physical_plan: None,
exact_subquery_endpoint: None,
Expand Down Expand Up @@ -307,6 +341,10 @@ impl ASAPQueryEngine {
self.query_forwarding_policy = policy;
self
}
pub fn with_max_stored_input_lag_ms(mut self, max_lag_ms: Option<u64>) -> Self {
self.max_stored_input_lag_ms = max_lag_ms;
self
}
pub fn with_execution_limits(mut self, limits: asap_physical_operators::dag::Limits) -> Self {
self.execution_limits = limits;
self
Expand Down Expand Up @@ -535,13 +573,31 @@ impl ASAPQueryEngine {
.as_deref()
.filter(|_| self.query_forwarding_policy.allows_external_queries())
.map(|endpoint| (&self.exact_subquery_client, endpoint));
super::logical_dag::native_values::execute_stored(
entry,
let (plan_id, plan_version) = (
physical.query_plan.plan_id,
physical.query_plan.plan_version,
index,
);
super::logical_dag::native_values::execute_stored(
entry,
raw_endpoint,
self.max_stored_input_lag_ms,
at,
&mut |binding, window, schema, max_bytes| {
let store = index.ok_or("summary store unavailable")?;
let address = asap_types::sds::StoredSummaryKey {
plan_id,
plan_version,
stored_output_id: binding.stored_output_reference.stored_output_id,
population: std::collections::BTreeMap::new(),
window,
};
store.read_bound_native_summary(
&address,
&binding.stored_output_reference,
schema,
max_bytes,
)
},
)
} else {
super::logical_dag::execute_installed(entry, leaves, at, |root, evaluation_ms| {
Expand Down Expand Up @@ -737,6 +793,7 @@ impl ASAPQueryEngine {
total.remote_evaluations += stats.remote_evaluations;
total.remote_rpcs += stats.remote_rpcs;
total.remote_branch_evaluations += stats.remote_branch_evaluations;
total.stored_input_lag_ms = total.stored_input_lag_ms.max(stats.stored_input_lag_ms);
let _step_result = crate::query_engines::request::reserve(result.retained_bytes())?;
let QueryResult::Vector(result) = result else {
return Err(EngineError::capability_miss(
Expand Down Expand Up @@ -1059,6 +1116,9 @@ fn annotate_logical_execution(
.into(),
);
}
if let Some(lag) = stats.stored_input_lag_ms {
warnings.push(format!("asap_stored_input_lag_ms:{lag}"));
}
warnings.push(format!(
"asap_logical_stats:raw={},summary={},memo_hits={},remote={},remote_rpcs={},remote_branches={}",
stats.raw_scan_evaluations,
Expand Down
3 changes: 3 additions & 0 deletions data_plane/src/query_engines/asap_query_engine/logical_dag.rs
Original file line number Diff line number Diff line change
Expand Up @@ -57,6 +57,9 @@ pub struct ExecutionStats {
pub remote_evaluations: usize,
pub remote_rpcs: usize,
pub remote_branch_evaluations: usize,
/// Largest `t_q - t_s` over stored inputs bound beside query-time raw
/// inputs; None when the program has no such mix.
pub stored_input_lag_ms: Option<u64>,
}
fn miss(detail: impl Into<String>) -> EngineError {
EngineError::capability_miss("installed_logical_dag", detail)
Expand Down
Loading
Loading