#![forbid(unsafe_code)]
use std::sync::Arc;
use std::time::Instant;
use tempfile::TempDir;
use crate::optimizer::foreground::ForegroundPolicy;
use crate::optimizer::policy::OptimizeOptions;
use crate::store::transaction::CrashHooks;
use crate::store::{NewEntry, Store, StoreConfig};
const LIVENESS_DEADLINE_SECS: u64 = 60;
#[test]
fn pool_write_path_saturation_stays_live() {
let dir = TempDir::new().unwrap();
let cfg = StoreConfig {
segment_size: 128 * 1024 * 1024,
..Default::default()
};
let _pool_guard = crate::store::workers::tests::POOL_LOCK
.lock()
.expect("pool test lock poisoned");
let store = Arc::new(Store::create(dir.path(), &cfg, [0x44; 16]).unwrap());
crate::store::workers::POOL.enable(8, 8); crate::store::workers::POOL.bind(&store);
store.enable_worker_pool();
let hooks = &CrashHooks::none();
store.set_semantic_mode(crate::dsfb::semantics::SemanticMode::Combined);
let fg = ForegroundPolicy {
pressure_enter: 0.80,
pressure_leave: 0.60,
pressure_defer_configurational: true,
pressure_max_deferred_bytes: 1024 * 1024 * 1024,
..ForegroundPolicy::focused()
};
let opts = OptimizeOptions::default();
let root = store.current_root().root_dir_ino;
let dir_ino = store
.epoch_create(root, b"sat", NewEntry::dir(0o755, 1000, 1000), hooks)
.unwrap();
let t = 16usize;
let rounds = 3usize;
let mut handles = Vec::new();
let done = Arc::new(std::sync::atomic::AtomicUsize::new(0));
for w in 0..t {
let store = Arc::clone(&store);
let done = Arc::clone(&done);
handles.push(std::thread::spawn(move || {
let hooks = &CrashHooks::none();
let mut state: u64 = 1000 + w as u64;
for r in 0..rounds {
let mut b = Vec::with_capacity(512 * 1024);
while b.len() < 512 * 1024 {
state = state
.wrapping_mul(6364136223846793005)
.wrapping_add(1442695040888963407);
b.push(
b"abcdefghijklmnopqrstuvwxyz0123456789{}();,= \n"
[((state >> 33) as usize) % 45],
);
}
let name = format!("w{w}-{r}");
let ino = store
.epoch_create(
dir_ino,
name.as_bytes(),
NewEntry::file(0o644, 1000, 1000),
hooks,
)
.unwrap();
store
.epoch_write_semantic(ino, 0, &b, opts, fg, None, hooks)
.unwrap();
}
done.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
}));
}
let deadline = Instant::now() + std::time::Duration::from_secs(LIVENESS_DEADLINE_SECS);
while done.load(std::sync::atomic::Ordering::Relaxed) < t {
assert!(
Instant::now() < deadline,
"pool write-path deadlock: {}/{} writers completed within the liveness deadline",
done.load(std::sync::atomic::Ordering::Relaxed),
t
);
std::thread::sleep(std::time::Duration::from_millis(50));
}
for h in handles {
h.join().unwrap();
}
crate::store::workers::POOL.disable();
let (de, db, _) = store.deferred_debt();
assert!(
de > 0 && db > 0,
"the pressure gate must have deferred under saturation (de={de}, db={db})"
);
}