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
4 changes: 2 additions & 2 deletions control_plane/src/opamp/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1027,9 +1027,9 @@ mod tests {
abstract_window_framework:
planner_types::post_asap::SummaryWindowFramework::Tumbling,
window_implementation_id: "collector-tumbling-v1".into(),
pane_secs: 60,
slide_secs: 60,
pane_origin_ms: Some(0),
state_layout: "anchored-pane-v1".into(),
window_layout: asap_types::WindowMaterializationLayout::Pane { pane_secs: 60 },
evidence_source: None,
lifecycle: crate::physical::compiler::CollectorLifecycle {
kind: "continuously_maintained".into(),
Expand Down
341 changes: 217 additions & 124 deletions control_plane/src/physical/compiler.rs

Large diffs are not rendered by default.

113 changes: 113 additions & 0 deletions crates/asap_types/src/aggregation_config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,79 @@ use crate::utils::normalize_spatial_filter;
use crate::AggregationType;
use crate::KeyByLabelNames;

/// Physical maintenance layout for one semantic windowed summary.
///
/// `window_size` and `slide_interval` on [`PrecomputeMaterialization`] retain
/// the query's window and evaluation cadence. This enum describes how that
/// semantic window is represented in storage; it must never be inferred by
/// overloading either semantic duration.
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq, Hash)]
#[serde(tag = "kind", rename_all = "snake_case", deny_unknown_fields)]
pub enum WindowMaterializationLayout {
/// Store disjoint mergeable states and compose a query window at read time.
Pane { pane_secs: u64 },
/// Maintain one complete state for every evaluation point.
FullWindow,
/// Store base panes plus coarser mergeable rollups. Every level is a
/// duration in seconds and is an integer multiple of its predecessor.
HierarchicalRollup {
base_pane_secs: u64,
levels_secs: Vec<u64>,
},
}

impl WindowMaterializationLayout {
pub fn base_pane_secs(&self) -> u64 {
match self {
Self::Pane { pane_secs } => *pane_secs,
Self::FullWindow => 0,
Self::HierarchicalRollup { base_pane_secs, .. } => *base_pane_secs,
}
}

pub fn validate(&self, window_secs: u64, slide_secs: u64) -> Result<(), String> {
if window_secs == 0 || slide_secs == 0 || slide_secs > window_secs {
return Err(
"window and slide must be positive and slide must not exceed window".into(),
);
}
match self {
Self::Pane { pane_secs } => {
if *pane_secs == 0 || window_secs % pane_secs != 0 || slide_secs % pane_secs != 0 {
return Err("pane size must divide both window size and slide".into());
}
}
Self::FullWindow => {}
Self::HierarchicalRollup {
base_pane_secs,
levels_secs,
} => {
if *base_pane_secs == 0
|| window_secs % base_pane_secs != 0
|| slide_secs % base_pane_secs != 0
|| levels_secs.is_empty()
{
return Err(
"rollup base pane must divide window and slide, with at least one level"
.into(),
);
}
let mut previous = *base_pane_secs;
for level in levels_secs {
if *level <= previous || *level % previous != 0 || window_secs % level != 0 {
return Err(
"rollup levels must increase by integral factors and divide the window"
.into(),
);
}
previous = *level;
}
}
}
Ok(())
}
}

