From f66d26b713b46af58e252a7f8e37e5eabb044245 Mon Sep 17 00:00:00 2001 From: zz_y Date: Thu, 10 Sep 2026 12:49:55 -0600 Subject: [PATCH] Persist summary coordination decisions --- Cargo.lock | 21 +- data_plane/Cargo.toml | 1 + .../precompute_engine/coordination_journal.rs | 475 ++++++++++++++++++ data_plane/src/precompute_engine/mod.rs | 1 + 4 files changed, 493 insertions(+), 5 deletions(-) create mode 100644 data_plane/src/precompute_engine/coordination_journal.rs diff --git a/Cargo.lock b/Cargo.lock index 2c52490b8..15298edbe 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1206,6 +1206,7 @@ dependencies = [ "dashmap 5.5.3", "flate2", "form_urlencoded", + "fs2", "futures", "hex", "http-body-util", @@ -1803,6 +1804,16 @@ dependencies = [ "percent-encoding", ] +[[package]] +name = "fs2" +version = "0.4.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9564fc758e15025b46aa6643b1b77d047d1a56a1aea6e01002ac0c7026876213" +dependencies = [ + "libc", + "winapi", +] + [[package]] name = "futures" version = "0.3.32" @@ -2261,7 +2272,7 @@ dependencies = [ "libc", "percent-encoding", "pin-project-lite", - "socket2 0.5.10", + "socket2 0.6.3", "tokio", "tower-service", "tracing", @@ -3434,7 +3445,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "be769465445e8c1474e9c5dac2018218498557af32d9ed057325ec9a41ae81bf" dependencies = [ "heck 0.5.0", - "itertools 0.10.5", + "itertools 0.13.0", "log", "multimap", "once_cell", @@ -3454,7 +3465,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "8a56d757972c98b346a9b766e3f02746cde6dd1cd1d1d563472929fdd74bec4d" dependencies = [ "anyhow", - "itertools 0.10.5", + "itertools 0.13.0", "proc-macro2", "quote", "syn 2.0.117", @@ -3552,7 +3563,7 @@ dependencies = [ "quinn-udp", "rustc-hash", "rustls 0.23.40", - "socket2 0.5.10", + "socket2 0.6.3", "thiserror 2.0.18", "tokio", "tracing", @@ -3589,7 +3600,7 @@ dependencies = [ "cfg_aliases", "libc", "once_cell", - "socket2 0.5.10", + "socket2 0.6.3", "tracing", "windows-sys 0.52.0", ] diff --git a/data_plane/Cargo.toml b/data_plane/Cargo.toml index 5df1d0c9a..460cb1993 100644 --- a/data_plane/Cargo.toml +++ b/data_plane/Cargo.toml @@ -123,6 +123,7 @@ asap-precompute-rs = { path = "../../ASAPCollector/asap-precompute-rs" } moka = { version = "0.12", features = ["sync"] } memmap2 = "0.9" crc32fast = "1.4" +fs2 = "0.4" # NOTE: the `asap-gorilla` path-dep + the `s3` (rust-s3) / `lru` # crates were dropped when the superseded GORILLA1 read side # (`gorilla_object_store`: GorillaS3Store + decode_block + the diff --git a/data_plane/src/precompute_engine/coordination_journal.rs b/data_plane/src/precompute_engine/coordination_journal.rs new file mode 100644 index 000000000..1f2cef2a0 --- /dev/null +++ b/data_plane/src/precompute_engine/coordination_journal.rs @@ -0,0 +1,475 @@ +//! 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::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 journal 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 JournalDocument { + schema_version: u32, + revision: u64, + staged: Vec, + watermarks: Vec, + published: Vec, +} + +impl Default for JournalDocument { + fn default() -> Self { + Self { + schema_version: SCHEMA_VERSION, + revision: 0, + staged: Vec::new(), + watermarks: Vec::new(), + published: Vec::new(), + } + } +} + +/// A small crash-safe snapshot journal. Mutations become visible in memory +/// only after the replacement file has been flushed and atomically renamed. +pub struct SummaryCoordinationJournal { + path: PathBuf, + /// Held for the journal lifetime; prevents two processes from replacing + /// the same snapshot concurrently. + _lock_file: File, + document: Mutex, +} + +impl SummaryCoordinationJournal { + 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 journal 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 journal: {error}")))?, + Err(error) if error.kind() == io::ErrorKind::NotFound => JournalDocument::default(), + Err(error) => return Err(error), + }; + validate_document(&document)?; + Ok(Self { + path, + _lock_file: lock_file, + document: Mutex::new(document), + }) + } + + 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) + }) + } + + 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()) + } + + fn mutate( + &self, + update: impl FnOnce(&mut JournalDocument) -> io::Result, + ) -> io::Result { + let mut guard = self.lock()?; + let mut next = guard.clone(); + let result = update(&mut next)?; + next.revision = next + .revision + .checked_add(1) + .ok_or_else(|| invalid("coordination journal revision overflow"))?; + persist_atomically(&self.path, &next)?; + *guard = next; + Ok(result) + } + + fn lock(&self) -> io::Result> { + self.document + .lock() + .map_err(|_| io::Error::other("coordination journal lock poisoned")) + } +} + +fn validate_document(document: &JournalDocument) -> io::Result<()> { + if document.schema_version != SCHEMA_VERSION { + return Err(invalid("unsupported coordination journal schema")); + } + 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: &JournalDocument) -> 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::SummaryDefinitionId, 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 { + summary_definition_id: SummaryDefinitionId(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 journal = SummaryCoordinationJournal::open(&path).unwrap(); + assert!(journal.stage_if_absent(staged(1)).unwrap()); + assert!(!journal.stage_if_absent(staged(1)).unwrap()); + drop(journal); + let recovered = SummaryCoordinationJournal::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 journal = SummaryCoordinationJournal::open(dir.path().join("journal.json")).unwrap(); + assert!(journal.stage_if_absent(staged(1)).unwrap()); + assert!(journal.stage_if_absent(staged(2)).unwrap()); + let mut conflicting = staged(1); + conflicting + .coordinates + .group_values + .insert("job".into(), "other".into()); + assert!(journal.stage_if_absent(conflicting).is_err()); + } + + #[test] + fn watermark_rejects_regression_but_new_epoch_starts_fresh() { + let dir = tempfile::tempdir().unwrap(); + let journal = SummaryCoordinationJournal::open(dir.path().join("journal.json")).unwrap(); + let barrier = |epoch, sequence, watermark_ms| SummaryWatermarkBarrier { + catalog_generation: generation(), + source: source(epoch), + sequence, + watermark_ms, + }; + assert!(journal.advance_watermark(barrier(1, 2, 20)).unwrap()); + assert!(!journal.advance_watermark(barrier(1, 2, 20)).unwrap()); + assert!(journal.advance_watermark(barrier(1, 2, 21)).is_err()); + assert!(journal.advance_watermark(barrier(1, 1, 30)).is_err()); + assert!(journal.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("journal.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 journal = SummaryCoordinationJournal::open(&path).unwrap(); + assert!(journal.publish_if_absent(key.clone()).unwrap()); + drop(journal); + let recovered = SummaryCoordinationJournal::open(&path).unwrap(); + assert!(!recovered.publish_if_absent(key.clone()).unwrap()); + assert!( + SummaryCoordinationJournal::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!(!SummaryCoordinationJournal::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("journal.json"); + fs::write(&path, br#"{"schema_version":1,"revision":0,"staged":[],"watermarks":[],"published":[],"unknown":true}"#).unwrap(); + assert!(SummaryCoordinationJournal::open(path).is_err()); + } + + #[test] + fn duplicate_primary_keys_on_disk_fail_closed() { + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join("journal.json"); + let duplicate = staged(1); + let document = JournalDocument { + schema_version: SCHEMA_VERSION, + revision: 1, + staged: vec![duplicate.clone(), duplicate], + watermarks: Vec::new(), + published: Vec::new(), + }; + fs::write(&path, serde_json::to_vec(&document).unwrap()).unwrap(); + assert!(SummaryCoordinationJournal::open(path).is_err()); + } +} diff --git a/data_plane/src/precompute_engine/mod.rs b/data_plane/src/precompute_engine/mod.rs index fee2f9b90..b86cde155 100644 --- a/data_plane/src/precompute_engine/mod.rs +++ b/data_plane/src/precompute_engine/mod.rs @@ -1,5 +1,6 @@ pub mod accumulator_factory; pub mod config; +pub mod coordination_journal; mod engine; pub mod frame_lineage; pub mod group_key;