sqlx-turso-driver 0.0.1

An asynchronous SQLx driver for embedded Turso databases
//! Re-executes only this integration harness's exact child test. No Cargo recursion or helper binary.
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 {
            // A file can be visible before its PID write finishes. Accept only
            // the complete expected child PID, otherwise keep the bounded wait.
            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) {
        // Only this Command-owned PID is ever killed; always reap on failure/cancellation.
        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?;
    // The pinned SDK itself rejects another process at open. Separately verify
    // the driver's real OS guard, not only the SDK's engine file exclusion.
    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?; // process::exit bypasses destructors: actual OS release
    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(())
}