From 60882600cd483a0573a6a617e61fcb8ca8246abf Mon Sep 17 00:00:00 2001 From: David Spencer <1526975+DecisionNerd@users.noreply.github.com> Date: Sat, 19 Sep 2026 15:51:13 +0000 Subject: [PATCH 1/2] fix(storage): bound contiguous partition routing buffers (#1445) --- .../src/graph_construction/shape.rs | 25 +++++- .../src/graph_construction/shape/tests.rs | 79 +++++++++++++++++++ 2 files changed, 100 insertions(+), 4 deletions(-) diff --git a/crates/graphforge-storage/src/graph_construction/shape.rs b/crates/graphforge-storage/src/graph_construction/shape.rs index b89dcf71c..cf45e7496 100644 --- a/crates/graphforge-storage/src/graph_construction/shape.rs +++ b/crates/graphforge-storage/src/graph_construction/shape.rs @@ -1098,9 +1098,11 @@ pub(super) fn route_fixed_run( combine_cache_cleanup(routed, released, "construction partition source") } -/// Accumulates one contiguous same-partition run of wire bytes so the -/// caller writes it in a single call instead of once per record (#1439 -/// follow-up). +/// Maximum wire bytes retained by one routing accumulator. +const PARTITION_RUN_BYTES: usize = 1024 * 1024; + +/// Accumulates bounded same-partition slices of wire bytes, amortizing writes +/// without retaining a whole partition run (#1445). /// /// The staged runs routed through [`route_fixed_run`] and /// [`route_identity_run`] are UUID-sorted and the partition function is @@ -1109,6 +1111,7 @@ pub(super) fn route_fixed_run( /// [`RowRangePartitioner::route_batch`] already exploits for Arrow rows; /// this is its fixed-width-record equivalent. struct PartitionRun { + bound: usize, partition: Option, bytes: Vec, records: u64, @@ -1116,7 +1119,12 @@ struct PartitionRun { impl PartitionRun { fn new() -> Self { + Self::with_bound(PARTITION_RUN_BYTES) + } + + fn with_bound(bound: usize) -> Self { Self { + bound, partition: None, bytes: Vec::new(), records: 0, @@ -1130,9 +1138,18 @@ impl PartitionRun { target: &mut FixedRangePartitioner, evidence: &mut GraphConstructionEvidence, ) -> Result<(), GfError> { - if self.partition.is_some_and(|current| current != partition) { + if self.partition.is_some_and(|current| current != partition) + || wire.len() > self.bound - self.bytes.len() + { self.flush(target, evidence)?; } + // Keep an oversized record whole without growing the accumulator. + if wire.len() >= self.bound { + return target.route_slice(partition, wire, 1, evidence); + } + if self.bytes.capacity() == 0 { + self.bytes.reserve_exact(self.bound); + } self.partition = Some(partition); self.bytes.extend_from_slice(wire); self.records = self diff --git a/crates/graphforge-storage/src/graph_construction/shape/tests.rs b/crates/graphforge-storage/src/graph_construction/shape/tests.rs index 5b4fa86f3..a233ebf53 100644 --- a/crates/graphforge-storage/src/graph_construction/shape/tests.rs +++ b/crates/graphforge-storage/src/graph_construction/shape/tests.rs @@ -710,3 +710,82 @@ fn resolved_endpoints_never_cost_one_write_submission_each() { "shaping submitted {submissions} writes for {resolved_endpoints} resolved endpoints" ); } + +/// A single partition larger than every candidate buffer must stream without +/// growing the accumulator or changing its authenticated bytes/row counts. +#[test] +fn partition_run_bounds_fixed_records_without_changing_output() { + let records = (1_u128..=131_073) + .map(u128::to_be_bytes) + .collect::>(); + for bound in [64 * 1024, 256 * 1024, 1024 * 1024] { + assert_bounded_partition_run(&records, None, bound); + } +} + +#[test] +fn partition_run_keeps_compact_records_whole_across_flushes() { + let records = (1_u128..=512) + .map(|key| { + let mut record = [0; NODE_DETAIL_WIDTH]; + record[..16].copy_from_slice(&key.to_be_bytes()); + let length = if key % 2 == 0 { 255 } else { 1 }; + record[16] = length; + record[17..17 + usize::from(length)].fill(b'x'); + record + }) + .collect::>(); + // The 64-byte case also exercises a record larger than the buffer. + for bound in [64, 1024, 64 * 1024] { + assert_bounded_partition_run(&records, Some(DetailCodec::Compact), bound); + } +} + +fn assert_bounded_partition_run( + records: &[[u8; N]], + codec: Option, + bound: usize, +) { + let root = TempDir::new().unwrap(); + crate::open_or_initialize_project(root.path()).unwrap(); + let mut session = open(&root, 0x1445); + let mut target = + FixedRangePartitioner::::new(&session.root, PartitionFamily::Identities, 2, codec, true) + .unwrap(); + let evidence = &mut session.checkpoint.evidence; + let mut run = PartitionRun::with_bound(bound); + let mut expected = Vec::new(); + // Two equally populated partitions exercise partition changes as well as + // byte-triggered flushes. Partition changes must retain partial buffers. + for (index, record) in records.iter().enumerate() { + let wire = codec + .map_or(Ok(record.as_slice()), |codec| codec.bytes(record)) + .unwrap(); + expected.extend_from_slice(wire); + run.push( + index / records.len().div_ceil(2), + wire, + &mut target, + evidence, + ) + .unwrap(); + assert!(run.bytes.len() <= bound); + assert!(run.bytes.capacity() <= bound); + } + run.flush(&mut target, evidence).unwrap(); + run.flush(&mut target, evidence).unwrap(); // Empty flush cannot duplicate rows. + assert_eq!(target.balance().total(), records.len() as u64); + let output = target + .finish_optional("staged-identities.run", &mut || false, evidence) + .unwrap() + .unwrap(); + let mut actual = Vec::new(); + session + .root + .open_child_file(OsStr::new(&output)) + .unwrap() + .read_to_end(&mut actual) + .unwrap(); + assert_eq!(actual, expected, "bound={bound}"); + assert_eq!(evidence.partition_rows, records.len() as u64); +} From 7eb882f7eee843732d02186a85b8a191e49da93b Mon Sep 17 00:00:00 2001 From: David Spencer <1526975+DecisionNerd@users.noreply.github.com> Date: Sat, 19 Sep 2026 17:06:34 +0000 Subject: [PATCH 2/2] perf(storage): select 64 KiB routing bound from measured curve --- .../src/graph_construction/shape.rs | 2 +- .../evidence/partition-run-bound-1445.json | 182 ++++++++++++++++++ .../evidence/partition-run-bound-1445.md | 96 +++++++++ 3 files changed, 279 insertions(+), 1 deletion(-) create mode 100644 docs/development/evidence/partition-run-bound-1445.json create mode 100644 docs/development/evidence/partition-run-bound-1445.md diff --git a/crates/graphforge-storage/src/graph_construction/shape.rs b/crates/graphforge-storage/src/graph_construction/shape.rs index cf45e7496..cdcfcf142 100644 --- a/crates/graphforge-storage/src/graph_construction/shape.rs +++ b/crates/graphforge-storage/src/graph_construction/shape.rs @@ -1099,7 +1099,7 @@ pub(super) fn route_fixed_run( } /// Maximum wire bytes retained by one routing accumulator. -const PARTITION_RUN_BYTES: usize = 1024 * 1024; +const PARTITION_RUN_BYTES: usize = 64 * 1024; /// Accumulates bounded same-partition slices of wire bytes, amortizing writes /// without retaining a whole partition run (#1445). diff --git a/docs/development/evidence/partition-run-bound-1445.json b/docs/development/evidence/partition-run-bound-1445.json new file mode 100644 index 000000000..d0b5c8479 --- /dev/null +++ b/docs/development/evidence/partition-run-bound-1445.json @@ -0,0 +1,182 @@ +{ + "base_sha": "51a0d6b984b3fab960a9c1116b7e1b4940e681d8", + "sha256": { + "/home/ubuntu/gf-1445-evidence/bin/bound-256k": "d5df2ddda8731795f626b6b7d09467c3ebd0e125b98e4bc375eb8a33d140fd1b", + "/home/ubuntu/gf-1445-evidence/bin/baseline": "dc10f5717d425611fc01742192384eaa8321b2551669b02fca9c86371e489bde", + "/home/ubuntu/gf-1445-evidence/bin/bound-64k": "00f681220c83e92a694957ff058baeb3128bc89f8f9531bf77349aa639a400f1", + "/home/ubuntu/gf-1445-evidence/bin/bound-1m": "2b2094da75352ea97ac0130cfadb9bf77848e93fb6b61178f60fa3c0f921977e", + "/home/ubuntu/gf-1445-evidence/candidate.patch": "7c1977d0ae00a8236103b4e27c91279e26acf20e2d4c2ddcdbec329d64c85f1e", + "/home/ubuntu/gf-1445-evidence/histogram.patch": "bd9ab301acc6672e36bc144bae7b68bd968fb745c012047604afb4ab4a00ec4c", + "/home/ubuntu/gf-1445-evidence/ingest.py": "485e1b968ecfda7fc9693e88712dcfdbf2a318b2ec47f5d0aab56d4bccde0a0a", + "/home/ubuntu/gf-1445-evidence/run-measured.sh": "6ca8c521b5956682c7f9abf8ab1113c81a82232b4de611ce63458d620e0aef56", + "/home/ubuntu/graphforge-ladder/cargo-bench/release/graphforge-benchmark-graph500-generator": "b1ddbd8bbcac15e50d530d7599f706603ab5a2e357a64730fb5ffa969d00e968", + "/home/ubuntu/gf-1445-evidence/quiet.py": "a6e86d2461a6bac9626d225d30c6ac17c8b2f73def00539117d9494aa04a54c5", + "/home/ubuntu/gf-1445-evidence/resume-curve.py": "9a8037abc4245652e52c3f13157373ac062c74b9abec4ea9a6e7fd62a5237bf8", + "/home/ubuntu/gf-1445-evidence/summarize-curve.py": "78baf16a2564a12db25e26fd77b8c72b6292b67ea59c3a54b8d3b48ec511278f", + "/home/ubuntu/gf-1445-evidence/measure-floors.py": "cac4be0cba0e780a4ad7262c2fa2269a9c2c0fabbb08ebcdca61f6c56c8a2a9d", + "/home/ubuntu/gf-1445-evidence/measurement-plan.txt": "69c3f83b6b471eaa2b9c712c3774fc94189a877ae1909ba9fedca7a8a7351ff3" + }, + "seed": 13907095936298285200, + "edge_factor": 16, + "input_sha256": { + "/home/ubuntu/gf-1445-evidence/input/edges.parquet": "30f197e5813adf8e0352d71d4451061686dfb96d267e681524c30344a96a52f4", + "/home/ubuntu/gf-1445-evidence/input/nodes.parquet": "44c9dfd9325013d0f6ea2f03bd86b00d2bea01265254cb28f70f0881a6478075", + "/home/ubuntu/gf-1445-evidence/input-s19/edges.parquet": "bd71780610259b1e6c721097926d930e20c434963ed1b3aa658802bdd2be91a0", + "/home/ubuntu/gf-1445-evidence/input-s19/nodes.parquet": "86491d7aa07331f25bbc788dc8aade06f97e7979f013e496e1970eee23cc1d98", + "/home/ubuntu/gf-1445-evidence/input-s20/edges.parquet": "c8e6b66f3d56b31f284a1288621d5536b7e5e347b6050e0ed3b0c8fed60e8953", + "/home/ubuntu/gf-1445-evidence/input-s20/nodes.parquet": "5792da943d39a3ec0cfe48c375fef1b078ae31f2c086af5d9bcfd417c74f24aa" + }, + "curve": { + "base_sha": "51a0d6b984b3fab960a9c1116b7e1b4940e681d8", + "selection_limit_wall_s": 30.78607031319989, + "selected": "bound-64k", + "curve": { + "baseline": { + "samples": [ + { + "label": "baseline-s18-guarded-round1", + "wall_s": 29.823490793118253, + "cpu_s": 25.896252, + "cgroup_peak_bytes": 792477696 + }, + { + "label": "baseline-s18-guarded-round2-attempt1", + "wall_s": 29.143522070953622, + "cpu_s": 25.505372, + "cgroup_peak_bytes": 821202944 + }, + { + "label": "baseline-s18-guarded-round3-attempt1", + "wall_s": 29.320066964952275, + "cpu_s": 25.723618, + "cgroup_peak_bytes": 763043840 + } + ], + "wall_s": 29.320066964952275, + "cpu_s": 25.723618, + "cgroup_peak_bytes": 792477696 + }, + "bound-64k": { + "samples": [ + { + "label": "bound-64k-s18-guarded-round1-attempt2", + "wall_s": 29.374653718201444, + "cpu_s": 25.724124, + "cgroup_peak_bytes": 799113216 + }, + { + "label": "bound-64k-s18-guarded-round2-attempt1", + "wall_s": 28.60761404898949, + "cpu_s": 25.612138, + "cgroup_peak_bytes": 790102016 + }, + { + "label": "bound-64k-s18-guarded-round3-attempt1", + "wall_s": 28.79827523091808, + "cpu_s": 25.088846, + "cgroup_peak_bytes": 791957504 + } + ], + "wall_s": 28.79827523091808, + "cpu_s": 25.612138, + "cgroup_peak_bytes": 791957504 + }, + "bound-256k": { + "samples": [ + { + "label": "bound-256k-s18-guarded-round1-attempt1", + "wall_s": 29.13570951204747, + "cpu_s": 25.580735, + "cgroup_peak_bytes": 845651968 + }, + { + "label": "bound-256k-s18-guarded-round2-attempt1", + "wall_s": 28.282239037798718, + "cpu_s": 25.22293, + "cgroup_peak_bytes": 784314368 + }, + { + "label": "bound-256k-s18-guarded-round3-attempt1", + "wall_s": 29.86603315989487, + "cpu_s": 25.689847, + "cgroup_peak_bytes": 803442688 + } + ], + "wall_s": 29.13570951204747, + "cpu_s": 25.580735, + "cgroup_peak_bytes": 803442688 + }, + "bound-1m": { + "samples": [ + { + "label": "bound-1m-s18-guarded-round1-attempt1", + "wall_s": 29.29350886703469, + "cpu_s": 25.370921, + "cgroup_peak_bytes": 840220672 + }, + { + "label": "bound-1m-s18-guarded-round2-attempt1", + "wall_s": 28.540352625073865, + "cpu_s": 24.991708, + "cgroup_peak_bytes": 791252992 + }, + { + "label": "bound-1m-s18-guarded-round3-attempt1", + "wall_s": 29.30297133116983, + "cpu_s": 25.272261, + "cgroup_peak_bytes": 791572480 + } + ], + "wall_s": 29.29350886703469, + "cpu_s": 25.272261, + "cgroup_peak_bytes": 791572480 + } + } + }, + "floors": { + "selected": "bound-64k", + "results": [ + { + "scale": 18, + "wall_s": 28.79827523091808, + "edges_per_second": 145644.2778731748, + "floor": 101175, + "basis": "three-run S18 median" + }, + { + "scale": 19, + "label": "bound-64k-s19-floor-attempt1", + "wall_s": 57.463328373851255, + "cpu_s": 50.350206, + "cgroup_peak_bytes": 1402183680, + "edges_per_second": 145981.93730485762, + "floor": 105415, + "basis": "one qualified complete-ingest observation" + }, + { + "scale": 20, + "label": "bound-64k-s20-floor-attempt1", + "wall_s": 113.88196149794385, + "cpu_s": 99.884978, + "cgroup_peak_bytes": 2673209344, + "edges_per_second": 147321.10142222053, + "floor": 105802, + "basis": "one qualified complete-ingest observation" + } + ] + }, + "excluded_attempts": [ + "baseline-s18-round1", + "baseline-s18-round2", + "bound-1m-s18-round1", + "bound-1m-s18-round2", + "bound-1m-s18-round3", + "bound-256k-s18-round1", + "bound-256k-s18-round2", + "bound-256k-s18-round3", + "bound-64k-s18-guarded-round1-attempt1", + "bound-64k-s18-guarded-round1", + "bound-64k-s18-round1", + "bound-64k-s18-round2" + ] +} diff --git a/docs/development/evidence/partition-run-bound-1445.md b/docs/development/evidence/partition-run-bound-1445.md new file mode 100644 index 000000000..d2d1d6172 --- /dev/null +++ b/docs/development/evidence/partition-run-bound-1445.md @@ -0,0 +1,96 @@ +# Partition routing buffer bound (#1445) + +`PartitionRun` previously accumulated an entire same-partition run within one +staged chunk. The repair bounds retained wire bytes, flushes before a record +would exceed the bound, and sends an oversized record directly. Records remain +whole; partition order, row counts and the publication protocol are unchanged. +Only one routing accumulator is live per routing call; this allocation is not +multiplied by the number of partitions. + +## Baseline histogram + +Measured on OVHC-AGENCY using release CLI code from `51a0d6b9`, plus temporary +flush-size logging. Graph500 S18, edge factor 16, seed `13907095936298285200`. +The logging patch is retained with the raw evidence and is absent from the PR. + +| Family | Flushes | Mean wire bytes | Maximum wire bytes | +| --- | ---: | ---: | ---: | +| Identities (26-byte records) | 323 | 358,723 | 456,144 | +| Node details (compact, padded width 272) | 259 | 21,255 | 21,672 | +| Edge details (compact, padded width 304) | 304 | 731,244 | 922,624 | +| Endpoints (33-byte records) | 16,365 | 16,916 | 511,863 | + +This does not support the historical hypothesis of multi-megabyte **mean** +flushes at S18. Routing restarts for each staged chunk. The remaining defect is +lack of an explicit byte bound, as the maintainer's updated #1445 acceptance +comment states; the earlier throughput regression was already repaired. + +## Bound-selection protocol + +BenchExec measures the complete command tree for `begin`, registration of both +Parquet sources, `validate`, and `commit`, including Python orchestration and CLI +launch costs. Each observation uses a fresh project, the same generated input, +16 logical CPUs and a 4 GiB cgroup limit. Reported memory is peak **cgroup memory**, +including cache; it is not process RSS. No cache dropping is performed. + +Compare 64 KiB, 256 KiB and 1 MiB with the baseline in three rotating-order S18 +rounds. The predeclared selection rule is the smallest bound with median wall +at most 5% above the matched baseline. This is a selection tolerance, not a +statistical confidence interval or an improvement claim. The baseline binary +retains the temporary logging branch but runs with logging disabled. Candidate +binaries contain no logging branch and differ only in the bound constant. + +Require the quiet-host guard before and after each observation, with one-second +checks during execution. The local guard also checks executable paths for renamed +GraphForge benchmark binaries, which the shared name-based guard misses. +Retain failed attempts and invalidate any contended +observation. The selected candidate must also protect #1445's recorded S18/S19/S20 +throughput floors: 101,175 / 105,415 / 105,802 edges/s. + +The initial curve was stopped after a concurrent #1452 benchmark with a renamed +executable was observed during round three. Its timing observations are retained +but excluded from bound selection. A fresh curve uses the stronger guard and +requires an initial 30-second quiet interval. Three clean rounds completed; the table reports medians: + +| Bound | Complete ingest wall (s) | CPU (s) | Peak cgroup memory (MiB) | +| --- | ---: | ---: | ---: | +| Baseline, unbounded | 29.320 | 25.724 | 755.77 | +| 64 KiB | 28.798 | 25.612 | 755.27 | +| 256 KiB | 29.136 | 25.581 | 766.22 | +| 1 MiB | 29.294 | 25.272 | 754.90 | + +**Selected: 64 KiB**, the smallest tested bound within the baseline +5% threshold +(30.786 s). This result supports bounding the accumulator without a material +whole-ingest regression under the stated selection rule; it does not establish a +throughput improvement. Two additional attempted 64 KiB observations were rejected +by the stronger guard when builds overlapped; neither enters these medians. + +All three historical floors pass under the complete-ingest measurement boundary: + +| Scale | Complete ingest wall (s) | Edges/s | Recorded floor | +| --- | ---: | ---: | ---: | +| S18 | 28.798 | 145,644 | 101,175 | +| S19 | 57.463 | 145,982 | 105,415 | +| S20 | 113.882 | 147,321 | 105,802 | + +S18 uses the three-run median; S19 and S20 each use one qualified observation. +Both larger ingests committed successfully, processed the expected input row +counts, and passed the before/during/after quiet checks. Peak cgroup memory was +1,337.23 MiB at S19 and 2,549.37 MiB at S20. + +## Correctness evidence + +`cargo test --release -p graphforge-storage --lib` passed 1,202 tests, with six +existing ignored tests. This includes determinism across sessions and partition +counts, interrupted shaping/recovery, corruption refusal and the new accumulator +regressions. The latter compare exact emitted bytes and row counts while checking +capacity at all three candidate bounds, partition transitions, compact records, +oversized records and final/empty flushes. + +`cargo test --release -p graphforge-api --test bdd` passed all 118 required API +scenarios and all 3,897 openCypher scenarios (zero regressions). Timing warnings +from that concurrent validation run are diagnostics, not benchmark evidence. + +Raw inputs, binaries, hashes, logging patch, command wrappers, histograms and +BenchExec output are retained at `/home/ubuntu/gf-1445-evidence/` on OVHC-AGENCY. +This report does not qualify S22–S26 or establish #1387's 1M edges/s floor.