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
20 changes: 13 additions & 7 deletions crates/graphforge-storage/src/construction_detail_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -75,7 +75,7 @@ mod compact_details {
)
.unwrap();
session
.append(ConstructionChunkKind::Edge, "edges", &edge_batch(100, 2))
.append(ConstructionChunkKind::Edge, "edges", &edge_batch(100, 1, 3, 2))
.unwrap();
session.seal().unwrap();
drop(session);
Expand Down Expand Up @@ -153,7 +153,7 @@ mod compact_details {
)
.unwrap();
initial
.append(ConstructionChunkKind::Edge, "edges", &edge_batch(100, 2))
.append(ConstructionChunkKind::Edge, "edges", &edge_batch(100, 1, 3, 2))
.unwrap();
initial.seal().unwrap();
let encoded = initial.prepare_canonical_encoding(1).unwrap();
Expand Down Expand Up @@ -295,7 +295,7 @@ mod compact_details {
.append(
ConstructionChunkKind::Edge,
&format!("edges-{chunk}"),
&edge_batch(10_000 + chunk * 128, 128),
&edge_batch(10_000 + chunk * 128, 1 + (chunk * 128) % 1024, 1024, 128),
)
.unwrap();
}
Expand Down Expand Up @@ -410,12 +410,18 @@ mod compact_details {
record[17..23].copy_from_slice(b"Person");
expected.push(record);
}
// src/dst walk the full 1024-node range across the 8 chunks and
// wrap modulo 1024 (#1439), rather than restarting at node 1 on
// every chunk: `edge_batch`'s `node_start` argument varies per
// chunk to match.
for chunk in 0_u128..8 {
for row in 0_u128..128 {
let mut record = vec![0; EDGE_DETAIL_WIDTH];
record[..16].copy_from_slice(&(10_000 + chunk * 128 + row).to_be_bytes());
record[16..32].copy_from_slice(&(row + 1).to_be_bytes());
record[32..48].copy_from_slice(&(row + 2).to_be_bytes());
let src = 1 + chunk * 128 + row;
let dst = 1 + (chunk * 128 + row + 1) % 1024;
record[16..32].copy_from_slice(&src.to_be_bytes());
record[32..48].copy_from_slice(&dst.to_be_bytes());
record[48] = 1;
record[49] = b'R';
expected.push(record);
Expand Down Expand Up @@ -496,7 +502,7 @@ mod compact_details {
.append(ConstructionChunkKind::Node, "nodes", &node_batch(1, 2))
.unwrap();
session
.append(ConstructionChunkKind::Edge, "edges", &edge_batch(100, 1))
.append(ConstructionChunkKind::Edge, "edges", &edge_batch(100, 1, 2, 1))
.unwrap();
session.seal().unwrap();
let encoded = session.prepare_canonical_encoding(1).unwrap();
Expand Down Expand Up @@ -567,7 +573,7 @@ mod compact_details {
}
if resumed.accepted_chunks() == 1 {
resumed
.append(ConstructionChunkKind::Edge, "edges", &edge_batch(100, 1))
.append(ConstructionChunkKind::Edge, "edges", &edge_batch(100, 1, 2, 1))
.unwrap();
}
resumed.seal().unwrap();
Expand Down
67 changes: 55 additions & 12 deletions crates/graphforge-storage/src/construction_determinism_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,17 @@
// That is a pre-existing durable-format defect, it is not introduced here, and
// fixing it is a separate concern.
//
// The table above is a historical measurement on `13632d4b`, predating S5's
// external-merge-tree deletion, not a pinned expectation: no test asserts
// those literal values. Read it for what it demonstrates -- self-consistency
// across two sessions -- not as a checksum to reproduce. It is stale for two
// rows as of #1439's follow-up (per-family splitters): routing node-keyed
// families with node-only splitters changes how edge details and edge
// endpoints partition, so `edge details` and `edge endpoints` legitimately
// produce different bytes now. `shaped-identities.run` and `node details`
// are unchanged, because their routing did not change. The claim this file's
// tests actually enforce -- stated precisely below -- still holds.
//
// So the claim these tests make is precise:
//
// * **Within a fixed set of recorded session parameters** — the session clock
Expand Down Expand Up @@ -108,16 +119,25 @@ mod determinism {
.unwrap()
}

fn edge_rows(edges: &[[u8; 16]], nodes: &[[u8; 16]]) -> RecordBatch {
/// `start` is this window's offset into the full edge sequence (#1439).
/// Every call site passes one chunk at a time, and indexing `nodes` from
/// a per-call-local `0` on every chunk -- the previous behaviour --
/// referenced only the first `chunk`-many nodes as an endpoint no matter
/// how many nodes or chunks existed, concentrating every edge onto a
/// narrow node-UUID band. That was invisible while endpoints were routed
/// with the joint identity splitters, which never resolve node UUIDs
/// finely enough to notice; routing them with node-only splitters
/// surfaces it as a real (fixture-caused, not production) skew.
fn edge_rows(start: usize, edges: &[[u8; 16]], nodes: &[[u8; 16]]) -> RecordBatch {
let src = edges
.iter()
.enumerate()
.map(|(index, _)| nodes[index % nodes.len()])
.map(|(index, _)| nodes[(start + index) % nodes.len()])
.collect::<Vec<_>>();
let dst = edges
.iter()
.enumerate()
.map(|(index, _)| nodes[(index + 1) % nodes.len()])
.map(|(index, _)| nodes[(start + index + 1) % nodes.len()])
.collect::<Vec<_>>();
RecordBatch::try_new(
CONSTRUCTION_EDGE_SCHEMA.clone(),
Expand Down Expand Up @@ -200,7 +220,7 @@ mod determinism {
.append(
ConstructionChunkKind::Edge,
&format!("edges-{index}"),
&edge_rows(window, nodes),
&edge_rows(index * chunk, window, nodes),
)
.unwrap();
}
Expand Down Expand Up @@ -451,7 +471,18 @@ mod determinism {
let edges = edge_ids(1_024);
let mut reference: Option<Fingerprint> = None;
let mut observed = Vec::new();
for partition_count in [1_u32, 4, 64, 256] {
// #1439 extends the range up to MAX_PARTITION_COUNT, the new default
// ceiling: at this fixture's small record count both 256 and 4_096
// clip to the same recorded-count-bounded cut, so this proves raising
// the ceiling all the way to its new default introduces no
// logical-result variance either.
for partition_count in [
1_u32,
4,
64,
256,
crate::graph_construction::partition::MAX_PARTITION_COUNT,
] {
let root = TempDir::new().unwrap();
let (fingerprint, layout) = ingest(&root, partition_count, &nodes, &edges, 128);
assert_eq!(layout.recorded_partition_count, u64::from(partition_count));
Expand Down Expand Up @@ -482,19 +513,31 @@ mod determinism {
}
}
// The layouts genuinely differ; the results do not.
// The cut is the recorded count bounded by the recorded record count, so
// that a small graph does not pay a large graph's durability price.
// 2048 identities admit at most 2048/16 = 128 partitions.
// The cut is bounded by the recorded record count on two independent
// terms, `min`-combined (#1439): a small-scale floor of
// `identities / 16` (unchanged since before #1439) and a data-driven
// target floored at `DEFAULT_PARTITION_COUNT`. 2048 identities admit
// at most 2048/16 = 128 partitions from the small-scale floor, which
// is far below the data-driven target's floor of 256, so the
// small-scale floor is what binds at this fixture's size for every
// requested count above it -- including the raised MAX_PARTITION_COUNT
// default.
let partitions = observed
.iter()
.map(|(_, partitions, ..)| *partitions)
.collect::<Vec<_>>();
let identities = observed[0].3;
assert_eq!(identities, (nodes.len() + edges.len()) as u64);
let expected = [1_u32, 4, 64, 256]
.into_iter()
.map(|requested| u64::from(requested).min(identities / 16))
.collect::<Vec<_>>();
let expected = [
1_u32,
4,
64,
256,
crate::graph_construction::partition::MAX_PARTITION_COUNT,
]
.into_iter()
.map(|requested| u64::from(requested).min(identities / 16))
.collect::<Vec<_>>();
assert_eq!(partitions, expected);
let peaks = observed
.iter()
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -529,7 +529,7 @@ mod lifecycle_budget {
.append(
ConstructionChunkKind::Edge,
&format!("edges-{chunk}"),
&edge_batch(1_000_000 + u128::from(chunk) * 4096, 4096),
&edge_batch(1_000_000 + u128::from(chunk) * 4096, 1 + (u128::from(chunk) * 4096) % 8192, 8192, 4096),
)
.unwrap();
}
Expand Down Expand Up @@ -1079,7 +1079,7 @@ mod lifecycle_budget {
.append(
ConstructionChunkKind::Edge,
&format!("edges-{chunk}"),
&edge_batch(1_000_000 + u128::from(chunk) * 4096, 4096),
&edge_batch(1_000_000 + u128::from(chunk) * 4096, 1 + (u128::from(chunk) * 4096) % 8192, 8192, 4096),
)
.unwrap();
}
Expand Down
10 changes: 10 additions & 0 deletions crates/graphforge-storage/src/graph_construction.rs
Original file line number Diff line number Diff line change
Expand Up @@ -822,6 +822,16 @@ struct ShapeIntent {
/// partition writes a byte.
#[serde(default)]
splitters: Vec<String>,
/// Recorded range-partition splitters over the staged **node** identity
/// domain only (#1439), canonical lower hex, strictly increasing.
/// Endpoints, node details and node-kind rows are keyed by node UUID and
/// are routed with these instead of `splitters`: Graph500-shaped input
/// puts nodes and edges in disjoint UUID bands, so the joint splitters
/// above route almost every node-keyed record into a handful of
/// partitions. A pure function of the same recorded chunk receipts as
/// `splitters`, so R1 holds for this set too.
#[serde(default)]
node_splitters: Vec<String>,
/// Measured identity rows per effective partition, in partition order.
#[serde(default)]
partition_identity_rows: Vec<u64>,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -239,7 +239,11 @@ fn schema_group_admission_is_constant_and_budgeted() {
.append(ConstructionChunkKind::Node, "nodes", &node_batch(1, 2))
.unwrap();
let error = session
.append(ConstructionChunkKind::Edge, "edges", &edge_batch(100, 1))
.append(
ConstructionChunkKind::Edge,
"edges",
&edge_batch(100, 1, 2, 1),
)
.unwrap_err();
assert!(error.to_string().contains("schema-group budget"));
assert_eq!(session.checkpoint.node_schema_sha256.len(), 1);
Expand Down
13 changes: 13 additions & 0 deletions crates/graphforge-storage/src/graph_construction/io_evidence.rs
Original file line number Diff line number Diff line change
Expand Up @@ -346,6 +346,19 @@ pub struct GraphConstructionEvidence {
/// Identity keys retained as splitter sample points.
#[serde(default)]
pub splitter_sample_records: u64,
/// Effective range partitions produced by the node-only splitter set
/// (#1439), used to route node-keyed families (node details, endpoints,
/// node rows) instead of the joint identity splitters.
#[serde(default)]
pub shape_node_partitions: u64,
/// Node identity records observed by the node-only splitter sampling
/// pass (#1439).
#[serde(default)]
pub node_splitter_sampled_source_records: u64,
/// Node identity keys retained as node-only splitter sample points
/// (#1439).
#[serde(default)]
pub node_splitter_sample_records: u64,
/// Rows routed into the largest identity partition.
#[serde(default)]
pub max_partition_identity_rows: u64,
Expand Down
23 changes: 22 additions & 1 deletion crates/graphforge-storage/src/graph_construction/partition.rs
Original file line number Diff line number Diff line change
Expand Up @@ -316,12 +316,23 @@ impl PartitionBalance {
/// Returns an error when the partition is out of range or the count
/// overflows.
pub(crate) fn record(&mut self, partition: usize) -> Result<(), GfError> {
self.record_many(partition, 1)
}

/// Charge `count` rows to `partition` in one step (#1439 follow-up: the
/// caller batching a contiguous same-partition run charges it once
/// rather than once per record).
///
/// # Errors
/// Returns an error when the partition is out of range or the count
/// overflows.
pub(crate) fn record_many(&mut self, partition: usize, count: u64) -> Result<(), GfError> {
let slot = self
.rows
.get_mut(partition)
.ok_or_else(|| storage("partition index is out of range"))?;
*slot = slot
.checked_add(1)
.checked_add(count)
.ok_or_else(|| storage("partition row count overflows"))?;
Ok(())
}
Expand All @@ -348,6 +359,16 @@ impl PartitionBalance {
/// where the quantile estimate has fewer points than partitions and the
/// ratio carries no information.
///
/// "Rows" are whatever the caller recorded. The splitters are quantiles
/// of the *key* domain, so the quantity they promise to balance is the
/// number of distinct keys per partition, not records: a range partition
/// cannot split one key, and a family with several records per key
/// (staged endpoints, keyed by node, hold one record per incident edge)
/// is expected to be row-skewed exactly as much as its key degrees are.
/// Such a family records its distinct-key counts here instead of its row
/// counts (`FixedRangePartitioner::finish_optional`, #1439); for a
/// unique-key family the two are the same number.
///
/// # Errors
/// Returns an error when the partitioning is skewed beyond the tolerance.
pub(crate) fn assert_balanced(&self, context: &str) -> Result<(), GfError> {
Expand Down
Loading