#![cfg(feature = "failpoints")]
use std::sync::Arc;
use std::sync::mpsc::sync_channel;
use std::time::Duration;
use tsoracle_consensus::{ConsensusDriver, ConsensusError};
use tsoracle_core::Epoch;
use tsoracle_driver_file::FileDriver;
static FAILPOINT_TEST_SERIAL: tokio::sync::Mutex<()> = tokio::sync::Mutex::const_new(());
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn load_is_not_blocked_by_in_flight_persist() {
let _serial = FAILPOINT_TEST_SERIAL.lock().await;
let _scenario = tsoracle_failpoint::fail::FailScenario::setup();
let dir = tempfile::tempdir().unwrap();
let driver = FileDriver::open_or_init(dir.path()).unwrap();
let (entered_tx, entered_rx) = sync_channel::<()>(1);
let (release_tx, release_rx) = sync_channel::<()>(1);
let release_rx = std::sync::Mutex::new(release_rx);
tsoracle_failpoint::fail::cfg_callback("file_driver::write_blocked", move || {
entered_tx.send(()).unwrap();
release_rx.lock().unwrap().recv().unwrap();
})
.unwrap();
let writer = {
let driver = Arc::clone(&driver);
tokio::spawn(async move { driver.persist_high_water(100, Epoch::ZERO).await })
};
tokio::task::spawn_blocking(move || entered_rx.recv().unwrap())
.await
.unwrap();
let read = tokio::time::timeout(Duration::from_millis(500), driver.load_high_water())
.await
.expect("load_high_water was blocked by an in-flight persist");
assert_eq!(
read.unwrap(),
0,
"writer has not published yet; reader should see the prior value"
);
release_tx.send(()).unwrap();
let persisted = writer.await.unwrap().unwrap();
assert_eq!(persisted, 100);
assert_eq!(driver.load_high_water().await.unwrap(), 100);
}
#[tokio::test]
async fn reopen_after_crash_before_write_returns_prior_high_water() {
let _serial = FAILPOINT_TEST_SERIAL.lock().await;
let _scenario = tsoracle_failpoint::fail::FailScenario::setup();
let dir = tempfile::tempdir().unwrap();
{
let driver = FileDriver::open_or_init(dir.path()).unwrap();
driver.persist_high_water(100, Epoch::ZERO).await.unwrap();
tsoracle_failpoint::fail::cfg("file_driver::before_write", "return(io)").unwrap();
let result = driver.persist_high_water(200, Epoch::ZERO).await;
assert!(
matches!(result, Err(ConsensusError::PermanentDriver(_))),
"expected PermanentDriver, got {result:?}"
);
tsoracle_failpoint::fail::cfg("file_driver::before_write", "off").unwrap();
}
let reopened = FileDriver::open_or_init(dir.path()).unwrap();
assert_eq!(
reopened.load_high_water().await.unwrap(),
100,
"reopen should see prior high-water; the failed persist must not have rewritten state"
);
}
#[tokio::test]
async fn reopen_after_crash_between_tmp_fsync_and_rename_returns_prior_high_water() {
let _serial = FAILPOINT_TEST_SERIAL.lock().await;
let _scenario = tsoracle_failpoint::fail::FailScenario::setup();
let dir = tempfile::tempdir().unwrap();
{
let driver = FileDriver::open_or_init(dir.path()).unwrap();
driver.persist_high_water(100, Epoch::ZERO).await.unwrap();
tsoracle_failpoint::fail::cfg("file_driver::after_tmp_fsync_before_rename", "return(io)")
.unwrap();
let result = driver.persist_high_water(200, Epoch::ZERO).await;
assert!(
matches!(result, Err(ConsensusError::PermanentDriver(_))),
"expected PermanentDriver, got {result:?}"
);
tsoracle_failpoint::fail::cfg("file_driver::after_tmp_fsync_before_rename", "off").unwrap();
}
let reopened = FileDriver::open_or_init(dir.path()).unwrap();
assert_eq!(reopened.load_high_water().await.unwrap(), 100);
}
#[tokio::test]
async fn reopen_after_crash_between_rename_and_dir_fsync_is_monotonic() {
let _serial = FAILPOINT_TEST_SERIAL.lock().await;
let _scenario = tsoracle_failpoint::fail::FailScenario::setup();
let dir = tempfile::tempdir().unwrap();
{
let driver = FileDriver::open_or_init(dir.path()).unwrap();
driver.persist_high_water(100, Epoch::ZERO).await.unwrap();
tsoracle_failpoint::fail::cfg("file_driver::after_rename_before_dir_fsync", "panic")
.unwrap();
let result = driver.persist_high_water(200, Epoch::ZERO).await;
assert!(
matches!(result, Err(ConsensusError::PermanentDriver(_))),
"spawn_blocking panic should surface as PermanentDriver, got {result:?}"
);
tsoracle_failpoint::fail::cfg("file_driver::after_rename_before_dir_fsync", "off").unwrap();
}
let reopened = FileDriver::open_or_init(dir.path()).unwrap();
let loaded = reopened.load_high_water().await.unwrap();
assert!(
(100..=200).contains(&loaded),
"expected loaded value in 100..=200, got {loaded}"
);
}