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
12 changes: 6 additions & 6 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

11 changes: 6 additions & 5 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -16,10 +16,10 @@ version = "0.1.0"
[workspace.dependencies]
# Keep Planner frontends, selection, and IR on the same immutable revision.
# Alias upstream asap-types because this workspace also defines asap_types.
planner-types = { package = "asap-types", git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "9a13ae43eecc8487fe6a25a98f6ec09399d2073e" }
asap-aware-mapping = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "9a13ae43eecc8487fe6a25a98f6ec09399d2073e" }
asap-frontend-promql = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "9a13ae43eecc8487fe6a25a98f6ec09399d2073e" }
asap-frontend-sql = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "9a13ae43eecc8487fe6a25a98f6ec09399d2073e" }
planner-types = { package = "asap-types", git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "ca422df610009bdcc4f05b9becf4b2f5ee0c56b2" }
asap-aware-mapping = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "ca422df610009bdcc4f05b9becf4b2f5ee0c56b2" }
asap-frontend-promql = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "ca422df610009bdcc4f05b9becf4b2f5ee0c56b2" }
asap-frontend-sql = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "ca422df610009bdcc4f05b9becf4b2f5ee0c56b2" }

# Shared external deps (used by 2+ crates)
serde = { version = "1.0", features = ["derive"] }
Expand All @@ -39,9 +39,10 @@ arc-swap = "1.7"
reqwest = { version = "0.12", default-features = false, features = ["json", "rustls-tls"] }

# Internal crates
asap-physical-operators = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "9a13ae43eecc8487fe6a25a98f6ec09399d2073e" }
asap-physical-operators = { git = "https://github.com/ProjectASAP/ASAPPlanner", rev = "ca422df610009bdcc4f05b9becf4b2f5ee0c56b2" }
asap_sketch_codec = { path = "crates/asap_sketch_codec" }
asap_summary_state = { path = "crates/asap_summary_state" }
asap_types = { path = "crates/asap_types" }
asap_otel_proto = { path = "crates/asap_otel_proto" }
indexmap = { version = "2.0", features = ["serde"] }

