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
23 changes: 23 additions & 0 deletions crates/asap_types/src/query_plan/native.rs
Original file line number Diff line number Diff line change
Expand Up @@ -111,7 +111,30 @@ 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) {
return Err(invalid(
"query-time raw inputs cannot be mixed with installed state",
));
}
for input in inputs {
if raw(input) {
continue;
}
if let Some(QueryPlanNode::ReadMaterialization { binding }) = self.nodes.get(input) {
if binding.readout_lookback_ms != Some(self.instant.lookback_ms)
|| binding.window_ms != self.instant.lookback_ms
Expand Down
20 changes: 15 additions & 5 deletions data_plane/src/query_engines/asap_query_engine/engine.rs
Original file line number Diff line number Diff line change
Expand Up @@ -519,19 +519,29 @@ impl ASAPQueryEngine {
inputs.iter().all(|id| {
matches!(
entry.nodes.get(id),
Some(asap_types::query_plan::QueryPlanNode::ReadMaterialization { .. })
Some(
asap_types::query_plan::QueryPlanNode::ReadMaterialization { .. }
| asap_types::query_plan::QueryPlanNode::Logical {
operator: asap_types::query_plan::query_time::QueryTimeOperator::Scan { .. },
..
}
)
)
})
});
let result = if native_stored {
let store = index.ok_or_else(|| {
EngineError::capability_miss("native_stored", "summary store unavailable")
})?;
// Query-time raw reads are external queries under the forwarding policy.
let raw_endpoint = self
.exact_subquery_endpoint
.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,
physical.query_plan.plan_id,
physical.query_plan.plan_version,
store,
index,
raw_endpoint,
at,
)
} else {
Expand Down
19 changes: 17 additions & 2 deletions data_plane/src/query_engines/asap_query_engine/exact_subqueries.rs
Original file line number Diff line number Diff line change
Expand Up @@ -87,10 +87,25 @@ fn leaves(
},
// A join is a typed composition node rather than a Logical
// wrapper, but its value input can still be a Prometheus leaf.
QueryPlanNode::PhysicalFragment { inputs, .. }
| QueryPlanNode::Physical { inputs, .. } => {
QueryPlanNode::PhysicalFragment { inputs, .. } => {
pending.extend(inputs.iter().map(|input| (*input, at)));
}
// Raw Scan inputs of a physical program are bound as lazy physical
// sources at execution, not prepared as exact leaves.
QueryPlanNode::Physical { inputs, .. } => pending.extend(
inputs
.iter()
.filter(|input| {
!matches!(
entry.nodes.get(input),
Some(QueryPlanNode::Logical {
operator: QueryTimeOperator::Scan { .. },
..
})
)
})
.map(|input| (*input, at)),
),

QueryPlanNode::ExternalExact { request, inputs } => {
pending.extend(inputs.iter().map(|input| (*input, at)));
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -587,11 +587,11 @@ mod tests {
#[test]
fn selected_candidate_input_budget_is_a_terminal_error() {
let plan = CompiledPhysicalDag::decode(&sorted()).unwrap();
let error = execute_batches(&plan, 1, 1, 42, |_, schema| {
Ok(Batch::try_new(
schema.clone(),
let error = execute_batches(&plan, 1, 1, 42, |_, contract| {
Ok(BoundInput::Rows(Batch::try_new(
contract.schema.clone(),
vec![vec![Value::Float64(1.)]],
)?)
)?))
})
.expect_err("input must exceed the byte budget");
assert!(
Expand Down Expand Up @@ -767,7 +767,8 @@ where
max_bytes,
bindings.len(),
at,
|input_id, schema| {
|input_id, contract| {
let schema = &contract.schema;
let values = super::vector(super::from_result(callback(bindings[&input_id], at)?)?)?;
let rows = values
.into_iter()
Expand All @@ -781,7 +782,9 @@ where
.map_err(|e| miss(e.to_string()))
})
.collect::<Result<Vec<_>, _>>()?;
Batch::try_new(schema.clone(), rows).map_err(EngineError::from)
Batch::try_new(schema.clone(), rows)
.map(BoundInput::Rows)
.map_err(EngineError::from)
},
)
}
Expand All @@ -790,7 +793,8 @@ pub(in crate::query_engines::asap_query_engine) fn execute_stored(
entry: &asap_types::query_plan::QueryPlanEntry,
plan_id: u64,
plan_version: u64,
store: &crate::storage_engines::sketch_db::index::SketchStore,
store: Option<&crate::storage_engines::sketch_db::index::SketchStore>,
raw_endpoint: Option<(&reqwest::Client, &str)>,
at: u64,
) -> Result<
(
Expand All @@ -805,16 +809,28 @@ pub(in crate::query_engines::asap_query_engine) fn execute_stored(
let program = entry
.recover_vector_physical_dag()
.map_err(|e| miss(e.to_string()))?;
execute_batches(&program, max_bytes, inputs.len(), at, |id, schema| {
execute_batches(&program, max_bytes, inputs.len(), at, |id, contract| {
let schema = &contract.schema;
let index = sources
.iter()
.position(|source| *source == id)
.ok_or_else(|| miss("native source is unbound"))?;
let Some(asap_types::query_plan::QueryPlanNode::ReadMaterialization { binding }) =
entry.nodes.get(&inputs[index])
let node = entry.nodes.get(&inputs[index]);
if let Some(asap_types::query_plan::QueryPlanNode::Logical { operator, .. }) = node {
let (client, endpoint) = raw_endpoint
.ok_or_else(|| miss("query-time raw input has no Prometheus endpoint"))?;
let at = i64::try_from(at).map_err(|_| miss("native timestamp overflow"))?;
return Ok(BoundInput::Lazy(super::super::raw_source::bind(
contract, operator, at, client, endpoint,
)?));
}
let Some(asap_types::query_plan::QueryPlanNode::ReadMaterialization { binding }) = node
else {
return Err(miss("native stored source has no deployed summary binding"));
};
let store = store.ok_or_else(|| {
EngineError::capability_miss("native_stored", "summary store unavailable")
})?;
let end = i64::try_from(at).map_err(|_| miss("native timestamp overflow"))?;
let start = at
.checked_sub(binding.window_ms)
Expand All @@ -836,6 +852,7 @@ pub(in crate::query_engines::asap_query_engine) fn execute_stored(
schema.clone(),
max_bytes as usize,
)
.map(BoundInput::Rows)
.map_err(|error| match error {
crate::storage_engines::sketch_db::index::NativeReadError::Unavailable(message) => {
miss(message)
Expand All @@ -847,12 +864,22 @@ pub(in crate::query_engines::asap_query_engine) fn execute_stored(
})
}

/// Stored and protocol inputs are read before execution; query-time raw
/// sources open lazily inside the run and charge its budget as they stream.
enum BoundInput {
Rows(Batch),
Lazy(asap_physical_operators::physical_planner::Source<'static>),
}

fn execute_batches(
program: &asap_physical_operators::physical_planner::CompiledPhysicalDag,
max_bytes: u64,
input_count: usize,
at: u64,
mut input_batch: impl FnMut(u64, &Schema) -> Result<Batch, EngineError>,
mut input_batch: impl FnMut(
u64,
&asap_physical_operators::physical_planner::InputContract,
) -> Result<BoundInput, EngineError>,
) -> Result<
(
crate::query_engines::query_result::QueryResult,
Expand Down Expand Up @@ -887,7 +914,13 @@ fn execute_batches(
let mut input_bytes = 0usize;
for (input_id, input) in program.input_contracts() {
crate::query_engines::request::check()?;
let batch = input_batch(input_id, &input.schema)?;
let batch = match input_batch(input_id, input)? {
BoundInput::Rows(batch) => batch,
BoundInput::Lazy(source) => {
sources.insert(input_id, source);
continue;
}
};
input_bytes = input_bytes
.checked_add(batch.bytes())
.ok_or(asap_physical_operators::Error::MemoryLimit)?;
Expand Down Expand Up @@ -996,6 +1029,7 @@ mod request_contract_tests {
)?;
execute_batches(&program, 64 * 1024, 1, 0, |_, _| {
Batch::try_new(schema.clone(), vec![vec![Value::Float64(1.0)]])
.map(BoundInput::Rows)
.map_err(EngineError::from)
})
},
Expand Down
1 change: 1 addition & 0 deletions data_plane/src/query_engines/asap_query_engine/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ mod exact_subqueries;
pub mod live_serve;
pub mod logical_dag;
pub mod post_asap_readout;
mod raw_source;
pub mod summary_executor;

pub use crate::storage_engines::sketch_db::query as asap_tier;
Expand Down
Loading
Loading