Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 5 additions & 4 deletions crates/graphforge-storage/src/graph_construction.rs
Original file line number Diff line number Diff line change
Expand Up @@ -23,10 +23,10 @@ use io_evidence::{
account_merge_read_bytes, account_merge_write, account_merge_write_bytes, account_probe_work,
account_sequential_read, account_sequential_write, checked_category_remove,
checked_evidence_sum, combine_cache_cleanup, combine_secondary_cleanup, copy_post_shape_io,
merge_cache_release_evidence, open_counted_fixed_reader, read_run_record,
record_active_identity_install, record_active_identity_remove, record_category_install,
record_encoded_active_artifacts, record_encoding_io_evidence, record_shape_artifact_install,
release_counted_reader_cache,
injected_input_release_failure, merge_cache_release_evidence, open_counted_fixed_reader,
open_fixed_reader, read_run_record, record_active_identity_install,
record_active_identity_remove, record_category_install, record_encoded_active_artifacts,
record_encoding_io_evidence, record_shape_artifact_install, release_counted_reader_cache,
};
pub(crate) struct AuthenticatedShapeSource {
pub(crate) file: File,
Expand Down Expand Up @@ -68,6 +68,7 @@ mod partition;
use catalog::{
build_runtime_catalog, load_parent_runtime_catalog, load_parent_runtime_catalog_from_compact,
};
mod partition_load;
mod partition_shaping;
mod supersession;

Expand Down
43 changes: 36 additions & 7 deletions crates/graphforge-storage/src/graph_construction/io_evidence.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1352,9 +1352,28 @@ pub(super) fn open_counted_fixed_reader(
IoCounter,
),
GfError,
> {
let (reader, counter, length) = open_fixed_reader(root, name)?;
account_sequential_read(length, evidence)?;
Ok((reader, counter))
}

/// Open a fixed-width run for counted reading without touching the shared
/// evidence; the third value is the file length the caller credits later.
/// This is the worker-side half of [`open_counted_fixed_reader`].
pub(super) fn open_fixed_reader(
root: &StableDirectory,
name: &str,
) -> Result<
(
BufReader<CountingRead<graphforge_filesystem::FileCacheReleasingReader>>,
IoCounter,
u64,
),
GfError,
> {
let file = root.open_child_file(OsStr::new(name)).map_err(storage)?;
account_sequential_read(file.metadata().map_err(storage)?.len(), evidence)?;
let length = file.metadata().map_err(storage)?.len();
let counter = IoCounter::default();
let cache_window =
graphforge_filesystem::cache_release_window_for_streams(4).map_err(storage)?;
Expand All @@ -1372,15 +1391,16 @@ pub(super) fn open_counted_fixed_reader(
},
),
counter,
length,
))
}

