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
192 changes: 191 additions & 1 deletion crates/asap_types/src/sds.rs
Original file line number Diff line number Diff line change
Expand Up @@ -64,6 +64,14 @@ impl SummaryInstanceId {
}
Ok(Self(value))
}

pub fn validate(&self) -> Result<(), SdsError> {
if self.0.trim().is_empty() {
Err(SdsError("summary instance ID must not be empty".into()))
} else {
Ok(())
}
}
}

/// Immutable identity of the desired catalog generation used to create state.
Expand All @@ -77,13 +85,118 @@ pub struct CatalogGeneration {
pub snapshot_sha256: String,
}

#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct HalfOpenTimeRange {
pub start_ms: i64,
pub end_ms: i64,
}

/// One ordered producer partition. Completion and watermark claims are scoped
/// to this identity; a maximum timestamp observed by an unrelated worker is
/// never a source-completeness signal.
#[derive(Debug, Clone, PartialEq, Eq, Hash, PartialOrd, Ord, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct SummarySourcePartition {
pub producer_id: String,
pub partition_id: String,
/// Changes whenever a producer restarts or loses its sequence state.
pub producer_epoch: u64,
}

/// A closed materialization bucket, including named group values so a
/// downstream maintenance DAG can project or shuffle groups without parsing a
/// routing key.
#[derive(Debug, Clone, PartialEq, Eq, Hash, PartialOrd, Ord, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct SummaryInstanceCoordinates {
pub summary_definition_id: SummaryDefinitionId,
pub time_range: HalfOpenTimeRange,
pub group_values: BTreeMap<String, String>,
}

#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct SummaryWindowCompletion {
pub catalog_generation: CatalogGeneration,
pub source: SummarySourcePartition,
/// Concrete SDS instance identity allocated by the storage engine.
pub instance_id: SummaryInstanceId,
pub coordinates: SummaryInstanceCoordinates,
/// Collision-free opaque producer lineage. Retries repeat these bytes.
pub input_lineage: Vec<u8>,
}

impl HalfOpenTimeRange {
pub fn validate(self) -> Result<(), SdsError> {
if self.start_ms >= self.end_ms {
Err(SdsError(
"time range must be non-empty and half-open".into(),
))
} else {
Ok(())
}
}
}

impl SummarySourcePartition {
pub fn validate(&self) -> Result<(), SdsError> {
if self.producer_id.trim().is_empty() || self.partition_id.trim().is_empty() {
return Err(SdsError(
"source producer and partition must be non-empty".into(),
));
}
if self.producer_epoch == 0 {
return Err(SdsError("source producer epoch must be positive".into()));
}
Ok(())
}
}

impl SummaryWindowCompletion {
pub fn validate(&self) -> Result<(), SdsError> {
self.source.validate()?;
self.instance_id.validate()?;
self.coordinates.time_range.validate()?;
if self.input_lineage.is_empty() {
return Err(SdsError(
"summary completion lineage must not be empty".into(),
));
}
validate_catalog_generation(&self.catalog_generation)
}
}

#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct SummaryWatermarkBarrier {
pub catalog_generation: CatalogGeneration,
pub source: SummarySourcePartition,
/// Monotonic sequence within the producer partition.
pub sequence: u64,
/// Every event in this partition at or before this event-time watermark
/// was published before the barrier.
pub watermark_ms: i64,
}

impl SummaryWatermarkBarrier {
pub fn validate(&self) -> Result<(), SdsError> {
self.source.validate()?;
if self.sequence == 0 {
return Err(SdsError("watermark sequence must be positive".into()));
}
validate_catalog_generation(&self.catalog_generation)
}
}

fn validate_catalog_generation(generation: &CatalogGeneration) -> Result<(), SdsError> {
if generation.schema_version == 0 || generation.snapshot_sha256.trim().is_empty() {
Err(SdsError("catalog generation identity is invalid".into()))
} else {
Ok(())
}
}

