From 4a666cb9506b1aa2e3682433406fff423790db2a Mon Sep 17 00:00:00 2001 From: zz_y Date: Thu, 10 Sep 2026 07:46:07 -0600 Subject: [PATCH] refactor: reuse lifecycle types in DAG exports --- .../src/summary_maintenance_dag_export.rs | 151 +----------------- .../src/summary_maintenance_lifecycle.rs | 3 +- .../src/post_asap/summary_maintenance.rs | 3 +- .../summary_maintenance_lifecycle.rs | 17 +- 4 files changed, 25 insertions(+), 149 deletions(-) diff --git a/crates/asap-aware-mapping/src/summary_maintenance_dag_export.rs b/crates/asap-aware-mapping/src/summary_maintenance_dag_export.rs index 3d3a5b8a..266fb02f 100644 --- a/crates/asap-aware-mapping/src/summary_maintenance_dag_export.rs +++ b/crates/asap-aware-mapping/src/summary_maintenance_dag_export.rs @@ -12,9 +12,8 @@ use serde::Serialize; use asap_types::dag_export::{self, SummaryDagGraph}; use asap_types::post_asap::{ - EvaluationSchedule, OutputRepresentation, ResultGuarantee, SummaryExpr, - SummaryMaintenanceLifecycle, SummaryMaintenanceLifecycleGuarantee, SummaryMaintenanceMode, - SummaryNode, SummaryWindowFramework, + ResultGuarantee, SummaryExpr, SummaryMaintenanceLifecycle, + SummaryMaintenanceLifecycleGuarantee, SummaryNode, SummaryWindowFramework, }; use crate::summary_maintenance_lifecycle::{ @@ -50,71 +49,14 @@ pub struct SummaryMaintenanceDeploymentExport { #[derive(Debug, Clone, Serialize)] pub struct SummaryMaintenanceLifecycleAlternativeExport { - pub lifecycle: SummaryMaintenanceLifecycleExport, + pub lifecycle: SummaryMaintenanceLifecycle, pub total_cost: Option, #[serde(skip_serializing_if = "Option::is_none")] - pub rejection: Option, + pub rejection: Option, pub assumptions: Vec, } -#[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 { - 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 type SummaryMaintenanceLifecycleGuaranteeExport = SummaryMaintenanceLifecycleGuarantee; pub fn export_summary_maintenance_plan( plan: &SummaryMaintenanceLifecyclePlan, @@ -128,14 +70,14 @@ pub fn export_summary_maintenance_plan( selected: deployment .summary_maintenance_lifecycle_guarantee .as_ref() - .map(export_guarantee), + .cloned(), alternatives: deployment .alternatives .iter() .map(|alternative| SummaryMaintenanceLifecycleAlternativeExport { - lifecycle: export_lifecycle(&alternative.summary_maintenance_lifecycle), + lifecycle: alternative.summary_maintenance_lifecycle.clone(), total_cost: alternative.total_cost.map(|cost| cost.0), - rejection: alternative.rejection.as_ref().map(export_rejection), + rejection: alternative.rejection.clone(), assumptions: alternative.assumptions.clone(), }) .collect(), @@ -215,80 +157,3 @@ fn summary_children(expr: &SummaryExpr) -> Vec<&Rc> { SummaryExpr::SummaryMerge { children } => children.iter().collect(), } } - -fn export_guarantee( - guarantee: &SummaryMaintenanceLifecycleGuarantee, -) -> 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, - 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/asap-aware-mapping/src/summary_maintenance_lifecycle.rs b/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs index cc909c59..91f21eda 100644 --- a/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs +++ b/crates/asap-aware-mapping/src/summary_maintenance_lifecycle.rs @@ -119,7 +119,8 @@ pub struct SummaryMaintenanceLifecycleCostInputs { pub retirement_cost: Option, } -#[derive(Debug, Clone, PartialEq, Eq)] +#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)] +#[serde(rename_all = "snake_case")] pub enum SummaryMaintenanceLifecycleRejection { UnsupportedByRuntime, RequiresPredictableOneTimeQuery, diff --git a/crates/types/src/post_asap/summary_maintenance.rs b/crates/types/src/post_asap/summary_maintenance.rs index 7526b3fe..d1e50d7e 100644 --- a/crates/types/src/post_asap/summary_maintenance.rs +++ b/crates/types/src/post_asap/summary_maintenance.rs @@ -7,7 +7,8 @@ //! compilation chooses its concrete implementation. /// How a summary deployment obtains its state, independent of implementation. -#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)] +#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, serde::Serialize, serde::Deserialize)] +#[serde(rename_all = "snake_case")] pub enum SummaryMaintenanceMode { /// Rebuild the summary from its complete input when the deployment needs /// a value. No update stream is required. diff --git a/crates/types/src/post_asap/summary_maintenance_lifecycle.rs b/crates/types/src/post_asap/summary_maintenance_lifecycle.rs index a2eb09c0..798d861f 100644 --- a/crates/types/src/post_asap/summary_maintenance_lifecycle.rs +++ b/crates/types/src/post_asap/summary_maintenance_lifecycle.rs @@ -12,7 +12,8 @@ use crate::workload::{DurationMs, TimestampMs}; /// When an operator is evaluated. This is independent of whether it owns /// state and how long that state is retained. -#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)] +#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, serde::Serialize, serde::Deserialize)] +#[serde(rename_all = "snake_case")] pub enum EvaluationSchedule { OneShot, PerUpdate, @@ -24,7 +25,8 @@ pub enum EvaluationSchedule { /// summary's output; for example, `Estimate` is the consumer in /// `SummaryAgg -> Estimate`. The exposed result is ordinary rows, reusable /// summary state, or a finalized value. -#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)] +#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, serde::Serialize, serde::Deserialize)] +#[serde(rename_all = "snake_case")] pub enum OutputRepresentation { PlainRows, SummaryState, @@ -38,14 +40,18 @@ pub enum OutputRepresentation { /// provides the expected number and timing of reads; data arrival provides the /// expected state-update demand. The planner combines those quantities with /// costs and runtime capabilities to compare these policies. -#[derive(Debug, Clone, PartialEq, Eq, Hash)] +#[derive(Debug, Clone, PartialEq, Eq, Hash, serde::Serialize, serde::Deserialize)] +#[serde(tag = "kind", rename_all = "snake_case")] pub enum SummaryMaintenanceLifecycle { Ephemeral, Prepared { + #[serde(rename = "activate_at_ms")] activate_at: TimestampMs, + #[serde(rename = "retire_at_ms")] retire_at: TimestampMs, }, Shared { + #[serde(rename = "retention_ms")] retention: DurationMs, }, ContinuouslyMaintained, @@ -56,9 +62,12 @@ pub enum SummaryMaintenanceLifecycle { /// This names the summary-maintenance promise explicitly so consumers do not /// confuse it with guarantees about the broader data lifecycle. Accuracy is a /// separate [`super::ResultGuarantee`]. -#[derive(Debug, Clone, PartialEq, Eq, Hash)] +#[derive(Debug, Clone, PartialEq, Eq, Hash, serde::Serialize, serde::Deserialize)] +#[serde(deny_unknown_fields)] pub struct SummaryMaintenanceLifecycleGuarantee { + #[serde(rename = "lifecycle")] pub summary_maintenance_lifecycle: SummaryMaintenanceLifecycle, + #[serde(rename = "maintenance_mode")] pub summary_maintenance_mode: SummaryMaintenanceMode, pub evaluation_schedule: EvaluationSchedule, pub output_representation: OutputRepresentation,