use std::ops::Range;
use std::time::Instant;
use spg_storage::{
Catalog, ColumnSchema, DataType, FreezeSlice, IndexKey, Row, StorageError, TableSchema, Value,
};
const POPULATE_ROWS: usize = 500_000;
const PAYLOAD_LEN: usize = 1_024;
const BATCH_ROWS: usize = POPULATE_ROWS;
const PREPARE_REPS: usize = 5;
fn wide_users_schema() -> TableSchema {
TableSchema {
name: "users".to_string(),
columns: vec![
ColumnSchema::new("id".to_string(), DataType::BigInt, false),
ColumnSchema::new("payload".to_string(), DataType::Text, false),
],
hot_tier_bytes: None,
foreign_keys: Vec::new(),
uniqueness_constraints: Vec::new(),
checks: Vec::new(),
partition_role: None,
policies: Vec::new(),
row_security: false,
force_row_security: false,
exclusion_constraints: Vec::new(),
acl: Vec::new(),
owner: None,
}
}
fn build_populated_catalog() -> Catalog {
let mut cat = Catalog::new();
cat.create_table(wide_users_schema()).unwrap();
let t = cat.get_mut("users").unwrap();
let payload: String = "x".repeat(PAYLOAD_LEN);
for id in 0..POPULATE_ROWS as i64 {
let row = Row::new(vec![Value::BigInt(id), Value::text(payload.clone())]);
t.insert(row).unwrap();
}
t.add_index("by_id".to_string(), "id").unwrap();
cat
}
fn partition_range(n: usize, parts: usize) -> Vec<Range<usize>> {
let mut out = Vec::with_capacity(parts);
let base = n / parts;
let extra = n % parts;
let mut start = 0;
for i in 0..parts {
let len = base + usize::from(i < extra);
out.push(start..start + len);
start += len;
}
out
}
fn freeze_once(cat: &mut Catalog, workers: usize) -> Result<(), StorageError> {
let ranges = partition_range(BATCH_ROWS, workers);
let prep_t0 = Instant::now();
let slices: Vec<FreezeSlice> = if workers == 1 {
ranges
.into_iter()
.map(|r| cat.prepare_freeze_slice("users", "by_id", r))
.collect::<Result<_, _>>()?
} else {
let cat_ref: &Catalog = cat;
std::thread::scope(|s| {
let handles: Vec<_> = ranges
.into_iter()
.map(|r| s.spawn(move || cat_ref.prepare_freeze_slice("users", "by_id", r)))
.collect();
handles
.into_iter()
.map(|h| h.join().expect("worker panicked"))
.collect::<Result<Vec<_>, _>>()
})?
};
let prep_wall = prep_t0.elapsed();
let commit_t0 = Instant::now();
cat.commit_freeze_slices("users", "by_id", slices)?;
let commit_wall = commit_t0.elapsed();
println!(" workers={workers}: prepare={prep_wall:?}, commit={commit_wall:?}");
Ok(())
}
fn measure_prepare_wall(cat: &Catalog, workers: usize, reps: usize) -> std::time::Duration {
let mut best = std::time::Duration::from_secs(u64::MAX);
for _ in 0..reps {
let ranges = partition_range(BATCH_ROWS, workers);
let t0 = Instant::now();
let _slices: Vec<FreezeSlice> = if workers == 1 {
ranges
.into_iter()
.map(|r| cat.prepare_freeze_slice("users", "by_id", r).unwrap())
.collect()
} else {
std::thread::scope(|s| {
let handles: Vec<_> = ranges
.into_iter()
.map(|r| {
s.spawn(move || cat.prepare_freeze_slice("users", "by_id", r).unwrap())
})
.collect();
handles.into_iter().map(|h| h.join().unwrap()).collect()
})
};
let elapsed = t0.elapsed();
if elapsed < best {
best = elapsed;
}
}
best
}
#[test]
#[ignore = "CPU-topology-sensitive — full tier; see gate comment"]
fn four_worker_prepare_speedup_scales() {
let _lock = crate::perf_lock();
let base = build_populated_catalog();
let t_single = measure_prepare_wall(&base, 1, PREPARE_REPS);
let t_quad = measure_prepare_wall(&base, 4, PREPARE_REPS);
let speedup = t_single.as_secs_f64() / t_quad.as_secs_f64().max(1e-9);
println!(
"perf_parallel_freezer prepare phase: \
t_single={t_single:?}, t_quad={t_quad:?}, speedup={speedup:.2}×"
);
let mut c1 = base.clone();
freeze_once(&mut c1, 1).expect("1-worker freeze");
let mut c4 = base.clone();
freeze_once(&mut c4, 4).expect("4-worker freeze");
let single_seg = c1
.cold_segment(0)
.expect("seg 0 on serial freeze")
.bytes()
.to_vec();
let quad_seg = c4
.cold_segment(0)
.expect("seg 0 on parallel freeze")
.bytes()
.to_vec();
assert_eq!(
single_seg, quad_seg,
"parallel freeze produced different segment bytes than serial freeze"
);
for id in [0i64, 1, BATCH_ROWS as i64 / 2, (BATCH_ROWS - 1) as i64] {
assert_eq!(
c1.lookup_by_pk("users", "by_id", &IndexKey::Int(id)),
c4.lookup_by_pk("users", "by_id", &IndexKey::Int(id))
);
}
let cores = std::thread::available_parallelism()
.map(|n| n.get())
.unwrap_or(2);
let threshold = if cores >= 8 {
1.4
} else if cores >= 4 {
1.15
} else {
1.05
};
assert!(
speedup >= threshold,
"speedup {speedup:.2}× < required {threshold}× on a {cores}-core host \
(t_single={t_single:?}, t_quad={t_quad:?})"
);
}