20 changes: 20 additions & 0 deletions control_plane/src/physical/executable_binding.rs
Original file line number Diff line number Diff line change
Expand Up @@ -100,6 +100,7 @@ pub(super) fn compile_precompute_programs(
use asap_types::executable_plan::BackendNodeBinding;
let dag = installed.document.decode()?;
let mut groups = std::collections::BTreeMap::<_, Vec<u64>>::new();
let mut raw = Vec::new();
for sink in &installed.binding.precompute_sinks {
if installed.native_programs.contains_key(sink) {
continue;
Expand All @@ -114,6 +115,24 @@ pub(super) fn compile_precompute_programs(
.find(|c| c.policy_fingerprint() == stored_output.fingerprint())
.ok_or("precompute output configuration missing")?;
let Some(derived) = &config.derived_input else {
// Raw ingest: Planner compiles the summary over its raw sample boundary.
let raw_inputs = dag
.edges
.iter()
.filter(|edge| edge.consumer == *sink)
.map(|edge| u64::from(edge.producer.0))
.collect::<Vec<_>>();
let program = asap_physical_operators::physical_planner::precompute::compile(
&dag,
&raw_inputs,
&[u64::from(sink.0)],
)
.map_err(|e| format!("raw precompute output {}: {e}", sink.0))?;
raw.push((
*sink,
serde_json::from_slice(&program.encode().map_err(|e| e.to_string())?)
.map_err(|e| e.to_string())?,
));
continue;
};
let frontiers = installed.binding.nodes.iter().filter_map(|(id, binding)| {
Expand All @@ -133,6 +152,7 @@ pub(super) fn compile_precompute_programs(
.or_default()
.push(u64::from(sink.0));
}
installed.native_programs.extend(raw);
for ((frontiers, _, _, _, _), roots) in groups {
let program = asap_physical_operators::physical_planner::precompute::compile(
&dag, &frontiers, &roots,
Expand Down
6 changes: 5 additions & 1 deletion control_plane/src/physical/plan_dot.rs
Original file line number Diff line number Diff line change
Expand Up @@ -71,7 +71,11 @@ pub fn render(plan: &CompiledPhysicalPlan) -> String {
render_native(
&mut dot,
&format!("maintenance_{dag_index}_{}", sink.0),
"Planner maintenance operators",
if installed.reads_raw_samples(*sink) {
"Planner raw precompute operators"
} else {
"Planner maintenance operators"
},
program,
);
}
Expand Down
4 changes: 4 additions & 0 deletions control_plane/src/physical/workload_cost.rs
Original file line number Diff line number Diff line change
Expand Up @@ -215,6 +215,10 @@ pub fn manifest(
.iter()
.find(|m| m.policy_fingerprint() == stored_output.fingerprint())
.ok_or_else(|| invalid("native maintenance config is absent"))?;
// Raw ingest programs are priced by their stored state components.
if config.derived_input.is_none() {
continue;
}
add(
format!("maintenance:{}", stored_output.0),
json!({"physical_program":program,"stored_output":stored_output,"window_ms":config.stored_window_ms(),"interval_ms":config.slide_interval.saturating_mul(1000)}),
Expand Down
34 changes: 18 additions & 16 deletions control_plane/tests/native_rate_topk.rs
Original file line number Diff line number Diff line change
Expand Up @@ -55,12 +55,11 @@ fn rate_heap_candidates_bind_durable_counter_windows() {
let Ok(plan) = DeploymentPlanCompiler.compile_promql(candidate, environment.clone()) else {
continue;
};
if plan
.precompute_plan
.executable_dags
.values()
.any(|dag| !dag.native_programs.is_empty())
{
if plan.precompute_plan.executable_dags.values().any(|dag| {
dag.native_programs
.keys()
.any(|sink| !dag.reads_raw_samples(*sink))
}) {
continue;
}
let entry = plan.query_plan.entries.values().next().unwrap();
Expand Down Expand Up @@ -151,11 +150,11 @@ fn rate_heap_costs_can_select_each_compiled_candidate() {
.map(ToString::to_string)
.unwrap_or_default();
let preferred_plan = entry.physical_vector_binding().is_some()
&& !plan
.precompute_plan
.executable_dags
.values()
.any(|dag| !dag.native_programs.is_empty())
&& !plan.precompute_plan.executable_dags.values().any(|dag| {
dag.native_programs
.keys()
.any(|sink| !dag.reads_raw_samples(*sink))
})
&& if preferred == "exact" {
!program.contains("WithHeap")
} else {
Expand Down Expand Up @@ -218,6 +217,9 @@ fn fixed_window_rate_heap_candidates_install_both_physical_graphs() {
};
for installed in plan.precompute_plan.executable_dags.values() {
for sink in installed.native_programs.keys() {
if installed.reads_raw_samples(*sink) {
continue;
}
let program = installed.native_program(*sink).unwrap().unwrap();
let encoded = String::from_utf8(program.encode().unwrap()).unwrap();
for family in ["CmsWithHeap", "CountSketchWithHeap"] {
Expand Down Expand Up @@ -278,11 +280,11 @@ fn grouped_rate_placement_follows_summary_store_cost() {
continue;
}
entry.recover_vector_physical_dag().unwrap();
let stored = plan
.precompute_plan
.executable_dags
.values()
.any(|dag| !dag.native_programs.is_empty());
let stored = plan.precompute_plan.executable_dags.values().any(|dag| {
dag.native_programs
.keys()
.any(|sink| !dag.reads_raw_samples(*sink))
});
let rebuilt = program.to_string().contains("SummaryBuild");
assert!(!(stored && rebuilt));
// Planner also offers a relational Sum over the readouts, which has
Expand Down
50 changes: 48 additions & 2 deletions crates/asap_types/src/executable_plan.rs
Original file line number Diff line number Diff line change
Expand Up @@ -202,6 +202,20 @@ impl InstalledPostAsapDag {
Ok(())
}

/// Whether `sink`'s Planner program reads raw ingested samples rather
/// than stored states (a raw-ingest output, not a derived maintenance one).
pub fn reads_raw_samples(&self, sink: PostAsapNodeId) -> bool {
let raw = asap_physical_operators::physical_planner::precompute::raw_sample_schema();
self.native_program(sink)
.ok()
.flatten()
.is_some_and(|program| {
program
.input_contracts()
.all(|(_, contract)| contract.schema == raw)
})
}

/// Recovery validates the physical producer's typed storage boundaries;
/// it never lowers the semantic provenance document again.
pub fn native_program(
Expand Down Expand Up @@ -235,14 +249,39 @@ impl InstalledPostAsapDag {
}
}
let dag = self.document.decode()?;
let mut raw_inputs = 0;
for (id, contract) in program.input_contracts() {
let id = PostAsapNodeId(u32::try_from(id).map_err(|_| "physical source id overflow")?);
let node = dag.nodes.iter().find(|n| n.id == id);
// A raw sample boundary is fed by ingestion, not a stored state.
let raw = node.is_some_and(|n| {
asap_physical_operators::physical_planner::precompute::boundary_schema(n)
.is_ok_and(|schema| {
schema == contract.schema
&& schema
== asap_physical_operators::physical_planner::precompute::raw_sample_schema()
})
});
if raw {
if program.roots().contains(&u64::from(id.0))
|| !matches!(
self.binding.node(id),
Some(BackendNodeBinding::MaintenanceInput)
)
{
return Err(
"native raw sample source differs from its ingestion binding".into(),
);
}
raw_inputs += 1;
continue;
}
if program.roots().contains(&u64::from(id.0))
|| !matches!(
self.binding.node(id),
Some(BackendNodeBinding::Materialization { .. })
)
|| dag.nodes.iter().find(|n| n.id == id).is_none_or(|n| {
|| node.is_none_or(|n| {
n.output_schema != *contract.schema
&& asap_physical_operators::physical_planner::precompute::source_schema(
&n.output_schema,
Expand All @@ -258,14 +297,21 @@ impl InstalledPostAsapDag {
if program.input_contracts().count() == 0 {
return Err("native precompute program has no bound inputs".into());
}
// Raw summaries are stored as their family's population state; their
// semantic schema may also carry source columns that ingestion drops.
let raw_program = raw_inputs == program.input_contracts().count();
for root in program.roots() {
let output = program.output_contract(*root).map_err(|e| e.to_string())?;
if dag
.nodes
.iter()
.find(|n| u64::from(n.id.0) == *root)
.is_none_or(|n| {
n.output_schema != *output.schema
let raw_summary = raw_program
&& matches!(&n.payload, planner_types::post_asap::PostAsapOperatorPayload::SummaryAgg { family, .. }
if output.schema == asap_physical_operators::physical_planner::precompute::population_schema(family.clone()));
!raw_summary
&& n.output_schema != *output.schema
&& asap_physical_operators::physical_planner::precompute::source_schema(
&n.output_schema,
)
Expand Down
Loading
Loading