/// Per-aggregation policy carried in the streaming config.
///
/// **PR 5 (merged-sid-identity refactor)** retired the
Expand All @@ -34,6 +107,7 @@ pub struct PrecomputeMaterialization {
pub window_size: u64, // Window size in seconds (e.g., 900s for 15m)
pub slide_interval: u64, // Slide/hop interval in seconds (e.g., 30s)
pub window_type: WindowKind, // Tumbling or Sliding
pub window_layout: WindowMaterializationLayout,
/// Unix millisecond timestamp on the pane-boundary grid selected from
/// the consuming query workload. Missing on legacy definitions, which
/// must not be used for certified pane-only reads.
Expand Down Expand Up @@ -121,6 +195,13 @@ impl PrecomputeMaterialization {
window_size,
slide_interval,
window_type,
window_layout: WindowMaterializationLayout::Pane {
pane_secs: if slide_interval == 0 {
window_size
} else {
slide_interval
},
},
pane_origin_ms: None,
spatial_filter,
spatial_filter_normalized,
Expand Down Expand Up @@ -420,6 +501,38 @@ impl SerializableToSink for PrecomputeMaterialization {
}
}

#[cfg(test)]
mod window_layout_tests {
use super::WindowMaterializationLayout;

#[test]
fn validates_multiple_slides_and_rejects_uncomposable_panes() {
for slide in [5, 10, 15, 30] {
WindowMaterializationLayout::Pane { pane_secs: 5 }
.validate(60, slide)
.unwrap();
WindowMaterializationLayout::FullWindow
.validate(60, slide)
.unwrap();
}
assert!(WindowMaterializationLayout::Pane { pane_secs: 7 }
.validate(60, 10)
.is_err());
assert!(WindowMaterializationLayout::HierarchicalRollup {
base_pane_secs: 5,
levels_secs: vec![10, 30],
}
.validate(60, 10)
.is_ok());
assert!(WindowMaterializationLayout::HierarchicalRollup {
base_pane_secs: 5,
levels_secs: vec![12],
}
.validate(60, 10)
.is_err());
}
}

#[cfg(test)]
mod tests {
use super::*;
Expand Down
28 changes: 28 additions & 0 deletions crates/asap_types/src/policy_fingerprint.rs
Original file line number Diff line number Diff line change
Expand Up @@ -148,6 +148,16 @@ impl PolicyFingerprint {
);
buf.push(0);

// Physical representation is part of state identity. A full overlapping
// window and a mergeable pane layout may share semantic descriptors but
// never share payload instances or lifecycle accounting.
buf.extend_from_slice(
serde_json::to_string(&cfg.window_layout)
.unwrap_or_default()
.as_bytes(),
);
buf.push(0);

// 9. pane origin. Presence is explicit so a legacy definition with
// unknown phase cannot alias an epoch-aligned definition.
match cfg.pane_origin_ms {
Expand Down Expand Up @@ -239,6 +249,24 @@ mod tests {
);
}

#[test]
fn physical_window_layout_is_part_of_state_identity() {
let panes = cfg(
"http_lat",
AggregationType::Sum,
HashMap::new(),
vec!["zone"],
60,
"",
);
let mut full = panes.clone();
full.window_layout = crate::WindowMaterializationLayout::FullWindow;
assert_ne!(
PolicyFingerprint::from_config(&panes),
PolicyFingerprint::from_config(&full)
);
}

#[test]
fn different_metric_yields_different_fingerprint() {
let a = cfg(
Expand Down
37 changes: 32 additions & 5 deletions crates/asap_types/src/summary_catalog.rs
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ use crate::sds::{
SummaryDescriptor, SummaryDescriptorId, ValueProjectionIdentity,
};
use crate::PolicyFingerprint;
use crate::WindowMaterializationLayout;
use serde::{Deserialize, Serialize};

pub const SUMMARY_CATALOG_SCHEMA_VERSION: u32 = 2;
Expand All @@ -21,6 +22,10 @@ pub const SUMMARY_CATALOG_SCHEMA_VERSION: u32 = 2;
pub struct SummaryDefinitionIdentity {
pub summary_descriptor_id: SummaryDescriptorId,
pub data_descriptor_id: DataDescriptorId,
/// Backend-selected physical representation. It is catalog-visible so
/// producers, readers, lifecycle management, and recovery agree on the
/// concrete state being referenced.
pub window_layout: WindowMaterializationLayout,
/// Pane boundary selected from the shared consumer workload. Legacy
/// snapshots deserialize as unknown and fail closed at pane-only reads.
#[serde(
Expand Down Expand Up @@ -123,6 +128,7 @@ impl SummaryCatalog {
config.policy_fingerprint(),
summary,
data,
config.window_layout.clone(),
config.pane_origin_ms,
))
})
Expand All @@ -133,14 +139,23 @@ impl SummaryCatalog {
pub fn build(
plan_id: u64,
plan_version: u64,
entries: impl IntoIterator<Item = (PolicyFingerprint, SummaryDescriptor, DataDescriptor)>,
entries: impl IntoIterator<
Item = (
PolicyFingerprint,
SummaryDescriptor,
DataDescriptor,
WindowMaterializationLayout,
),
>,
) -> Result<Self, SummaryCatalogError> {
Self::build_with_origins(
plan_id,
plan_version,
entries
.into_iter()
.map(|(fingerprint, summary, data)| (fingerprint, summary, data, None)),
.map(|(fingerprint, summary, data, layout)| {
(fingerprint, summary, data, layout, None)
}),
)
}

Expand All @@ -152,6 +167,7 @@ impl SummaryCatalog {
PolicyFingerprint,
SummaryDescriptor,
DataDescriptor,
WindowMaterializationLayout,
Option<i64>,
),
>,
Expand All @@ -164,7 +180,7 @@ impl SummaryCatalog {
data_descriptors: BTreeMap::new(),
materializations: BTreeMap::new(),
};
for (fingerprint, summary, data, pane_origin_ms) in entries {
for (fingerprint, summary, data, window_layout, pane_origin_ms) in entries {
let materialization = SummaryDefinitionId::from(fingerprint);
summary
.validate()
Expand All @@ -174,6 +190,7 @@ impl SummaryCatalog {
let binding = SummaryDefinitionIdentity {
summary_descriptor_id: summary.id().clone(),
data_descriptor_id: data.id().clone(),
window_layout,
pane_origin_ms,
};
if catalog
Expand Down Expand Up @@ -374,8 +391,18 @@ mod tests {
1,
1,
[
(first.policy_fingerprint(), summary.clone(), data[0].clone()),
(first.policy_fingerprint(), summary, data[1].clone()),
(
first.policy_fingerprint(),
summary.clone(),
data[0].clone(),
first.window_layout.clone(),
),
(
first.policy_fingerprint(),
summary,
data[1].clone(),
first.window_layout.clone(),
),
],
)
.unwrap_err();
Expand Down
3 changes: 2 additions & 1 deletion data_plane/src/drivers/ingest/prometheus_remote_write.rs
Original file line number Diff line number Diff line change
Expand Up @@ -365,7 +365,6 @@ impl PrometheusRemoteWriteReceiver {
.ingest
.samples_ingested
.fetch_add(new_samples.len() as u64, Ordering::Relaxed);
crate::precompute_engine::metrics::record_accepted_samples(new_samples.len() as u64);
Ok(())
}
}
Expand Down Expand Up @@ -840,6 +839,7 @@ mod tests {
window_size: 60,
slide_interval: 60,
window_type: WindowKind::Tumbling,
window_layout: asap_types::WindowMaterializationLayout::Pane { pane_secs: 60 },
pane_origin_ms: None,
spatial_filter: String::new(),
spatial_filter_normalized: String::new(),
Expand Down Expand Up @@ -892,6 +892,7 @@ mod tests {
window_size: 60,
slide_interval: 60,
window_type: WindowKind::Tumbling,
window_layout: asap_types::WindowMaterializationLayout::Pane { pane_secs: 60 },
pane_origin_ms: Some(0),
spatial_filter: String::new(),
spatial_filter_normalized: String::new(),
Expand Down
Loading
Loading