Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
30 commits
Select commit Hold shift + click to select a range
2752c31
Acquire complete immutable input cohorts before maintenance publication
zzylol Sep 11, 2026
190964c
Hold cohort lifetimes through immutable registration and publication
zzylol Sep 11, 2026
537bc86
Bind immutable receipts to canonical complete input lineage
zzylol Sep 11, 2026
b761197
feat: validate shared derived source window cohorts
zzylol Sep 11, 2026
0e96c09
fix: retain legacy nonoverlapping origin contract
zzylol Sep 11, 2026
91f380b
test: distinguish explicit sliding and legacy origins
zzylol Sep 11, 2026
8db3319
Resolve frozen maintenance inputs at each materialized frontier
zzylol Sep 11, 2026
c71c801
Execute complete aligned source cohorts through the installed mainten…
zzylol Sep 11, 2026
be1ec6c
refactor: share the data-plane Float64 arithmetic kernel
zzylol Sep 11, 2026
4c0b4ac
Preserve binary operand roles in maintenance scheduling
zzylol Sep 11, 2026
a38262e
Evaluate aligned frozen Float64 rows with explicit timestamp provenance
zzylol Sep 11, 2026
f39609f
fix: consume explicit binary timing in query compilation
zzylol Sep 11, 2026
91402c8
fix: consume explicit binary timing in query compilation
zzylol Sep 11, 2026
37e2464
Execute explicitly timed binary maintenance over frozen rows
zzylol Sep 11, 2026
c0cd2ea
merge: preserve timed maintenance regressions on current main
zzylol Sep 11, 2026
dfcd7b1
test: preserve read-time binary evidence fixture
zzylol Sep 11, 2026
cf775c4
Merge commit 'dfcd7b19' into feat/compile-multi-source-maintenance
zzylol Sep 11, 2026
6407992
Schedule complete aligned source cohorts under the captured catalog g…
zzylol Sep 11, 2026
675f9d1
Merge commit 'dfcd7b19' into feat/maintenance-finite-cohort
zzylol Sep 11, 2026
6dc0567
fix: retain binary timing in calibration candidate exports
zzylol Sep 11, 2026
8ec8b41
feat: bind actual finite multi-source maintenance plans
zzylol Sep 11, 2026
0245316
Merge commit '6dc0567c' into feat/compile-multi-source-maintenance
zzylol Sep 11, 2026
2271bc0
docs: describe finite arithmetic input cohorts
zzylol Sep 11, 2026
c8a9467
test: assert empty local routing errors without warm provenance
zzylol Sep 11, 2026
124bfe8
Merge remote-tracking branch 'origin/main' into feat/compile-multi-so…
zzylol Sep 11, 2026
ecc7810
Version population key encoding in shared materialization identity
zzylol Sep 11, 2026
0a66643
Keep legacy population encoding in existing test fixtures
zzylol Sep 11, 2026
4475320
Merge remote-tracking branch 'origin/main' into feat/population-key-e…
zzylol Sep 11, 2026
d5c90a6
Merge branch 'feat/compile-multi-source-maintenance' into feat/popula…
zzylol Sep 11, 2026
6d41ee0
Merge remote-tracking branch 'origin/main' into feat/population-key-e…
zzylol Sep 11, 2026
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
19 changes: 19 additions & 0 deletions crates/asap_types/src/aggregation_config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -101,6 +101,11 @@ pub struct PrecomputeMaterialization {
pub parameters: HashMap<String, Value>,
#[serde(serialize_with = "crate::grouping_projection::serialize_config_grouping")]
pub grouping_labels: crate::GroupingProjection,
#[serde(
default,
skip_serializing_if = "crate::grouping_projection::PopulationKeyEncoding::is_legacy"
)]
pub population_key_encoding: crate::grouping_projection::PopulationKeyEncoding,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub partitioning: Option<crate::sds::PopulationPartitioning>,
#[serde(default, skip_serializing_if = "Option::is_none")]
Expand Down Expand Up @@ -286,6 +291,7 @@ impl PrecomputeMaterialization {
aggregation_sub_type,
parameters,
grouping_labels: grouping_labels.into(),
population_key_encoding: Default::default(),
partitioning: None,
derived_input: None,
aggregated_labels,
Expand Down Expand Up @@ -433,6 +439,11 @@ impl PrecomputeMaterialization {
table_name,
value_column,
);
config.population_key_encoding = data
.get("population_key_encoding")
.map(|value| serde_json::from_value(value.clone()))
.transpose()?
.unwrap_or_default();
config.derived_input = data
.get("derived_input")
.filter(|v| !v.is_null())
Expand Down Expand Up @@ -633,6 +644,11 @@ impl PrecomputeMaterialization {
table_name,
value_column,
);
config.population_key_encoding = aggregation_data
.get("population_key_encoding")
.map(|value| serde_yaml::from_value(value.clone()))
.transpose()?
.unwrap_or_default();
config.derived_input = aggregation_data
.get("derived_input")
.filter(|v| !v.is_null())
Expand Down Expand Up @@ -683,6 +699,9 @@ impl SerializableToSink for PrecomputeMaterialization {
"metric": self.metric,
});

