From deebfd5c2566ac33dc6dd80bd229f6e28ff12da3 Mon Sep 17 00:00:00 2001 From: zz_y Date: Sat, 29 Aug 2026 14:46:55 -0600 Subject: [PATCH 1/2] feat(export): include summary lifecycle decisions and rejections --- Cargo.lock | 2 + crates/asap-aware-mapping/Cargo.toml | 1 + crates/asap-aware-mapping/src/lib.rs | 5 + .../src/lifecycle_dag_export.rs | 207 ++++++++++++++++++ crates/integration-tests/Cargo.toml | 1 + .../tests/workload_lifecycle_e2e.rs | 20 +- 6 files changed, 235 insertions(+), 1 deletion(-) create mode 100644 crates/asap-aware-mapping/src/lifecycle_dag_export.rs diff --git a/Cargo.lock b/Cargo.lock index b7dfffce..34f106da 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -307,6 +307,7 @@ name = "asap-aware-mapping" version = "0.1.0" dependencies = [ "asap-types", + "serde", "serde_json", "thiserror", ] @@ -355,6 +356,7 @@ dependencies = [ "asap-frontend-promql", "asap-frontend-sql", "asap-types", + "serde_json", "tokio", ] diff --git a/crates/asap-aware-mapping/Cargo.toml b/crates/asap-aware-mapping/Cargo.toml index f7283834..352ed1cf 100644 --- a/crates/asap-aware-mapping/Cargo.toml +++ b/crates/asap-aware-mapping/Cargo.toml @@ -10,4 +10,5 @@ edition = "2021" [dependencies] asap-types = { path = "../types" } thiserror = "2" +serde = { version = "1", features = ["derive"] } serde_json = "1" diff --git a/crates/asap-aware-mapping/src/lib.rs b/crates/asap-aware-mapping/src/lib.rs index ed1570e4..3c0c0538 100644 --- a/crates/asap-aware-mapping/src/lib.rs +++ b/crates/asap-aware-mapping/src/lib.rs @@ -186,6 +186,7 @@ pub mod accuracy_reconciliation; pub mod cost_model; pub mod explanation; pub mod grouping; +pub mod lifecycle_dag_export; pub mod recurrence; pub mod replacement; pub mod rewrite; @@ -204,6 +205,10 @@ pub use explanation::{ explain_replacements, explain_replacements_with, ExplanationKind, ReplacementExplanation, }; pub use grouping::{has_subpopulations, HydraGroupingStrategy}; +pub use lifecycle_dag_export::{ + export_summary_maintenance_plan, SummaryMaintenanceDagExport, + SummaryMaintenanceDeploymentExport, SummaryMaintenanceLifecycleAlternativeExport, +}; pub use recurrence::{ evaluation_rate_of, total_cost, update_rate_from_data_workload, CostRate, EvaluationRate, Horizon, RecurrenceCostExplanation, RecurrenceError, RecurrenceProfile, RootRecurrence, diff --git a/crates/asap-aware-mapping/src/lifecycle_dag_export.rs b/crates/asap-aware-mapping/src/lifecycle_dag_export.rs new file mode 100644 index 00000000..a79b977a --- /dev/null +++ b/crates/asap-aware-mapping/src/lifecycle_dag_export.rs @@ -0,0 +1,207 @@ +//! Serializable DAG export for a materialized summary-maintenance plan. +//! +//! `asap-types::dag_export` owns the crate-neutral post-ASAP graph shape. This +//! adapter lives in the mapping layer, where lifecycle alternatives and their +//! typed rejection reasons are available, and emits both views together. + +use serde::Serialize; + +use asap_types::dag_export::{self, SummaryDagGraph}; +use asap_types::post_asap::{ + EvaluationSchedule, OutputRepresentation, SummaryMaintenanceLifecycle, + SummaryMaintenanceLifecycleGuarantee, +}; + +use crate::summary_maintenance_lifecycle::{ + SummaryMaintenanceLifecyclePlan, SummaryMaintenanceLifecycleRejection, +}; + +#[derive(Debug, Clone, Serialize)] +pub struct SummaryMaintenanceDagExport { + pub graph: SummaryDagGraph, + pub deployments: Vec, + pub horizon_seconds: Option, + pub evaluation_rate_per_second: Option, + pub update_rate_per_second: Option, + pub expected_reads: Option, + pub selected_raw_recompute: bool, + pub summary_total_cost: Option, + pub raw_recompute_total_cost: Option, +} + +#[derive(Debug, Clone, Serialize)] +pub struct SummaryMaintenanceDeploymentExport { + pub summary_index: usize, + #[serde(skip_serializing_if = "Option::is_none")] + pub selected: Option, + pub alternatives: Vec, +} + +#[derive(Debug, Clone, Serialize)] +pub struct SummaryMaintenanceLifecycleAlternativeExport { + pub lifecycle: SummaryMaintenanceLifecycleExport, + pub total_cost: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub rejection: Option, + pub assumptions: Vec, +} + +#[derive(Debug, Clone, Serialize)] +pub struct SummaryMaintenanceLifecycleGuaranteeExport { + pub lifecycle: SummaryMaintenanceLifecycleExport, + pub evaluation_schedule: EvaluationScheduleExport, + pub output_representation: OutputRepresentationExport, +} + +#[derive(Debug, Clone, Serialize)] +#[serde(tag = "kind", rename_all = "snake_case")] +pub enum SummaryMaintenanceLifecycleExport { + Ephemeral, + Prepared { + activate_at_ms: u64, + retire_at_ms: u64, + }, + Shared { + retention_ms: u64, + }, + ContinuouslyMaintained, +} + +#[derive(Debug, Clone, Copy, Serialize)] +#[serde(rename_all = "snake_case")] +pub enum EvaluationScheduleExport { + OneShot, + PerUpdate, + OnRead, +} + +#[derive(Debug, Clone, Copy, Serialize)] +#[serde(rename_all = "snake_case")] +pub enum OutputRepresentationExport { + PlainRows, + SummaryState, + FinalizedValue, +} + +#[derive(Debug, Clone, Copy, Serialize)] +#[serde(rename_all = "snake_case")] +pub enum SummaryMaintenanceLifecycleRejectionExport { + UnsupportedByRuntime, + RequiresPredictableOneTimeQuery, + RequiresMultipleReads, + RequiresHorizon, + RequiresContinuousData, + MissingOrStaleIngestionRate, + SummaryDoesNotSupportIncrementalUpdates, + SummaryDoesNotSupportDeletion, + MissingCostEvidence, +} + +pub fn export_summary_maintenance_plan( + plan: &SummaryMaintenanceLifecyclePlan, +) -> SummaryMaintenanceDagExport { + SummaryMaintenanceDagExport { + graph: dag_export::export_summary(&plan.root), + deployments: plan + .deployments + .iter() + .map(|deployment| SummaryMaintenanceDeploymentExport { + summary_index: deployment.summary_index, + selected: deployment + .summary_maintenance_lifecycle_guarantee + .as_ref() + .map(export_guarantee), + alternatives: deployment + .alternatives + .iter() + .map(|alternative| SummaryMaintenanceLifecycleAlternativeExport { + lifecycle: export_lifecycle(&alternative.summary_maintenance_lifecycle), + total_cost: alternative.total_cost.map(|cost| cost.0), + rejection: alternative.rejection.as_ref().map(export_rejection), + assumptions: alternative.assumptions.clone(), + }) + .collect(), + }) + .collect(), + horizon_seconds: plan.horizon.map(|horizon| horizon.0), + evaluation_rate_per_second: plan.evaluation_rate.map(|rate| rate.0), + update_rate_per_second: plan.update_rate.map(|rate| rate.0), + expected_reads: plan.expected_reads, + selected_raw_recompute: plan.selected_raw_recompute, + summary_total_cost: plan.summary_total_cost.map(|cost| cost.0), + raw_recompute_total_cost: plan.raw_recompute_total_cost.map(|cost| cost.0), + } +} + +fn export_guarantee( + guarantee: &SummaryMaintenanceLifecycleGuarantee, +) -> SummaryMaintenanceLifecycleGuaranteeExport { + SummaryMaintenanceLifecycleGuaranteeExport { + lifecycle: export_lifecycle(&guarantee.summary_maintenance_lifecycle), + evaluation_schedule: match guarantee.evaluation_schedule { + EvaluationSchedule::OneShot => EvaluationScheduleExport::OneShot, + EvaluationSchedule::PerUpdate => EvaluationScheduleExport::PerUpdate, + EvaluationSchedule::OnRead => EvaluationScheduleExport::OnRead, + }, + output_representation: match guarantee.output_representation { + OutputRepresentation::PlainRows => OutputRepresentationExport::PlainRows, + OutputRepresentation::SummaryState => OutputRepresentationExport::SummaryState, + OutputRepresentation::FinalizedValue => OutputRepresentationExport::FinalizedValue, + }, + } +} + +fn export_lifecycle(lifecycle: &SummaryMaintenanceLifecycle) -> SummaryMaintenanceLifecycleExport { + match lifecycle { + SummaryMaintenanceLifecycle::Ephemeral => SummaryMaintenanceLifecycleExport::Ephemeral, + SummaryMaintenanceLifecycle::Prepared { + activate_at, + retire_at, + } => SummaryMaintenanceLifecycleExport::Prepared { + activate_at_ms: activate_at.0, + retire_at_ms: retire_at.0, + }, + SummaryMaintenanceLifecycle::Shared { retention } => { + SummaryMaintenanceLifecycleExport::Shared { + retention_ms: retention.0, + } + } + SummaryMaintenanceLifecycle::ContinuouslyMaintained => { + SummaryMaintenanceLifecycleExport::ContinuouslyMaintained + } + } +} + +fn export_rejection( + rejection: &SummaryMaintenanceLifecycleRejection, +) -> SummaryMaintenanceLifecycleRejectionExport { + match rejection { + SummaryMaintenanceLifecycleRejection::UnsupportedByRuntime => { + SummaryMaintenanceLifecycleRejectionExport::UnsupportedByRuntime + } + SummaryMaintenanceLifecycleRejection::RequiresPredictableOneTimeQuery => { + SummaryMaintenanceLifecycleRejectionExport::RequiresPredictableOneTimeQuery + } + SummaryMaintenanceLifecycleRejection::RequiresMultipleReads => { + SummaryMaintenanceLifecycleRejectionExport::RequiresMultipleReads + } + SummaryMaintenanceLifecycleRejection::RequiresHorizon => { + SummaryMaintenanceLifecycleRejectionExport::RequiresHorizon + } + SummaryMaintenanceLifecycleRejection::RequiresContinuousData => { + SummaryMaintenanceLifecycleRejectionExport::RequiresContinuousData + } + SummaryMaintenanceLifecycleRejection::MissingOrStaleIngestionRate => { + SummaryMaintenanceLifecycleRejectionExport::MissingOrStaleIngestionRate + } + SummaryMaintenanceLifecycleRejection::SummaryDoesNotSupportIncrementalUpdates => { + SummaryMaintenanceLifecycleRejectionExport::SummaryDoesNotSupportIncrementalUpdates + } + SummaryMaintenanceLifecycleRejection::SummaryDoesNotSupportDeletion => { + SummaryMaintenanceLifecycleRejectionExport::SummaryDoesNotSupportDeletion + } + SummaryMaintenanceLifecycleRejection::MissingCostEvidence => { + SummaryMaintenanceLifecycleRejectionExport::MissingCostEvidence + } + } +} diff --git a/crates/integration-tests/Cargo.toml b/crates/integration-tests/Cargo.toml index 11873eff..878d8236 100644 --- a/crates/integration-tests/Cargo.toml +++ b/crates/integration-tests/Cargo.toml @@ -10,4 +10,5 @@ asap-frontend-sql = { path = "../frontend-sql" } asap-aware-mapping = { path = "../asap-aware-mapping" } [dev-dependencies] +serde_json = "1" tokio = { version = "1", features = ["rt", "macros", "rt-multi-thread"] } diff --git a/crates/integration-tests/tests/workload_lifecycle_e2e.rs b/crates/integration-tests/tests/workload_lifecycle_e2e.rs index 66e7a414..5d23c433 100644 --- a/crates/integration-tests/tests/workload_lifecycle_e2e.rs +++ b/crates/integration-tests/tests/workload_lifecycle_e2e.rs @@ -7,7 +7,7 @@ use std::rc::Rc; use asap_aware_mapping::cost_model::Cost; use asap_aware_mapping::CostRate; use asap_aware_mapping::{ - global_selection_with_summary_maintenance_lifecycles, + export_summary_maintenance_plan, global_selection_with_summary_maintenance_lifecycles, materialize_with_summary_maintenance_lifecycles, search_workload_with, CostModel, Horizon, SummaryMaintenanceCapabilities, SummaryMaintenanceLifecycleCapabilities, SummaryMaintenanceLifecycleCostInputs, SummaryMaintenanceLifecycleRejection, WorkloadDemand, @@ -176,4 +176,22 @@ fn promql_dashboard_materializes_continuous_summary_with_explained_rejections() ) && alternative.rejection == Some(SummaryMaintenanceLifecycleRejection::UnsupportedByRuntime) })); + + let exported = serde_json::to_value(export_summary_maintenance_plan(&plan)).unwrap(); + assert_eq!( + exported["deployments"][0]["selected"]["lifecycle"]["kind"], + "continuously_maintained" + ); + let alternatives = exported["deployments"][0]["alternatives"] + .as_array() + .expect("exported lifecycle alternatives"); + assert!(alternatives.iter().any(|alternative| { + alternative["lifecycle"]["kind"] == "prepared" + && alternative["rejection"] == "requires_predictable_one_time_query" + })); + assert!(alternatives.iter().any(|alternative| { + alternative["lifecycle"]["kind"] == "shared" + && alternative["rejection"] == "unsupported_by_runtime" + })); + assert!(exported["graph"]["nodes"].as_array().is_some()); } From 505c10414a9672f93379c4ac5575279305736089 Mon Sep 17 00:00:00 2001 From: zz_y Date: Sun, 30 Aug 2026 15:21:20 -0600 Subject: [PATCH 2/2] feat(export): include selected summary maintenance mode --- .../asap-aware-mapping/src/lifecycle_dag_export.rs | 14 +++++++++++++- .../tests/workload_lifecycle_e2e.rs | 4 ++++ 2 files changed, 17 insertions(+), 1 deletion(-) diff --git a/crates/asap-aware-mapping/src/lifecycle_dag_export.rs b/crates/asap-aware-mapping/src/lifecycle_dag_export.rs index a79b977a..420ace85 100644 --- a/crates/asap-aware-mapping/src/lifecycle_dag_export.rs +++ b/crates/asap-aware-mapping/src/lifecycle_dag_export.rs @@ -9,7 +9,7 @@ use serde::Serialize; use asap_types::dag_export::{self, SummaryDagGraph}; use asap_types::post_asap::{ EvaluationSchedule, OutputRepresentation, SummaryMaintenanceLifecycle, - SummaryMaintenanceLifecycleGuarantee, + SummaryMaintenanceLifecycleGuarantee, SummaryMaintenanceMode, }; use crate::summary_maintenance_lifecycle::{ @@ -49,10 +49,18 @@ pub struct SummaryMaintenanceLifecycleAlternativeExport { #[derive(Debug, Clone, Serialize)] pub struct SummaryMaintenanceLifecycleGuaranteeExport { pub lifecycle: SummaryMaintenanceLifecycleExport, + pub maintenance_mode: SummaryMaintenanceModeExport, pub evaluation_schedule: EvaluationScheduleExport, pub output_representation: OutputRepresentationExport, } +#[derive(Debug, Clone, Copy, Serialize)] +#[serde(rename_all = "snake_case")] +pub enum SummaryMaintenanceModeExport { + DirectBuild, + Incremental, +} + #[derive(Debug, Clone, Serialize)] #[serde(tag = "kind", rename_all = "snake_case")] pub enum SummaryMaintenanceLifecycleExport { @@ -138,6 +146,10 @@ fn export_guarantee( ) -> SummaryMaintenanceLifecycleGuaranteeExport { SummaryMaintenanceLifecycleGuaranteeExport { lifecycle: export_lifecycle(&guarantee.summary_maintenance_lifecycle), + maintenance_mode: match guarantee.summary_maintenance_mode { + SummaryMaintenanceMode::DirectBuild => SummaryMaintenanceModeExport::DirectBuild, + SummaryMaintenanceMode::Incremental => SummaryMaintenanceModeExport::Incremental, + }, evaluation_schedule: match guarantee.evaluation_schedule { EvaluationSchedule::OneShot => EvaluationScheduleExport::OneShot, EvaluationSchedule::PerUpdate => EvaluationScheduleExport::PerUpdate, diff --git a/crates/integration-tests/tests/workload_lifecycle_e2e.rs b/crates/integration-tests/tests/workload_lifecycle_e2e.rs index 5d23c433..7fa5eb2b 100644 --- a/crates/integration-tests/tests/workload_lifecycle_e2e.rs +++ b/crates/integration-tests/tests/workload_lifecycle_e2e.rs @@ -182,6 +182,10 @@ fn promql_dashboard_materializes_continuous_summary_with_explained_rejections() exported["deployments"][0]["selected"]["lifecycle"]["kind"], "continuously_maintained" ); + assert_eq!( + exported["deployments"][0]["selected"]["maintenance_mode"], + "incremental" + ); let alternatives = exported["deployments"][0]["alternatives"] .as_array() .expect("exported lifecycle alternatives");