Skip to content
Merged
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
151 changes: 8 additions & 143 deletions crates/asap-aware-mapping/src/summary_maintenance_dag_export.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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::{
Expand Down Expand Up @@ -50,71 +49,14 @@ pub struct SummaryMaintenanceDeploymentExport {

#[derive(Debug, Clone, Serialize)]
pub struct SummaryMaintenanceLifecycleAlternativeExport {
pub lifecycle: SummaryMaintenanceLifecycleExport,
pub lifecycle: SummaryMaintenanceLifecycle,
pub total_cost: Option<f64>,
#[serde(skip_serializing_if = "Option::is_none")]
pub rejection: Option<SummaryMaintenanceLifecycleRejectionExport>,
pub rejection: Option<SummaryMaintenanceLifecycleRejection>,
pub assumptions: Vec<String>,
}

#[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,
Expand All @@ -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(),
Expand Down Expand Up @@ -215,80 +157,3 @@ fn summary_children(expr: &SummaryExpr) -> Vec<&Rc<SummaryNode>> {
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
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -119,7 +119,8 @@ pub struct SummaryMaintenanceLifecycleCostInputs {
pub retirement_cost: Option<Cost>,
}

#[derive(Debug, Clone, PartialEq, Eq)]
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
#[serde(rename_all = "snake_case")]
pub enum SummaryMaintenanceLifecycleRejection {
UnsupportedByRuntime,
RequiresPredictableOneTimeQuery,
Expand Down
3 changes: 2 additions & 1 deletion crates/types/src/post_asap/summary_maintenance.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
17 changes: 13 additions & 4 deletions crates/types/src/post_asap/summary_maintenance_lifecycle.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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,
Expand All @@ -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,
Expand All @@ -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,
Expand Down
Loading