From 98066ceb4bc4efed1bf0573edfe3ba6132183cde Mon Sep 17 00:00:00 2001 From: zzylol Date: Tue, 29 Sep 2026 20:21:06 +0000 Subject: [PATCH] refactor: delete unused precompute coordinator and series buffer `multisource_coordinator`, `coordination_checkpoint` and `series_buffer` have no callers outside themselves and `mod.rs`; production precompute does not construct them. Remove them (1.8k lines) and their design-doc sections. Co-Authored-By: Claude Opus 5.5 --- .../coordination_checkpoint.rs | 548 -------- data_plane/src/precompute_engine/mod.rs | 3 - .../multisource_coordinator.rs | 1104 ----------------- .../precompute_engine_design_doc.md | 27 +- .../src/precompute_engine/series_buffer.rs | 163 --- 5 files changed, 4 insertions(+), 1841 deletions(-) delete mode 100644 data_plane/src/precompute_engine/coordination_checkpoint.rs delete mode 100644 data_plane/src/precompute_engine/multisource_coordinator.rs delete mode 100644 data_plane/src/precompute_engine/series_buffer.rs diff --git a/data_plane/src/precompute_engine/coordination_checkpoint.rs b/data_plane/src/precompute_engine/coordination_checkpoint.rs deleted file mode 100644 index 873725505..000000000 --- a/data_plane/src/precompute_engine/coordination_checkpoint.rs +++ /dev/null @@ -1,548 +0,0 @@ -//! Durable staging and publication decisions for cross-source maintenance. - -use asap_types::sds::{ - CatalogGeneration, SummaryInstanceCoordinates, SummaryInstanceId, SummarySourcePartition, - SummaryStateReference, SummaryWatermarkBarrier, -}; -use fs2::FileExt; -use serde::{Deserialize, Serialize}; -use std::fs::{self, File, OpenOptions}; -use std::io::{self, Write}; -use std::path::{Path, PathBuf}; -use std::sync::atomic::{AtomicBool, Ordering}; -use std::sync::Mutex; - -const SCHEMA_VERSION: u32 = 1; - -#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] -#[serde(deny_unknown_fields)] -pub struct StagedSummaryInput { - pub catalog_generation: CatalogGeneration, - pub dag_id: String, - pub consumer_node_id: String, - pub input_node_id: String, - pub source: SummarySourcePartition, - pub instance_id: SummaryInstanceId, - pub coordinates: SummaryInstanceCoordinates, - pub input_lineage: Vec, - /// Durable SummaryStore payload; the checkpoint store never duplicates state bytes. - pub state_reference: SummaryStateReference, -} - -impl StagedSummaryInput { - fn validate(&self) -> io::Result<()> { - if self.dag_id.trim().is_empty() - || self.consumer_node_id.trim().is_empty() - || self.input_node_id.trim().is_empty() - { - return Err(invalid("staged input contains an empty required field")); - } - let completion = asap_types::sds::SummaryWindowCompletion { - catalog_generation: self.catalog_generation.clone(), - source: self.source.clone(), - instance_id: self.instance_id.clone(), - coordinates: self.coordinates.clone(), - input_lineage: self.input_lineage.clone(), - }; - completion - .validate() - .map_err(|error| invalid(error.to_string()))?; - validate_state_reference(&self.state_reference) - } -} - -#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] -#[serde(deny_unknown_fields)] -pub struct AtomicPublicationKey { - pub catalog_generation: CatalogGeneration, - pub dag_id: String, - pub sink_node_id: String, - pub instance_id: SummaryInstanceId, - pub coordinates: SummaryInstanceCoordinates, - pub output_lineage: Vec, - pub state_reference: SummaryStateReference, -} - -impl AtomicPublicationKey { - fn validate(&self) -> io::Result<()> { - if self.dag_id.trim().is_empty() - || self.sink_node_id.trim().is_empty() - || self.output_lineage.is_empty() - { - return Err(invalid("publication key contains an empty required field")); - } - self.instance_id - .validate() - .map_err(|error| invalid(error.to_string()))?; - self.coordinates - .time_range - .validate() - .map_err(|error| invalid(error.to_string()))?; - validate_state_reference(&self.state_reference) - } -} - -#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] -#[serde(deny_unknown_fields)] -struct CheckpointDocument { - schema_version: u32, - revision: u64, - #[serde(default, skip_serializing_if = "Option::is_none")] - coordinator_scope: Option, - staged: Vec, - watermarks: Vec, - published: Vec, -} - -impl Default for CheckpointDocument { - fn default() -> Self { - Self { - schema_version: SCHEMA_VERSION, - revision: 0, - coordinator_scope: None, - staged: Vec::new(), - watermarks: Vec::new(), - published: Vec::new(), - } - } -} - -/// A small crash-safe checkpoint snapshot store. Mutations become visible in memory -/// only after the replacement file has been flushed and atomically renamed. -pub struct SummaryCoordinationCheckpointStore { - path: PathBuf, - /// Held for the checkpoint store lifetime; prevents two processes from replacing - /// the same snapshot concurrently. - _lock_file: File, - document: Mutex, - persistence_uncertain: AtomicBool, -} - -impl SummaryCoordinationCheckpointStore { - pub fn open(path: impl Into) -> io::Result { - let path = path.into(); - let parent = path.parent().unwrap_or_else(|| Path::new(".")); - fs::create_dir_all(parent)?; - let lock_path = path.with_extension("lock"); - let lock_file = OpenOptions::new() - .create(true) - .read(true) - .write(true) - .open(lock_path)?; - lock_file.try_lock_exclusive().map_err(|error| { - io::Error::new( - io::ErrorKind::WouldBlock, - format!("coordination checkpoint store already has a writer: {error}"), - ) - })?; - let document = match fs::read(&path) { - Ok(bytes) => serde_json::from_slice::(&bytes).map_err(|error| { - invalid(format!("invalid coordination checkpoint store: {error}")) - })?, - Err(error) if error.kind() == io::ErrorKind::NotFound => CheckpointDocument::default(), - Err(error) => return Err(error), - }; - validate_document(&document)?; - Ok(Self { - path, - _lock_file: lock_file, - document: Mutex::new(document), - persistence_uncertain: AtomicBool::new(false), - }) - } - - pub(crate) fn bind_coordinator_scope( - &self, - spec: &super::multisource_coordinator::MultiSourceNodeSpec, - ) -> io::Result { - super::multisource_coordinator::validate_spec(spec)?; - self.mutate(|document| { - if let Some(existing) = &document.coordinator_scope { - return if existing == spec { - Ok(false) - } else { - Err(invalid("persisted coordinator scope changed")) - }; - } - if !document.staged.is_empty() - || !document.watermarks.is_empty() - || !document.published.is_empty() - { - return Err(invalid( - "nonempty legacy checkpoint has no authoritative coordinator scope", - )); - } - document.coordinator_scope = Some(spec.clone()); - Ok(true) - }) - } - - pub fn stage_if_absent(&self, input: StagedSummaryInput) -> io::Result { - input.validate()?; - self.mutate(|document| { - if let Some(existing) = document - .staged - .iter() - .find(|item| item.instance_id == input.instance_id) - { - return if existing == &input { - Ok(false) - } else { - Err(invalid( - "summary instance ID was reused with different metadata", - )) - }; - } - document.staged.push(input); - Ok(true) - }) - } - - pub fn advance_watermark(&self, barrier: SummaryWatermarkBarrier) -> io::Result { - barrier - .validate() - .map_err(|error| invalid(error.to_string()))?; - self.mutate(|document| { - if let Some(existing) = document.watermarks.iter_mut().find(|item| { - item.catalog_generation == barrier.catalog_generation - && item.source == barrier.source - }) { - if barrier.sequence == existing.sequence && *existing != barrier { - return Err(invalid("equal watermark sequence changed its claim")); - } - if barrier.sequence < existing.sequence - || barrier.watermark_ms < existing.watermark_ms - { - return Err(invalid( - "watermark or sequence regressed within a producer epoch", - )); - } - if *existing == barrier { - return Ok(false); - } - *existing = barrier; - } else { - document.watermarks.push(barrier); - } - Ok(true) - }) - } - - #[cfg(test)] - pub fn publish_if_absent(&self, key: AtomicPublicationKey) -> io::Result { - key.validate()?; - self.mutate(|document| { - if let Some(existing) = document - .published - .iter() - .find(|item| item.instance_id == key.instance_id) - { - return if existing == &key { - Ok(false) - } else { - Err(invalid( - "published summary instance ID was reused with different metadata", - )) - }; - } - document.published.push(key); - Ok(true) - }) - } - - pub fn staged(&self) -> io::Result> { - Ok(self.lock()?.staged.clone()) - } - - pub fn watermarks(&self) -> io::Result> { - Ok(self.lock()?.watermarks.clone()) - } - - pub fn is_published(&self, key: &AtomicPublicationKey) -> io::Result { - Ok(self.lock()?.published.contains(key)) - } - - fn mutate( - &self, - update: impl FnOnce(&mut CheckpointDocument) -> io::Result, - ) -> io::Result { - let mut guard = self.lock()?; - let mut next = guard.clone(); - let result = update(&mut next)?; - if next == *guard { - return Ok(result); - } - next.revision = next - .revision - .checked_add(1) - .ok_or_else(|| invalid("coordination checkpoint store revision overflow"))?; - if let Err(error) = persist_atomically(&self.path, &next) { - // Rename may have succeeded before directory fsync failed. Do not - // overwrite that possible newer checkpoint from stale live state. - self.persistence_uncertain.store(true, Ordering::Release); - return Err(error); - } - *guard = next; - Ok(result) - } - - fn lock(&self) -> io::Result> { - let guard = self - .document - .lock() - .map_err(|_| io::Error::other("coordination checkpoint store lock poisoned"))?; - if self.persistence_uncertain.load(Ordering::Acquire) { - return Err(io::Error::other( - "checkpoint persistence is uncertain; reopen before continuing", - )); - } - Ok(guard) - } -} - -fn validate_document(document: &CheckpointDocument) -> io::Result<()> { - if document.schema_version != SCHEMA_VERSION { - return Err(invalid("unsupported coordination checkpoint store schema")); - } - if let Some(scope) = &document.coordinator_scope { - super::multisource_coordinator::validate_spec(scope)?; - } - let mut staged_ids = std::collections::BTreeSet::new(); - for input in &document.staged { - input.validate()?; - if !staged_ids.insert(input.instance_id.canonical()) { - return Err(invalid("duplicate staged summary instance ID")); - } - } - let mut watermark_ids = std::collections::BTreeSet::new(); - for barrier in &document.watermarks { - barrier - .validate() - .map_err(|error| invalid(error.to_string()))?; - let identity = ( - barrier.catalog_generation.schema_version, - barrier.catalog_generation.plan_id, - barrier.catalog_generation.plan_version, - &barrier.catalog_generation.snapshot_sha256, - &barrier.source, - ); - if !watermark_ids.insert(identity) { - return Err(invalid("duplicate source-epoch watermark")); - } - } - let mut published_ids = std::collections::BTreeSet::new(); - for key in &document.published { - key.validate()?; - if !published_ids.insert(key.instance_id.canonical()) { - return Err(invalid("duplicate published summary instance ID")); - } - } - Ok(()) -} - -fn validate_state_reference(reference: &SummaryStateReference) -> io::Result<()> { - if reference.store.trim().is_empty() - || reference.key.trim().is_empty() - || reference.state_schema_version == 0 - || reference - .checksum - .as_deref() - .is_none_or(|checksum| checksum.trim().is_empty()) - { - Err(invalid( - "coordination state reference must be durable and checksummed", - )) - } else { - Ok(()) - } -} - -fn persist_atomically(path: &Path, document: &CheckpointDocument) -> io::Result<()> { - let parent = path.parent().unwrap_or_else(|| Path::new(".")); - fs::create_dir_all(parent)?; - let tmp = path.with_extension("tmp"); - let bytes = serde_json::to_vec(document).map_err(io::Error::other)?; - let mut file = File::create(&tmp)?; - file.write_all(&bytes)?; - file.sync_all()?; - fs::rename(&tmp, path)?; - File::open(parent)?.sync_all()?; - Ok(()) -} - -fn invalid(message: impl Into) -> io::Error { - io::Error::new(io::ErrorKind::InvalidData, message.into()) -} - -#[cfg(test)] -mod tests { - use super::*; - use asap_types::{sds::HalfOpenTimeRange, sds::StoredOutputId, PolicyFingerprint}; - use std::collections::BTreeMap; - - fn generation() -> CatalogGeneration { - CatalogGeneration { - schema_version: 1, - plan_id: 1, - plan_version: 2, - snapshot_sha256: "sha".into(), - } - } - - fn source(epoch: u64) -> SummarySourcePartition { - SummarySourcePartition { - producer_id: "producer".into(), - partition_id: "0".into(), - producer_epoch: epoch, - } - } - - fn coordinates() -> SummaryInstanceCoordinates { - SummaryInstanceCoordinates { - stored_output_id: StoredOutputId::from(PolicyFingerprint(7)), - time_range: HalfOpenTimeRange { - start_ms: 0, - end_ms: 10, - }, - group_values: BTreeMap::from([("job".into(), "api".into())]), - } - } - - fn staged(epoch: u64) -> StagedSummaryInput { - StagedSummaryInput { - catalog_generation: generation(), - dag_id: "dag".into(), - consumer_node_id: "join".into(), - input_node_id: "left".into(), - source: source(epoch), - instance_id: SummaryInstanceId::new(format!("instance-{epoch}")).unwrap(), - coordinates: coordinates(), - input_lineage: vec![epoch as u8], - state_reference: state_reference(format!("state-{epoch}")), - } - } - - fn state_reference(key: String) -> SummaryStateReference { - SummaryStateReference { - store: "summary-store".into(), - key, - state_schema_version: 1, - generation: 1, - sequence: 1, - checksum: Some("sha256:abc".into()), - } - } - - #[test] - fn restart_recovers_staging_and_idempotence() { - let dir = tempfile::tempdir().unwrap(); - let path = dir.path().join("coordination.json"); - let checkpoint_store = SummaryCoordinationCheckpointStore::open(&path).unwrap(); - assert!(checkpoint_store.stage_if_absent(staged(1)).unwrap()); - assert!(!checkpoint_store.stage_if_absent(staged(1)).unwrap()); - drop(checkpoint_store); - let recovered = SummaryCoordinationCheckpointStore::open(&path).unwrap(); - assert_eq!(recovered.staged().unwrap(), vec![staged(1)]); - assert!(!recovered.stage_if_absent(staged(1)).unwrap()); - } - - #[test] - fn epochs_are_distinct_and_instance_id_equivocation_is_rejected() { - let dir = tempfile::tempdir().unwrap(); - let checkpoint_store = - SummaryCoordinationCheckpointStore::open(dir.path().join("checkpoint.json")).unwrap(); - assert!(checkpoint_store.stage_if_absent(staged(1)).unwrap()); - assert!(checkpoint_store.stage_if_absent(staged(2)).unwrap()); - let mut conflicting = staged(1); - conflicting - .coordinates - .group_values - .insert("job".into(), "other".into()); - assert!(checkpoint_store.stage_if_absent(conflicting).is_err()); - } - - #[test] - fn watermark_rejects_regression_but_new_epoch_starts_fresh() { - let dir = tempfile::tempdir().unwrap(); - let checkpoint_store = - SummaryCoordinationCheckpointStore::open(dir.path().join("checkpoint.json")).unwrap(); - let barrier = |epoch, sequence, watermark_ms| SummaryWatermarkBarrier { - catalog_generation: generation(), - source: source(epoch), - sequence, - watermark_ms, - }; - assert!(checkpoint_store - .advance_watermark(barrier(1, 2, 20)) - .unwrap()); - assert!(!checkpoint_store - .advance_watermark(barrier(1, 2, 20)) - .unwrap()); - assert!(checkpoint_store - .advance_watermark(barrier(1, 2, 21)) - .is_err()); - assert!(checkpoint_store - .advance_watermark(barrier(1, 1, 30)) - .is_err()); - assert!(checkpoint_store - .advance_watermark(barrier(2, 1, 5)) - .unwrap()); - } - - #[test] - fn publication_key_is_durable_and_idempotent() { - let dir = tempfile::tempdir().unwrap(); - let path = dir.path().join("checkpoint.json"); - let key = AtomicPublicationKey { - catalog_generation: generation(), - dag_id: "dag".into(), - sink_node_id: "sink".into(), - instance_id: SummaryInstanceId::new("output").unwrap(), - coordinates: coordinates(), - output_lineage: vec![1], - state_reference: state_reference("output-state".into()), - }; - let checkpoint_store = SummaryCoordinationCheckpointStore::open(&path).unwrap(); - assert!(checkpoint_store.publish_if_absent(key.clone()).unwrap()); - drop(checkpoint_store); - let recovered = SummaryCoordinationCheckpointStore::open(&path).unwrap(); - assert!(!recovered.publish_if_absent(key.clone()).unwrap()); - assert!( - SummaryCoordinationCheckpointStore::open(&path).is_err(), - "writer lock remains held" - ); - let mut conflicting = key.clone(); - conflicting.output_lineage = vec![2]; - assert!(recovered.publish_if_absent(conflicting).is_err()); - drop(recovered); - assert!(!SummaryCoordinationCheckpointStore::open(&path) - .unwrap() - .publish_if_absent(key) - .unwrap()); - } - - #[test] - fn corrupt_or_unknown_wire_data_fails_closed() { - let dir = tempfile::tempdir().unwrap(); - let path = dir.path().join("checkpoint.json"); - fs::write(&path, br#"{"schema_version":1,"revision":0,"staged":[],"watermarks":[],"published":[],"unknown":true}"#).unwrap(); - assert!(SummaryCoordinationCheckpointStore::open(path).is_err()); - } - - #[test] - fn duplicate_primary_keys_on_disk_fail_closed() { - let dir = tempfile::tempdir().unwrap(); - let path = dir.path().join("checkpoint.json"); - let duplicate = staged(1); - let document = CheckpointDocument { - schema_version: SCHEMA_VERSION, - revision: 1, - coordinator_scope: None, - staged: vec![duplicate.clone(), duplicate], - watermarks: Vec::new(), - published: Vec::new(), - }; - fs::write(&path, serde_json::to_vec(&document).unwrap()).unwrap(); - assert!(SummaryCoordinationCheckpointStore::open(path).is_err()); - } -} diff --git a/data_plane/src/precompute_engine/mod.rs b/data_plane/src/precompute_engine/mod.rs index 1aaa257e4..2a684d470 100644 --- a/data_plane/src/precompute_engine/mod.rs +++ b/data_plane/src/precompute_engine/mod.rs @@ -1,5 +1,4 @@ pub mod config; -pub mod coordination_checkpoint; mod engine; pub mod erp_observer; pub mod frame_lineage; @@ -7,10 +6,8 @@ pub mod group_key; pub mod ingest_handler; pub mod maintenance_runtime; pub(crate) mod metrics; -pub mod multisource_coordinator; pub mod output_sink; pub mod raw_dag; -pub mod series_buffer; pub mod series_router; pub mod subdag_scheduler; pub mod window_manager; diff --git a/data_plane/src/precompute_engine/multisource_coordinator.rs b/data_plane/src/precompute_engine/multisource_coordinator.rs deleted file mode 100644 index 1afe3a407..000000000 --- a/data_plane/src/precompute_engine/multisource_coordinator.rs +++ /dev/null @@ -1,1104 +0,0 @@ -//! Keyed, watermark-gated staging for multi-source maintenance DAG nodes. - -use super::coordination_checkpoint::{ - AtomicPublicationKey, StagedSummaryInput, SummaryCoordinationCheckpointStore, -}; -use asap_types::sds::{ - CatalogGeneration, HalfOpenTimeRange, SummaryInstanceCoordinates, SummaryInstanceId, - SummarySourcePartition, SummaryStateReference, SummaryWatermarkBarrier, -}; -use asap_types::PolicyFingerprint; -use serde::{Deserialize, Serialize}; -use std::collections::{BTreeMap, BTreeSet}; -use std::io; -use std::sync::{Arc, Mutex}; - -#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)] -#[serde(deny_unknown_fields)] -pub struct LogicalSourcePartition { - pub producer_id: String, - pub partition_id: String, -} - -#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] -#[serde(deny_unknown_fields)] -pub struct CoordinatedInput { - pub input_node_id: String, - pub stored_output_id: asap_types::sds::StoredOutputId, - pub partitions: BTreeSet, -} - -#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] -#[serde(deny_unknown_fields)] -pub struct MultiSourceNodeSpec { - pub catalog_generation: CatalogGeneration, - pub dag_id: String, - pub consumer_node_id: String, - /// The content-addressed installed output binds its source/window contract. - #[serde(default, skip_serializing_if = "Option::is_none")] - pub output_definition: Option, - pub inputs: Vec, - /// Named output grouping. An empty projection represents one global group. - pub output_grouping: Vec, -} - -#[derive(Debug, Clone, PartialEq, Eq)] -pub struct ReadyInputBatch { - pub time_range: HalfOpenTimeRange, - pub group_values: BTreeMap, - /// Inputs are ordered by the spec's input order, then source partition. - pub inputs: Vec, -} - -/// Serializes stage/barrier/readiness transitions around the durable checkpoint store. -/// This is coordination, not operator execution: family-specific joins consume -/// a `ReadyInputBatch` through the typed maintenance operator registry. -pub struct MultiSourceCoordinator { - spec: MultiSourceNodeSpec, - installed_plan: Option>, - checkpoint_store: SummaryCoordinationCheckpointStore, - transition: Mutex<()>, -} - -impl MultiSourceCoordinator { - pub fn new( - spec: MultiSourceNodeSpec, - checkpoint_store: SummaryCoordinationCheckpointStore, - ) -> io::Result { - if spec.output_definition.is_some() { - return Err(invalid( - "installed coordinator requires its authoritative plan", - )); - } - Self::from_scope(spec, checkpoint_store) - } - - fn from_scope( - spec: MultiSourceNodeSpec, - checkpoint_store: SummaryCoordinationCheckpointStore, - ) -> io::Result { - validate_spec(&spec)?; - checkpoint_store.bind_coordinator_scope(&spec)?; - Ok(Self { - spec, - installed_plan: None, - checkpoint_store, - transition: Mutex::new(()), - }) - } - - #[cfg(test)] - /// Bind the complete installed producer roster before accepting any input - /// or barrier. This validates scope; transport authentication and durable - /// source payload publication remain caller obligations. - pub fn for_installed_plan( - spec: MultiSourceNodeSpec, - plan: Arc, - checkpoint_store: SummaryCoordinationCheckpointStore, - ) -> io::Result { - plan.validate() - .map_err(|error| invalid(error.to_string()))?; - if plan.summary_catalog.as_ref() != Some(&spec.catalog_generation) - || plan.envelope.plan_id != spec.catalog_generation.plan_id - || plan.envelope.plan_version != spec.catalog_generation.plan_version - { - return Err(invalid( - "coordinator generation differs from installed plan", - )); - } - let target = spec - .output_definition - .ok_or_else(|| invalid("installed coordinator needs an output definition"))?; - let config = plan - .materializations - .iter() - .find(|config| config.policy_fingerprint() == target.fingerprint()) - .ok_or_else(|| invalid("coordinator output is not installed"))?; - let expected = &config - .derived_input - .as_ref() - .ok_or_else(|| invalid("coordinator output needs derived input"))? - .inputs; - let supplied: BTreeSet<_> = spec - .inputs - .iter() - .map(|input| input.stored_output_id) - .collect(); - if &supplied != expected - || supplied.len() != spec.inputs.len() - || config.partitioning != Some(asap_types::sds::PopulationPartitioning::Grouped) - || spec - .output_grouping - .iter() - .cloned() - .collect::>() - != config.grouping_labels.names().into_iter().collect() - { - return Err(invalid( - "coordinator source set or output grouping differs from installed target", - )); - } - let mut sources = Vec::new(); - for input in &spec.inputs { - let source = plan - .materializations - .iter() - .find(|config| config.policy_fingerprint() == input.stored_output_id.fingerprint()) - .ok_or_else(|| invalid("coordinator input is not installed"))?; - sources.push(source); - let producers: Vec<_> = plan - .producers - .iter() - .filter(|producer| producer.materialization == input.stored_output_id) - .collect(); - if producers.is_empty() - || producers - .iter() - .any(|producer| producer.partition_ids.is_empty()) - { - return Err(invalid( - "coordinator input has no complete authoritative partition roster", - )); - } - let expected_partitions = producers - .into_iter() - .flat_map(|producer| { - producer - .partition_ids - .iter() - .map(move |partition| LogicalSourcePartition { - producer_id: producer.producer_id.clone(), - partition_id: partition.clone(), - }) - }) - .collect::>(); - if input.partitions != expected_partitions { - return Err(invalid( - "coordinator input omits or adds installed producer partitions", - )); - } - } - asap_types::precompute_plan::validated_source_window_cohort(config, &sources) - .map_err(|error| invalid(error.to_string()))?; - let mut coordinator = Self::from_scope(spec, checkpoint_store)?; - coordinator.installed_plan = Some(plan); - for input in coordinator.checkpoint_store.staged()? { - coordinator.validate_input(&input)?; - } - for barrier in coordinator.checkpoint_store.watermarks()? { - if barrier.catalog_generation != coordinator.spec.catalog_generation - || !coordinator - .logical_partitions() - .contains(&logical(&barrier.source)) - { - return Err(invalid("restored barrier is outside installed scope")); - } - for input in &coordinator.spec.inputs { - if input.partitions.contains(&logical(&barrier.source)) { - coordinator - .installed_plan - .as_ref() - .unwrap() - .validate_watermark_scope(input.stored_output_id, &barrier) - .map_err(|error| invalid(error.to_string()))?; - } - } - } - Ok(coordinator) - } - - pub fn stage(&self, input: StagedSummaryInput) -> io::Result { - let _transition = self - .transition - .lock() - .map_err(|_| io::Error::other("multi-source coordinator lock poisoned"))?; - self.validate_input(&input)?; - let staged = self.checkpoint_store.staged()?; - let watermarks = self.checkpoint_store.watermarks()?; - let already_staged = staged - .iter() - .any(|existing| same_input_identity(existing, &input)); - if !already_staged { - if self - .active_epochs(&staged, &watermarks) - .get(&logical(&input.source)) - .is_some_and(|epoch| input.source.producer_epoch < *epoch) - { - return Err(invalid("new input belongs to a superseded producer epoch")); - } - if watermarks.iter().any(|barrier| { - barrier.catalog_generation == input.catalog_generation - && barrier.source == input.source - && barrier.watermark_ms >= input.coordinates.time_range.end_ms - }) { - return Err(invalid( - "new input arrived after its source epoch completed the window", - )); - } - } - self.checkpoint_store.stage_if_absent(input) - } - - pub fn advance_watermark(&self, barrier: SummaryWatermarkBarrier) -> io::Result { - let _transition = self - .transition - .lock() - .map_err(|_| io::Error::other("multi-source coordinator lock poisoned"))?; - if barrier.catalog_generation != self.spec.catalog_generation - || !self - .logical_partitions() - .contains(&logical(&barrier.source)) - { - return Err(invalid( - "watermark does not belong to this maintenance node", - )); - } - barrier - .validate() - .map_err(|error| invalid(error.to_string()))?; - if let Some(plan) = &self.installed_plan { - for input in self - .spec - .inputs - .iter() - .filter(|input| input.partitions.contains(&logical(&barrier.source))) - { - plan.validate_watermark_scope(input.stored_output_id, &barrier) - .map_err(|error| invalid(error.to_string()))?; - } - } - let staged = self.checkpoint_store.staged()?; - let watermarks = self.checkpoint_store.watermarks()?; - if self - .active_epochs(&staged, &watermarks) - .get(&logical(&barrier.source)) - .is_some_and(|epoch| barrier.source.producer_epoch < *epoch) - { - return Err(invalid("watermark belongs to a superseded producer epoch")); - } - self.checkpoint_store.advance_watermark(barrier) - } - - #[cfg(test)] - pub fn ready_batches(&self) -> io::Result> { - let _transition = self - .transition - .lock() - .map_err(|_| io::Error::other("multi-source coordinator lock poisoned"))?; - let staged = self.checkpoint_store.staged()?; - let watermarks = self.checkpoint_store.watermarks()?; - let active_epochs = self.active_epochs(&staged, &watermarks); - let mut buckets = - BTreeMap::<(i64, i64, Vec<(String, String)>), Vec>::new(); - for input in staged.into_iter().filter(|input| { - input.catalog_generation == self.spec.catalog_generation - && input.dag_id == self.spec.dag_id - && input.consumer_node_id == self.spec.consumer_node_id - && active_epochs.get(&logical(&input.source)) == Some(&input.source.producer_epoch) - }) { - self.validate_input(&input)?; - let projected = - project_group(&input.coordinates.group_values, &self.spec.output_grouping)?; - buckets - .entry(( - input.coordinates.time_range.start_ms, - input.coordinates.time_range.end_ms, - projected.into_iter().collect(), - )) - .or_default() - .push(input); - } - - let mut ready = Vec::new(); - for ((start_ms, end_ms, group), inputs) in buckets { - if self.complete(&inputs, &watermarks, &active_epochs, end_ms) { - let mut ordered = Vec::new(); - for requirement in &self.spec.inputs { - let mut matching = inputs - .iter() - .filter(|input| input.input_node_id == requirement.input_node_id) - .cloned() - .collect::>(); - matching.sort_by(|a, b| a.source.cmp(&b.source)); - ordered.extend(matching); - } - ready.push(ReadyInputBatch { - time_range: HalfOpenTimeRange { start_ms, end_ms }, - group_values: group.into_iter().collect(), - inputs: ordered, - }); - } - } - Ok(ready) - } - - fn validate_input(&self, input: &StagedSummaryInput) -> io::Result<()> { - if input.catalog_generation != self.spec.catalog_generation - || input.dag_id != self.spec.dag_id - || input.consumer_node_id != self.spec.consumer_node_id - { - return Err(invalid("staged input belongs to another plan or node")); - } - let requirement = self - .spec - .inputs - .iter() - .find(|requirement| requirement.input_node_id == input.input_node_id) - .ok_or_else(|| invalid("staged input node is not required"))?; - if input.coordinates.stored_output_id != requirement.stored_output_id - || !requirement.partitions.contains(&logical(&input.source)) - { - return Err(invalid( - "staged input source does not match its requirement", - )); - } - if let Some(plan) = &self.installed_plan { - let config = plan - .materializations - .iter() - .find(|config| { - config.policy_fingerprint() == requirement.stored_output_id.fingerprint() - }) - .ok_or_else(|| invalid("staged input definition is not installed"))?; - let window = input.coordinates.time_range; - let width = i64::try_from(config.stored_window_ms()) - .map_err(|_| invalid("source window width overflow"))?; - let slide = config - .slide_interval - .checked_mul(1000) - .filter(|slide| *slide > 0) - .ok_or_else(|| invalid("source window slide is invalid"))?; - if window.end_ms.checked_sub(window.start_ms) != Some(width) - || (i128::from(window.start_ms) - i128::from(config.pane_origin_ms.unwrap_or(0))) - .rem_euclid(i128::from(slide)) - != 0 - { - return Err(invalid( - "staged input window differs from installed source contract", - )); - } - } - project_group(&input.coordinates.group_values, &self.spec.output_grouping)?; - Ok(()) - } - - fn logical_partitions(&self) -> BTreeSet { - self.spec - .inputs - .iter() - .flat_map(|input| input.partitions.iter().cloned()) - .collect() - } - - /// A staged input already observes a new epoch; waiting until its first - /// watermark would allow the previous epoch's barrier to authorize work. - /// Historical staged metadata remains available for audit/idempotent retry. - fn active_epochs( - &self, - staged: &[StagedSummaryInput], - watermarks: &[SummaryWatermarkBarrier], - ) -> BTreeMap { - let partitions = self.logical_partitions(); - let mut active = BTreeMap::::new(); - let sources = staged - .iter() - .filter(|input| input.catalog_generation == self.spec.catalog_generation) - .map(|input| &input.source) - .chain( - watermarks - .iter() - .filter(|barrier| barrier.catalog_generation == self.spec.catalog_generation) - .map(|barrier| &barrier.source), - ); - for source in sources { - let partition = logical(source); - if partitions.contains(&partition) { - let epoch = active.entry(partition).or_default(); - *epoch = (*epoch).max(source.producer_epoch); - } - } - active - } - - fn complete( - &self, - inputs: &[StagedSummaryInput], - watermarks: &[SummaryWatermarkBarrier], - active_epochs: &BTreeMap, - end_ms: i64, - ) -> bool { - self.spec.inputs.iter().all(|requirement| { - requirement.partitions.iter().all(|partition| { - let active_epoch = active_epochs.get(partition).copied(); - active_epoch.is_some_and(|epoch| { - let source = SummarySourcePartition { - producer_id: partition.producer_id.clone(), - partition_id: partition.partition_id.clone(), - producer_epoch: epoch, - }; - watermarks.iter().any(|barrier| { - barrier.catalog_generation == self.spec.catalog_generation - && barrier.source == source - && barrier.watermark_ms >= end_ms - }) && inputs.iter().any(|input| { - input.input_node_id == requirement.input_node_id && input.source == source - }) - }) - }) - }) - } -} - -pub(crate) fn validate_spec(spec: &MultiSourceNodeSpec) -> io::Result<()> { - if spec.dag_id.trim().is_empty() - || spec.consumer_node_id.trim().is_empty() - || spec.inputs.len() < 2 - { - return Err(invalid( - "multi-source spec needs IDs and at least two inputs", - )); - } - let mut nodes = BTreeSet::new(); - for input in &spec.inputs { - if input.input_node_id.trim().is_empty() - || input.partitions.is_empty() - || !nodes.insert(input.input_node_id.as_str()) - || input - .partitions - .iter() - .any(|p| p.producer_id.trim().is_empty() || p.partition_id.trim().is_empty()) - { - return Err(invalid("multi-source input requirement is invalid")); - } - } - let mut grouping = BTreeSet::new(); - if spec - .output_grouping - .iter() - .any(|key| key.trim().is_empty() || !grouping.insert(key)) - { - return Err(invalid( - "output grouping contains an empty or duplicate key", - )); - } - Ok(()) -} - -fn project_group( - values: &BTreeMap, - keys: &[String], -) -> io::Result> { - keys.iter() - .map(|key| { - values - .get(key) - .cloned() - .map(|value| (key.clone(), value)) - .ok_or_else(|| invalid(format!("input is missing grouping label {key}"))) - }) - .collect() -} - -fn logical(source: &SummarySourcePartition) -> LogicalSourcePartition { - LogicalSourcePartition { - producer_id: source.producer_id.clone(), - partition_id: source.partition_id.clone(), - } -} - -fn same_input_identity(a: &StagedSummaryInput, b: &StagedSummaryInput) -> bool { - a.instance_id == b.instance_id -} - -fn invalid(message: impl Into) -> io::Error { - io::Error::new(io::ErrorKind::InvalidData, message.into()) -} - -#[cfg(test)] -mod tests { - use super::*; - use tempfile::tempdir; - - fn generation() -> CatalogGeneration { - CatalogGeneration { - schema_version: 1, - plan_id: 1, - plan_version: 1, - snapshot_sha256: "sha".into(), - } - } - fn partition(id: &str) -> LogicalSourcePartition { - LogicalSourcePartition { - producer_id: "p".into(), - partition_id: id.into(), - } - } - fn spec() -> MultiSourceNodeSpec { - MultiSourceNodeSpec { - catalog_generation: generation(), - dag_id: "dag".into(), - consumer_node_id: "join".into(), - output_definition: None, - inputs: vec![ - CoordinatedInput { - input_node_id: "left".into(), - stored_output_id: PolicyFingerprint(1).into(), - partitions: BTreeSet::from([partition("0")]), - }, - CoordinatedInput { - input_node_id: "right".into(), - stored_output_id: PolicyFingerprint(2).into(), - partitions: BTreeSet::from([partition("1")]), - }, - ], - output_grouping: vec!["job".into()], - } - } - fn input(node: &str, definition: u64, partition_id: &str, epoch: u64) -> StagedSummaryInput { - StagedSummaryInput { - catalog_generation: generation(), - dag_id: "dag".into(), - consumer_node_id: "join".into(), - input_node_id: node.into(), - source: SummarySourcePartition { - producer_id: "p".into(), - partition_id: partition_id.into(), - producer_epoch: epoch, - }, - instance_id: SummaryInstanceId::new(format!("{node}-{epoch}")).unwrap(), - coordinates: SummaryInstanceCoordinates { - stored_output_id: PolicyFingerprint(definition).into(), - time_range: HalfOpenTimeRange { - start_ms: 0, - end_ms: 10, - }, - group_values: BTreeMap::from([ - ("job".into(), "api".into()), - ("instance".into(), node.into()), - ]), - }, - input_lineage: vec![definition as u8, epoch as u8], - state_reference: SummaryStateReference { - store: "summary-store".into(), - key: format!("{node}/{epoch}"), - state_schema_version: 1, - generation: 1, - sequence: 1, - checksum: Some(format!("sha256:{definition}")), - }, - } - } - fn barrier(partition_id: &str, epoch: u64, watermark_ms: i64) -> SummaryWatermarkBarrier { - SummaryWatermarkBarrier { - catalog_generation: generation(), - source: SummarySourcePartition { - producer_id: "p".into(), - partition_id: partition_id.into(), - producer_epoch: epoch, - }, - sequence: 1, - watermark_ms, - } - } - - fn installed_fixture() -> ( - Arc, - MultiSourceNodeSpec, - ) { - let mut wire: serde_json::Value = serde_json::from_str(include_str!( - "../../../docs/examples/asapquery-compatibility-demo-snapshot.json" - )) - .unwrap(); - let mut query = wire["query_workload"]["repeating_queries"][3].clone(); - query["query"] = "quantile(0.9, sum_over_time(m[1m]) + sum_over_time(n[1m]))".into(); - query["demand"]["fixed_interval_at"]["interval"] = 60000.into(); - query["demand"]["fixed_interval_at"]["evaluation_phase"] = 0.into(); - wire["query_workload"]["repeating_queries"] = serde_json::json!([query]); - let snapshot: control_plane::physical::compiler::BackendLocalPlanningInput = - serde_json::from_value(wire).unwrap(); - let mut plan = crate::tests::test_utilities::planning::quoted_snapshot(snapshot, false) - .compile_promql() - .unwrap() - .precompute_plan; - let target = plan - .materializations - .iter() - .find(|config| config.derived_input.is_some()) - .unwrap(); - let mut spec = MultiSourceNodeSpec { - catalog_generation: plan.summary_catalog.clone().unwrap(), - dag_id: "selected-maintenance".into(), - consumer_node_id: "global".into(), - output_definition: Some(target.policy_fingerprint().into()), - inputs: Vec::new(), - output_grouping: Vec::new(), - }; - let definitions = target.derived_input.as_ref().unwrap().inputs.clone(); - plan.producers.clear(); - for (ordinal, definition) in definitions.into_iter().enumerate() { - let producer_id = format!("producer-{ordinal}"); - let partition_ids = BTreeSet::from(["east".into(), "west".into()]); - plan.producers - .push(asap_types::precompute_plan::ProducerContract { - producer_id: producer_id.clone(), - collector_id: producer_id.clone(), - materialization: definition, - schema_id: plan - .schemas - .iter() - .find(|schema| schema.materialization == definition) - .unwrap() - .schema_id - .clone(), - partition_ids: partition_ids.clone(), - }); - spec.inputs.push(CoordinatedInput { - input_node_id: format!("input-{ordinal}"), - stored_output_id: definition, - partitions: partition_ids - .into_iter() - .map(|partition_id| LogicalSourcePartition { - producer_id: producer_id.clone(), - partition_id, - }) - .collect(), - }); - } - plan.validate().unwrap(); - (Arc::new(plan), spec) - } - - fn roster_inputs( - spec: &MultiSourceNodeSpec, - epoch: u64, - start: i64, - ) -> Vec { - spec.inputs - .iter() - .flat_map(|requirement| { - requirement.partitions.iter().map(move |partition| { - let key = format!( - "{}-{}-{epoch}-{start}", - partition.producer_id, partition.partition_id - ); - StagedSummaryInput { - catalog_generation: spec.catalog_generation.clone(), - dag_id: spec.dag_id.clone(), - consumer_node_id: spec.consumer_node_id.clone(), - input_node_id: requirement.input_node_id.clone(), - source: SummarySourcePartition { - producer_id: partition.producer_id.clone(), - partition_id: partition.partition_id.clone(), - producer_epoch: epoch, - }, - instance_id: SummaryInstanceId::new(key.clone()).unwrap(), - coordinates: SummaryInstanceCoordinates { - stored_output_id: requirement.stored_output_id, - time_range: HalfOpenTimeRange { - start_ms: start, - end_ms: start + 60000, - }, - group_values: BTreeMap::new(), - }, - input_lineage: key.as_bytes().to_vec(), - state_reference: SummaryStateReference { - store: "durable-summary-store".into(), - key: key.clone(), - state_schema_version: 1, - generation: 1, - sequence: 1, - checksum: Some(format!("sha256:{key}")), - }, - } - }) - }) - .collect() - } - - fn input_barrier( - input: &StagedSummaryInput, - sequence: u64, - watermark_ms: i64, - ) -> SummaryWatermarkBarrier { - SummaryWatermarkBarrier { - catalog_generation: input.catalog_generation.clone(), - source: input.source.clone(), - sequence, - watermark_ms, - } - } - - #[test] - fn installed_roster_barriers_require_every_producer_partition_and_survive_restart() { - let (plan, spec) = installed_fixture(); - let temp = tempdir().unwrap(); - let path = temp.path().join("coordinator.json"); - let coordinator = MultiSourceCoordinator::for_installed_plan( - spec.clone(), - Arc::clone(&plan), - SummaryCoordinationCheckpointStore::open(&path).unwrap(), - ) - .unwrap(); - let inputs = roster_inputs(&spec, 1, 0); - assert_eq!(inputs.len(), 4); - for input in &inputs { - coordinator.stage(input.clone()).unwrap(); - } - for input in &inputs[..3] { - coordinator - .advance_watermark(input_barrier(input, 1, 60000)) - .unwrap(); - } - assert!(coordinator.ready_batches().unwrap().is_empty()); - coordinator - .advance_watermark(input_barrier(&inputs[3], 1, 60000)) - .unwrap(); - assert_eq!(coordinator.ready_batches().unwrap()[0].inputs.len(), 4); - let before = std::fs::read(&path).unwrap(); - assert!(!coordinator.stage(inputs[0].clone()).unwrap()); - assert!(!coordinator - .advance_watermark(input_barrier(&inputs[0], 1, 60000)) - .unwrap()); - assert_eq!(std::fs::read(&path).unwrap(), before); - for case in 0..7 { - let mut wrong = input_barrier(&inputs[0], 1, 60000); - match case { - 0 => wrong.source.producer_id = "foreign".into(), - 1 => wrong.source.partition_id = "missing".into(), - 2 => wrong.catalog_generation.snapshot_sha256 = "other".into(), - 3 => wrong.source.producer_epoch = 0, - 4 => wrong.sequence = 0, - 5 => wrong.watermark_ms += 1, - _ => wrong.watermark_ms -= 1, - } - assert!(coordinator.advance_watermark(wrong).is_err()); - assert_eq!(std::fs::read(&path).unwrap(), before); - } - let mut changed = inputs[0].clone(); - changed.state_reference.checksum = Some("changed".into()); - assert!(coordinator.stage(changed).is_err()); - let mut wrong_window = roster_inputs(&spec, 1, 1)[0].clone(); - wrong_window.coordinates.time_range.end_ms = 60001; - assert!(coordinator.stage(wrong_window).is_err()); - assert_eq!(std::fs::read(&path).unwrap(), before); - coordinator - .advance_watermark(input_barrier(&inputs[0], 2, 120000)) - .unwrap(); - let progress = std::fs::read(&path).unwrap(); - assert!(coordinator - .advance_watermark(input_barrier(&inputs[0], 1, 60000)) - .is_err()); - assert!(coordinator - .advance_watermark(input_barrier(&inputs[0], 3, 60000)) - .is_err()); - assert_eq!(std::fs::read(&path).unwrap(), progress); - let newer = roster_inputs(&spec, 2, 60000)[0].clone(); - coordinator.stage(newer.clone()).unwrap(); - assert!(coordinator.ready_batches().unwrap().is_empty()); - let epoch_progress = std::fs::read(&path).unwrap(); - assert!(coordinator - .advance_watermark(input_barrier(&inputs[0], 99, 999999)) - .is_err()); - assert_eq!(std::fs::read(&path).unwrap(), epoch_progress); - drop(coordinator); - let reopened = MultiSourceCoordinator::for_installed_plan( - spec.clone(), - Arc::clone(&plan), - SummaryCoordinationCheckpointStore::open(&path).unwrap(), - ) - .unwrap(); - assert!(reopened.ready_batches().unwrap().is_empty()); - assert!(reopened - .advance_watermark(input_barrier(&inputs[0], 100, 999999)) - .is_err()); - assert!(!reopened.stage(newer).unwrap()); - drop(reopened); - let mut reduced_plan = (*plan).clone(); - reduced_plan.producers[0].partition_ids.remove("west"); - let mut reduced_spec = spec.clone(); - reduced_spec.inputs[0] - .partitions - .retain(|partition| partition.partition_id != "west"); - assert!(MultiSourceCoordinator::for_installed_plan( - reduced_spec, - Arc::new(reduced_plan), - SummaryCoordinationCheckpointStore::open(&path).unwrap() - ) - .is_err()); - assert_eq!(std::fs::read(&path).unwrap(), epoch_progress); - } - - #[test] - fn installed_checkpoint_cannot_downgrade_or_restore_foreign_windows() { - let (plan, spec) = installed_fixture(); - let dir = tempdir().unwrap(); - let path = dir.path().join("checkpoint.json"); - let coordinator = MultiSourceCoordinator::for_installed_plan( - spec.clone(), - plan.clone(), - SummaryCoordinationCheckpointStore::open(&path).unwrap(), - ) - .unwrap(); - drop(coordinator); - let before = std::fs::read(&path).unwrap(); - assert!(MultiSourceCoordinator::new( - spec.clone(), - SummaryCoordinationCheckpointStore::open(&path).unwrap() - ) - .is_err()); - assert_eq!(before, std::fs::read(&path).unwrap()); - let store = SummaryCoordinationCheckpointStore::open(&path).unwrap(); - let mut malformed = roster_inputs(&spec, 1, 0).remove(0); - malformed.coordinates.time_range.end_ms -= 1; - store.stage_if_absent(malformed).unwrap(); - drop(store); - let before = std::fs::read(&path).unwrap(); - assert!(MultiSourceCoordinator::for_installed_plan( - spec, - plan, - SummaryCoordinationCheckpointStore::open(&path).unwrap() - ) - .is_err()); - assert_eq!(before, std::fs::read(&path).unwrap()); - } - - #[test] - fn installed_roster_rejects_omissions_and_nonempty_unscoped_checkpoints() { - let (plan, spec) = installed_fixture(); - let temp = tempdir().unwrap(); - let path = temp.path().join("coordinator.json"); - let mut incomplete = spec.clone(); - incomplete.inputs[0].partitions.pop_last(); - assert!(MultiSourceCoordinator::for_installed_plan( - incomplete, - Arc::clone(&plan), - SummaryCoordinationCheckpointStore::open(&path).unwrap() - ) - .is_err()); - assert!(!path.exists()); - let input = roster_inputs(&spec, 1, 0)[0].clone(); - let legacy = SummaryCoordinationCheckpointStore::open(&path).unwrap(); - legacy - .advance_watermark(input_barrier(&input, 1, 60000)) - .unwrap(); - drop(legacy); - let bytes = std::fs::read(&path).unwrap(); - assert!(MultiSourceCoordinator::for_installed_plan( - spec, - plan, - SummaryCoordinationCheckpointStore::open(&path).unwrap() - ) - .is_err()); - assert_eq!(std::fs::read(&path).unwrap(), bytes); - } - - #[test] - fn failed_checkpoint_persistence_stops_barrier_and_readiness_until_reopen() { - let (plan, spec) = installed_fixture(); - let temp = tempdir().unwrap(); - let path = temp.path().join("coordinator.json"); - let coordinator = MultiSourceCoordinator::for_installed_plan( - spec.clone(), - Arc::clone(&plan), - SummaryCoordinationCheckpointStore::open(&path).unwrap(), - ) - .unwrap(); - let input = roster_inputs(&spec, 1, 0)[0].clone(); - let before = std::fs::read(&path).unwrap(); - std::fs::create_dir(path.with_extension("tmp")).unwrap(); - assert!(coordinator - .advance_watermark(input_barrier(&input, 1, 60000)) - .is_err()); - assert!(coordinator.ready_batches().is_err()); - assert!(coordinator.stage(input.clone()).is_err()); - std::fs::remove_dir(path.with_extension("tmp")).unwrap(); - assert!(coordinator - .advance_watermark(input_barrier(&input, 1, 60000)) - .is_err()); - assert_eq!(std::fs::read(&path).unwrap(), before); - drop(coordinator); - let reopened = MultiSourceCoordinator::for_installed_plan( - spec, - plan, - SummaryCoordinationCheckpointStore::open(&path).unwrap(), - ) - .unwrap(); - assert!(reopened.checkpoint_store.watermarks().unwrap().is_empty()); - assert!(reopened - .advance_watermark(input_barrier(&input, 1, 60000)) - .unwrap()); - } - - #[test] - fn waits_for_every_input_and_watermark_then_projects_group() { - let dir = tempdir().unwrap(); - let coordinator = MultiSourceCoordinator::new( - spec(), - SummaryCoordinationCheckpointStore::open(dir.path().join("checkpoint.json")).unwrap(), - ) - .unwrap(); - coordinator.stage(input("left", 1, "0", 1)).unwrap(); - coordinator.advance_watermark(barrier("0", 1, 10)).unwrap(); - assert!(coordinator.ready_batches().unwrap().is_empty()); - coordinator.stage(input("right", 2, "1", 1)).unwrap(); - coordinator.advance_watermark(barrier("1", 1, 9)).unwrap(); - assert!(coordinator.ready_batches().unwrap().is_empty()); - coordinator - .advance_watermark(SummaryWatermarkBarrier { - sequence: 2, - ..barrier("1", 1, 10) - }) - .unwrap(); - let ready = coordinator.ready_batches().unwrap(); - assert_eq!(ready.len(), 1); - assert_eq!( - ready[0].group_values, - BTreeMap::from([("job".into(), "api".into())]) - ); - assert_eq!( - ready[0] - .inputs - .iter() - .map(|i| i.input_node_id.as_str()) - .collect::>(), - vec!["left", "right"] - ); - } - - #[test] - fn restart_preserves_readiness_and_publication_identity() { - let dir = tempdir().unwrap(); - let path = dir.path().join("checkpoint.json"); - let coordinator = MultiSourceCoordinator::new( - spec(), - SummaryCoordinationCheckpointStore::open(&path).unwrap(), - ) - .unwrap(); - coordinator.stage(input("left", 1, "0", 1)).unwrap(); - coordinator.stage(input("right", 2, "1", 1)).unwrap(); - coordinator.advance_watermark(barrier("0", 1, 10)).unwrap(); - coordinator.advance_watermark(barrier("1", 1, 10)).unwrap(); - drop(coordinator); - let recovered = MultiSourceCoordinator::new( - spec(), - SummaryCoordinationCheckpointStore::open(path).unwrap(), - ) - .unwrap(); - assert_eq!(recovered.ready_batches().unwrap().len(), 1); - } - - #[test] - fn late_new_input_is_rejected_but_exact_retry_is_idempotent() { - let dir = tempdir().unwrap(); - let coordinator = MultiSourceCoordinator::new( - spec(), - SummaryCoordinationCheckpointStore::open(dir.path().join("checkpoint.json")).unwrap(), - ) - .unwrap(); - let left = input("left", 1, "0", 1); - coordinator.stage(left.clone()).unwrap(); - coordinator.advance_watermark(barrier("0", 1, 10)).unwrap(); - assert!(!coordinator.stage(left).unwrap()); - let mut late = input("left", 1, "0", 1); - late.instance_id = SummaryInstanceId::new("late").unwrap(); - late.input_lineage = vec![9]; - assert!(coordinator.stage(late).is_err()); - } - - #[test] - fn old_epoch_cannot_complete_new_epoch_input() { - let dir = tempdir().unwrap(); - let coordinator = MultiSourceCoordinator::new( - spec(), - SummaryCoordinationCheckpointStore::open(dir.path().join("checkpoint.json")).unwrap(), - ) - .unwrap(); - coordinator.stage(input("left", 1, "0", 2)).unwrap(); - coordinator.stage(input("right", 2, "1", 2)).unwrap(); - assert!(coordinator.advance_watermark(barrier("0", 1, 10)).is_err()); - assert!(coordinator.advance_watermark(barrier("1", 1, 10)).is_err()); - assert!(coordinator.ready_batches().unwrap().is_empty()); - coordinator.advance_watermark(barrier("0", 2, 10)).unwrap(); - coordinator.advance_watermark(barrier("1", 2, 10)).unwrap(); - assert_eq!(coordinator.ready_batches().unwrap().len(), 1); - } - - #[test] - fn restarted_partitions_wait_for_new_barriers_and_exclude_old_epoch_inputs() { - let dir = tempdir().unwrap(); - let path = dir.path().join("checkpoint.json"); - let coordinator = MultiSourceCoordinator::new( - spec(), - SummaryCoordinationCheckpointStore::open(path.clone()).unwrap(), - ) - .unwrap(); - for (node, definition, partition) in [("left", 1, "0"), ("right", 2, "1")] { - coordinator - .stage(input(node, definition, partition, 1)) - .unwrap(); - coordinator - .advance_watermark(barrier(partition, 1, 10)) - .unwrap(); - } - assert_eq!(coordinator.ready_batches().unwrap()[0].inputs.len(), 2); - for (node, definition, partition) in [("left", 1, "0"), ("right", 2, "1")] { - coordinator - .stage(input(node, definition, partition, 2)) - .unwrap(); - } - // Observing a restarted partition invalidates its old completion proof - // even before the new epoch has emitted its first barrier. - assert!(coordinator.ready_batches().unwrap().is_empty()); - coordinator.advance_watermark(barrier("0", 2, 10)).unwrap(); - assert!(coordinator.ready_batches().unwrap().is_empty()); - coordinator.advance_watermark(barrier("1", 2, 10)).unwrap(); - let ready = coordinator.ready_batches().unwrap(); - assert_eq!(ready.len(), 1); - assert_eq!(ready[0].inputs.len(), 2); - assert!(ready[0] - .inputs - .iter() - .all(|input| input.source.producer_epoch == 2)); - drop(coordinator); - let coordinator = MultiSourceCoordinator::new( - spec(), - SummaryCoordinationCheckpointStore::open(path).unwrap(), - ) - .unwrap(); - assert_eq!(coordinator.ready_batches().unwrap(), ready); - let before = coordinator.checkpoint_store.staged().unwrap(); - // Retain old metadata for replay/audit; do not turn an identical retry - // into a new contribution or erase uncommitted historical inputs. - assert_eq!(before.len(), 4); - assert!(!coordinator.stage(input("left", 1, "0", 1)).unwrap()); - let mut late = input("left", 1, "0", 1); - late.instance_id = SummaryInstanceId::new("late-old-epoch").unwrap(); - late.coordinates.time_range = HalfOpenTimeRange { - start_ms: 20, - end_ms: 30, - }; - assert!(coordinator.stage(late).is_err()); - assert_eq!(coordinator.checkpoint_store.staged().unwrap(), before); - assert_eq!(coordinator.ready_batches().unwrap(), ready); - } - - #[test] - fn differing_projected_groups_do_not_join() { - let dir = tempdir().unwrap(); - let coordinator = MultiSourceCoordinator::new( - spec(), - SummaryCoordinationCheckpointStore::open(dir.path().join("checkpoint.json")).unwrap(), - ) - .unwrap(); - coordinator.stage(input("left", 1, "0", 1)).unwrap(); - let mut right = input("right", 2, "1", 1); - right - .coordinates - .group_values - .insert("job".into(), "worker".into()); - coordinator.stage(right).unwrap(); - coordinator.advance_watermark(barrier("0", 1, 10)).unwrap(); - coordinator.advance_watermark(barrier("1", 1, 10)).unwrap(); - assert!(coordinator.ready_batches().unwrap().is_empty()); - } -} diff --git a/data_plane/src/precompute_engine/precompute_engine_design_doc.md b/data_plane/src/precompute_engine/precompute_engine_design_doc.md index b5dd19891..123d8837a 100644 --- a/data_plane/src/precompute_engine/precompute_engine_design_doc.md +++ b/data_plane/src/precompute_engine/precompute_engine_design_doc.md @@ -247,7 +247,6 @@ struct Worker { **Per-series state:** ```rust struct SeriesState { - buffer: SeriesBuffer, // sorted sample buffer previous_watermark_ms: i64, // last-seen watermark aggregations: Vec, // one per matching config } @@ -375,24 +374,7 @@ from `active_panes`. Remaining panes are read non-destructively via When `pass_raw_samples = true`, the entire aggregation pipeline is bypassed. Each sample is emitted as a `SumAccumulator::with_sum(value)` with point-window bounds `[ts, ts]` and the configured `raw_mode_aggregation_id`. -### 3.5 SeriesBuffer (`series_buffer.rs`) - -Per-series in-memory buffer backed by `BTreeMap`. - -```rust -struct SeriesBuffer { - samples: BTreeMap, // timestamp_ms → value - watermark_ms: i64, // max timestamp ever seen (monotonic) - max_buffer_size: usize, -} -``` - -- Samples are automatically sorted by timestamp. -- Watermark only advances forward (monotonic). -- When the buffer exceeds `max_buffer_size`, the oldest samples are evicted. -- Supports range reads (`read_range`) and destructive drains (`drain_up_to`). - -### 3.6 WindowManager (`window_manager.rs`) +### 3.5 WindowManager (`window_manager.rs`) Handles both tumbling and sliding window semantics. @@ -457,7 +439,7 @@ The worker calls `window_starts_containing(ts)` for each incoming sample and fee the value into the accumulator for every matching window. When `closed_windows()` fires, each closed window's accumulator is extracted and emitted independently. -### 3.7 AccumulatorUpdater (`accumulator_factory.rs`) +### 3.6 AccumulatorUpdater (`accumulator_factory.rs`) Trait-based interface for feeding samples into sketch accumulators: @@ -492,7 +474,7 @@ The factory function `create_accumulator_updater(config)` dispatches on | MultipleSubpopulation | CMS | CmsAccumulatorUpdater | | MultipleSubpopulation | HydraKLL | HydraKllAccumulatorUpdater | -### 3.8 OutputSink (`output_sink.rs`) +### 3.7 OutputSink (`output_sink.rs`) ```rust trait OutputSink: Send + Sync { @@ -1131,7 +1113,7 @@ store with the Kafka consumer path. | `test_late_data_drop` | Sample behind the event watermark with `Drop` policy -> 0 emits and records the action | | `test_late_data_forward_to_store` | Late sample for evicted pane with `ForwardToStore` -> 1 emit as mini-accumulator with correct window bounds and sum | -- **Unit tests -- other modules**: `window_manager.rs` (tumbling/sliding arithmetic, pane enumeration, closure detection), `series_buffer.rs` (ordering, watermark), `accumulator_factory.rs` (updater creation and reset), `series_router.rs` (consistent hash routing), `config.rs` (defaults). +- **Unit tests -- other modules**: `window_manager.rs` (tumbling/sliding arithmetic, pane enumeration, closure detection), `accumulator_factory.rs` (updater creation and reset), `series_router.rs` (consistent hash routing), `config.rs` (defaults). - **E2E coverage**: end-to-end paths now run through the OTLP receiver driving the same `IngestState` (`tests/component_process_e2e.rs`; the runnable multi-node demo @@ -1184,7 +1166,6 @@ A lighter alternative: **periodic pane snapshots** written to disk at each flush | `precompute_engine/config.rs` | `PrecomputeEngineConfig`, `LateDataPolicy` | | `precompute_engine/worker.rs` | Per-shard processing, aggregation, window management | | `precompute_engine/series_router.rs` | Hash-based series → worker routing | -| `precompute_engine/series_buffer.rs` | Per-series BTreeMap sample buffer | | `precompute_engine/window_manager.rs` | Tumbling/sliding window logic | | `precompute_engine/accumulator_factory.rs` | `AccumulatorUpdater` trait + factory | | `precompute_engine/output_sink.rs` | `OutputSink` trait + `StoreOutputSink`, `NoopOutputSink`, `CapturingOutputSink` (testing) | diff --git a/data_plane/src/precompute_engine/series_buffer.rs b/data_plane/src/precompute_engine/series_buffer.rs deleted file mode 100644 index 2d8d8e48d..000000000 --- a/data_plane/src/precompute_engine/series_buffer.rs +++ /dev/null @@ -1,163 +0,0 @@ -use std::collections::BTreeMap; - -/// Per-series sample buffer backed by a `BTreeMap` for automatic -/// ordering by timestamp. Tracks a per-series watermark. -pub struct SeriesBuffer { - /// Samples keyed by timestamp_ms. BTreeMap keeps them sorted. - samples: BTreeMap, - /// High-watermark: the maximum timestamp seen so far for this series. - watermark_ms: i64, - /// Maximum number of samples to retain. When exceeded, oldest are evicted. - max_buffer_size: usize, -} - -impl SeriesBuffer { - pub fn new(max_buffer_size: usize) -> Self { - Self { - samples: BTreeMap::new(), - watermark_ms: i64::MIN, - max_buffer_size, - } - } - - /// Insert a sample. Updates the watermark if `timestamp_ms` is the new max. - /// Returns `true` if the sample was actually inserted (not a duplicate timestamp - /// with the same value). - pub fn insert(&mut self, timestamp_ms: i64, value: f64) -> bool { - if timestamp_ms > self.watermark_ms { - self.watermark_ms = timestamp_ms; - } - self.samples.insert(timestamp_ms, value); - - // Enforce max buffer size by evicting oldest entries - while self.samples.len() > self.max_buffer_size { - self.samples.pop_first(); - } - - true - } - - /// Return current watermark. - pub fn watermark_ms(&self) -> i64 { - self.watermark_ms - } - - #[cfg(test)] - /// Read all samples in `[start_ms, end_ms)` — inclusive start, exclusive end. - /// Returns them in timestamp order. - pub fn read_range(&self, start_ms: i64, end_ms: i64) -> Vec<(i64, f64)> { - self.samples - .range(start_ms..end_ms) - .map(|(&ts, &val)| (ts, val)) - .collect() - } - - #[cfg(test)] - /// Drain (remove and return) all samples with `timestamp_ms < up_to_ms`. - pub fn drain_up_to(&mut self, up_to_ms: i64) -> Vec<(i64, f64)> { - let mut drained = Vec::new(); - // split_off returns everything >= up_to_ms; we keep that part - let remaining = self.samples.split_off(&up_to_ms); - // self.samples now contains everything < up_to_ms - drained.extend(self.samples.iter().map(|(&ts, &val)| (ts, val))); - self.samples = remaining; - drained - } - - /// Number of buffered samples. - pub fn len(&self) -> usize { - self.samples.len() - } - - /// Whether the buffer is empty. - pub fn is_empty(&self) -> bool { - self.samples.is_empty() - } -} - -#[cfg(test)] -mod tests { - use super::*; - - #[test] - fn test_insert_and_watermark() { - let mut buf = SeriesBuffer::new(100); - assert_eq!(buf.watermark_ms(), i64::MIN); - - buf.insert(1000, 1.0); - assert_eq!(buf.watermark_ms(), 1000); - - buf.insert(500, 0.5); // out-of-order - assert_eq!(buf.watermark_ms(), 1000); // watermark should not go back - - buf.insert(2000, 2.0); - assert_eq!(buf.watermark_ms(), 2000); - } - - #[test] - fn test_sorted_order() { - let mut buf = SeriesBuffer::new(100); - buf.insert(3000, 3.0); - buf.insert(1000, 1.0); - buf.insert(2000, 2.0); - - let all = buf.read_range(0, 4000); - assert_eq!(all, vec![(1000, 1.0), (2000, 2.0), (3000, 3.0)]); - } - - #[test] - fn test_read_range() { - let mut buf = SeriesBuffer::new(100); - for t in [1000, 2000, 3000, 4000, 5000] { - buf.insert(t, t as f64); - } - - // [2000, 4000) should return 2000, 3000 - let range = buf.read_range(2000, 4000); - assert_eq!(range, vec![(2000, 2000.0), (3000, 3000.0)]); - } - - #[test] - fn test_drain_up_to() { - let mut buf = SeriesBuffer::new(100); - for t in [1000, 2000, 3000, 4000, 5000] { - buf.insert(t, t as f64); - } - - let drained = buf.drain_up_to(3000); - assert_eq!(drained, vec![(1000, 1000.0), (2000, 2000.0)]); - assert_eq!(buf.len(), 3); // 3000, 4000, 5000 remain - } - - #[test] - fn test_max_buffer_enforcement() { - let mut buf = SeriesBuffer::new(3); - buf.insert(1000, 1.0); - buf.insert(2000, 2.0); - buf.insert(3000, 3.0); - buf.insert(4000, 4.0); // should evict 1000 - assert_eq!(buf.len(), 3); - - let all = buf.read_range(0, 5000); - assert_eq!(all, vec![(2000, 2.0), (3000, 3.0), (4000, 4.0)]); - } - - #[test] - fn test_dedup_by_timestamp() { - let mut buf = SeriesBuffer::new(100); - buf.insert(1000, 1.0); - buf.insert(1000, 2.0); // same timestamp, overwrites - assert_eq!(buf.len(), 1); - - let all = buf.read_range(0, 2000); - assert_eq!(all, vec![(1000, 2.0)]); - } - - #[test] - fn test_empty_operations() { - let buf = SeriesBuffer::new(100); - assert!(buf.is_empty()); - assert_eq!(buf.len(), 0); - assert_eq!(buf.read_range(0, 1000), vec![]); - } -}