if !self.population_key_encoding.is_legacy() {
json["population_key_encoding"] = serde_json::json!(self.population_key_encoding);
}
if let Some(input) = &self.derived_input {
json["derived_input"] = serde_json::json!(input);
}
Expand Down
15 changes: 15 additions & 0 deletions crates/asap_types/src/grouping_projection.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,21 @@ use crate::KeyByLabelNames;
use planner_types::pre_asap::{Column, DataType};
use serde::{Deserialize, Deserializer, Serialize, Serializer};

/// Version of the population-key routing contract. Choosing a new version
/// changes materialization and data identity; it never migrates old SIDs.
#[derive(Debug, Default, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum PopulationKeyEncoding {
#[default]
LegacyDelimited,
CanonicalLabelsV1,
}
impl PopulationKeyEncoding {
pub fn is_legacy(&self) -> bool {
*self == Self::LegacyDelimited
}
}

#[derive(Debug, Clone, PartialEq, Eq)]
pub struct GroupingProjection(Vec<Column>);

Expand Down
2 changes: 1 addition & 1 deletion crates/asap_types/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,7 @@ pub use accuracy::{AccuracyKind, AccuracyProfile};
pub use aggregation_config::*;
pub use aggregation_type::AggregationType;
pub use enums::*;
pub use grouping_projection::GroupingProjection;
pub use grouping_projection::{GroupingProjection, PopulationKeyEncoding};
pub use key_by_label_names::KeyByLabelNames;
pub use monitor_spec::{MonitorFunctional, MonitorSpec};
pub use policy_fingerprint::PolicyFingerprint;
Expand Down
43 changes: 43 additions & 0 deletions crates/asap_types/src/policy_fingerprint.rs
Original file line number Diff line number Diff line change
Expand Up @@ -89,6 +89,10 @@ impl PolicyFingerprint {
pub fn from_config(cfg: &AggregationConfig) -> Self {
let mut buf: Vec<u8> = Vec::with_capacity(512);

if !cfg.population_key_encoding.is_legacy() {
// UTF-8 raw metrics cannot alias this nonlegacy domain prefix.
buf.extend_from_slice(b"\xffpopulation-key-canonical-labels-v1\0");
}
// 1. metric name
if cfg.derived_input.is_some() {
// Raw policies start with UTF-8 metric bytes; 0xff is impossible
Expand Down Expand Up @@ -285,6 +289,45 @@ mod tests {
)
}

#[test]
fn population_encoding_preserves_legacy_wire_and_separates_identity() {
use crate::grouping_projection::PopulationKeyEncoding;
let legacy = cfg(
"m",
AggregationType::Sum,
HashMap::new(),
vec!["host"],
60,
"",
);
let wire = serde_json::to_value(&legacy).unwrap();
assert!(wire.get("population_key_encoding").is_none());
let decoded: AggregationConfig = serde_json::from_value(wire).unwrap();
assert!(decoded.population_key_encoding.is_legacy());
assert_eq!(legacy.policy_fingerprint(), decoded.policy_fingerprint());
let mut canonical = legacy.clone();
canonical.population_key_encoding = PopulationKeyEncoding::CanonicalLabelsV1;
assert_ne!(legacy.policy_fingerprint(), canonical.policy_fingerprint());
let wire = serde_json::to_value(&canonical).unwrap();
assert_eq!(wire["population_key_encoding"], "canonical_labels_v1");
let decoded: AggregationConfig = serde_json::from_value(wire).unwrap();
assert_eq!(decoded.policy_fingerprint(), canonical.policy_fingerprint());
use crate::traits::SerializableToSink;
let mut sink = canonical.serialize_to_json();
// Transport wrappers supply the three label projections separately.
sink["groupingLabels"] = serde_json::to_value(&canonical.grouping_labels).unwrap();
sink["aggregatedLabels"] =
serde_json::to_value(&canonical.aggregated_labels.labels).unwrap();
sink["rollupLabels"] = serde_json::to_value(&canonical.rollup_labels.labels).unwrap();
let decoded = AggregationConfig::deserialize_from_json(&sink).unwrap();
assert_eq!(
decoded.population_key_encoding,
canonical.population_key_encoding
);

assert!(serde_json::from_str::<PopulationKeyEncoding>("\"canonical_labels_v2\"").is_err());
}

#[test]
fn same_config_yields_same_fingerprint() {
let a = cfg(
Expand Down
28 changes: 28 additions & 0 deletions crates/asap_types/src/precompute_plan.rs
Original file line number Diff line number Diff line change
Expand Up @@ -410,6 +410,11 @@ impl PrecomputePlan {
return Err(PrecomputePlanError::UnsupportedIngestEndpoint);
}
for config in &self.materializations {
if !config.population_key_encoding.is_legacy() {
return Err(PrecomputePlanError::CatalogContract(
"population key encoding is not supported by the installed runtime".into(),
));
}
let Some(derived) = &config.derived_input else {
continue;
};
Expand Down Expand Up @@ -824,6 +829,29 @@ mod source_window_cohort_tests {
}))
.unwrap()
}
#[test]
fn canonical_population_encoding_is_not_yet_installable() {
let mut config = full_window();
config.slide_interval = config.window_size;
config.population_key_encoding =
crate::grouping_projection::PopulationKeyEncoding::CanonicalLabelsV1;
let envelope = PlanEnvelope {
plan_id: 1,
plan_version: 1,
generated_at_unix_ms: 0,
activation_unix_ms: 0,
expiry_unix_ms: None,
backend_compat: "test".into(),
planner_revision: "test".into(),
capability_snapshot_id: "test".into(),
};
let error = PrecomputePlan::build_backend_local(envelope, vec![config]).unwrap_err();
assert!(
error.to_string().contains("population key encoding"),
"{error}"
);
}

#[test]
fn full_sliding_cohort_preserves_explicit_windows_and_identity() {
let target = full_window();
Expand Down
3 changes: 2 additions & 1 deletion crates/asap_types/src/precompute_plan/catalog.rs
Original file line number Diff line number Diff line change
Expand Up @@ -87,7 +87,8 @@ impl PrecomputePlan {
.map_err(invalid)?;
}
let expected_projection = config.effective_value_projection();
if data.partitioning != config.partitioning
if data.population_key_encoding != config.population_key_encoding
|| data.partitioning != config.partitioning
|| data.timestamp_column != config.table_timestamp_column
|| data.source != expected_source
|| &data.value_projection != expected_projection
Expand Down
41 changes: 41 additions & 0 deletions crates/asap_types/src/sds.rs
Original file line number Diff line number Diff line change
Expand Up @@ -842,6 +842,11 @@ pub struct DataDescriptor {
pub timestamp_column: Option<String>,
pub population_filter_canonical: String,
pub group_by_keys: crate::GroupingProjection,
#[serde(
default,
skip_serializing_if = "crate::grouping_projection::PopulationKeyEncoding::is_legacy"
)]
pub population_key_encoding: crate::grouping_projection::PopulationKeyEncoding,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub partitioning: Option<PopulationPartitioning>,
/// Versioned contract for timestamp interpretation and
Expand Down Expand Up @@ -905,6 +910,7 @@ impl DataDescriptor {
&observation_semantics,
None,
None,
Default::default(),
);
Self {
id,
Expand All @@ -914,6 +920,7 @@ impl DataDescriptor {
population_filter_canonical,
group_by_keys,
partitioning: None,
population_key_encoding: Default::default(),
observation_semantics,
}
}
Expand All @@ -927,6 +934,7 @@ impl DataDescriptor {
&self.observation_semantics,
self.partitioning,
self.timestamp_column.as_deref(),
self.population_key_encoding,
);
self
}
Expand All @@ -940,9 +948,18 @@ impl DataDescriptor {
&self.observation_semantics,
partitioning,
self.timestamp_column.as_deref(),
self.population_key_encoding,
);
self
}
pub fn with_population_key_encoding(
mut self,
encoding: crate::grouping_projection::PopulationKeyEncoding,
) -> Self {
self.population_key_encoding = encoding;
let partitioning = self.partitioning;
self.with_partitioning(partitioning)
}
pub fn with_timestamp_column(mut self, column: Option<String>) -> Self {
self.timestamp_column = column;
let partitioning = self.partitioning;
Expand Down Expand Up @@ -987,13 +1004,15 @@ impl DataDescriptor {
&self.observation_semantics,
self.partitioning,
self.timestamp_column.as_deref(),
self.population_key_encoding,
)
{
return Err(SdsError("data descriptor ID/content mismatch".into()));
}
Ok(())
}
}
#[allow(clippy::too_many_arguments)]
fn data_descriptor_id(
source: &DataSourceIdentity,
value_projection: &ValueProjectionIdentity,
Expand All @@ -1002,6 +1021,7 @@ fn data_descriptor_id(
observation_semantics: &str,
partitioning: Option<PopulationPartitioning>,
timestamp_column: Option<&str>,
population_key_encoding: crate::grouping_projection::PopulationKeyEncoding,
) -> DataDescriptorId {
// Length framing keeps distinct typed sources, projections, predicates,
// and grouping keys collision-free in the content identity.
Expand Down Expand Up @@ -1032,6 +1052,9 @@ fn data_descriptor_id(
"|{}:{observation_semantics}",
observation_semantics.len()
));
if !population_key_encoding.is_legacy() {
key = format!("data:population-key:canonical-labels-v1|{key}");
}
DataDescriptorId(key)
}

