use anyhow::{Context, Result};
use async_trait::async_trait;
use serde::{Deserialize, Serialize};
use std::collections::HashSet;
use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::time::Duration;
use tokio_rusqlite::{Connection, Row, rusqlite};
const INIT_SCHEMA_SQL: &str = "CREATE TABLE IF NOT EXISTS job_runs (
id TEXT PRIMARY KEY,
job_name TEXT NOT NULL,
parameters TEXT NOT NULL,
status TEXT NOT NULL,
created_at TEXT NOT NULL,
started_at TEXT,
finished_at TEXT,
rows_written INTEGER,
snapshot_id TEXT,
error TEXT,
submission_id TEXT
);
CREATE INDEX IF NOT EXISTS idx_job_runs_name_created
ON job_runs (job_name, created_at DESC);
CREATE INDEX IF NOT EXISTS idx_job_runs_status
ON job_runs (status);";
const SUBMISSION_ID_INDEX_SQL: &str = "DROP INDEX IF EXISTS idx_job_runs_submission_id;
CREATE UNIQUE INDEX IF NOT EXISTS idx_job_runs_submission_id_unique
ON job_runs (submission_id) WHERE submission_id IS NOT NULL;";
const ADDED_COLUMNS: [&str; 1] = ["submission_id"];
const JOBS_BUSY_TIMEOUT: Duration = Duration::from_secs(5);
fn is_duplicate_column(e: &rusqlite::Error) -> bool {
e.to_string().contains("duplicate column name")
}
fn ensure_added_columns(conn: &rusqlite::Connection, columns: &[&str]) -> rusqlite::Result<()> {
let mut stmt = conn.prepare("PRAGMA table_info(job_runs)")?;
let existing: HashSet<String> = stmt
.query_map([], |row| row.get::<_, String>(1))?
.collect::<std::result::Result<_, _>>()?;
for col in columns {
if existing.contains(*col) {
continue;
}
match conn.execute(&format!("ALTER TABLE job_runs ADD COLUMN {col} TEXT"), []) {
Ok(_) => {}
Err(e) if is_duplicate_column(&e) => {}
Err(e) => return Err(e),
}
}
Ok(())
}
fn create_private_file(path: &Path) -> Result<()> {
if path.exists() {
return Ok(());
}
let mut options = std::fs::OpenOptions::new();
options.write(true).create_new(true);
#[cfg(unix)]
{
use std::os::unix::fs::OpenOptionsExt;
options.mode(0o600);
}
match options.open(path) {
Ok(_) => Ok(()),
Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => Ok(()),
Err(e) => {
Err(e).with_context(|| format!("Failed to create jobs.db file: {}", path.display()))
}
}
}
fn restrict_permissions(path: &Path) -> Result<()> {
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o600)).with_context(
|| {
format!(
"Failed to restrict jobs.db file permissions: {}",
path.display()
)
},
)?;
}
#[cfg(not(unix))]
let _ = path;
Ok(())
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum JobRunStatus {
Pending,
Running,
Succeeded,
Failed,
Cancelled,
}
impl JobRunStatus {
pub fn as_str(&self) -> &'static str {
match self {
Self::Pending => "pending",
Self::Running => "running",
Self::Succeeded => "succeeded",
Self::Failed => "failed",
Self::Cancelled => "cancelled",
}
}
pub fn from_str(s: &str) -> Result<Self> {
Ok(match s {
"pending" => Self::Pending,
"running" => Self::Running,
"succeeded" => Self::Succeeded,
"failed" => Self::Failed,
"cancelled" => Self::Cancelled,
other => anyhow::bail!("unknown job status: {other}"),
})
}
pub fn is_terminal(&self) -> bool {
matches!(self, Self::Succeeded | Self::Failed | Self::Cancelled)
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct JobRun {
pub id: String,
pub job_name: String,
pub parameters: String,
pub status: JobRunStatus,
pub created_at: String,
pub started_at: Option<String>,
pub finished_at: Option<String>,
pub rows_written: Option<u64>,
pub snapshot_id: Option<String>,
pub error: Option<String>,
pub submission_id: Option<String>,
}
#[async_trait]
pub trait JobStore: Send + Sync {
async fn create_run(&self, run: &JobRun) -> Result<()>;
async fn get_run(&self, run_id: &str) -> Result<Option<JobRun>>;
async fn get_run_by_submission_id(&self, submission_id: &str) -> Result<Option<JobRun>>;
async fn list_runs(&self, job_name: Option<&str>, limit: usize) -> Result<Vec<JobRun>>;
async fn update_status(
&self,
run_id: &str,
status: JobRunStatus,
started_at: Option<String>,
finished_at: Option<String>,
rows_written: Option<u64>,
snapshot_id: Option<String>,
error: Option<String>,
) -> Result<()>;
async fn reconcile_orphaned(&self, reason: &str) -> Result<usize>;
}
pub struct SqliteJobStore {
conn: Arc<Connection>,
path: PathBuf,
}
impl std::fmt::Debug for SqliteJobStore {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("SqliteJobStore")
.field("path", &self.path)
.finish()
}
}
impl SqliteJobStore {
pub async fn open(path: impl AsRef<Path>) -> Result<Self> {
let path = path.as_ref().to_path_buf();
if let Some(parent) = path.parent() {
std::fs::create_dir_all(parent)
.with_context(|| format!("Failed to create jobs.db parent dir: {parent:?}"))?;
}
create_private_file(&path)?;
let conn = Connection::open(&path)
.await
.with_context(|| format!("Failed to open jobs.db: {path:?}"))?;
conn.call(|conn| -> std::result::Result<(), rusqlite::Error> {
conn.pragma_update(None, "journal_mode", "WAL")?;
conn.busy_timeout(JOBS_BUSY_TIMEOUT)?;
Ok(())
})
.await
.with_context(|| format!("Failed to configure jobs.db: {}", path.display()))?;
let store = Self {
conn: Arc::new(conn),
path,
};
store.ensure_schema().await?;
restrict_permissions(&store.path)?;
for suffix in ["-wal", "-shm"] {
let mut sidecar = store.path.clone().into_os_string();
sidecar.push(suffix);
let sidecar = PathBuf::from(sidecar);
if sidecar.exists() {
restrict_permissions(&sidecar)?;
}
}
Ok(store)
}
pub async fn open_in_memory() -> Result<Self> {
let conn = Connection::open(":memory:")
.await
.context("Failed to open in-memory jobs.db")?;
let store = Self {
conn: Arc::new(conn),
path: PathBuf::from(":memory:"),
};
store.ensure_schema().await?;
Ok(store)
}
pub fn path(&self) -> &Path {
&self.path
}
async fn ensure_schema(&self) -> Result<()> {
self.conn
.call(|conn| -> std::result::Result<(), rusqlite::Error> {
conn.execute_batch(INIT_SCHEMA_SQL)?;
Ok(())
})
.await
.map_err(|e| anyhow::anyhow!("Failed to initialize jobs.db schema: {e}"))?;
self.conn
.call(|conn| -> std::result::Result<(), rusqlite::Error> {
ensure_added_columns(conn, &ADDED_COLUMNS)?;
conn.execute_batch(SUBMISSION_ID_INDEX_SQL)?;
Ok(())
})
.await
.map_err(|e| {
anyhow::anyhow!("Failed to migrate jobs.db schema (submission_id): {e}")
})?;
Ok(())
}
}
fn row_to_job_run(row: &Row<'_>) -> rusqlite::Result<JobRun> {
let status_str: String = row.get("status")?;
let status = JobRunStatus::from_str(&status_str).map_err(|e| {
rusqlite::Error::FromSqlConversionFailure(
3,
rusqlite::types::Type::Text,
Box::new(std::io::Error::other(e.to_string())),
)
})?;
Ok(JobRun {
id: row.get("id")?,
job_name: row.get("job_name")?,
parameters: row.get("parameters")?,
status,
created_at: row.get("created_at")?,
started_at: row.get("started_at")?,
finished_at: row.get("finished_at")?,
rows_written: row.get::<_, Option<i64>>("rows_written")?.map(|v| v as u64),
snapshot_id: row.get("snapshot_id")?,
error: row.get("error")?,
submission_id: row.get("submission_id")?,
})
}
#[async_trait]
impl JobStore for SqliteJobStore {
async fn create_run(&self, run: &JobRun) -> Result<()> {
let run = run.clone();
self.conn
.call(move |conn| -> std::result::Result<(), rusqlite::Error> {
conn.execute(
"INSERT INTO job_runs
(id, job_name, parameters, status, created_at, started_at,
finished_at, rows_written, snapshot_id, error, submission_id)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11)",
rusqlite::params![
run.id,
run.job_name,
run.parameters,
run.status.as_str(),
run.created_at,
run.started_at,
run.finished_at,
run.rows_written.map(|v| v as i64),
run.snapshot_id,
run.error,
run.submission_id,
],
)?;
Ok(())
})
.await
.map_err(|e| anyhow::anyhow!("create_run failed: {e}"))?;
Ok(())
}
async fn get_run(&self, run_id: &str) -> Result<Option<JobRun>> {
let run_id = run_id.to_string();
let row = self
.conn
.call(
move |conn| -> std::result::Result<Option<JobRun>, rusqlite::Error> {
let mut stmt = conn.prepare(
"SELECT id, job_name, parameters, status, created_at, started_at,
finished_at, rows_written, snapshot_id, error, submission_id
FROM job_runs
WHERE id = ?1",
)?;
let mut rows = stmt.query(rusqlite::params![run_id])?;
match rows.next()? {
Some(row) => Ok(Some(row_to_job_run(row)?)),
None => Ok(None),
}
},
)
.await
.map_err(|e| anyhow::anyhow!("get_run failed: {e}"))?;
Ok(row)
}
async fn get_run_by_submission_id(&self, submission_id: &str) -> Result<Option<JobRun>> {
let submission_id = submission_id.to_string();
let row = self
.conn
.call(
move |conn| -> std::result::Result<Option<JobRun>, rusqlite::Error> {
let mut stmt = conn.prepare(
"SELECT id, job_name, parameters, status, created_at, started_at,
finished_at, rows_written, snapshot_id, error, submission_id
FROM job_runs
WHERE submission_id = ?1
ORDER BY created_at DESC
LIMIT 1",
)?;
let mut rows = stmt.query(rusqlite::params![submission_id])?;
match rows.next()? {
Some(row) => Ok(Some(row_to_job_run(row)?)),
None => Ok(None),
}
},
)
.await
.map_err(|e| anyhow::anyhow!("get_run_by_submission_id failed: {e}"))?;
Ok(row)
}
async fn list_runs(&self, job_name: Option<&str>, limit: usize) -> Result<Vec<JobRun>> {
let job_name = job_name.map(|s| s.to_string());
let limit = limit as i64;
let rows = self
.conn
.call(
move |conn| -> std::result::Result<Vec<JobRun>, rusqlite::Error> {
let (sql, params): (&str, Vec<rusqlite::types::Value>) = match &job_name {
Some(name) => (
"SELECT id, job_name, parameters, status, created_at, started_at,
finished_at, rows_written, snapshot_id, error, submission_id
FROM job_runs
WHERE job_name = ?1
ORDER BY created_at DESC
LIMIT ?2",
vec![
rusqlite::types::Value::Text(name.clone()),
rusqlite::types::Value::Integer(limit),
],
),
None => (
"SELECT id, job_name, parameters, status, created_at, started_at,
finished_at, rows_written, snapshot_id, error, submission_id
FROM job_runs
ORDER BY created_at DESC
LIMIT ?1",
vec![rusqlite::types::Value::Integer(limit)],
),
};
let mut stmt = conn.prepare(sql)?;
let rows = stmt
.query_map(rusqlite::params_from_iter(params), |row| {
row_to_job_run(row)
})?
.collect::<rusqlite::Result<Vec<_>>>()?;
Ok(rows)
},
)
.await
.map_err(|e| anyhow::anyhow!("list_runs failed: {e}"))?;
Ok(rows)
}
async fn update_status(
&self,
run_id: &str,
status: JobRunStatus,
started_at: Option<String>,
finished_at: Option<String>,
rows_written: Option<u64>,
snapshot_id: Option<String>,
error: Option<String>,
) -> Result<()> {
let run_id = run_id.to_string();
let status_str = status.as_str().to_string();
self.conn
.call(move |conn| -> std::result::Result<(), rusqlite::Error> {
conn.execute(
"UPDATE job_runs
SET status = ?2,
started_at = COALESCE(?3, started_at),
finished_at = COALESCE(?4, finished_at),
rows_written = COALESCE(?5, rows_written),
snapshot_id = COALESCE(?6, snapshot_id),
error = COALESCE(?7, error)
WHERE id = ?1",
rusqlite::params![
run_id,
status_str,
started_at,
finished_at,
rows_written.map(|v| v as i64),
snapshot_id,
error,
],
)?;
Ok(())
})
.await
.map_err(|e| anyhow::anyhow!("update_status failed: {e}"))?;
Ok(())
}
async fn reconcile_orphaned(&self, reason: &str) -> Result<usize> {
let reason = reason.to_string();
let updated = self
.conn
.call(move |conn| -> std::result::Result<usize, rusqlite::Error> {
let now = chrono::Utc::now().to_rfc3339();
let n = conn.execute(
"UPDATE job_runs
SET status = 'failed',
finished_at = ?1,
error = ?2
WHERE status IN ('pending', 'running')",
rusqlite::params![now, reason],
)?;
Ok(n)
})
.await
.map_err(|e| anyhow::anyhow!("reconcile_orphaned failed: {e}"))?;
Ok(updated)
}
}
#[cfg(test)]
mod tests {
use super::*;
fn sample_run(id: &str, job_name: &str) -> JobRun {
sample_run_with_submission(id, job_name, None)
}
fn sample_run_with_submission(id: &str, job_name: &str, submission_id: Option<&str>) -> JobRun {
JobRun {
id: id.to_string(),
job_name: job_name.to_string(),
parameters: r#"{"from_date":"2026-01-01"}"#.to_string(),
status: JobRunStatus::Pending,
created_at: "2026-04-21T00:00:00Z".to_string(),
started_at: None,
finished_at: None,
rows_written: None,
snapshot_id: None,
error: None,
submission_id: submission_id.map(str::to_string),
}
}
#[tokio::test]
async fn status_round_trip_strings() {
for status in [
JobRunStatus::Pending,
JobRunStatus::Running,
JobRunStatus::Succeeded,
JobRunStatus::Failed,
JobRunStatus::Cancelled,
] {
assert_eq!(JobRunStatus::from_str(status.as_str()).unwrap(), status);
}
assert!(JobRunStatus::from_str("not-a-status").is_err());
}
#[tokio::test]
async fn create_get_list_round_trip() {
let store = SqliteJobStore::open_in_memory().await.unwrap();
store.create_run(&sample_run("r1", "ingest")).await.unwrap();
store.create_run(&sample_run("r2", "ingest")).await.unwrap();
store.create_run(&sample_run("r3", "other")).await.unwrap();
let got = store.get_run("r1").await.unwrap().unwrap();
assert_eq!(got.job_name, "ingest");
assert_eq!(got.status, JobRunStatus::Pending);
let by_name = store.list_runs(Some("ingest"), 10).await.unwrap();
assert_eq!(by_name.len(), 2);
assert!(by_name.iter().all(|r| r.job_name == "ingest"));
let all = store.list_runs(None, 10).await.unwrap();
assert_eq!(all.len(), 3);
}
#[tokio::test]
async fn update_status_progresses_row() {
let store = SqliteJobStore::open_in_memory().await.unwrap();
store.create_run(&sample_run("r1", "ingest")).await.unwrap();
store
.update_status(
"r1",
JobRunStatus::Running,
Some("2026-04-21T00:01:00Z".to_string()),
None,
None,
None,
None,
)
.await
.unwrap();
let got = store.get_run("r1").await.unwrap().unwrap();
assert_eq!(got.status, JobRunStatus::Running);
assert_eq!(got.started_at.as_deref(), Some("2026-04-21T00:01:00Z"));
store
.update_status(
"r1",
JobRunStatus::Succeeded,
None,
Some("2026-04-21T00:02:00Z".to_string()),
Some(123),
Some("7".to_string()),
None,
)
.await
.unwrap();
let got = store.get_run("r1").await.unwrap().unwrap();
assert_eq!(got.status, JobRunStatus::Succeeded);
assert_eq!(got.rows_written, Some(123));
assert_eq!(got.snapshot_id.as_deref(), Some("7"));
}
#[tokio::test]
async fn reconcile_orphaned_marks_non_terminal_rows_failed() {
let store = SqliteJobStore::open_in_memory().await.unwrap();
let mut r_pending = sample_run("r-pending", "ingest");
r_pending.status = JobRunStatus::Pending;
store.create_run(&r_pending).await.unwrap();
let mut r_running = sample_run("r-running", "ingest");
r_running.status = JobRunStatus::Running;
store.create_run(&r_running).await.unwrap();
let mut r_done = sample_run("r-done", "ingest");
r_done.status = JobRunStatus::Succeeded;
store.create_run(&r_done).await.unwrap();
let updated = store
.reconcile_orphaned("server restart mid-run")
.await
.unwrap();
assert_eq!(updated, 2, "pending + running rows should be reconciled");
assert_eq!(
store.get_run("r-pending").await.unwrap().unwrap().status,
JobRunStatus::Failed
);
assert_eq!(
store.get_run("r-running").await.unwrap().unwrap().status,
JobRunStatus::Failed
);
assert_eq!(
store.get_run("r-done").await.unwrap().unwrap().status,
JobRunStatus::Succeeded
);
}
#[tokio::test]
async fn persists_across_opens() {
let tmp = tempfile::NamedTempFile::new().unwrap();
let p = tmp.path().to_path_buf();
drop(tmp);
{
let s = SqliteJobStore::open(&p).await.unwrap();
s.create_run(&sample_run("rX", "ingest")).await.unwrap();
}
{
let s = SqliteJobStore::open(&p).await.unwrap();
let got = s.get_run("rX").await.unwrap().unwrap();
assert_eq!(got.job_name, "ingest");
}
}
#[tokio::test]
async fn submission_id_round_trips_and_looks_up_in_reverse() {
let s = SqliteJobStore::open_in_memory().await.unwrap();
s.create_run(&sample_run_with_submission(
"r1",
"ingest",
Some("audit-abc"),
))
.await
.unwrap();
assert_eq!(
s.get_run("r1").await.unwrap().unwrap().submission_id,
Some("audit-abc".to_string())
);
let found = s
.get_run_by_submission_id("audit-abc")
.await
.unwrap()
.expect("submission token did not resolve to its run");
assert_eq!(found.id, "r1");
assert!(
s.get_run_by_submission_id("audit-nope")
.await
.unwrap()
.is_none()
);
}
#[tokio::test]
async fn unattributed_runs_keep_a_null_submission_id() {
let s = SqliteJobStore::open_in_memory().await.unwrap();
s.create_run(&sample_run("r1", "ingest")).await.unwrap();
s.create_run(&sample_run("r2", "ingest")).await.unwrap();
assert!(
s.get_run("r1")
.await
.unwrap()
.unwrap()
.submission_id
.is_none()
);
assert_eq!(s.list_runs(None, 10).await.unwrap().len(), 2);
}
#[tokio::test]
async fn open_migrates_a_ledger_written_before_submission_id() {
let tmp = tempfile::NamedTempFile::new().unwrap();
let path = tmp.path().to_path_buf();
drop(tmp);
{
let conn = rusqlite::Connection::open(&path).unwrap();
conn.execute_batch(
"CREATE TABLE job_runs (
id TEXT PRIMARY KEY, job_name TEXT NOT NULL,
parameters TEXT NOT NULL, status TEXT NOT NULL,
created_at TEXT NOT NULL, started_at TEXT, finished_at TEXT,
rows_written INTEGER, snapshot_id TEXT, error TEXT);
INSERT INTO job_runs (id, job_name, parameters, status, created_at)
VALUES ('old-run', 'ingest', '{}', 'succeeded', '2026-01-01T00:00:00Z');",
)
.unwrap();
}
let s = SqliteJobStore::open(&path).await.unwrap();
let old = s.get_run("old-run").await.unwrap().unwrap();
assert_eq!(old.job_name, "ingest");
assert!(old.submission_id.is_none());
s.create_run(&sample_run_with_submission(
"new-run",
"ingest",
Some("audit-1"),
))
.await
.unwrap();
assert_eq!(
s.get_run_by_submission_id("audit-1")
.await
.unwrap()
.unwrap()
.id,
"new-run"
);
let (has_unique, has_legacy) = s
.conn
.call(
|conn| -> std::result::Result<(bool, bool), rusqlite::Error> {
let unique = conn
.prepare(
"SELECT 1 FROM sqlite_master
WHERE type = 'index'
AND name = 'idx_job_runs_submission_id_unique'",
)?
.exists([])?;
let legacy = conn
.prepare(
"SELECT 1 FROM sqlite_master
WHERE type = 'index' AND name = 'idx_job_runs_submission_id'",
)?
.exists([])?;
Ok((unique, legacy))
},
)
.await
.unwrap();
assert!(
has_unique,
"migrated ledger is missing the unique submission_id index"
);
assert!(
!has_legacy,
"the pre-unique index survived the migration, so uniqueness is not enforced"
);
}
#[tokio::test]
async fn a_duplicate_submission_id_is_rejected_rather_than_shadowing_a_run() {
let s = SqliteJobStore::open_in_memory().await.unwrap();
s.create_run(&sample_run_with_submission("r1", "ingest", Some("audit-1")))
.await
.unwrap();
let err = s
.create_run(&sample_run_with_submission("r2", "ingest", Some("audit-1")))
.await
.unwrap_err();
assert!(
err.to_string().contains("UNIQUE constraint failed"),
"expected a uniqueness failure, got {err}"
);
assert_eq!(
s.get_run_by_submission_id("audit-1")
.await
.unwrap()
.unwrap()
.id,
"r1"
);
}
#[tokio::test]
async fn many_unattributed_runs_do_not_collide_under_the_unique_index() {
let s = SqliteJobStore::open_in_memory().await.unwrap();
for i in 0..5 {
s.create_run(&sample_run_with_submission(
&format!("r{i}"),
"ingest",
None,
))
.await
.unwrap();
}
assert_eq!(s.list_runs(None, 10).await.unwrap().len(), 5);
}
#[tokio::test]
async fn open_is_idempotent_across_the_unique_index_rename() {
let tmp = tempfile::NamedTempFile::new().unwrap();
let path = tmp.path().to_path_buf();
drop(tmp);
{
let conn = rusqlite::Connection::open(&path).unwrap();
conn.execute_batch(
"CREATE TABLE job_runs (
id TEXT PRIMARY KEY, job_name TEXT NOT NULL,
parameters TEXT NOT NULL, status TEXT NOT NULL,
created_at TEXT NOT NULL, started_at TEXT, finished_at TEXT,
rows_written INTEGER, snapshot_id TEXT, error TEXT,
submission_id TEXT);
CREATE INDEX idx_job_runs_submission_id
ON job_runs (submission_id) WHERE submission_id IS NOT NULL;",
)
.unwrap();
}
for _ in 0..2 {
let s = SqliteJobStore::open(&path).await.unwrap();
drop(s);
}
let s = SqliteJobStore::open(&path).await.unwrap();
s.create_run(&sample_run_with_submission("r1", "ingest", Some("audit-1")))
.await
.unwrap();
let err = s
.create_run(&sample_run_with_submission("r2", "ingest", Some("audit-1")))
.await
.unwrap_err();
assert!(
err.to_string().contains("UNIQUE constraint failed"),
"the rename left the ledger unconstrained: {err}"
);
}
#[tokio::test]
async fn open_restricts_permissions_and_sets_the_concurrency_pragmas() {
let tmp = tempfile::NamedTempFile::new().unwrap();
let path = tmp.path().to_path_buf();
drop(tmp);
let s = SqliteJobStore::open(&path).await.unwrap();
s.create_run(&sample_run("r1", "ingest")).await.unwrap();
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
let mode = std::fs::metadata(&path).unwrap().permissions().mode() & 0o777;
assert_eq!(mode, 0o600, "jobs.db is not owner-only: {mode:o}");
for suffix in ["-wal", "-shm"] {
let mut sidecar = path.clone().into_os_string();
sidecar.push(suffix);
let sidecar = PathBuf::from(sidecar);
if sidecar.exists() {
let mode = std::fs::metadata(&sidecar).unwrap().permissions().mode() & 0o777;
assert_eq!(mode, 0o600, "{sidecar:?} is not owner-only: {mode:o}");
}
}
}
let (journal_mode, busy_timeout) = s
.conn
.call(
|conn| -> std::result::Result<(String, i64), rusqlite::Error> {
let journal: String =
conn.query_row("PRAGMA journal_mode", [], |r| r.get(0))?;
let busy: i64 = conn.query_row("PRAGMA busy_timeout", [], |r| r.get(0))?;
Ok((journal, busy))
},
)
.await
.unwrap();
assert_eq!(journal_mode.to_lowercase(), "wal");
assert_eq!(busy_timeout, JOBS_BUSY_TIMEOUT.as_millis() as i64);
}
#[test]
fn duplicate_column_error_is_recognised() {
let conn = rusqlite::Connection::open_in_memory().unwrap();
conn.execute_batch("CREATE TABLE job_runs (id TEXT PRIMARY KEY);")
.unwrap();
conn.execute("ALTER TABLE job_runs ADD COLUMN submission_id TEXT", [])
.unwrap();
let err = conn
.execute("ALTER TABLE job_runs ADD COLUMN submission_id TEXT", [])
.expect_err("adding an existing column must fail");
assert!(is_duplicate_column(&err), "unrecognised: {err}");
let other = conn
.execute("ALTER TABLE nonexistent ADD COLUMN submission_id TEXT", [])
.expect_err("altering a missing table must fail");
assert!(!is_duplicate_column(&other), "over-matched: {other}");
}
}