Skip to content
Merged
388 changes: 236 additions & 152 deletions crates/asap-aware-mapping/src/analytical_cost.rs

Large diffs are not rendered by default.

913 changes: 913 additions & 0 deletions crates/asap-aware-mapping/src/empirical_comparison.rs

Large diffs are not rendered by default.

774 changes: 774 additions & 0 deletions crates/asap-aware-mapping/src/empirical_cost.rs

Large diffs are not rendered by default.

202 changes: 202 additions & 0 deletions crates/asap-aware-mapping/src/empirical_resources.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,202 @@
//! Measured resource payloads and compatibility with the v1 benchmark wire format.
//!
//! Physical dimensions live in `asap_types::resources`; flat wire structs below
//! exist only to keep archived artifacts readable and preserve their field names.

use asap_types::resources::PhysicalResources;
use serde::{Deserialize, Serialize};

pub use asap_types::resources::{MeasuredCpu, MeasuredResources, Measurement};

#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
#[serde(from = "SketchWire", into = "SketchWire")]
pub struct ResourceMeasurements {
pub resources: MeasuredResources,
}

#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
#[serde(from = "ExactWire", into = "ExactWire")]
pub struct ExactResourceMeasurements {
pub resources: MeasuredResources,
}

// Keep the archived flat v1 schema at the serialization boundary only. New
// optional dimensions are omitted when absent, so legacy snapshots round-trip.
#[derive(Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
struct SketchWire {
build_cpu_ns: Option<Measurement>,
update_cpu_ns: Option<Measurement>,
merge_cpu_ns: Option<Measurement>,
read_cpu_ns: Option<Measurement>,
#[serde(default, skip_serializing_if = "Option::is_none")]
prepare_cpu_ns: Option<Measurement>,
retained_bytes: Option<Measurement>,
peak_bytes: Option<Measurement>,
serialized_bytes: Option<Measurement>,
disk_bytes: Option<Measurement>,
#[serde(default, skip_serializing_if = "Option::is_none")]
scan_bytes: Option<Measurement>,
}

#[derive(Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
struct ExactWire {
empty_build_cpu_ns: Option<Measurement>,
update_cpu_ns: Option<Measurement>,
prepare_cpu_ns: Option<Measurement>,
read_cpu_ns: Option<Measurement>,
#[serde(default, skip_serializing_if = "Option::is_none")]
merge_cpu_ns: Option<Measurement>,
retained_bytes: Option<Measurement>,
peak_bytes: Option<Measurement>,
#[serde(default, skip_serializing_if = "Option::is_none")]
serialized_bytes: Option<Measurement>,
#[serde(default, skip_serializing_if = "Option::is_none")]
disk_bytes: Option<Measurement>,
#[serde(default, skip_serializing_if = "Option::is_none")]
scan_bytes: Option<Measurement>,
}

impl From<SketchWire> for ResourceMeasurements {
fn from(w: SketchWire) -> Self {
Self {
resources: PhysicalResources {
cpu: MeasuredCpu {
build_cpu_ns: w.build_cpu_ns,
update_cpu_ns: w.update_cpu_ns,
merge_cpu_ns: w.merge_cpu_ns,
prepare_cpu_ns: w.prepare_cpu_ns,
read_cpu_ns: w.read_cpu_ns,
},
retained_memory_bytes: w.retained_bytes,
peak_memory_bytes: w.peak_bytes,
serialized_bytes: w.serialized_bytes,
disk_bytes: w.disk_bytes,
scan_bytes: w.scan_bytes,
},
}
}
}

impl From<ResourceMeasurements> for SketchWire {
fn from(value: ResourceMeasurements) -> Self {
let r = value.resources;
Self {
build_cpu_ns: r.cpu.build_cpu_ns,
update_cpu_ns: r.cpu.update_cpu_ns,
merge_cpu_ns: r.cpu.merge_cpu_ns,
prepare_cpu_ns: r.cpu.prepare_cpu_ns,
read_cpu_ns: r.cpu.read_cpu_ns,
retained_bytes: r.retained_memory_bytes,
peak_bytes: r.peak_memory_bytes,
serialized_bytes: r.serialized_bytes,
disk_bytes: r.disk_bytes,
scan_bytes: r.scan_bytes,
}
}
}

