use std::fs;
use std::path::{Path, PathBuf};
use crate::__bypass::RawAccessExt as _;
use crate::context::{DjogiContext, PinnedCtx};
use crate::error::DjogiError;
use super::ledger::compute_checksum;
use super::projection::BucketKey;
use super::runner::{advisory_lock_key, release_advisory_lock};
pub const SEED_LEDGER_TABLE_DDL: &str = "\
CREATE TABLE IF NOT EXISTS djogi_seed_runs (
id BIGINT GENERATED BY DEFAULT AS IDENTITY PRIMARY KEY,
seed_name TEXT NOT NULL UNIQUE,
checksum_up VARCHAR(68) NOT NULL,
status TEXT NOT NULL DEFAULT 'applied',
applied_at TIMESTAMPTZ NOT NULL DEFAULT now(),
applied_by TEXT NOT NULL DEFAULT current_user,
failure_note TEXT
)";
const SEED_RUN_LOCK_APP_LABEL: &str = "__djogi_seed_run__";
pub const SEEDS_DIRNAME: &str = "seeds";
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct DiscoveredSeed {
pub path: PathBuf,
pub seed_name: String,
}
#[derive(Debug)]
pub enum SeedError {
Io {
path: PathBuf,
source: std::io::Error,
},
LedgerWrite { source: DjogiError },
LedgerDecode {
column: String,
source: tokio_postgres::Error,
},
ApplyFailed {
seed_name: String,
source: DjogiError,
},
ChecksumDrift {
seed_name: String,
ledger_checksum: String,
on_disk_checksum: String,
},
LocalhostGate { database_url: String },
MalformedApplicationUrl { application_url: String },
RunAlreadyInProgress { database: String },
StaleClaim {
seed_name: String,
checksum_up: String,
},
FailedClaim {
seed_name: String,
failure_note: Option<String>,
},
LedgerStatusUnknown { value: String },
LockQueryFailed {
database: String,
source: DjogiError,
},
AdvisoryUnlockReturnedFalse { database: String, key: i64 },
PinnedSessionCheckoutFailed { source: DjogiError },
}
impl std::fmt::Display for SeedError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
SeedError::Io { path, source } => {
write!(f, "seed I/O at {}: {source}", path.display())
}
SeedError::LedgerWrite { source } => {
write!(f, "seed ledger write failed: {source}")
}
SeedError::LedgerDecode { column, source } => write!(
f,
"seed ledger row decode failed for column `{column}`: {source}; \
the column is present in `djogi_seed_runs` but the bytes do not \
deserialise into the expected type — typically the row was \
written by an older Djogi or hand-edited",
),
SeedError::ApplyFailed { seed_name, source } => write!(
f,
"seed `{seed_name}` failed to apply: {source}; \
previously-applied seeds are NOT rolled back",
),
SeedError::ChecksumDrift {
seed_name,
ledger_checksum,
on_disk_checksum,
} => write!(
f,
"seed `{seed_name}` was edited after it was first applied: \
ledger has checksum `{ledger_checksum}` but the file on disk \
hashes to `{on_disk_checksum}`; either revert the edit or \
delete the matching `djogi_seed_runs` row before re-running",
),
SeedError::LocalhostGate { database_url } => write!(
f,
"djogi db seed refuses to run when DATABASE_URL is not \
localhost (got `{database_url}`); pass `--allow-non-localhost` \
to override (typical use: a CI integration suite seeding a \
remote test database)",
),
SeedError::MalformedApplicationUrl { application_url } => write!(
f,
"djogi db seed: application URL `{application_url}` has no \
database-name component the runner can splice `--database \
<name>` into; the URL must be of the form \
`postgres://<authority>/<database>` so per-database routing \
can derive the connection target",
),
SeedError::RunAlreadyInProgress { database } => write!(
f,
"djogi db seed refused to run for database `{database}` because another \
seed runner already holds the per-database advisory lock",
),
SeedError::StaleClaim {
seed_name,
checksum_up,
} => write!(
f,
"seed `{seed_name}` has a stale `running` claim with checksum \
`{checksum_up}`; the previous run may have executed SQL but did not \
durably finalise the ledger row, so automatic re-execution is refused",
),
SeedError::FailedClaim {
seed_name,
failure_note,
} => match failure_note {
Some(note) if !note.is_empty() => write!(
f,
"seed `{seed_name}` is blocked by a prior failed claim: {note}; \
repair or delete the claim row before retrying",
),
_ => write!(
f,
"seed `{seed_name}` is blocked by a prior failed claim; repair or \
delete the claim row before retrying",
),
},
SeedError::LedgerStatusUnknown { value } => write!(
f,
"seed ledger row carries unknown status `{value}`; the djogi_seed_runs \
table was likely hand-edited or written by incompatible code",
),
SeedError::LockQueryFailed { database, source } => write!(
f,
"seed-run advisory lock query failed for database `{database}`: {source}",
),
SeedError::AdvisoryUnlockReturnedFalse { database, key } => write!(
f,
"seed-run advisory unlock returned false for database `{database}` \
(key=0x{key:016x}); the lock was not held on this session",
),
SeedError::PinnedSessionCheckoutFailed { source } => write!(
f,
"seed runner failed to check out a pinned Postgres session from the pool: \
{source}",
),
}
}
}
impl std::error::Error for SeedError {
fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
match self {
SeedError::Io { source, .. } => Some(source),
SeedError::LedgerWrite { source } => Some(source),
SeedError::LedgerDecode { source, .. } => Some(source),
SeedError::ApplyFailed { source, .. } => Some(source),
SeedError::LockQueryFailed { source, .. } => Some(source),
SeedError::PinnedSessionCheckoutFailed { source } => Some(source),
_ => None,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum SeedLedgerStatus {
Running,
Applied,
Failed,
}
impl SeedLedgerStatus {
const fn as_db_str(self) -> &'static str {
match self {
SeedLedgerStatus::Running => "running",
SeedLedgerStatus::Applied => "applied",
SeedLedgerStatus::Failed => "failed",
}
}
#[allow(clippy::result_large_err)]
fn from_db_str(value: &str) -> Result<Self, SeedError> {
match value {
"running" => Ok(Self::Running),
"applied" => Ok(Self::Applied),
"failed" => Ok(Self::Failed),
other => Err(SeedError::LedgerStatusUnknown {
value: other.to_string(),
}),
}
}
}
#[derive(Debug, Clone)]
struct SeedLedgerRow {
checksum_up: String,
status: SeedLedgerStatus,
failure_note: Option<String>,
}
fn seed_run_lock_bucket(database: &str) -> BucketKey {
BucketKey {
database: database.to_string(),
app: SEED_RUN_LOCK_APP_LABEL.to_string(),
}
}
#[derive(Debug, Clone)]
pub struct SeedReport {
pub entries: Vec<SeedReportEntry>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct SeedReportEntry {
pub seed_name: String,
pub outcome: SeedOutcome,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum SeedOutcome {
Applied,
SkippedAlreadyApplied,
}
pub fn derive_per_database_url(application_url: &str, database: &str) -> Option<String> {
super::reset::replace_db_in_url(application_url, database)
}
#[allow(clippy::result_large_err)]
pub fn discover_seeds(
workspace_root: &Path,
database: &str,
) -> Result<Vec<DiscoveredSeed>, SeedError> {
let dir = workspace_root.join(SEEDS_DIRNAME).join(database);
let mut out: Vec<DiscoveredSeed> = Vec::new();
let entries = match fs::read_dir(&dir) {
Ok(e) => e,
Err(err) if err.kind() == std::io::ErrorKind::NotFound => return Ok(out),
Err(err) => {
return Err(SeedError::Io {
path: dir,
source: err,
});
}
};
for entry in entries {
let entry = entry.map_err(|err| SeedError::Io {
path: dir.clone(),
source: err,
})?;
let ft = entry.file_type().map_err(|err| SeedError::Io {
path: entry.path(),
source: err,
})?;
if !ft.is_file() {
continue;
}
let Some(name) = entry.file_name().to_str().map(str::to_string) else {
continue;
};
if name.as_bytes().first() == Some(&b'.') {
continue;
}
let Some(stem) = name.strip_suffix(".sql") else {
continue;
};
if stem.is_empty() {
continue;
}
out.push(DiscoveredSeed {
path: entry.path(),
seed_name: stem.to_string(),
});
}
out.sort_by(|a, b| a.seed_name.cmp(&b.seed_name));
Ok(out)
}
pub async fn bootstrap(ctx: &mut DjogiContext) -> Result<(), DjogiError> {
ctx.raw_ddl(SEED_LEDGER_TABLE_DDL).await?;
ctx.raw_ddl(
"ALTER TABLE djogi_seed_runs \
ADD COLUMN IF NOT EXISTS status TEXT NOT NULL DEFAULT 'applied'",
)
.await?;
ctx.raw_ddl(
"ALTER TABLE djogi_seed_runs \
ADD COLUMN IF NOT EXISTS failure_note TEXT",
)
.await?;
Ok(())
}
async fn fetch_seed_ledger_row(
ctx: &mut DjogiContext,
seed_name: &str,
) -> Result<Option<SeedLedgerRow>, SeedError> {
let row_opt = ctx
.query_opt(
"SELECT checksum_up, status, failure_note \
FROM djogi_seed_runs \
WHERE seed_name = $1",
&[&seed_name],
)
.await
.map_err(|e| SeedError::LedgerWrite { source: e })?;
let Some(row) = row_opt else {
return Ok(None);
};
let checksum_up: String = row
.try_get("checksum_up")
.map_err(|e| SeedError::LedgerDecode {
column: "checksum_up".to_string(),
source: e,
})?;
let status_str: String = row.try_get("status").map_err(|e| SeedError::LedgerDecode {
column: "status".to_string(),
source: e,
})?;
let failure_note: Option<String> =
row.try_get("failure_note")
.map_err(|e| SeedError::LedgerDecode {
column: "failure_note".to_string(),
source: e,
})?;
Ok(Some(SeedLedgerRow {
checksum_up,
status: SeedLedgerStatus::from_db_str(&status_str)?,
failure_note,
}))
}
pub async fn fetch_recorded_checksum(
ctx: &mut DjogiContext,
seed_name: &str,
) -> Result<Option<String>, SeedError> {
let row_opt = ctx
.query_opt(
"SELECT checksum_up FROM djogi_seed_runs WHERE seed_name = $1",
&[&seed_name],
)
.await
.map_err(|e| SeedError::LedgerWrite { source: e })?;
match row_opt {
Some(row) => {
let checksum: String =
row.try_get::<_, String>("checksum_up")
.map_err(|e| SeedError::LedgerDecode {
column: "checksum_up".to_string(),
source: e,
})?;
Ok(Some(checksum))
}
None => Ok(None),
}
}
pub async fn insert_recorded(
ctx: &mut DjogiContext,
seed_name: &str,
checksum: &str,
) -> Result<(), DjogiError> {
ctx.raw_execute(
"INSERT INTO djogi_seed_runs (seed_name, checksum_up) VALUES ($1, $2)",
&[&seed_name, &checksum],
)
.await?;
Ok(())
}
async fn insert_running_claim(
ctx: &mut DjogiContext,
seed_name: &str,
checksum: &str,
) -> Result<(), DjogiError> {
ctx.raw_execute(
"INSERT INTO djogi_seed_runs \
(seed_name, checksum_up, status, failure_note) \
VALUES ($1, $2, $3, NULL)",
&[
&seed_name,
&checksum,
&SeedLedgerStatus::Running.as_db_str(),
],
)
.await?;
Ok(())
}
async fn mark_seed_applied(
ctx: &mut DjogiContext,
seed_name: &str,
checksum: &str,
) -> Result<(), DjogiError> {
ctx.raw_execute(
"UPDATE djogi_seed_runs \
SET checksum_up = $2, \
status = $3, \
applied_at = now(), \
failure_note = NULL \
WHERE seed_name = $1",
&[
&seed_name,
&checksum,
&SeedLedgerStatus::Applied.as_db_str(),
],
)
.await?;
Ok(())
}
async fn mark_seed_failed(
ctx: &mut DjogiContext,
seed_name: &str,
note: &str,
) -> Result<(), DjogiError> {
ctx.raw_execute(
"UPDATE djogi_seed_runs \
SET status = $2, \
failure_note = $3 \
WHERE seed_name = $1",
&[&seed_name, &SeedLedgerStatus::Failed.as_db_str(), ¬e],
)
.await?;
Ok(())
}
pub fn compute_seed_checksum(body: &str) -> String {
compute_checksum([body])
}
async fn apply_seed_body(ctx: &mut DjogiContext, body: &str) -> Result<(), DjogiError> {
ctx.raw_ddl(body).await
}
async fn reset_after_seed_apply_failure(ctx: &mut DjogiContext) {
if let Err(source) = ctx.raw_ddl("ROLLBACK").await {
tracing::warn!(
?source,
"seed apply failed and best-effort ROLLBACK did not clear session state",
);
}
}
pub async fn run_seeds(
ctx: &mut DjogiContext,
workspace_root: &Path,
database: &str,
database_url: &str,
allow_non_localhost: bool,
) -> Result<SeedReport, SeedError> {
let mut pinned = ctx
.pin_for_migration()
.await
.map_err(|e| SeedError::PinnedSessionCheckoutFailed { source: e })?;
run_seeds_pinned(
&mut pinned,
workspace_root,
database,
database_url,
allow_non_localhost,
)
.await
}
async fn run_seeds_pinned(
ctx: &mut PinnedCtx<'_>,
workspace_root: &Path,
database: &str,
database_url: &str,
allow_non_localhost: bool,
) -> Result<SeedReport, SeedError> {
if !allow_non_localhost && !super::policy::is_localhost_connection(database_url) {
return Err(SeedError::LocalhostGate {
database_url: database_url.to_string(),
});
}
bootstrap(ctx)
.await
.map_err(|e| SeedError::LedgerWrite { source: e })?;
let lock_bucket = seed_run_lock_bucket(database);
let lock_key = advisory_lock_key(&lock_bucket);
try_acquire_seed_run_lock(ctx, database, lock_key).await?;
let result = run_seeds_inner(ctx, workspace_root, database).await;
let released = release_advisory_lock(ctx, lock_key).await;
let result = handle_seed_unlock(result, released, database, lock_key);
if result.is_ok() {
ctx.mark_clean();
}
result
}
async fn run_seeds_inner(
ctx: &mut DjogiContext,
workspace_root: &Path,
database: &str,
) -> Result<SeedReport, SeedError> {
let discovered = discover_seeds(workspace_root, database)?;
let mut entries: Vec<SeedReportEntry> = Vec::with_capacity(discovered.len());
for seed in discovered {
let body_bytes = fs::read(&seed.path).map_err(|err| SeedError::Io {
path: seed.path.clone(),
source: err,
})?;
let body = String::from_utf8_lossy(&body_bytes).into_owned();
let on_disk_checksum = compute_seed_checksum(&body);
match fetch_seed_ledger_row(ctx, &seed.seed_name).await? {
Some(row)
if row.status == SeedLedgerStatus::Applied
&& row.checksum_up == on_disk_checksum =>
{
entries.push(SeedReportEntry {
seed_name: seed.seed_name.clone(),
outcome: SeedOutcome::SkippedAlreadyApplied,
});
continue;
}
Some(row) if row.status == SeedLedgerStatus::Applied => {
return Err(SeedError::ChecksumDrift {
seed_name: seed.seed_name,
ledger_checksum: row.checksum_up,
on_disk_checksum,
});
}
Some(row) if row.status == SeedLedgerStatus::Running => {
return Err(SeedError::StaleClaim {
seed_name: seed.seed_name,
checksum_up: row.checksum_up,
});
}
Some(row) => {
return Err(SeedError::FailedClaim {
seed_name: seed.seed_name,
failure_note: row.failure_note,
});
}
None => {}
}
insert_running_claim(ctx, &seed.seed_name, &on_disk_checksum)
.await
.map_err(|e| SeedError::LedgerWrite { source: e })?;
if let Err(e) = apply_seed_body(ctx, &body).await {
reset_after_seed_apply_failure(ctx).await;
let note = format!("seed apply failed: {e}");
let _ = mark_seed_failed(ctx, &seed.seed_name, ¬e).await;
return Err(SeedError::ApplyFailed {
seed_name: seed.seed_name,
source: e,
});
}
mark_seed_applied(ctx, &seed.seed_name, &on_disk_checksum)
.await
.map_err(|e| SeedError::LedgerWrite { source: e })?;
entries.push(SeedReportEntry {
seed_name: seed.seed_name,
outcome: SeedOutcome::Applied,
});
}
Ok(SeedReport { entries })
}
async fn try_acquire_seed_run_lock(
ctx: &mut PinnedCtx<'_>,
database: &str,
key: i64,
) -> Result<(), SeedError> {
assert!(
!ctx.is_pool_backed(),
"seed-run advisory lock called on a pool-backed context — \
the lock would be acquired on an arbitrary pool connection \
and subsequent operations would run on different connections. \
Callers must use ctx.pin_for_migration() first (GH #274 / #331).",
);
let row = ctx
.query_one("SELECT pg_try_advisory_lock($1)", &[&key])
.await
.map_err(|e| SeedError::LockQueryFailed {
database: database.to_string(),
source: e,
})?;
let acquired: bool = row.try_get(0).map_err(|e| SeedError::LockQueryFailed {
database: database.to_string(),
source: DjogiError::from(e),
})?;
if acquired {
Ok(())
} else {
Err(SeedError::RunAlreadyInProgress {
database: database.to_string(),
})
}
}
#[allow(clippy::result_large_err)]
fn handle_seed_unlock(
result: Result<SeedReport, SeedError>,
released: bool,
database: &str,
key: i64,
) -> Result<SeedReport, SeedError> {
match (result, released) {
(Ok(report), true) => Ok(report),
(Ok(_), false) => Err(SeedError::AdvisoryUnlockReturnedFalse {
database: database.to_string(),
key,
}),
(Err(err), true) => Err(err),
(Err(err), false) => {
tracing::error!(
database,
key = format_args!("0x{key:016x}"),
"seed-run advisory unlock returned false while propagating an earlier error",
);
Err(err)
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::atomic::{AtomicUsize, Ordering};
fn temp_root(tag: &str) -> PathBuf {
static COUNTER: AtomicUsize = AtomicUsize::new(0);
let n = COUNTER.fetch_add(1, Ordering::SeqCst);
let nanos = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos();
let p = std::env::temp_dir().join(format!("djogi-seed-{tag}-{nanos}-{n}"));
fs::create_dir_all(&p).unwrap();
p
}
#[test]
fn discover_seeds_returns_empty_when_dir_missing() {
let root = temp_root("empty");
let seeds = discover_seeds(&root, "main").expect("ok");
assert!(seeds.is_empty());
let _ = fs::remove_dir_all(&root);
}
#[test]
fn discover_seeds_filters_to_sql_files_in_alphabetical_order() {
let root = temp_root("filter");
let dir = root.join("seeds/main");
fs::create_dir_all(&dir).unwrap();
fs::write(dir.join("02_data.sql"), "INSERT INTO foo VALUES (1);\n").unwrap();
fs::write(dir.join("01_init.sql"), "INSERT INTO foo VALUES (0);\n").unwrap();
fs::write(dir.join("readme.md"), "# notes").unwrap();
fs::write(dir.join(".gitkeep"), "").unwrap();
fs::create_dir_all(dir.join("subdir")).unwrap();
let seeds = discover_seeds(&root, "main").expect("ok");
let names: Vec<&str> = seeds.iter().map(|s| s.seed_name.as_str()).collect();
assert_eq!(names, vec!["01_init", "02_data"]);
let _ = fs::remove_dir_all(&root);
}
#[test]
fn discover_seeds_rejects_zero_length_stem() {
let root = temp_root("empty_stem");
let dir = root.join("seeds/main");
fs::create_dir_all(&dir).unwrap();
fs::write(dir.join(".sql"), "noop").unwrap();
let seeds = discover_seeds(&root, "main").expect("ok");
assert!(seeds.is_empty(), "got {seeds:?}");
let _ = fs::remove_dir_all(&root);
}
#[test]
fn compute_seed_checksum_is_v1_prefixed_and_stable() {
let a = compute_seed_checksum("INSERT INTO foo VALUES (1)");
let b = compute_seed_checksum("INSERT INTO foo VALUES (1)");
assert_eq!(a, b, "checksum is deterministic");
assert!(a.starts_with("V1:"), "format: V1:<sha256-hex>");
assert_eq!(a.len(), 3 + 64, "V1 + 64-char sha-256 hex digest");
}
#[test]
fn checksum_drift_message_names_both_sides() {
let e = SeedError::ChecksumDrift {
seed_name: "01_init".to_string(),
ledger_checksum: "V1:aaaa".to_string(),
on_disk_checksum: "V1:bbbb".to_string(),
};
let msg = e.to_string();
assert!(msg.contains("01_init"));
assert!(msg.contains("V1:aaaa"));
assert!(msg.contains("V1:bbbb"));
}
#[test]
fn localhost_gate_message_names_url() {
let e = SeedError::LocalhostGate {
database_url: "postgres://prod.example.com/main".to_string(),
};
let msg = e.to_string();
assert!(msg.contains("postgres://prod.example.com/main"));
assert!(msg.contains("--allow-non-localhost"));
}
#[test]
fn derive_per_database_url_replaces_path_component() {
assert_eq!(
derive_per_database_url("postgres://localhost/main", "crud_log"),
Some("postgres://localhost/crud_log".to_string())
);
assert_eq!(
derive_per_database_url("postgres://user:pass@localhost:5432/main", "event_log"),
Some("postgres://user:pass@localhost:5432/event_log".to_string())
);
assert_eq!(
derive_per_database_url("postgresql://localhost/main?sslmode=disable", "audit"),
Some("postgresql://localhost/audit?sslmode=disable".to_string())
);
}
#[test]
fn derive_per_database_url_returns_none_for_pathless_urls() {
assert_eq!(
derive_per_database_url("postgres://localhost", "crud_log"),
None
);
assert_eq!(
derive_per_database_url("mysql://localhost/main", "crud_log"),
None
);
}
#[test]
fn ledger_decode_variant_threads_error_source() {
let write = SeedError::LedgerWrite {
source: DjogiError::Db(crate::error::DbError::other("boom".to_string())),
};
match write {
SeedError::LedgerWrite { .. } => {}
other => panic!("expected LedgerWrite, got {other:?}"),
}
}
#[test]
fn seed_ledger_status_round_trips_through_db_strings() {
for status in [
SeedLedgerStatus::Running,
SeedLedgerStatus::Applied,
SeedLedgerStatus::Failed,
] {
let db = status.as_db_str();
assert_eq!(
SeedLedgerStatus::from_db_str(db).expect("round-trip"),
status
);
}
}
#[test]
fn stale_and_failed_claim_messages_name_the_seed() {
let stale = SeedError::StaleClaim {
seed_name: "01_seed".to_string(),
checksum_up: "V1:abcd".to_string(),
};
let stale_msg = stale.to_string();
assert!(stale_msg.contains("01_seed"));
assert!(stale_msg.contains("V1:abcd"));
let failed = SeedError::FailedClaim {
seed_name: "02_seed".to_string(),
failure_note: Some("seed apply failed: duplicate key".to_string()),
};
let failed_msg = failed.to_string();
assert!(failed_msg.contains("02_seed"));
assert!(failed_msg.contains("duplicate key"));
}
}