pub(super) fn release_counted_reader_cache(
reader: &mut BufReader<CountingRead<graphforge_filesystem::FileCacheReleasingReader>>,
evidence: &mut GraphConstructionEvidence,
) -> Result<(), GfError> {
let released = reader.get_mut().inner.finish().map_err(storage)?;
account_cache_release(released, evidence)?;
/// Test-only injected failure for the input-release step. It is a
/// thread-local, so it is consumed on the thread that armed it: a partition
/// loaded on a worker thread reports it through the coordinator's merge step
/// rather than from the worker.
#[allow(clippy::unnecessary_wraps)]
pub(super) fn injected_input_release_failure() -> Result<(), GfError> {
#[cfg(test)]
SHAPE_CLEANUP_FAILURES.with(|failures| {
if failures.borrow().input_release {
Expand All @@ -1392,6 +1412,15 @@ pub(super) fn release_counted_reader_cache(
Ok(())
}

pub(super) fn release_counted_reader_cache(
reader: &mut BufReader<CountingRead<graphforge_filesystem::FileCacheReleasingReader>>,
evidence: &mut GraphConstructionEvidence,
) -> Result<(), GfError> {
let released = reader.get_mut().inner.finish().map_err(storage)?;
account_cache_release(released, evidence)?;
injected_input_release_failure()
}

pub(super) fn combine_cache_cleanup<T>(
primary: Result<T, GfError>,
cleanup: Result<(), GfError>,
Expand Down
307 changes: 307 additions & 0 deletions crates/graphforge-storage/src/graph_construction/partition_load.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,307 @@
//! Worker-local partition loads with coordinator-owned publication (#1456,
//! work item 1).
//!
//! Every parallelism attempt on the construction path (#1429, #1448) died on
//! the same shape: `&mut GraphConstructionEvidence` threaded through every
//! stage, and a whole-struct clone / mutate / write-back across an fsync in
//! `finish_optional`. Concurrent lanes lose updates even under a mutex.
//!
//! This module separates the two halves that were conflated:
//!
//! * **Workers return their own results and local counters.** A load never
//! sees the shared evidence. It returns the sorted partition plus a
//! [`PartitionLoadCounters`], whose merge into the evidence is explicit
//! and total ([`PartitionLoadCounters::merge_into`]): plain sums add,
//! per-partition maxima take `max`.
//! * **One coordinator owns the checkpoint, the allocation ledger and
//! publication.** [`consume_in_partition_order`] runs the loads on a bounded
//! pool and hands each result to the calling thread in canonical partition
//! index order, so scheduling cannot change output. The ordered critical
//! section is the coordinator's consume step; decoding, sorting and the bulk
//! read happen outside it. No worker touches the allocation ledger, so the
//! exact coexistence peak stays what it was.
//!
//! Two figures that are not the same figure:
//!
//! * The **observed** allocation peak is evidence. It is exact and depends on
//! interleaving, which is why the ledger append stays ordered on the
//! coordinator (two workers each installing and removing 100 bytes peak at
//! 100 or 200 depending on overlap; per-worker maxima are 100 and 100 either
//! way and cannot reconstruct which).
//! * A **reservation bound** is what an admission gate needs. It is
//! order-independent by construction. For materialized partitions it is
//! [`materialized_records_bound`]: `min(workers, partitions)` times the
//! largest partition, which the scheduler's window guarantees is never
//! exceeded.
//!
//! Cancellation needs no `Send` bound on the caller's `FnMut() -> bool`: the
//! callback stays on the coordinator, which polls it in the consume step, and
//! the pool observes a stop flag. Workers poll the flag between records so an
//! abandoned load exits promptly.

use super::{GraphConstructionEvidence, account_cache_release, account_sequential_read, storage};
use graphforge_core::GfError;
use std::collections::BTreeMap;
use std::num::NonZeroUsize;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Condvar, Mutex, MutexGuard, PoisonError};

/// Concurrent partition jobs per fixed-width family finish: one partition
/// loading and sorting while the coordinator writes the previous one.
///
/// This is a scheduling decision, not a recorded format parameter: the
/// durable partition count and splitters are unchanged by it, and the
/// schedule-independence tests hold the evidence and output bytes equal
/// across worker counts. It bounds the materialized partitions, see
/// [`consume_in_partition_order`].
pub(super) const PARTITION_LOAD_WORKERS: NonZeroUsize = NonZeroUsize::new(2).unwrap();

/// How often a load polls the stop flag, in records.
const STOP_POLL_RECORDS: usize = 4096;

/// Evidence a worker accumulates while loading one partition, kept apart from
/// the shared [`GraphConstructionEvidence`] until the coordinator merges it.
///
/// Every field here is either a plain sum or a per-partition maximum, so the
/// merge is commutative: the same set of loads produces the same evidence in
/// every consume order. The allocation ledger is deliberately absent; a load
/// installs and removes nothing.
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
pub(super) struct PartitionLoadCounters {
/// Spill length on disk, credited to `merge_read_blocks` as whole blocks.
pub(super) spill_bytes: u64,
/// Records decoded from the spill; also this partition's size for the
/// `peak_partition_records` maximum.
pub(super) records: u64,
/// Payload bytes actually read from the descriptor.
pub(super) read_bytes: u64,
/// Non-empty read submissions actually completed.
pub(super) read_operations: u64,
/// Page-cache release boundaries observed while reading.
pub(super) cache_release: graphforge_filesystem::FileCacheReleaseEvidence,
}

impl PartitionLoadCounters {
/// Fold this load into the shared evidence. Total: every counter a load
/// accumulates lands here, and nothing else is touched.
///
/// Fields written: `merge_read_blocks`, `merge_read_records`,
/// `merge_read_bytes`, `merge_read_operations`, `cache_release_operations`,
/// `cache_release_unsupported_operations`, `cache_released_bytes`,
/// `peak_cache_release_window_bytes` (max) and `peak_partition_records`
/// (max). The test `merge_is_total_and_touches_nothing_else` pins that list.
pub(super) fn merge_into(
self,
evidence: &mut GraphConstructionEvidence,
) -> Result<(), GfError> {
account_sequential_read(self.spill_bytes, evidence)?;
evidence.merge_read_records = evidence
.merge_read_records
.checked_add(self.records)
.ok_or_else(|| storage("merge read record count overflows"))?;
if (self.read_bytes == 0) != (self.read_operations == 0) {
return Err(storage("fixed-run read bytes and submissions disagree"));
}
evidence.merge_read_bytes = evidence
.merge_read_bytes
.checked_add(self.read_bytes)
.ok_or_else(|| storage("merge read byte count overflows"))?;
evidence.merge_read_operations = evidence
.merge_read_operations
.checked_add(self.read_operations)
.ok_or_else(|| storage("merge read operation count overflows"))?;
account_cache_release(self.cache_release, evidence)?;
evidence.peak_partition_records = evidence.peak_partition_records.max(self.records);
Ok(())
}
}

/// Reservation bound on simultaneously materialized fixed-width records under
/// [`consume_in_partition_order`]: never more than `min(workers, partitions)`
/// partitions exist at once, and none is larger than the largest.
///
/// This is the order-independent figure an admission gate may reserve
/// against. It is a bound, not an observation: the observed maximum is at
/// most this in every schedule, and equal to it only when the largest
/// partitions happen to coincide. `None` on overflow.
///
/// Test-only for now: no production budget governs materialized partition
/// records (`max_run_records` bounds staged runs, not partitions), so there
/// is no admission gate to hand this to yet. It is defined here, next to the
/// window that guarantees it, so the gate that arrives computes this figure
/// rather than a smaller one.
#[cfg(test)]
pub(super) fn materialized_records_bound(
workers: NonZeroUsize,
partitions: usize,
peak_partition_records: u64,
) -> Option<u64> {
let slots = u64::try_from(workers.get().min(partitions)).ok()?;
slots.checked_mul(peak_partition_records)
}

struct State<T> {
/// Next partition index to hand to a worker. Dispatch is in index order.
next: usize,
/// Partitions the coordinator has consumed and dropped.
released: usize,
/// Loaded partitions awaiting the coordinator, keyed by index.
ready: BTreeMap<usize, Result<T, GfError>>,
/// Workers that have not exited. The coordinator refuses to wait on a
/// partition no worker can deliver.
live_workers: usize,
}

struct Shared<T> {
state: Mutex<State<T>>,
changed: Condvar,
/// Set once by the coordinator when it stops consuming, for any reason.
stop: AtomicBool,
}

impl<T> Shared<T> {
fn lock(&self) -> MutexGuard<'_, State<T>> {
self.state.lock().unwrap_or_else(PoisonError::into_inner)
}

fn wait<'a>(&'a self, guard: MutexGuard<'a, State<T>>) -> MutexGuard<'a, State<T>> {
self.changed
.wait(guard)
.unwrap_or_else(PoisonError::into_inner)
}
}

/// Decrements the live-worker count however the worker exits.
struct LiveWorker<'a, T>(&'a Shared<T>);

impl<T> Drop for LiveWorker<'_, T> {
fn drop(&mut self) {
let mut state = self.0.lock();
state.live_workers -= 1;
drop(state);
self.0.changed.notify_all();
}
}

fn worker<T, L>(shared: &Shared<T>, partitions: usize, window: usize, load: &L)
where
L: Fn(usize, &AtomicBool) -> Result<T, GfError>,
{
let _live = LiveWorker(shared);
loop {
let index = {
let mut state = shared.lock();
loop {
if shared.stop.load(Ordering::Acquire) || state.next >= partitions {
return;
}
// Backpressure: running loads plus loaded-but-unconsumed
// results plus the partition being consumed never exceed the
// window. A slow first partition therefore holds the pool at
// `window` materialized partitions, not `partitions - 1`
// (#1448 measured exactly that defect).
if state.next < state.released + window {
break;
}
state = shared.wait(state);
}
let index = state.next;
state.next += 1;
index
};
let result = load(index, &shared.stop);
let mut state = shared.lock();
state.ready.insert(index, result);
drop(state);
shared.changed.notify_all();
}
}

/// Load `partitions` on `workers` threads and consume every result on the
/// calling thread in partition index order.
///
/// `load(index, stop)` runs on a worker; it may poll `stop` and abandon the
/// load once it is set. `consume(index, value)` runs on the calling thread and
/// may therefore borrow non-`Send` state: the shared evidence, the output
/// writer and the caller's cancellation callback. `value` is dropped when
/// `consume` returns and only then is its slot released to the pool.
///
/// Bound: at any instant at most `workers` partitions are materialized,
/// counting loads in flight, loaded results not yet consumed, and the one
/// being consumed. This is the `threads` bound #1448's reorder buffer claimed
/// and did not deliver.
///
/// The first error, in partition order for loads or immediately for
/// `consume`, stops dispatch, sets `stop`, and is returned after every worker
/// has exited. A load that fails after an earlier partition already failed is
/// never observed.
pub(super) fn consume_in_partition_order<T, L, C>(
partitions: usize,
workers: NonZeroUsize,
load: L,
mut consume: C,
) -> Result<(), GfError>
where
T: Send,
L: Fn(usize, &AtomicBool) -> Result<T, GfError> + Sync,
C: FnMut(usize, T) -> Result<(), GfError>,
{
if partitions == 0 {
return Ok(());
}
let window = workers.get();
let shared = Shared {
state: Mutex::new(State {
next: 0,
released: 0,
ready: BTreeMap::new(),
live_workers: window,
}),
changed: Condvar::new(),
stop: AtomicBool::new(false),
};
std::thread::scope(|scope| {
for _ in 0..window {
scope.spawn(|| worker(&shared, partitions, window, &load));
}
let outcome = (|| -> Result<(), GfError> {
for index in 0..partitions {
let loaded = {
let mut state = shared.lock();
loop {
if let Some(loaded) = state.ready.remove(&index) {
break loaded;
}
if state.live_workers == 0 {
return Err(storage(
"partition loaders exited before every partition was delivered",
));
}
state = shared.wait(state);
}
};
let consumed = consume(index, loaded?);
let mut state = shared.lock();
state.released += 1;
drop(state);
shared.changed.notify_all();
consumed?;
}
Ok(())
})();
shared.stop.store(true, Ordering::Release);
shared.changed.notify_all();
outcome
})
}

/// Poll the stop flag once every [`STOP_POLL_RECORDS`] records of a load.
pub(super) fn abandon_if_stopped(records: usize, stop: &AtomicBool) -> Result<(), GfError> {
if records.is_multiple_of(STOP_POLL_RECORDS) && stop.load(Ordering::Acquire) {
return Err(storage("partition load abandoned after coordinator stop"));
}
Ok(())
}

#[cfg(test)]
mod tests;
Loading