#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum InstanceCompleteness {
Expand Down Expand Up @@ -712,6 +825,83 @@ fn data_descriptor_id(
mod tests {
use super::*;

fn completion(epoch: u64) -> SummaryWindowCompletion {
SummaryWindowCompletion {
catalog_generation: CatalogGeneration {
schema_version: 1,
plan_id: 4,
plan_version: 2,
snapshot_sha256: "snapshot".into(),
},
source: SummarySourcePartition {
producer_id: "producer-a".into(),
partition_id: "partition-0".into(),
producer_epoch: epoch,
},
instance_id: SummaryInstanceId::new("instance-1").unwrap(),
coordinates: SummaryInstanceCoordinates {
summary_definition_id: SummaryDefinitionId(crate::PolicyFingerprint(7)),
time_range: HalfOpenTimeRange {
start_ms: 1_000,
end_ms: 2_000,
},
group_values: BTreeMap::from([("job".into(), "api".into())]),
},
input_lineage: vec![1, 2, 3],
}
}

#[test]
fn completion_contract_validates_identity_epoch_range_and_lineage() {
completion(1).validate().unwrap();

let mut invalid = completion(1);
invalid.coordinates.time_range.end_ms = invalid.coordinates.time_range.start_ms;
assert!(invalid.validate().is_err());
invalid = completion(1);
invalid.source.producer_id.clear();
assert!(invalid.validate().is_err());
invalid = completion(0);
assert!(invalid.validate().is_err());
invalid = completion(1);
invalid.input_lineage.clear();
assert!(invalid.validate().is_err());

let mut wire = serde_json::to_value(completion(1)).unwrap();
wire.as_object_mut()
.unwrap()
.insert("future_field".into(), json!(true));
assert!(serde_json::from_value::<SummaryWindowCompletion>(wire).is_err());
}

#[test]
fn producer_epoch_separates_restart_incarnations() {
let first = completion(1).source;
let second = completion(2).source;
assert_ne!(first, second);
assert_eq!(BTreeSet::from([first, second]).len(), 2);
}

#[test]
fn watermark_contract_rejects_zero_sequence_and_unknown_fields() {
let completion = completion(3);
let mut barrier = SummaryWatermarkBarrier {
catalog_generation: completion.catalog_generation,
source: completion.source,
sequence: 1,
watermark_ms: 2_000,
};
barrier.validate().unwrap();
barrier.sequence = 0;
assert!(barrier.validate().is_err());

let mut wire = serde_json::to_value(&barrier).unwrap();
wire.as_object_mut()
.unwrap()
.insert("ignored".into(), json!(1));
assert!(serde_json::from_value::<SummaryWatermarkBarrier>(wire).is_err());
}

fn observed_instance(lifecycle: InstanceLifecycle) -> SummaryInstance {
SummaryInstance {
instance_id: SummaryInstanceId::new("instance-1").unwrap(),
Expand Down
51 changes: 51 additions & 0 deletions data_plane/src/precompute_engine/output_sink.rs
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,18 @@ pub trait OutputSink: Send + Sync {
&self,
outputs: Vec<(PrecomputedOutput, Box<dyn AggregateCore>)>,
) -> Result<(), Box<dyn std::error::Error + Send + Sync>>;

/// Publish an explicit source-partition event-time barrier after every
/// preceding window from that partition has reached this sink.
fn advance_summary_watermark(
&self,
_barrier: &asap_types::sds::SummaryWatermarkBarrier,
) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
Err(Box::new(std::io::Error::new(
std::io::ErrorKind::Unsupported,
"output sink does not support summary watermark barriers",
)))
}
}

fn consume_in_order<T>(items: Vec<T>, mut persist: impl FnMut(&T) -> bool) -> usize {
Expand Down Expand Up @@ -268,6 +280,45 @@ mod tests {
use std::collections::HashMap;
use std::sync::atomic::{AtomicUsize, Ordering};

struct EmitOnlySink;

impl OutputSink for EmitOnlySink {
fn emit_batch(
&self,
_outputs: Vec<(PrecomputedOutput, Box<dyn AggregateCore>)>,
) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
Ok(())
}
}

#[test]
fn watermark_support_must_be_explicit() {
let completion_source = asap_types::sds::SummarySourcePartition {
producer_id: "producer".into(),
partition_id: "partition".into(),
producer_epoch: 1,
};
let barrier = asap_types::sds::SummaryWatermarkBarrier {
catalog_generation: asap_types::sds::CatalogGeneration {
schema_version: 1,
plan_id: 1,
plan_version: 1,
snapshot_sha256: "snapshot".into(),
},
source: completion_source,
sequence: 1,
watermark_ms: 10,
};
let err = EmitOnlySink
.advance_summary_watermark(&barrier)
.expect_err("sink must reject an unsupported barrier");
assert_eq!(
err.downcast_ref::<std::io::Error>()
.map(std::io::Error::kind),
Some(std::io::ErrorKind::Unsupported)
);
}

fn sum_agg_config(_id: u64, metric: &str, grouping_keys: &[&str]) -> AggregationConfig {
// `_id` is unused after PR 5 — identity is content-addressed
// via `PolicyFingerprint::from_config`. Callers obtain the id
Expand Down
Loading