From c6766e451f9e676a9d26fb08ca3e81b515a2e416 Mon Sep 17 00:00:00 2001 From: David Spencer <1526975+DecisionNerd@users.noreply.github.com> Date: Sun, 20 Sep 2026 03:32:55 +0000 Subject: [PATCH 1/3] fix(storage): admit adaptive partitions against recorded memory budgets (#1439) --- .../src/construction_determinism_tests.rs | 1 + .../src/construction_lifecycle_tests.rs | 1 + .../src/graph_construction.rs | 37 +++- .../src/graph_construction/partition.rs | 56 +++++- .../src/graph_construction/partition/tests.rs | 62 ++++++ .../graph_construction/partition_memory.rs | 177 ++++++++++++++++++ .../graph_construction/partition_shaping.rs | 107 ++++++++++- .../partition_shaping/tests.rs | 116 ++++++++++++ .../src/graph_construction/recovery.rs | 25 +++ .../src/graph_construction/shape.rs | 26 ++- .../src/graph_construction/tests.rs | 68 ++++++- docs/book/architecture/resumable-import.md | 32 ++++ 12 files changed, 682 insertions(+), 26 deletions(-) create mode 100644 crates/graphforge-storage/src/graph_construction/partition_memory.rs diff --git a/crates/graphforge-storage/src/construction_determinism_tests.rs b/crates/graphforge-storage/src/construction_determinism_tests.rs index 409f74f31..0771cda8f 100644 --- a/crates/graphforge-storage/src/construction_determinism_tests.rs +++ b/crates/graphforge-storage/src/construction_determinism_tests.rs @@ -508,6 +508,7 @@ mod determinism { max_batch_rows: 8_192, max_run_records: 262_144, partition_count, + target_partition_records: 16, ..GraphConstructionBudgets::default() }; let (fingerprint, layout) = diff --git a/crates/graphforge-storage/src/construction_lifecycle_tests.rs b/crates/graphforge-storage/src/construction_lifecycle_tests.rs index 96da3b909..52899e389 100644 --- a/crates/graphforge-storage/src/construction_lifecycle_tests.rs +++ b/crates/graphforge-storage/src/construction_lifecycle_tests.rs @@ -1552,6 +1552,7 @@ mod lifecycle_budget { Some("endpoints.run"), None, window_rows, + partition::default_materialization_bytes(), &mut || false, &mut evidence, ) diff --git a/crates/graphforge-storage/src/graph_construction.rs b/crates/graphforge-storage/src/graph_construction.rs index 1a4bcf34f..f32ae7d35 100644 --- a/crates/graphforge-storage/src/graph_construction.rs +++ b/crates/graphforge-storage/src/graph_construction.rs @@ -70,6 +70,7 @@ use catalog::{ load_parent_runtime_catalog_from_compact, }; mod partition_load; +mod partition_memory; mod partition_records; pub(crate) mod partition_shaping; mod supersession; @@ -473,12 +474,26 @@ pub struct GraphConstructionBudgets { pub max_catalog_decoded_bytes: usize, /// Maximum UTF-8 identifier bytes retained by the complete runtime catalog. pub max_catalog_identifier_bytes: usize, - /// Range partitions shaping cuts the identity key space into. + /// Maximum range partitions shaping cuts the identity key space into. /// /// This is a recorded format parameter, never derived from the machine. It /// must not be tied to `available_parallelism()` or to a thread count: the /// same logical input has to stage identically on hosts of different sizes. pub partition_count: u32, + /// Desired identities per sampled range; a planning target, not a byte bound. + #[serde( + default = "partition::default_target_records", + skip_serializing_if = "partition::is_default_target_records" + )] + pub target_partition_records: u64, + /// Maximum accounted materialization bytes per partition load. Fixed-family + /// workers admit at most two loads; Arrow rows are loaded serially. This is + /// separate from staging windows, writer buffers and whole-process memory. + #[serde( + default = "partition::default_materialization_bytes", + skip_serializing_if = "partition::is_default_materialization_bytes" + )] + pub max_partition_bytes: u64, } impl Default for GraphConstructionBudgets { @@ -495,6 +510,8 @@ impl Default for GraphConstructionBudgets { max_catalog_decoded_bytes: 256 << 20, max_catalog_identifier_bytes: 64 << 20, partition_count: partition::DEFAULT_PARTITION_COUNT, + target_partition_records: partition::default_target_records(), + max_partition_bytes: partition::default_materialization_bytes(), } } } @@ -511,6 +528,8 @@ impl GraphConstructionBudgets { || self.max_catalog_entries == 0 || self.max_catalog_decoded_bytes == 0 || self.max_catalog_identifier_bytes == 0 + || self.target_partition_records == 0 + || self.max_partition_bytes == 0 || self.partition_count == 0 || self.partition_count > partition::MAX_PARTITION_COUNT { @@ -1288,6 +1307,22 @@ impl GraphConstructionSession { Err(error) if error.kind() == std::io::ErrorKind::NotFound => None, Err(error) => return Err(storage(error)), }; + // The former public default recorded a 256-partition maximum. Opening + // that exact legacy default with today's default retains its authority; + // explicit non-default changes still fail validate_checkpoint below. + let legacy_defaults = GraphConstructionBudgets { + partition_count: 256, + ..GraphConstructionBudgets::default() + }; + let budgets = if budgets == GraphConstructionBudgets::default() + && recovered_checkpoint + .as_ref() + .is_some_and(|checkpoint| checkpoint.budgets == legacy_defaults) + { + legacy_defaults + } else { + budgets + }; if let Some(checkpoint) = recovered_checkpoint.as_mut() && (DetailCodec::from_version(checkpoint.format_version).is_err() || checkpoint.operation_uuid != operation_uuid diff --git a/crates/graphforge-storage/src/graph_construction/partition.rs b/crates/graphforge-storage/src/graph_construction/partition.rs index acd2491ba..c081bd8fb 100644 --- a/crates/graphforge-storage/src/graph_construction/partition.rs +++ b/crates/graphforge-storage/src/graph_construction/partition.rs @@ -33,16 +33,42 @@ use super::storage; use graphforge_core::GfError; -/// Number of range partitions requested by a construction session. +/// Maximum range partitions requested by a construction session. /// /// This is a recorded format parameter, never derived from the machine: it must /// not depend on `available_parallelism()` or on thread count, or the same /// logical input would stage differently on differently-sized hosts. -pub(crate) const DEFAULT_PARTITION_COUNT: u32 = 256; +pub(crate) const DEFAULT_PARTITION_COUNT: u32 = 4_096; /// Upper bound on the recorded partition count. pub(crate) const MAX_PARTITION_COUNT: u32 = 4_096; +/// Deterministic planning target; not a byte bound or an optimality claim. +/// Admission checks actual bytes independently; skew can exceed this target. +pub(super) const fn default_target_records() -> u64 { + 16_384 +} +pub(super) fn is_default_target_records(value: &u64) -> bool { + *value == default_target_records() +} +pub(super) const fn default_materialization_bytes() -> u64 { + 256 << 20 +} +pub(super) fn is_default_materialization_bytes(value: &u64) -> bool { + *value == default_materialization_bytes() +} + +/// Fail before allocating the retained partition or its sort/reorder buffers. +pub(super) fn admit_materialization(bytes: Option, limit: u64) -> Result { + let bytes = bytes.ok_or_else(|| storage("partition materialization byte count overflows"))?; + if bytes > limit { + return Err(storage(format!( + "partition materialization requires {bytes} bytes, exceeds recorded budget {limit}" + ))); + } + Ok(bytes) +} + /// Sample points drawn per requested partition. Oversampling the quantile /// estimate is what keeps the partitions close to balanced; 64 points per /// partition is the usual sample-sort factor. @@ -62,12 +88,9 @@ const BALANCE_MIN_MEAN_ROWS: u64 = 16; /// Never cut more partitions than the balance check can meaningfully assert on. /// -/// Each partition costs a durable spill artifact, with its own barrier and -/// writer receipt, in every family. Cutting the full recorded count regardless -/// of input size makes that cost a constant: a two-thousand-row graph would pay -/// the same durability price as a sixty-seven-million-row one, and shaping's -/// barrier count would stop tracking the work it protects. Bounding the cut by -/// the recorded record count keeps the cost proportional. +/// Each partition retains spill names and writer state in every family. +/// Bounding the cut by the recorded record count avoids charging tiny inputs +/// for the maximum cut. Directory synchronization is batched across spills. /// /// This does not make the partition count machine-dependent, which is the thing /// R1 forbids. The bound is a pure function of the staged record count recorded @@ -216,6 +239,23 @@ pub(crate) struct IdentitySampler { } impl IdentitySampler { + /// Adapt within the recorded maximum, retaining the small-input balance floor. + pub(crate) fn with_target( + partition_count: u32, + total_records: u64, + target: u64, + ) -> Result { + validate_partition_count(partition_count)?; + if target == 0 { + return Err(storage("partition target must be positive")); + } + let desired = total_records + .div_ceil(target) + .max(256) + .min(u64::from(partition_count)); + Self::new(u32::try_from(desired).map_err(storage)?, total_records) + } + /// Build a sampler for `total_records` staged identities. /// /// # Errors diff --git a/crates/graphforge-storage/src/graph_construction/partition/tests.rs b/crates/graphforge-storage/src/graph_construction/partition/tests.rs index a907eae7a..328724252 100644 --- a/crates/graphforge-storage/src/graph_construction/partition/tests.rs +++ b/crates/graphforge-storage/src/graph_construction/partition/tests.rs @@ -236,3 +236,65 @@ fn every_partition_is_reachable_for_a_well_sampled_ingest() { balance.rows() ); } + +#[test] +fn adaptive_cut_grows_from_recorded_data_and_stops_at_recorded_ceiling() { + for (records, expected) in [ + (0, 1), + (32, 2), + (1 << 20, 256), + (1 << 23, 512), + (1 << 24, 1024), + (1 << 26, 4096), + (1 << 30, 4096), + ] { + let sampler = + IdentitySampler::with_target(4096, records, default_target_records()).unwrap(); + assert_eq!(sampler.cut(), expected); + assert_eq!( + sampler.positions().collect::>(), + IdentitySampler::with_target(4096, records, default_target_records()) + .unwrap() + .positions() + .collect::>() + ); + } + assert_eq!( + IdentitySampler::with_target(64, 1 << 30, 16).unwrap().cut(), + 64 + ); + assert!(IdentitySampler::with_target(4096, 100, 0).is_err()); +} + +#[test] +fn materialization_budget_refuses_overflow_and_one_byte_over_limit() { + assert_eq!(admit_materialization(Some(1024), 1024).unwrap(), 1024); + assert!(admit_materialization(Some(1025), 1024).is_err()); + assert!(admit_materialization(None, u64::MAX).is_err()); +} + +#[test] +fn legacy_budget_serialization_is_unchanged_and_new_limits_round_trip() { + let legacy = super::super::GraphConstructionBudgets { + partition_count: 256, + ..Default::default() + }; + let encoded = serde_json::to_value(legacy).unwrap(); + assert!(encoded.get("max_partition_bytes").is_none()); + assert!(encoded.get("target_partition_records").is_none()); + assert_eq!( + serde_json::from_value::(encoded).unwrap(), + legacy + ); + let custom = super::super::GraphConstructionBudgets { + max_partition_bytes: 12345, + target_partition_records: 123, + ..legacy + }; + let encoded = serde_json::to_value(custom).unwrap(); + assert_eq!(encoded["max_partition_bytes"], 12345); + assert_eq!( + serde_json::from_value::(encoded).unwrap(), + custom + ); +} diff --git a/crates/graphforge-storage/src/graph_construction/partition_memory.rs b/crates/graphforge-storage/src/graph_construction/partition_memory.rs new file mode 100644 index 000000000..9b0d3f511 --- /dev/null +++ b/crates/graphforge-storage/src/graph_construction/partition_memory.rs @@ -0,0 +1,177 @@ +//! Conservative admission for Arrow partition buffers and reorder scratch. +//! This is separate from allocator metadata, routing and process RSS. +use super::{partition::admit_materialization, storage}; +use arrow::{array::ArrayData, datatypes::DataType, record_batch::RecordBatch}; +use graphforge_core::GfError; + +#[derive(Clone, Copy, Default)] +pub(super) struct RowReservation { + buffers: u64, + elements: u64, + rows: u64, + batches: u64, +} + +fn add(a: u64, b: u64) -> Result { + a.checked_add(b) + .ok_or_else(|| storage("partition memory accounting overflows")) +} +fn mul(a: u64, b: u64) -> Result { + a.checked_mul(b) + .ok_or_else(|| storage("partition memory accounting overflows")) +} +fn aligned(bytes: u64) -> Result { + mul(add(bytes, 63)? / 64, 64) +} + +// Capacity floors matter even for empty nested lists: MutableArrayData reserves +// child buffers using its parent's requested capacity. Full child data is a +// conservative charge for a sliced list; fixed/text leaves use slice lengths. +fn array_reservation(data: &ArrayData, floor: u64) -> Result<(u64, u64), GfError> { + let n = (data.len() as u64).max(floor); + let mut bytes = aligned(add(n, 7)? / 8)?; + let mut elements = n; + let mut child_floor = n; + let own = if let Some(width) = data.data_type().primitive_width() { + mul(n, width as u64)? + } else { + match data.data_type() { + DataType::Null => 0, + DataType::Boolean => add(n, 7)? / 8, + DataType::FixedSizeBinary(width) => mul(n, u64::try_from(*width).map_err(storage)?)?, + DataType::Utf8 | DataType::Binary | DataType::LargeUtf8 | DataType::LargeBinary => { + let width = if matches!( + data.data_type(), + DataType::LargeUtf8 | DataType::LargeBinary + ) { + 8 + } else { + 4 + }; + // Slice memory includes selected values, but not the terminal + // offset. Additional capacity covers generic byte preallocation. + add( + data.get_slice_memory_size().map_err(storage)? as u64, + add(mul(add(n, 1)?, width)?, n)?, + )? + } + DataType::List(_) | DataType::Map(_, _) => mul(add(n, 1)?, 4)?, + DataType::LargeList(_) => mul(add(n, 1)?, 8)?, + DataType::Struct(_) => 0, + DataType::FixedSizeList(_, width) => { + child_floor = mul(n, u64::try_from(*width).map_err(storage)?)?; + 0 + } + other => { + return Err(storage(format!( + "partition memory accounting does not support {other:?}" + ))); + } + } + }; + bytes = add(bytes, aligned(own)?)?; + for child in data.child_data() { + let (child_bytes, child_elements) = array_reservation(child, child_floor)?; + bytes = add(bytes, child_bytes)?; + elements = add(elements, child_elements)?; + } + Ok((bytes, elements)) +} + +impl RowReservation { + pub(super) fn with_batch(self, batch: &RecordBatch) -> Result { + let mut next = self; + for column in batch.columns() { + let (bytes, elements) = array_reservation(&column.to_data(), 0)?; + next.buffers = add(next.buffers, bytes)?; + next.elements = add(next.elements, elements)?; + } + next.rows = add(next.rows, batch.num_rows() as u64)?; + next.batches = add(next.batches, 1)?; + Ok(next) + } + + pub(super) fn admit(self, spill_bytes: u64, limit: u64) -> Result<(), GfError> { + // Retained IPC bodies and decoded buffers, concat/take output, allocator + // growth/alignment, UUID keys and permutation/nested take indexes. + let bytes = add(mul(spill_bytes, 2)?, mul(self.buffers, 6)?)?; + let bytes = add(bytes, mul(self.rows, 20)?)?; + let bytes = add(bytes, mul(self.elements, 8)?)?; + let bytes = add(bytes, mul(self.batches, 64)?)?; + admit_materialization(Some(bytes), limit)?; + Ok(()) + } +} + +#[cfg(test)] +mod tests { + use super::*; + use arrow::array::{Array, ArrayRef, FixedSizeBinaryArray, ListArray, StringArray}; + use arrow::buffer::{OffsetBuffer, ScalarBuffer}; + use arrow::datatypes::{Field, Schema}; + use std::sync::Arc; + + fn batch(array: ArrayRef) -> RecordBatch { + RecordBatch::try_new( + Arc::new(Schema::new(vec![Field::new( + "value", + array.data_type().clone(), + true, + )])), + vec![array], + ) + .unwrap() + } + + #[test] + fn sliced_strings_charge_selected_values_not_entire_backing_buffer() { + let array = StringArray::from(vec!["x".repeat(1 << 20), "small".to_owned()]); + let selected = batch(Arc::new(array.slice(1, 1))); + let reservation = RowReservation::default().with_batch(&selected).unwrap(); + reservation.admit(0, 4096).unwrap(); + assert!( + RowReservation::default() + .with_batch(&batch(Arc::new(array))) + .unwrap() + .admit(0, 4096) + .is_err() + ); + } + + #[test] + fn empty_nested_lists_charge_wide_child_capacity_and_tiny_batches_accumulate() { + let values = FixedSizeBinaryArray::new( + 4096, + ScalarBuffer::::from(Vec::new()).into_inner(), + None, + ); + let field = Arc::new(Field::new("item", DataType::FixedSizeBinary(4096), true)); + let list = ListArray::new( + field, + OffsetBuffer::new(ScalarBuffer::from(vec![0_i32, 0, 0])), + Arc::new(values), + None, + ); + let nested = batch(Arc::new(list)); + let reservation = RowReservation::default().with_batch(&nested).unwrap(); + assert!( + reservation.buffers >= 2 * 4096, + "empty children still inherit parent capacity" + ); + assert!(reservation.admit(0, 8192).is_err()); + let tiny = batch(Arc::new(StringArray::from(vec!["a"]))); + let mut accumulated = RowReservation::default(); + for _ in 0..100 { + accumulated = accumulated.with_batch(&tiny).unwrap(); + } + assert!(accumulated.admit(0, 4096).is_err()); + assert!( + RowReservation { + buffers: u64::MAX, + ..Default::default() + } + .admit(0, u64::MAX) + .is_err() + ); + } +} diff --git a/crates/graphforge-storage/src/graph_construction/partition_shaping.rs b/crates/graphforge-storage/src/graph_construction/partition_shaping.rs index 6bce6b60c..feec965f6 100644 --- a/crates/graphforge-storage/src/graph_construction/partition_shaping.rs +++ b/crates/graphforge-storage/src/graph_construction/partition_shaping.rs @@ -398,6 +398,7 @@ pub(super) struct FixedRangePartitioner<'a, const N: usize> { /// Concurrent partition loads while finishing; see /// [`super::partition_load::consume_in_partition_order`] for the bound. load_workers: NonZeroUsize, + max_partition_bytes: u64, } impl Drop for FixedRangePartitioner<'_, N> { @@ -434,9 +435,15 @@ impl<'a, const N: usize> FixedRangePartitioner<'a, N> { balance: PartitionBalance::new(partitions), records: 0, load_workers: PARTITION_LOAD_WORKERS, + max_partition_bytes: super::partition::default_materialization_bytes(), }) } + pub(super) fn with_materialization_limit(mut self, bytes: u64) -> Self { + self.max_partition_bytes = bytes; + self + } + /// Force the finish-time worker count. Scheduling only: the tests hold the /// evidence and output bytes equal across every value. #[cfg(test)] @@ -619,7 +626,14 @@ impl<'a, const N: usize> FixedRangePartitioner<'a, N> { let codec = self.codec; let load = |job: usize, stop: &AtomicBool| { let (_, name, expected) = &jobs[job]; - load_fixed_partition::(root, name, *expected, codec, stop) + load_fixed_partition::( + root, + name, + *expected, + codec, + self.max_partition_bytes, + stop, + ) }; let consume = |job: usize, (records, counters): (PartitionRecords, PartitionLoadCounters)| { @@ -830,6 +844,7 @@ fn load_fixed_partition( name: &str, expected_records: Option, codec: Option, + max_partition_bytes: u64, stop: &AtomicBool, ) -> Result<(PartitionRecords, PartitionLoadCounters), GfError> { let (mut reader, counter, spill_bytes) = open_fixed_reader(root, name)?; @@ -838,11 +853,50 @@ fn load_fixed_partition( ..PartitionLoadCounters::default() }; let loaded = (|| -> Result, GfError> { - let mut records = PartitionRecords::new(codec, expected_records, spill_bytes)?; + if let Some(codec) = codec { + codec.validate_size(N, 0, 0).map_err(super::storage)?; + } + let count = expected_records.unwrap_or( + spill_bytes + / if codec.is_some() { + (N - 255) as u64 + } else { + N as u64 + }, + ); + let retained = if codec.is_some() { + count + .checked_mul(std::mem::size_of::() as u64) + .and_then(|offsets| spill_bytes.checked_add(offsets)) + } else { + count.checked_mul(N as u64) + }; + super::partition::admit_materialization(retained, max_partition_bytes)?; + let mut records = PartitionRecords::new(codec, Some(count), spill_bytes)?; + let mut admitted_wire_bytes = 0_u64; while let Some(record) = read_run_record::(&mut reader, codec)? { + if records.len() as u64 >= count { + return Err(super::storage("partition exceeds admitted record count")); + } + let wire_bytes = if codec.is_some() { + N - 255 + usize::from(record[N - 256]) + } else { + N + }; + admitted_wire_bytes = admitted_wire_bytes + .checked_add(wire_bytes as u64) + .ok_or_else(|| super::storage("partition wire byte count overflows"))?; + if admitted_wire_bytes > spill_bytes { + return Err(super::storage("partition exceeds admitted wire bytes")); + } records.push(record); abandon_if_stopped(records.len(), stop)?; } + if expected_records.is_some_and(|expected| records.len() as u64 != expected) { + return Err(super::storage( + "partition differs from admitted record count", + )); + } Ok(records) })(); let released = reader @@ -874,6 +928,8 @@ pub(super) struct RowRangePartitioner<'a> { sealed: Vec>, balance: PartitionBalance, rows: u64, + reservations: Vec, + max_partition_bytes: u64, } struct RowSpill { @@ -926,9 +982,16 @@ impl<'a> RowRangePartitioner<'a> { sealed: (0..partitions).map(|_| None).collect(), balance: PartitionBalance::new(partitions), rows: 0, + reservations: vec![super::partition_memory::RowReservation::default(); partitions], + max_partition_bytes: super::partition::default_materialization_bytes(), }) } + pub(super) fn with_materialization_limit(mut self, bytes: u64) -> Self { + self.max_partition_bytes = bytes; + self + } + /// Route every row of one staged chunk Parquet into its partition. /// /// Staged chunks are already UUID-sorted and the partition function is @@ -1013,6 +1076,14 @@ impl<'a> RowRangePartitioner<'a> { partition: usize, evidence: &mut GraphConstructionEvidence, ) -> Result<(), GfError> { + let slice = batch.slice(offset, length); + let reservation = self + .reservations + .get_mut(partition) + .ok_or_else(|| super::storage("row reservation partition is out of range"))?; + let next = reservation.with_batch(&slice)?; + next.admit(0, self.max_partition_bytes)?; + *reservation = next; let slot = self .spills .get_mut(partition) @@ -1038,10 +1109,7 @@ impl<'a> RowRangePartitioner<'a> { let spill = slot .as_mut() .ok_or_else(|| super::storage("row partition spill is absent"))?; - spill - .writer - .write(&batch.slice(offset, length)) - .map_err(super::storage)?; + spill.writer.write(&slice).map_err(super::storage)?; for _ in 0..length { self.balance.record(partition)?; } @@ -1140,7 +1208,7 @@ impl<'a> RowRangePartitioner<'a> { continue; }; reject_cancelled(cancelled)?; - let sorted = self.load_partition(&name, &schema, evidence)?; + let sorted = self.load_partition(partition, &name, &schema, evidence)?; let uuids = key_column(&sorted)?; for row in 0..sorted.num_rows() { let uuid = uuid_value(uuids, row)?; @@ -1247,10 +1315,29 @@ impl<'a> RowRangePartitioner<'a> { /// Read one row partition spill into memory and sort it by UUID. fn load_partition( &self, + partition: usize, name: &str, schema: &SchemaRef, evidence: &mut GraphConstructionEvidence, ) -> Result { + let receipt = self + .sealed + .get(partition) + .and_then(Option::as_ref) + .ok_or_else(|| super::storage("row partition receipt is absent"))?; + self.reservations[partition].admit(receipt.bytes, self.max_partition_bytes)?; + // Authenticate with bounded buffers BEFORE Arrow reads bodyLength and + // allocates an IPC body. Post-decode checks cannot protect that step. + let work = super::recovery::authenticate_row_spill(self.root, receipt)?; + account_cache_release(work.cache_release, evidence)?; + evidence.merge_read_bytes = evidence + .merge_read_bytes + .checked_add(work.bytes) + .ok_or_else(|| super::storage("partition read bytes overflow"))?; + evidence.merge_read_operations = evidence + .merge_read_operations + .checked_add(work.operations) + .ok_or_else(|| super::storage("partition read operations overflow"))?; let file = self .root .open_child_file(OsStr::new(name)) @@ -1270,8 +1357,12 @@ impl<'a> RowRangePartitioner<'a> { let stream = arrow::ipc::reader::StreamReader::try_new(&mut reader, None) .map_err(super::storage)?; let mut batches = Vec::new(); + let mut decoded = super::partition_memory::RowReservation::default(); for batch in stream { - batches.push(batch.map_err(super::storage)?); + let batch = batch.map_err(super::storage)?; + decoded = decoded.with_batch(&batch)?; + decoded.admit(receipt.bytes, self.max_partition_bytes)?; + batches.push(batch); } let combined = concat_batches(schema, &batches).map_err(super::storage)?; drop(batches); diff --git a/crates/graphforge-storage/src/graph_construction/partition_shaping/tests.rs b/crates/graphforge-storage/src/graph_construction/partition_shaping/tests.rs index 5dc1bfc2c..2e83345bc 100644 --- a/crates/graphforge-storage/src/graph_construction/partition_shaping/tests.rs +++ b/crates/graphforge-storage/src/graph_construction/partition_shaping/tests.rs @@ -1,4 +1,5 @@ use super::*; +use std::sync::Arc; // BenchExec can run either fixture in a fresh test process with // GF_DETAIL_PARTITION_ROWS set to a fixed larger scale. No timing or resource @@ -36,6 +37,7 @@ fn detail_partition_load(family: PartitionFamily) { &fixed_spill_name(family, 0), Some(count), Some(codec), + super::super::partition::default_materialization_bytes(), &AtomicBool::new(false), ) .unwrap(); @@ -78,3 +80,117 @@ fn node_detail_partition_load_preserves_sorted_wire_bytes() { fn edge_detail_partition_load_preserves_sorted_wire_bytes() { detail_partition_load::<304>(PartitionFamily::EdgeDetails); } + +#[test] +fn concentrated_fixed_partition_is_refused_before_record_allocation() { + let root = tempfile::TempDir::new().unwrap(); + crate::open_or_initialize_project(root.path()).unwrap(); + let mut session = super::super::tests::open(&root, 0x1439); + let mut partitioner = + FixedRangePartitioner::<33>::new(&session.root, PartitionFamily::Endpoints, 1, None, false) + .unwrap(); + // Every endpoint has the same hub key. More UUID splitters cannot separate it. + let mut wire = Vec::new(); + for i in 0_u64..1024 { + let mut row = [0; 33]; + row[24..32].copy_from_slice(&i.to_be_bytes()); + wire.extend_from_slice(&row); + } + partitioner + .route_slice(0, &wire, 1024, &mut session.checkpoint.evidence) + .unwrap(); + partitioner.seal(&mut session.checkpoint.evidence).unwrap(); + let name = fixed_spill_name(PartitionFamily::Endpoints, 0); + let error = load_fixed_partition::<33>( + &session.root, + &name, + Some(1024), + None, + 1024 * 33 - 1, + &AtomicBool::new(false), + ) + .err() + .unwrap(); + assert!(error.to_string().contains("exceeds recorded budget")); + let (rows, _) = load_fixed_partition::<33>( + &session.root, + &name, + Some(1024), + None, + 1024 * 33, + &AtomicBool::new(false), + ) + .unwrap(); + assert_eq!(rows.len(), 1024); + // A corrupt routed count cannot make the Vec grow beyond its reservation. + assert!( + load_fixed_partition::<33>( + &session.root, + &name, + Some(1), + None, + 33, + &AtomicBool::new(false) + ) + .err() + .unwrap() + .to_string() + .contains("admitted record count") + ); +} + +#[test] +fn row_partition_admits_before_decode_and_authenticates_before_ipc_allocation() { + use arrow::{ + array::{FixedSizeBinaryArray, StringArray}, + datatypes::{DataType, Field, Schema}, + }; + let root = tempfile::TempDir::new().unwrap(); + crate::open_or_initialize_project(root.path()).unwrap(); + let mut session = super::super::tests::open(&root, 0x143a); + let schema = Arc::new(Schema::new(vec![ + Field::new("uuid", DataType::FixedSizeBinary(16), false), + Field::new("text", DataType::Utf8, false), + ])); + let ids = [[2_u8; 16], [1_u8; 16]]; + let batch = RecordBatch::try_new( + schema.clone(), + vec![ + Arc::new(FixedSizeBinaryArray::try_from_iter(ids.iter()).unwrap()), + Arc::new(StringArray::from(vec![ + "long property".repeat(1000), + "other".to_owned(), + ])), + ], + ) + .unwrap(); + let mut partitioner = RowRangePartitioner::new(&session.root, "node-budget-test", 1).unwrap(); + partitioner + .route_slice(&schema, &batch, 0, 2, 0, &mut session.checkpoint.evidence) + .unwrap(); + partitioner.seal(&mut session.checkpoint.evidence).unwrap(); + let name = partitioner.sealed[0].as_ref().unwrap().name.clone(); + let sorted = partitioner + .load_partition(0, &name, &schema, &mut session.checkpoint.evidence) + .unwrap(); + assert_eq!(uuid_value(key_column(&sorted).unwrap(), 0).unwrap(), ids[1]); + partitioner.max_partition_bytes = 1; + let error = partitioner + .load_partition(0, &name, &schema, &mut session.checkpoint.evidence) + .unwrap_err(); + assert!(error.to_string().contains("exceeds recorded budget")); + partitioner.max_partition_bytes = super::super::partition::default_materialization_bytes(); + let path = root + .path() + .join(".graphforge-construction") + .join(format!("{:032x}", 0x143a)) + .join(&name); + // Corrupt the IPC metadata length. Arrow must never get to allocate it. + let mut bytes = std::fs::read(&path).unwrap(); + bytes[4..8].copy_from_slice(&i32::MAX.to_le_bytes()); + std::fs::write(&path, bytes).unwrap(); + let error = partitioner + .load_partition(0, &name, &schema, &mut session.checkpoint.evidence) + .unwrap_err(); + assert!(error.to_string().contains("digest"), "{error}"); +} diff --git a/crates/graphforge-storage/src/graph_construction/recovery.rs b/crates/graphforge-storage/src/graph_construction/recovery.rs index f9a399006..3694477e2 100644 --- a/crates/graphforge-storage/src/graph_construction/recovery.rs +++ b/crates/graphforge-storage/src/graph_construction/recovery.rs @@ -803,6 +803,31 @@ pub(super) fn authenticate_artifact( let _diagnostic_scope = crate::graph_construction::diagnostics::Scope::start("artifact_authentication"); validate_artifact_name(receipt)?; + authenticate_artifact_contents(root, receipt, codec) +} + +/// Authenticate an internally generated, sealed row spill before IPC decoding. +/// These temporary spills are not published construction artifacts. +pub(super) fn authenticate_row_spill( + root: &StableDirectory, + receipt: &ArtifactReceipt, +) -> Result { + if !receipt.name.starts_with("part-rows-") + || !receipt.name.ends_with(".arrow") + || receipt.name.contains('/') + || receipt.name.contains('\\') + || !is_canonical_sha256(&receipt.sha256) + { + return Err(storage("invalid row partition spill receipt")); + } + authenticate_artifact_contents(root, receipt, DetailCodec::Compact) +} + +fn authenticate_artifact_contents( + root: &StableDirectory, + receipt: &ArtifactReceipt, + codec: DetailCodec, +) -> Result { let file = root .open_child_file(OsStr::new(&receipt.name)) .map_err(storage)?; diff --git a/crates/graphforge-storage/src/graph_construction/shape.rs b/crates/graphforge-storage/src/graph_construction/shape.rs index 508bbc027..18c18b85c 100644 --- a/crates/graphforge-storage/src/graph_construction/shape.rs +++ b/crates/graphforge-storage/src/graph_construction/shape.rs @@ -140,7 +140,11 @@ impl GraphConstructionSession { // empirically measured RSS-vs-buffer-overhead minimum into the // buffer-dominated side of that curve. Reverted: the cut is sized // from the sample domain here, same as the joint plan. - let mut sampler = IdentitySampler::new(partition_count, total)?; + let mut sampler = IdentitySampler::with_target( + partition_count, + total, + self.checkpoint.budgets.target_partition_records, + )?; let positions = sampler.positions().collect::>(); let mut cursor = 0_usize; for position in positions { @@ -266,28 +270,32 @@ impl GraphConstructionSession { partitions, None, true, - )?; + )? + .with_materialization_limit(self.checkpoint.budgets.max_partition_bytes); let mut node_details = FixedRangePartitioner::::new( &self.root, PartitionFamily::NodeDetails, node_partitions, Some(detail_codec), false, - )?; + )? + .with_materialization_limit(self.checkpoint.budgets.max_partition_bytes); let mut edge_details = FixedRangePartitioner::::new( &self.root, PartitionFamily::EdgeDetails, partitions, Some(detail_codec), false, - )?; + )? + .with_materialization_limit(self.checkpoint.budgets.max_partition_bytes); let mut endpoints = FixedRangePartitioner::::new( &self.root, PartitionFamily::Endpoints, node_partitions, None, false, - )?; + )? + .with_materialization_limit(self.checkpoint.budgets.max_partition_bytes); let mut row_groups: BTreeMap<(u8, String), RowRangePartitioner> = BTreeMap::new(); // #1455. A kind whose every staged chunk carries the bare canonical // schema has no property columns, so its shaped rows would hold @@ -392,7 +400,8 @@ impl GraphConstructionSession { &self.root, &format!("{kind}-{}", receipt.schema_sha256), row_plan.partitions(), - )?, + )? + .with_materialization_limit(self.checkpoint.budgets.max_partition_bytes), ); } row_groups @@ -547,6 +556,7 @@ impl GraphConstructionSession { endpoints.as_deref(), self.base_snapshot.as_mut(), self.checkpoint.budgets.max_batch_rows, + self.checkpoint.budgets.max_partition_bytes, &mut cancelled, &mut self.checkpoint.evidence, )?; @@ -1603,6 +1613,7 @@ pub(super) fn resolve_endpoint_surrogates( endpoints_name: Option<&str>, mut base: Option<&mut AuthenticatedUuidIndexSnapshot>, window_rows: usize, + max_partition_bytes: u64, cancelled: &mut impl FnMut() -> bool, evidence: &mut GraphConstructionEvidence, ) -> Result, GfError> { @@ -1620,7 +1631,8 @@ pub(super) fn resolve_endpoint_surrogates( plan.partitions(), None, false, - )?; + )? + .with_materialization_limit(max_partition_bytes); loop { let mut endpoint_window = Vec::with_capacity(window_rows); for _ in 0..window_rows { diff --git a/crates/graphforge-storage/src/graph_construction/tests.rs b/crates/graphforge-storage/src/graph_construction/tests.rs index 029d2839c..79207def0 100644 --- a/crates/graphforge-storage/src/graph_construction/tests.rs +++ b/crates/graphforge-storage/src/graph_construction/tests.rs @@ -438,8 +438,9 @@ fn million_chunk_shaping_retains_name_state_bounded_by_the_partition_count() { let slots = u64::from(budgets.partition_count) * families; // Name state is a function of the recorded partition count and the family // set, never of the staged chunk count. - assert_eq!(slots, 1_280, "retained slots: {slots}"); - assert!(slots < budgets.max_chunks / 100, "retained slots: {slots}"); + assert_eq!(slots, 20_480, "retained slots: {slots}"); + // Even at the maximum cut, name slots stay independent of chunk count. + assert!(slots < budgets.max_chunks / 40, "retained slots: {slots}"); assert_eq!(budgets.max_schema_groups, 256); assert_eq!( budgets.partition_count, @@ -1323,3 +1324,66 @@ fn seal_batches_report_one_directory_barrier_per_seal_not_per_spill() { "one directory barrier per seal batch" ); } + +#[test] +fn legacy_default_partition_budget_resumes_without_rewriting_authority() { + let root = TempDir::new().unwrap(); + let operation = Uuid::from_u128(0x1439); + let legacy = GraphConstructionBudgets { + partition_count: 256, + ..Default::default() + }; + let session = GraphConstructionSession::open(root.path(), operation, 0, legacy).unwrap(); + let before = serde_json::to_value(&session.checkpoint.budgets).unwrap(); + drop(session); + let resumed = GraphConstructionSession::open( + root.path(), + operation, + 0, + GraphConstructionBudgets::default(), + ) + .unwrap(); + assert_eq!(resumed.checkpoint.budgets, legacy); + assert_eq!( + serde_json::to_value(&resumed.checkpoint.budgets).unwrap(), + before + ); + drop(resumed); + let changed = GraphConstructionBudgets { + max_partition_bytes: 1024, + ..Default::default() + }; + assert!(GraphConstructionSession::open(root.path(), operation, 0, changed).is_err()); +} + +#[test] +fn recorded_partition_refusal_survives_reopen_with_no_completed_shape() { + let root = TempDir::new().unwrap(); + let operation = Uuid::from_u128(0x143b); + let budgets = GraphConstructionBudgets { + max_partition_bytes: 1, + ..Default::default() + }; + let mut session = GraphConstructionSession::open(root.path(), operation, 0, budgets).unwrap(); + session + .append(ConstructionChunkKind::Node, "node", &node_batch(1, 1)) + .unwrap(); + session.seal().unwrap(); + let error = session + .shape_canonical_with_cancellation(|| false) + .unwrap_err(); + assert!( + error.to_string().contains("exceeds recorded budget"), + "{error}" + ); + drop(session); + let mut resumed = GraphConstructionSession::open(root.path(), operation, 0, budgets).unwrap(); + assert!(resumed.checkpoint.shape_authority_sha256.is_none()); + let error = resumed + .shape_canonical_with_cancellation(|| false) + .unwrap_err(); + assert!( + error.to_string().contains("exceeds recorded budget"), + "{error}" + ); +} diff --git a/docs/book/architecture/resumable-import.md b/docs/book/architecture/resumable-import.md index 4735caadc..ea08eaacd 100644 --- a/docs/book/architecture/resumable-import.md +++ b/docs/book/architecture/resumable-import.md @@ -23,6 +23,38 @@ change establishes cross-process determinism for shaped payloads, not all encode metadata. Existing completed shapes remain authenticated and resumable with their recorded payloads; they are not rewritten to change historical catalog times. +Construction shaping also has recorded partition limits. New sessions allow up +to 4,096 ranges, with a target of 16,384 sampled identities per range. The cut +retains a 256-range floor where the input has enough identities for meaningful +balance checks, and never exceeds the recorded maximum or one range per 16 +identities. This is a deterministic planning policy, not a guarantee that skewed +or variable-width records fit. Recorded splitters remain the recovery authority; +worker count and host memory do not select the layout. + +`GraphConstructionBudgets::max_partition_bytes` defaults to 256 MiB of accounted +materialization per load. Fixed-width families admit the routed record count +and wire representation before reserving storage; compact details include their +sorting offsets. The production worker window holds at most two such loads. +Arrow property-row partitions load serially and charge decoded buffers, +concatenation/reordering capacity, nested child capacity, UUID keys and indexes. +Their conservative accounting can refuse a partition even when a less +conservative representation might fit. Unsupported Arrow representations are +refused rather than assigned an unproven estimate. Sealed IPC spills are +authenticated with bounded reads before Arrow can allocate from their metadata. + +A partition exceeding its recorded budget is refused before its retained load +is allocated. Increasing the cut cannot split one node's hub group; a large hub +can therefore be refused even at 4,096 ranges. The failed private shape remains +unpublished and reopening retains the same limits. The materialization budget +is separate from source decoding, routing/writer buffers, allocator metadata, +page cache and the process memory limit; it is not a whole-process RSS promise. + +The new budget fields have stable defaults for historical checkpoints and omit +default values when serialized. Opening a checkpoint with the exact former +256-partition default through today's default retains its recorded budgets; +other explicit budget mismatches remain errors. Completed shapes and their +recorded authority are not rewritten. + The resource envelope is explicit in `ImportSessionLimits`. Decoding is capped by `batch_rows`, source bytes and files have hard limits, and source readers are bounded by `io_concurrency`. Status reports accepted and rejected rows, bytes, From 6d6b7f17578508d7ecb4b2add9233e4fc0d43143 Mon Sep 17 00:00:00 2001 From: David Spencer <1526975+DecisionNerd@users.noreply.github.com> Date: Sun, 20 Sep 2026 03:39:06 +0000 Subject: [PATCH 2/3] style(storage): satisfy accounting and authentication lint contracts --- .../src/graph_construction/partition.rs | 8 ++++++++ .../src/graph_construction/partition_memory.rs | 3 +-- .../graphforge-storage/src/graph_construction/recovery.rs | 4 ++-- 3 files changed, 11 insertions(+), 4 deletions(-) diff --git a/crates/graphforge-storage/src/graph_construction/partition.rs b/crates/graphforge-storage/src/graph_construction/partition.rs index c081bd8fb..4fdc0208e 100644 --- a/crates/graphforge-storage/src/graph_construction/partition.rs +++ b/crates/graphforge-storage/src/graph_construction/partition.rs @@ -48,12 +48,20 @@ pub(crate) const MAX_PARTITION_COUNT: u32 = 4_096; pub(super) const fn default_target_records() -> u64 { 16_384 } +#[expect( + clippy::trivially_copy_pass_by_ref, + reason = "Serde skip_serializing_if requires a reference" +)] pub(super) fn is_default_target_records(value: &u64) -> bool { *value == default_target_records() } pub(super) const fn default_materialization_bytes() -> u64 { 256 << 20 } +#[expect( + clippy::trivially_copy_pass_by_ref, + reason = "Serde skip_serializing_if requires a reference" +)] pub(super) fn is_default_materialization_bytes(value: &u64) -> bool { *value == default_materialization_bytes() } diff --git a/crates/graphforge-storage/src/graph_construction/partition_memory.rs b/crates/graphforge-storage/src/graph_construction/partition_memory.rs index 9b0d3f511..870d339c7 100644 --- a/crates/graphforge-storage/src/graph_construction/partition_memory.rs +++ b/crates/graphforge-storage/src/graph_construction/partition_memory.rs @@ -36,7 +36,7 @@ fn array_reservation(data: &ArrayData, floor: u64) -> Result<(u64, u64), GfError mul(n, width as u64)? } else { match data.data_type() { - DataType::Null => 0, + DataType::Null | DataType::Struct(_) => 0, DataType::Boolean => add(n, 7)? / 8, DataType::FixedSizeBinary(width) => mul(n, u64::try_from(*width).map_err(storage)?)?, DataType::Utf8 | DataType::Binary | DataType::LargeUtf8 | DataType::LargeBinary => { @@ -57,7 +57,6 @@ fn array_reservation(data: &ArrayData, floor: u64) -> Result<(u64, u64), GfError } DataType::List(_) | DataType::Map(_, _) => mul(add(n, 1)?, 4)?, DataType::LargeList(_) => mul(add(n, 1)?, 8)?, - DataType::Struct(_) => 0, DataType::FixedSizeList(_, width) => { child_floor = mul(n, u64::try_from(*width).map_err(storage)?)?; 0 diff --git a/crates/graphforge-storage/src/graph_construction/recovery.rs b/crates/graphforge-storage/src/graph_construction/recovery.rs index 3694477e2..b947d8a3d 100644 --- a/crates/graphforge-storage/src/graph_construction/recovery.rs +++ b/crates/graphforge-storage/src/graph_construction/recovery.rs @@ -794,7 +794,6 @@ pub(super) struct ReadWork { pub(super) cache_release: graphforge_filesystem::FileCacheReleaseEvidence, } -#[allow(clippy::too_many_lines)] // One authentication pass validates format, digest, and cache cleanup together. pub(super) fn authenticate_artifact( root: &StableDirectory, receipt: &ArtifactReceipt, @@ -813,7 +812,7 @@ pub(super) fn authenticate_row_spill( receipt: &ArtifactReceipt, ) -> Result { if !receipt.name.starts_with("part-rows-") - || !receipt.name.ends_with(".arrow") + || std::path::Path::new(&receipt.name).extension() != Some(OsStr::new("arrow")) || receipt.name.contains('/') || receipt.name.contains('\\') || !is_canonical_sha256(&receipt.sha256) @@ -823,6 +822,7 @@ pub(super) fn authenticate_row_spill( authenticate_artifact_contents(root, receipt, DetailCodec::Compact) } +#[allow(clippy::too_many_lines)] // One authentication pass validates format, digest, and cache cleanup together. fn authenticate_artifact_contents( root: &StableDirectory, receipt: &ArtifactReceipt, From f34efd3671ec7b3f99f2906376ece1bc881e3d46 Mon Sep 17 00:00:00 2001 From: David Spencer <1526975+DecisionNerd@users.noreply.github.com> Date: Sun, 20 Sep 2026 04:03:13 +0000 Subject: [PATCH 3/3] fix(storage): decode row spills through authenticated descriptors --- .../graph_construction/partition_shaping.rs | 12 ++++--- .../partition_shaping/tests.rs | 34 +++++++++++++++++-- .../src/graph_construction/recovery.rs | 25 ++++++++++---- 3 files changed, 57 insertions(+), 14 deletions(-) diff --git a/crates/graphforge-storage/src/graph_construction/partition_shaping.rs b/crates/graphforge-storage/src/graph_construction/partition_shaping.rs index feec965f6..d440a54f5 100644 --- a/crates/graphforge-storage/src/graph_construction/partition_shaping.rs +++ b/crates/graphforge-storage/src/graph_construction/partition_shaping.rs @@ -1325,10 +1325,15 @@ impl<'a> RowRangePartitioner<'a> { .get(partition) .and_then(Option::as_ref) .ok_or_else(|| super::storage("row partition receipt is absent"))?; + if receipt.name != name { + return Err(super::storage( + "row partition name differs from its receipt", + )); + } self.reservations[partition].admit(receipt.bytes, self.max_partition_bytes)?; // Authenticate with bounded buffers BEFORE Arrow reads bodyLength and // allocates an IPC body. Post-decode checks cannot protect that step. - let work = super::recovery::authenticate_row_spill(self.root, receipt)?; + let (file, work) = super::recovery::authenticate_row_spill(self.root, receipt)?; account_cache_release(work.cache_release, evidence)?; evidence.merge_read_bytes = evidence .merge_read_bytes @@ -1338,10 +1343,7 @@ impl<'a> RowRangePartitioner<'a> { .merge_read_operations .checked_add(work.operations) .ok_or_else(|| super::storage("partition read operations overflow"))?; - let file = self - .root - .open_child_file(OsStr::new(name)) - .map_err(super::storage)?; + let counter = IoCounter::default(); let reader = super::CountingRead { inner: graphforge_filesystem::FileCacheReleasingReader::with_window_bytes( diff --git a/crates/graphforge-storage/src/graph_construction/partition_shaping/tests.rs b/crates/graphforge-storage/src/graph_construction/partition_shaping/tests.rs index 2e83345bc..0cb242f62 100644 --- a/crates/graphforge-storage/src/graph_construction/partition_shaping/tests.rs +++ b/crates/graphforge-storage/src/graph_construction/partition_shaping/tests.rs @@ -186,11 +186,41 @@ fn row_partition_admits_before_decode_and_authenticates_before_ipc_allocation() .join(format!("{:032x}", 0x143a)) .join(&name); // Corrupt the IPC metadata length. Arrow must never get to allocate it. - let mut bytes = std::fs::read(&path).unwrap(); + let original = std::fs::read(&path).unwrap(); + let mut bytes = original.clone(); bytes[4..8].copy_from_slice(&i32::MAX.to_le_bytes()); - std::fs::write(&path, bytes).unwrap(); + std::fs::write(&path, &bytes).unwrap(); let error = partitioner .load_partition(0, &name, &schema, &mut session.checkpoint.evidence) .unwrap_err(); assert!(error.to_string().contains("digest"), "{error}"); + + // Substituting the pathname after authentication cannot redirect decoding. + std::fs::write(&path, original).unwrap(); + let (file, _) = super::super::recovery::authenticate_row_spill( + &session.root, + partitioner.sealed[0].as_ref().unwrap(), + ) + .unwrap(); + let renamed = std::fs::rename(&path, path.with_extension("saved")); + #[cfg(windows)] + assert!( + renamed.is_err(), + "the retained handle denies delete sharing" + ); + #[cfg(not(windows))] + { + renamed.unwrap(); + std::fs::write(&path, bytes).unwrap(); + } + let mut reader = arrow::ipc::reader::StreamReader::try_new(file, None).unwrap(); + assert_eq!(reader.next().unwrap().unwrap().num_rows(), 2); + assert!(reader.next().is_none()); + #[cfg(not(windows))] + { + let error = partitioner + .load_partition(0, &name, &schema, &mut session.checkpoint.evidence) + .unwrap_err(); + assert!(error.to_string().contains("identity"), "{error}"); + } } diff --git a/crates/graphforge-storage/src/graph_construction/recovery.rs b/crates/graphforge-storage/src/graph_construction/recovery.rs index b947d8a3d..27d5f7d68 100644 --- a/crates/graphforge-storage/src/graph_construction/recovery.rs +++ b/crates/graphforge-storage/src/graph_construction/recovery.rs @@ -802,7 +802,10 @@ pub(super) fn authenticate_artifact( let _diagnostic_scope = crate::graph_construction::diagnostics::Scope::start("artifact_authentication"); validate_artifact_name(receipt)?; - authenticate_artifact_contents(root, receipt, codec) + let file = root + .open_child_file(OsStr::new(&receipt.name)) + .map_err(storage)?; + authenticate_artifact_contents(file, receipt, codec) } /// Authenticate an internally generated, sealed row spill before IPC decoding. @@ -810,7 +813,7 @@ pub(super) fn authenticate_artifact( pub(super) fn authenticate_row_spill( root: &StableDirectory, receipt: &ArtifactReceipt, -) -> Result { +) -> Result<(std::fs::File, ReadWork), GfError> { if !receipt.name.starts_with("part-rows-") || std::path::Path::new(&receipt.name).extension() != Some(OsStr::new("arrow")) || receipt.name.contains('/') @@ -819,18 +822,26 @@ pub(super) fn authenticate_row_spill( { return Err(storage("invalid row partition spill receipt")); } - authenticate_artifact_contents(root, receipt, DetailCodec::Compact) + let mut file = root + .open_child_file(OsStr::new(&receipt.name)) + .map_err(storage)?; + // Cloning retains this exact file identity rather than reopening its name. + // The duplicate shares its cursor, so rewind after the authentication pass. + let work = authenticate_artifact_contents( + file.try_clone().map_err(storage)?, + receipt, + DetailCodec::Compact, + )?; + std::io::Seek::rewind(&mut file).map_err(storage)?; + Ok((file, work)) } #[allow(clippy::too_many_lines)] // One authentication pass validates format, digest, and cache cleanup together. fn authenticate_artifact_contents( - root: &StableDirectory, + file: std::fs::File, receipt: &ArtifactReceipt, codec: DetailCodec, ) -> Result { - let file = root - .open_child_file(OsStr::new(&receipt.name)) - .map_err(storage)?; if !receipt .identity .matches(file_identity(&file).map_err(storage)?)