Expand All @@ -1040,6 +1063,24 @@ mod tests {
use super::*;
use std::collections::BTreeSet;

#[test]
fn population_encoding_is_part_of_data_identity_with_legacy_default() {
use crate::grouping_projection::PopulationKeyEncoding;
let legacy = DataDescriptor::new("m", "", ["host".to_string()]);
let wire = serde_json::to_value(&legacy).unwrap();
assert!(wire.get("population_key_encoding").is_none());
let decoded: DataDescriptor = serde_json::from_value(wire).unwrap();
assert_eq!(legacy.id(), decoded.id());
let canonical = legacy
.clone()
.with_population_key_encoding(PopulationKeyEncoding::CanonicalLabelsV1);
assert_ne!(legacy.id(), canonical.id());
canonical.validate().unwrap();
let reverted =
canonical.with_population_key_encoding(PopulationKeyEncoding::LegacyDelimited);
assert_eq!(legacy.id(), reverted.id());
}

fn completion(epoch: u64) -> SummaryWindowCompletion {
SummaryWindowCompletion {
catalog_generation: CatalogGeneration {
Expand Down
1 change: 1 addition & 0 deletions crates/asap_types/src/summary_catalog.rs
Original file line number Diff line number Diff line change
Expand Up @@ -120,6 +120,7 @@ impl SummaryCatalog {
)
.with_grouping_projection(config.grouping_labels.clone())
.with_partitioning(config.partitioning)
.with_population_key_encoding(config.population_key_encoding)
.with_timestamp_column(config.table_timestamp_column.clone());
Ok((
config.policy_fingerprint(),
Expand Down
2 changes: 2 additions & 0 deletions data_plane/src/drivers/ingest/prometheus_remote_write.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1013,6 +1013,7 @@ mod tests {
use asap_types::enums::WindowKind;
use asap_types::{AggregationConfig, AggregationType, KeyByLabelNames};
let aggregation = AggregationConfig {
population_key_encoding: Default::default(),
aggregation_type: AggregationType::Sum,
aggregation_sub_type: String::new(),
parameters: HashMap::new(),
Expand Down Expand Up @@ -1075,6 +1076,7 @@ mod tests {

let config =
|aggregation_type, grouping: Vec<String>, aggregated: Vec<String>| AggregationConfig {
population_key_encoding: Default::default(),
aggregation_type,
aggregation_sub_type: String::new(),
parameters: match aggregation_type {
Expand Down
1 change: 1 addition & 0 deletions data_plane/src/drivers/query/servers/http.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3704,6 +3704,7 @@ aggregations:
for marker in active_agg_ids {
let metric = format!("metric_{marker}");
let cfg = AggregationConfig {
population_key_encoding: Default::default(),
aggregation_type: AggregationType::Sum,
aggregation_sub_type: String::new(),
parameters: HashMap::new(),
Expand Down
1 change: 1 addition & 0 deletions data_plane/src/precompute_engine/output_sink.rs
Original file line number Diff line number Diff line change
Expand Up @@ -404,6 +404,7 @@ mod tests {
// via `PolicyFingerprint::from_config`. Callers obtain the id
// via `config.policy_fp_u64()`.
AggregationConfig {
population_key_encoding: Default::default(),
aggregation_type: AggregationType::Sum,
aggregation_sub_type: String::new(),
parameters: HashMap::new(),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -257,6 +257,7 @@ mod tests {

fn sum_agg_config(id: u64) -> AggregationConfig {
AggregationConfig {
population_key_encoding: Default::default(),
aggregation_type: AggregationType::Sum,
aggregation_sub_type: String::new(),
parameters: HashMap::new(),
Expand Down
Loading
Loading