use time::OffsetDateTime;
use crate::__bypass::RawAccessExt as _;
use crate::context::DjogiContext;
use crate::error::{DbError, DjogiError};
use crate::live_migrate::plan::PlanClassification;
use crate::types::HeerId;
pub const INSTALL_SQL: &str = r#"
CREATE TABLE IF NOT EXISTS djogi_live_plans (
plan_id BIGINT PRIMARY KEY,
slug TEXT NOT NULL,
plan_file_checksum VARCHAR(68) NOT NULL,
classification TEXT NOT NULL
CHECK (classification IN (
'online_safe', 'expand_contract', 'offline_only'
)),
status TEXT NOT NULL DEFAULT 'pending'
CHECK (status IN (
'pending', 'running', 'paused',
'validating', 'cutover', 'finalizing',
'complete', 'abandoned', 'failed',
'failed_retriable', 'failed_terminal'
)),
current_step TEXT,
current_step_index INTEGER NOT NULL DEFAULT 0,
backfill_rows_done BIGINT NOT NULL DEFAULT 0,
backfill_rows_total BIGINT,
started_at TIMESTAMPTZ,
last_progress_at TIMESTAMPTZ,
completed_at TIMESTAMPTZ,
last_error TEXT,
originating_migration TEXT NOT NULL,
target_database TEXT NOT NULL DEFAULT 'main',
app_label TEXT NOT NULL DEFAULT '',
daemon_session_token TEXT
);
CREATE UNIQUE INDEX IF NOT EXISTS djogi_live_plans_bucket_plan_id_uidx
ON djogi_live_plans (target_database, app_label, plan_id);
ALTER TABLE djogi_live_plans
ADD COLUMN IF NOT EXISTS claimed_by_pid BIGINT;
ALTER TABLE djogi_live_plans
ADD COLUMN IF NOT EXISTS claimed_by_host TEXT;
ALTER TABLE djogi_live_plans
ADD COLUMN IF NOT EXISTS claimed_at TIMESTAMPTZ;
ALTER TABLE djogi_live_plans
ADD COLUMN IF NOT EXISTS daemon_session_token TEXT;
"#;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
#[non_exhaustive]
pub enum PlanStatus {
Pending,
Running,
Paused,
Validating,
Cutover,
Finalizing,
Complete,
Abandoned,
Failed,
FailedRetriable,
FailedTerminal,
}
const PLAN_STATUS_LABELS: &[(PlanStatus, &str)] = &[
(PlanStatus::Pending, "pending"),
(PlanStatus::Running, "running"),
(PlanStatus::Paused, "paused"),
(PlanStatus::Validating, "validating"),
(PlanStatus::Cutover, "cutover"),
(PlanStatus::Finalizing, "finalizing"),
(PlanStatus::Complete, "complete"),
(PlanStatus::Abandoned, "abandoned"),
(PlanStatus::Failed, "failed"),
(PlanStatus::FailedRetriable, "failed_retriable"),
(PlanStatus::FailedTerminal, "failed_terminal"),
];
const _PLAN_STATUS_DRIFT_GUARD: () = {
const fn check(s: PlanStatus) -> u8 {
match s {
PlanStatus::Pending => 0,
PlanStatus::Running => 1,
PlanStatus::Paused => 2,
PlanStatus::Validating => 3,
PlanStatus::Cutover => 4,
PlanStatus::Finalizing => 5,
PlanStatus::Complete => 6,
PlanStatus::Abandoned => 7,
PlanStatus::Failed => 8,
PlanStatus::FailedRetriable => 9,
PlanStatus::FailedTerminal => 10,
}
}
let _ = check(PlanStatus::Pending);
};
impl PlanStatus {
pub fn as_db_str(self) -> &'static str {
PLAN_STATUS_LABELS
.iter()
.find_map(|(v, label)| (*v == self).then_some(*label))
.expect("invariant: PlanStatus variant missing from PLAN_STATUS_LABELS")
}
pub fn from_db_str(s: &str) -> Option<Self> {
PLAN_STATUS_LABELS
.iter()
.find_map(|(variant, label)| (*label == s).then_some(*variant))
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct LivePlanRow {
pub plan_id: HeerId,
pub slug: String,
pub plan_file_checksum: String,
pub classification: PlanClassification,
pub status: PlanStatus,
pub current_step: Option<String>,
pub current_step_index: i32,
pub backfill_rows_done: i64,
pub backfill_rows_total: Option<i64>,
pub started_at: Option<OffsetDateTime>,
pub last_progress_at: Option<OffsetDateTime>,
pub completed_at: Option<OffsetDateTime>,
pub last_error: Option<String>,
pub originating_migration: String,
pub target_database: String,
pub app_label: String,
pub daemon_session_token: Option<String>,
}
pub async fn install(ctx: &mut DjogiContext) -> Result<(), DjogiError> {
ctx.raw_ddl(INSTALL_SQL).await
}
pub async fn insert_row(ctx: &mut DjogiContext, row: &LivePlanRow) -> Result<(), DjogiError> {
let sql = "INSERT INTO djogi_live_plans \
(plan_id, slug, plan_file_checksum, classification, status, \
current_step, current_step_index, backfill_rows_done, \
backfill_rows_total, started_at, last_progress_at, completed_at, \
last_error, originating_migration, target_database, app_label, \
daemon_session_token) \
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16, $17)";
let plan_id = row.plan_id.as_i64();
let classification = row.classification.as_db_str();
let status = row.status.as_db_str();
ctx.execute(
sql,
&[
&plan_id,
&row.slug,
&row.plan_file_checksum,
&classification,
&status,
&row.current_step,
&row.current_step_index,
&row.backfill_rows_done,
&row.backfill_rows_total,
&row.started_at,
&row.last_progress_at,
&row.completed_at,
&row.last_error,
&row.originating_migration,
&row.target_database,
&row.app_label,
&row.daemon_session_token,
],
)
.await?;
Ok(())
}
pub async fn fetch_row_by_id(
ctx: &mut DjogiContext,
plan_id: HeerId,
target_database: &str,
app_label: &str,
) -> Result<Option<LivePlanRow>, DjogiError> {
let sql = "SELECT plan_id, slug, plan_file_checksum, classification, status, \
current_step, current_step_index, backfill_rows_done, \
backfill_rows_total, started_at, last_progress_at, completed_at, \
last_error, originating_migration, target_database, app_label, \
daemon_session_token \
FROM djogi_live_plans \
WHERE target_database = $1 AND app_label = $2 AND plan_id = $3";
let plan_id_i64 = plan_id.as_i64();
let row_opt = ctx
.query_opt(sql, &[&target_database, &app_label, &plan_id_i64])
.await?;
let Some(row) = row_opt else {
return Ok(None);
};
let parsed = row_to_live_plan_row(&row)?;
Ok(Some(parsed))
}
fn row_to_live_plan_row(row: &tokio_postgres::Row) -> Result<LivePlanRow, DjogiError> {
let plan_id_i64: i64 = row.try_get(0)?;
let plan_id = HeerId::from_i64(plan_id_i64)
.map_err(|e| DjogiError::Db(DbError::other(format!("invalid plan_id in row: {e}"))))?;
let slug: String = row.try_get(1)?;
let plan_file_checksum: String = row.try_get(2)?;
let classification_s: String = row.try_get(3)?;
let classification = PlanClassification::from_db_str(&classification_s).ok_or_else(|| {
DjogiError::Db(DbError::other(format!(
"unknown classification in djogi_live_plans: {classification_s:?}"
)))
})?;
let status_s: String = row.try_get(4)?;
let status = PlanStatus::from_db_str(&status_s).ok_or_else(|| {
DjogiError::Db(DbError::other(format!(
"unknown status in djogi_live_plans: {status_s:?}"
)))
})?;
let current_step: Option<String> = row.try_get(5)?;
let current_step_index: i32 = row.try_get(6)?;
let backfill_rows_done: i64 = row.try_get(7)?;
let backfill_rows_total: Option<i64> = row.try_get(8)?;
let started_at: Option<OffsetDateTime> = row.try_get(9)?;
let last_progress_at: Option<OffsetDateTime> = row.try_get(10)?;
let completed_at: Option<OffsetDateTime> = row.try_get(11)?;
let last_error: Option<String> = row.try_get(12)?;
let originating_migration: String = row.try_get(13)?;
let target_database: String = row.try_get(14)?;
let app_label: String = row.try_get(15)?;
let daemon_session_token: Option<String> = row.try_get(16)?;
Ok(LivePlanRow {
plan_id,
slug,
plan_file_checksum,
classification,
status,
current_step,
current_step_index,
backfill_rows_done,
backfill_rows_total,
started_at,
last_progress_at,
completed_at,
last_error,
originating_migration,
target_database,
app_label,
daemon_session_token,
})
}
pub async fn update_progress(
ctx: &mut DjogiContext,
plan_id: HeerId,
target_database: &str,
app_label: &str,
rows_done: i64,
) -> Result<(), DjogiError> {
let sql = "UPDATE djogi_live_plans \
SET backfill_rows_done = $4, last_progress_at = now() \
WHERE target_database = $1 AND app_label = $2 AND plan_id = $3";
let plan_id_i64 = plan_id.as_i64();
ctx.execute(
sql,
&[&target_database, &app_label, &plan_id_i64, &rows_done],
)
.await?;
Ok(())
}
pub async fn update_step_index(
ctx: &mut DjogiContext,
plan_id: HeerId,
target_database: &str,
app_label: &str,
new_index: i32,
new_step_label: Option<&str>,
) -> Result<(), DjogiError> {
let sql = "UPDATE djogi_live_plans \
SET current_step_index = $4, current_step = $5 \
WHERE target_database = $1 AND app_label = $2 AND plan_id = $3";
let plan_id_i64 = plan_id.as_i64();
ctx.execute(
sql,
&[
&target_database,
&app_label,
&plan_id_i64,
&new_index,
&new_step_label,
],
)
.await?;
Ok(())
}
pub async fn stamp_completed_at(
ctx: &mut DjogiContext,
plan_id: HeerId,
target_database: &str,
app_label: &str,
) -> Result<(), DjogiError> {
let sql = "UPDATE djogi_live_plans \
SET completed_at = now() \
WHERE target_database = $1 AND app_label = $2 AND plan_id = $3";
let plan_id_i64 = plan_id.as_i64();
ctx.execute(sql, &[&target_database, &app_label, &plan_id_i64])
.await?;
Ok(())
}
pub async fn update_status(
ctx: &mut DjogiContext,
plan_id: HeerId,
target_database: &str,
app_label: &str,
new_status: PlanStatus,
) -> Result<(), DjogiError> {
let sql = "UPDATE djogi_live_plans \
SET status = $4 \
WHERE target_database = $1 AND app_label = $2 AND plan_id = $3";
let plan_id_i64 = plan_id.as_i64();
let status = new_status.as_db_str();
ctx.execute(sql, &[&target_database, &app_label, &plan_id_i64, &status])
.await?;
Ok(())
}
pub async fn update_status_with_error(
ctx: &mut DjogiContext,
plan_id: HeerId,
target_database: &str,
app_label: &str,
status: PlanStatus,
error_msg: String,
) -> Result<(), DjogiError> {
let sql = "UPDATE djogi_live_plans \
SET status = $4, last_error = $5 \
WHERE target_database = $1 AND app_label = $2 AND plan_id = $3";
let plan_id_i64 = plan_id.as_i64();
let status_str = status.as_db_str();
ctx.execute(
sql,
&[
&target_database,
&app_label,
&plan_id_i64,
&status_str,
&error_msg,
],
)
.await?;
Ok(())
}
pub async fn record_failure(
ctx: &mut DjogiContext,
plan_id: HeerId,
target_database: &str,
app_label: &str,
error_msg: String,
retriable: bool,
) -> Result<(), DjogiError> {
let status = if retriable {
PlanStatus::FailedRetriable
} else {
PlanStatus::FailedTerminal
};
update_status_with_error(ctx, plan_id, target_database, app_label, status, error_msg).await
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn plan_status_as_db_str_matches_check_constraint_list() {
let pairs = [
(PlanStatus::Pending, "pending"),
(PlanStatus::Running, "running"),
(PlanStatus::Paused, "paused"),
(PlanStatus::Validating, "validating"),
(PlanStatus::Cutover, "cutover"),
(PlanStatus::Finalizing, "finalizing"),
(PlanStatus::Complete, "complete"),
(PlanStatus::Abandoned, "abandoned"),
(PlanStatus::Failed, "failed"),
(PlanStatus::FailedRetriable, "failed_retriable"),
(PlanStatus::FailedTerminal, "failed_terminal"),
];
for (variant, expected) in pairs {
assert_eq!(variant.as_db_str(), expected);
assert_eq!(PlanStatus::from_db_str(expected), Some(variant));
}
}
#[test]
fn plan_status_from_db_str_rejects_unknown() {
assert_eq!(PlanStatus::from_db_str("whatever"), None);
assert_eq!(PlanStatus::from_db_str(""), None);
assert_eq!(PlanStatus::from_db_str("Pending"), None);
}
#[test]
fn plan_status_label_slice_covers_every_variant() {
let exhaustive_list = [
PlanStatus::Pending,
PlanStatus::Running,
PlanStatus::Paused,
PlanStatus::Validating,
PlanStatus::Cutover,
PlanStatus::Finalizing,
PlanStatus::Complete,
PlanStatus::Abandoned,
PlanStatus::Failed,
PlanStatus::FailedRetriable,
PlanStatus::FailedTerminal,
];
assert_eq!(
PLAN_STATUS_LABELS.len(),
exhaustive_list.len(),
"PLAN_STATUS_LABELS must contain exactly one row per PlanStatus variant"
);
for variant in exhaustive_list {
let label = variant.as_db_str();
assert_eq!(
PlanStatus::from_db_str(label),
Some(variant),
"round-trip failed for {variant:?}",
);
}
}
#[test]
fn install_sql_lists_every_plan_status_variant() {
for (variant, label) in PLAN_STATUS_LABELS {
assert_eq!(
*label,
variant.as_db_str(),
"drift between PLAN_STATUS_LABELS and `as_db_str()` for {variant:?}",
);
let needle = format!("'{label}'");
assert!(
INSTALL_SQL.contains(&needle),
"INSTALL_SQL missing CHECK entry for status {needle}",
);
}
}
#[test]
fn install_sql_lists_every_plan_classification_variant() {
for c in [
PlanClassification::OnlineSafe,
PlanClassification::ExpandContract,
PlanClassification::OfflineOnly,
] {
let label = c.as_db_str();
let needle = format!("'{label}'");
assert!(
INSTALL_SQL.contains(&needle),
"INSTALL_SQL missing CHECK entry for classification {needle}",
);
}
}
#[test]
fn install_sql_declares_unique_index_on_bucket_plan_id() {
assert!(
INSTALL_SQL.contains("CREATE UNIQUE INDEX IF NOT EXISTS"),
"INSTALL_SQL must create the bucket+plan_id unique index",
);
assert!(
INSTALL_SQL.contains("(target_database, app_label, plan_id)"),
"unique index must be on (target_database, app_label, plan_id)",
);
}
}