impl From<ExactWire> for ExactResourceMeasurements {
fn from(w: ExactWire) -> Self {
Self {
resources: PhysicalResources {
cpu: MeasuredCpu {
build_cpu_ns: w.empty_build_cpu_ns,
update_cpu_ns: w.update_cpu_ns,
merge_cpu_ns: w.merge_cpu_ns,
prepare_cpu_ns: w.prepare_cpu_ns,
read_cpu_ns: w.read_cpu_ns,
},
retained_memory_bytes: w.retained_bytes,
peak_memory_bytes: w.peak_bytes,
serialized_bytes: w.serialized_bytes,
disk_bytes: w.disk_bytes,
scan_bytes: w.scan_bytes,
},
}
}
}

impl From<ExactResourceMeasurements> for ExactWire {
fn from(value: ExactResourceMeasurements) -> Self {
let r = value.resources;
Self {
empty_build_cpu_ns: r.cpu.build_cpu_ns,
update_cpu_ns: r.cpu.update_cpu_ns,
merge_cpu_ns: r.cpu.merge_cpu_ns,
prepare_cpu_ns: r.cpu.prepare_cpu_ns,
read_cpu_ns: r.cpu.read_cpu_ns,
retained_bytes: r.retained_memory_bytes,
peak_bytes: r.peak_memory_bytes,
serialized_bytes: r.serialized_bytes,
disk_bytes: r.disk_bytes,
scan_bytes: r.scan_bytes,
}
}
}

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

fn measured(value: f64) -> Option<Measurement> {
Some(Measurement {
value,
stddev: Some(0.5),
samples: 5,
method: Some("test only".into()),
})
}

/// The same resource dimensions retain uncertainty through either wire adapter.
#[test]
fn shared_resources_round_trip_without_losing_dimensions() {
let resources = PhysicalResources {
cpu: MeasuredCpu {
build_cpu_ns: measured(1.0),
update_cpu_ns: measured(2.0),
merge_cpu_ns: measured(3.0),
prepare_cpu_ns: measured(4.0),
read_cpu_ns: measured(5.0),
},
peak_memory_bytes: measured(100.0),
retained_memory_bytes: measured(70.0),
scan_bytes: measured(200.0),
serialized_bytes: measured(40.0),
disk_bytes: measured(4096.0),
};
let sketch = ResourceMeasurements {
resources: resources.clone(),
};
let exact = ExactResourceMeasurements { resources };
assert_eq!(
serde_json::from_value::<ResourceMeasurements>(serde_json::to_value(&sketch).unwrap())
.unwrap(),
sketch
);
assert_eq!(
serde_json::from_value::<ExactResourceMeasurements>(
serde_json::to_value(&exact).unwrap()
)
.unwrap(),
exact
);
}

