use super::{TABLE, TestResult, migrator};
use sqlx::migrate::Migrate;
use sqlx::{ConnectOptions, Connection};
use sqlx_turso_driver::TursoConnectOptions;
use std::{
path::{Path, PathBuf},
process::{Child, Command, Stdio},
time::{Duration, Instant},
};
const CHILD: &str = "migration_process::child";
const ENV: &str = "SQLX_TURSO_D5_OWNED_FIXTURE";
const MODE: &str = "SQLX_TURSO_D5_CHILD_MODE";
const DEADLINE: Duration = Duration::from_secs(7);
struct OwnedChild {
child: Child,
fixture: PathBuf,
}
impl OwnedChild {
fn spawn(fixture: &Path, mode: &str) -> std::io::Result<Self> {
let stdout = std::fs::File::create(fixture.join(format!("{mode}.stdout")))?;
let stderr = std::fs::File::create(fixture.join(format!("{mode}.stderr")))?;
let child = Command::new(std::env::current_exe()?)
.args([
"--exact",
CHILD,
"--ignored",
"--nocapture",
"--test-threads=1",
])
.env(ENV, fixture)
.env(MODE, mode)
.stdin(Stdio::null())
.stdout(stdout)
.stderr(stderr)
.spawn()?;
assert_ne!(child.id(), std::process::id());
Ok(Self {
child,
fixture: fixture.to_owned(),
})
}
async fn ready(&mut self, name: &str) -> TestResult {
let start = Instant::now();
loop {
if let Ok(pid) = std::fs::read_to_string(self.fixture.join(name))
&& pid == self.child.id().to_string()
{
return Ok(());
}
if let Some(status) = self.child.try_wait()? {
let mode = if name == "holder-ready" {
"hold"
} else if name == "contender-rejected" {
"contend"
} else {
"after-exit"
};
let stdout = std::fs::read_to_string(self.fixture.join(format!("{mode}.stdout")))?;
let stderr = std::fs::read_to_string(self.fixture.join(format!("{mode}.stderr")))?;
return Err(format!(
"child exited before {name}: {status}; stdout: {stdout}; stderr: {stderr}"
)
.into());
}
if start.elapsed() > DEADLINE {
return Err(format!("child readiness deadline: {name}").into());
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
}
async fn exit(&mut self, code: i32) -> TestResult {
let start = Instant::now();
loop {
if let Some(status) = self.child.try_wait()? {
assert_eq!(status.code(), Some(code), "child outcome");
return Ok(());
}
if start.elapsed() > DEADLINE {
return Err("child exit deadline".into());
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
}
}
impl Drop for OwnedChild {
fn drop(&mut self) {
if matches!(self.child.try_wait(), Ok(None)) {
let _ = self.child.kill();
}
let _ = self.child.wait();
}
}
pub(super) async fn exclusion_and_abnormal_exit() -> TestResult {
let directory = tempfile::tempdir()?;
let fixture = directory.path().canonicalize()?;
let options = TursoConnectOptions::file(fixture.join("process data.db"))?;
options.connect().await?.close().await?;
let mut holder = OwnedChild::spawn(&fixture, "hold")?;
holder.ready("holder-ready").await?;
let mut competitor = OwnedChild::spawn(&fixture, "contend")?;
competitor.ready("contender-rejected").await?;
competitor.exit(0).await?;
assert_os_guard_held(&fixture)?;
let error = options.connect().await.unwrap_err();
assert!(
matches!(error, sqlx::Error::Database(_)) && error.to_string().contains("Locking error"),
"{error}"
);
std::fs::write(fixture.join("exit-abnormally"), b"owned IPC")?;
holder.exit(23).await?; let mut after = OwnedChild::spawn(&fixture, "after-exit")?;
after.ready("after-applied").await?;
after.exit(0).await?;
let mut reopened = options.connect().await?;
assert_eq!(reopened.list_applied_migrations(TABLE).await?.len(), 1);
reopened.close().await?;
eprintln!(
"D5 children: contender=0 (rejected), holder=23 (abnormal exit), after-exit=0 (applied); all reaped"
);
Ok(())
}
fn assert_os_guard_held(fixture: &Path) -> TestResult {
let file = std::fs::OpenOptions::new()
.read(true)
.write(true)
.open(fixture.join("process data.db.sqlx-migrate.lock"))?;
assert!(
!fs4::fs_std::FileExt::try_lock_exclusive(&file)?,
"OS migration guard was not held"
);
Ok(())
}
#[tokio::test]
#[ignore = "re-executed only by parent with an owned fixture"]
async fn child() -> TestResult {
let Some(fixture) = std::env::var_os(ENV) else {
return Ok(());
};
let fixture = PathBuf::from(fixture);
assert!(fixture.is_absolute());
assert_eq!(fixture.canonicalize()?, fixture);
let mode = std::env::var(MODE)?;
let options = TursoConnectOptions::file(fixture.join("process data.db"))?;
let pid = std::process::id().to_string();
match mode.as_str() {
"hold" => {
let mut conn = options.connect().await?;
conn.lock().await?;
std::fs::write(fixture.join("holder-ready"), pid)?;
let start = Instant::now();
while !fixture.join("exit-abnormally").exists() {
if start.elapsed() > DEADLINE {
return Err("holder IPC deadline".into());
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
std::process::exit(23);
}
"contend" => {
let started = Instant::now();
assert_os_guard_held(&fixture)?;
let error = options.connect().await.unwrap_err();
assert!(
matches!(error, sqlx::Error::Database(_))
&& error.to_string().contains("Locking error"),
"{error}"
);
assert!(
started.elapsed() < Duration::from_secs(2),
"lock rejection must be bounded"
);
std::fs::write(fixture.join("contender-rejected"), pid)?;
}
"after-exit" => {
let mut conn = options.connect().await?;
migrator("CREATE TABLE child_applied(id INTEGER);")
.run(&mut conn)
.await?;
assert_eq!(conn.list_applied_migrations(TABLE).await?.len(), 1);
conn.close().await?;
std::fs::write(fixture.join("after-applied"), pid)?;
}
_ => return Err("unrecognized owned child mode".into()),
}
Ok(())
}