/// Archived names remain accepted, but physical fields use canonical names.
#[test]
fn legacy_fields_map_to_shared_resources() {
let value = serde_json::json!({"empty_build_cpu_ns": measured(3.0), "peak_bytes": measured(100.0), "retained_bytes": measured(70.0)});
let exact: ExactResourceMeasurements = serde_json::from_value(value).unwrap();
assert_eq!(exact.resources.cpu.build_cpu_ns, measured(3.0));
assert_eq!(exact.resources.peak_memory_bytes, measured(100.0));
assert_eq!(exact.resources.retained_memory_bytes, measured(70.0));
assert!(exact.resources.scan_bytes.is_none());
assert!(exact.resources.disk_bytes.is_none());
assert!(
serde_json::from_value::<ResourceMeasurements>(serde_json::json!({"cpu_ops": 42}))
.is_err()
);
}
}
3 changes: 3 additions & 0 deletions crates/asap-aware-mapping/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -185,6 +185,9 @@ pub mod accuracy;
pub mod accuracy_reconciliation;
pub mod analytical_cost;
pub mod cost_model;
pub mod empirical_comparison;
pub mod empirical_cost;
pub mod empirical_resources;
pub mod explanation;
pub mod grouping;
pub mod physical_operator_statistics;
Expand Down
10 changes: 5 additions & 5 deletions crates/asap-aware-mapping/src/query_physical_lowering.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1632,7 +1632,7 @@ mod tests {
let estimate =
estimate_physical_dag(&dag.nodes, &dag.root, &independent_scope, &dag.evidence)
.unwrap();
assert_eq!(estimate.scan_bytes, 1_600);
assert_eq!(estimate.scan_bytes(), 1_600);

let shared_provider = |request: PhysicalNodeRequest<'_>| {
let (physical_id, statistics) = match request.operator {
Expand Down Expand Up @@ -1668,8 +1668,8 @@ mod tests {
},
)
.unwrap();
assert_eq!(comparison.raw.scan_bytes, 1_600);
assert_eq!(comparison.candidate.scan_bytes, 800);
assert_eq!(comparison.raw.scan_bytes(), 1_600);
assert_eq!(comparison.candidate.scan_bytes(), 800);

let mut drifted_buffer = shared_dag.clone();
drifted_buffer.nodes[0].output_buffer_bytes += 1;
Expand Down Expand Up @@ -2167,8 +2167,8 @@ mod tests {
]);
let dag = lower_query_physical_dag(&root, &scope, &scripted(&provided)).unwrap();
let estimate = estimate_physical_dag(&dag.nodes, &dag.root, &scope, &dag.evidence).unwrap();
assert_eq!(estimate.cpu_ops, 200.0);
assert_eq!(estimate.scan_bytes, 800);
assert_eq!(estimate.cpu_ops(), 200.0);
assert_eq!(estimate.scan_bytes(), 800);
}

#[test]
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -545,18 +545,18 @@ pub(super) fn estimate_heterogeneous_summary(
if !cpu_ops.is_finite() {
return Err(AnalyticalCostError::Overflow);
}
Ok(ResourceEstimate {
Ok(ResourceEstimate::new(
cpu_ops,
peak_memory_bytes: persistent_bytes
persistent_bytes
.checked_add(transient_bytes)
.and_then(|bytes| bytes.checked_add(ephemeral_state_bytes))
.ok_or(AnalyticalCostError::Overflow)?,
scan_bytes: scans
scans
.values()
.try_fold(operator_io_bytes, |sum, (_, bytes)| {
sum.checked_add(*bytes).ok_or(AnalyticalCostError::Overflow)
})?,
})
))
}

fn add_operator_io(
Expand Down Expand Up @@ -1070,15 +1070,15 @@ pub(super) fn estimate_incremental_summary_maintenance_with_join(
.initial_input_bytes
.div_ceil(inputs.initial_input_rows)
};
Ok(ResourceEstimate {
Ok(ResourceEstimate::new(
cpu_ops,
peak_memory_bytes: retained_bytes
retained_bytes
.checked_add(transient_bytes)
.and_then(|bytes| bytes.checked_add(join_bytes))
.ok_or(AnalyticalCostError::Overflow)?
.max(bootstrap_row_buffer),
scan_bytes: inputs.initial_source_scan_bytes,
})
inputs.initial_source_scan_bytes,
))
}

pub(super) fn lifecycle_row_counts(
Expand Down
50 changes: 22 additions & 28 deletions crates/asap-aware-mapping/src/summary_maintenance_cost/model.rs
Original file line number Diff line number Diff line change
Expand Up @@ -637,28 +637,22 @@ impl SummaryMaintenanceCostModel {
let evidence = self.canonical_inputs(summary)?;
let inputs = evidence.inputs;
let insert = validated_operator_cpu("insert_cpu_ops", evidence.insert_cpu_ops).ok()?;
let build = self.calibrated(ResourceEstimate {
cpu_ops: inputs.initial_input_rows as f64
* inputs.bootstrap_window_count as f64
* insert,
peak_memory_bytes: 0,
scan_bytes: inputs.initial_source_scan_bytes,
})?;
let maintenance = self.calibrated(ResourceEstimate {
cpu_ops: inputs.active_window_count as f64 * insert,
peak_memory_bytes: 0,
scan_bytes: 0,
})?;
let build = self.calibrated(ResourceEstimate::new(
inputs.initial_input_rows as f64 * inputs.bootstrap_window_count as f64 * insert,
0,
inputs.initial_source_scan_bytes,
))?;
let maintenance = self.calibrated(ResourceEstimate::new(
inputs.active_window_count as f64 * insert,
0,
0,
))?;
let retained = inputs
.active_window_count
.checked_add(inputs.retained_window_count)?
.checked_mul(inputs.physical_summary_count)?
.checked_mul(inputs.state_bytes_per_summary)?;
let retention_total = self.calibrated(ResourceEstimate {
cpu_ops: 0.0,
peak_memory_bytes: retained,
scan_bytes: 0,
})?;
let retention_total = self.calibrated(ResourceEstimate::new(0.0, retained, 0))?;
let horizon_seconds = horizon.filter(|value| value.0 > 0.0)?.0;
Some(SummaryMaintenanceLifecycleCostInputs {
build_cost: Some(build),
Expand Down Expand Up @@ -1059,8 +1053,8 @@ mod tests {
)
.unwrap();
// 10 arrivals * 2 active windows * 2 insert ops + 5 reads * 2 summaries.
assert_eq!(estimate.cpu_ops, 50.0);
assert_eq!(estimate.scan_bytes, 0);
assert_eq!(estimate.cpu_ops(), 50.0);
assert_eq!(estimate.scan_bytes(), 0);
}

#[test]
Expand Down Expand Up @@ -1112,7 +1106,7 @@ mod tests {
)
.unwrap();
// 10 bootstrap rows * 3 windows * 2 insert ops + 5 reads * 2 summaries.
assert_eq!(estimate.cpu_ops, 70.0);
assert_eq!(estimate.cpu_ops(), 70.0);
}

#[test]
Expand Down Expand Up @@ -2604,9 +2598,9 @@ mod tests {
)
.unwrap();
// 10 bootstrap + 10 arrivals into two active windows; two states read 5 times.
assert_eq!(estimate.cpu_ops, 90.0);
assert_eq!(estimate.peak_memory_bytes, 1_000);
assert_eq!(estimate.scan_bytes, 640);
assert_eq!(estimate.cpu_ops(), 90.0);
assert_eq!(estimate.peak_memory_bytes(), 1_000);
assert_eq!(estimate.scan_bytes(), 640);
}

#[test]
Expand Down Expand Up @@ -2636,9 +2630,9 @@ mod tests {
},
)
.unwrap();
assert_eq!(estimate.cpu_ops, 21.0 + 20.0 + 30.0 + 200.0 + 70.0);
assert_eq!(estimate.cpu_ops(), 21.0 + 20.0 + 30.0 + 200.0 + 70.0);
// Three persistent windows plus one transient result, for two instances.
assert_eq!(estimate.peak_memory_bytes, 80);
assert_eq!(estimate.peak_memory_bytes(), 80);
}

#[test]
Expand Down Expand Up @@ -2762,7 +2756,7 @@ mod tests {
.unwrap();
// Two pre-activation arrivals join the bootstrap; eight more are
// maintained through the horizon; five reads are served.
assert_eq!(estimate.cpu_ops, 25.0);
assert_eq!(estimate.cpu_ops(), 25.0);
}

#[test]
Expand Down Expand Up @@ -2859,8 +2853,8 @@ mod tests {
}),
)
.unwrap();
assert_eq!(estimate.cpu_ops, 77.0);
assert_eq!(estimate.peak_memory_bytes, 64); // 4 persistent states + join memory.
assert_eq!(estimate.cpu_ops(), 77.0);
assert_eq!(estimate.peak_memory_bytes(), 64); // 4 persistent states + join memory.
}

fn summary_with_operations(merge: bool, subtract: bool, delete: bool) -> Rc<SummaryNode> {
Expand Down
Loading
Loading