Skip to main content

turso_orm_migration/
migrator.rs

1//! Applying and reverting migrations, modeled by [`MigratorTrait`].
2//!
3//! The migrator keeps a bookkeeping table, `turso_migrations` unless
4//! [`MigratorTrait::migration_table_name`] says otherwise, with one row per
5//! applied migration. Every operation starts by making sure that table
6//! exists and reading it, then walks the declared migrations in order (or in
7//! reverse for `down`) and skips the ones whose state already matches. The
8//! state is checked again once each migration holds the write lock, so two
9//! migrators started together never run the same migration twice.
10//!
11//! Each migration runs inside its own `BEGIN IMMEDIATE` transaction, and
12//! the bookkeeping insert or delete is issued on that same transaction
13//! before the commit. Taking the write lock up front avoids a busy error
14//! mid-migration, and bundling the version row with the schema change means
15//! a failure leaves neither a partial schema nor a misleading version row.
16//!
17//! The declared list is validated before anything runs. A duplicate name is
18//! always an error, since the bookkeeping table could not tell the two
19//! apart. A migration recorded as applied but no longer declared, or a
20//! pending one declared before an applied one, is a [`MigrationIssue`]:
21//! legitimate after a deployment is rolled back or two branches are merged,
22//! so `up` only warns about it unless [`MigratorTrait::strict`] says
23//! otherwise.
24
25use std::collections::HashSet;
26use std::fmt;
27
28use async_trait::async_trait;
29use turso_orm::sql::{ColumnDef, Expr, Order, Query, Table};
30use turso_orm::{ConnectionTrait, Database, DbErr, Statement, TransactionMode, TransactionTrait};
31
32use crate::MigrationTrait;
33use crate::manager::{SchemaManager, has_table};
34
35/// The default name of the bookkeeping table.
36const DEFAULT_TABLE: &str = "turso_migrations";
37
38/// The status of one declared migration.
39#[derive(Clone, Debug, PartialEq, Eq)]
40pub struct MigrationStatus {
41    /// The migration name.
42    pub name: String,
43    /// Whether the migration has been applied.
44    pub applied: bool,
45}
46
47/// A mismatch between the declared migrations and the bookkeeping table,
48/// reported by [`MigratorTrait::check`].
49#[derive(Clone, Debug, PartialEq, Eq)]
50#[non_exhaustive]
51pub enum MigrationIssue {
52    /// A migration recorded as applied that [`MigratorTrait::migrations`]
53    /// no longer declares, for example after a deployment was rolled back
54    /// to an older binary or a migration was renamed.
55    Unknown(String),
56    /// A pending migration declared before one that is already applied,
57    /// typically after two branches that each added a migration were
58    /// merged; `up` applies it after the later one.
59    OutOfOrder(String),
60}
61
62impl fmt::Display for MigrationIssue {
63    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
64        match self {
65            Self::Unknown(name) => write!(f, "applied migration `{name}` is not declared"),
66            Self::OutOfOrder(name) => write!(
67                f,
68                "pending migration `{name}` is declared before an applied one"
69            ),
70        }
71    }
72}
73
74/// Lists migrations and applies or reverts them in order.
75///
76/// Implement [`migrations`](Self::migrations) only; every other method has
77/// a default built on it.
78#[async_trait]
79pub trait MigratorTrait: Send {
80    /// Every migration, oldest first.
81    fn migrations() -> Vec<Box<dyn MigrationTrait>>;
82
83    /// The name of the bookkeeping table, `turso_migrations` by default.
84    ///
85    /// Override it to run several migrators against one database, or to
86    /// keep the table name of a schema that was migrated by another tool.
87    fn migration_table_name() -> &'static str {
88        DEFAULT_TABLE
89    }
90
91    /// Whether a [`MigrationIssue`] makes [`up`](Self::up) fail rather than
92    /// warn; `false` by default.
93    ///
94    /// Leave it off where an older binary may run against a database that
95    /// a newer one migrated, which leaves unknown migrations behind.
96    fn strict() -> bool {
97        false
98    }
99
100    /// Lists the mismatches between the declared migrations and the
101    /// bookkeeping table, without changing anything.
102    ///
103    /// Unknown migrations come first, in application order, then
104    /// out-of-order ones, in declaration order.
105    ///
106    /// # Errors
107    ///
108    /// Returns [`DbErr::Migration`] when two declared migrations share a
109    /// name; the errors of
110    /// [`get_applied_migrations`](Self::get_applied_migrations).
111    async fn check(db: &Database) -> Result<Vec<MigrationIssue>, DbErr> {
112        let migrations = Self::migrations();
113        ensure_unique(&migrations)?;
114        let applied = Self::get_applied_migrations(db).await?;
115        Ok(find_issues(&migrations, &applied))
116    }
117
118    /// Creates the bookkeeping table if it does not exist yet.
119    ///
120    /// # Errors
121    ///
122    /// Returns [`DbErr::Driver`] when the statement fails.
123    async fn install(db: &Database) -> Result<(), DbErr> {
124        let stmt = Table::create()
125            .table(Self::migration_table_name())
126            .if_not_exists()
127            .col(ColumnDef::text("version").primary_key().not_null())
128            .col(ColumnDef::integer("applied_at").not_null());
129        db.execute(turso_orm::Build::to_statement(&stmt)).await?;
130        Ok(())
131    }
132
133    /// The names of the applied migrations, in application order.
134    ///
135    /// Returns an empty list when the bookkeeping table does not exist, so
136    /// that status can be queried on a database that was never migrated.
137    ///
138    /// # Errors
139    ///
140    /// Returns [`DbErr::Migration`] when the catalog query returns no row;
141    /// [`DbErr::Driver`] when a query fails or a version cannot be decoded.
142    async fn get_applied_migrations(db: &Database) -> Result<Vec<String>, DbErr> {
143        if !has_table(db, Self::migration_table_name()).await? {
144            return Ok(Vec::new());
145        }
146        // The timestamp has second resolution, so the name breaks ties
147        // between migrations applied within the same second.
148        let stmt = Query::select()
149            .column("version")
150            .from(Self::migration_table_name())
151            .order_by("applied_at", Order::Asc)
152            .order_by("version", Order::Asc);
153        let rows = db.query_all(turso_orm::Build::to_statement(&stmt)).await?;
154        rows.iter()
155            .map(|r| r.get::<String>("version").map_err(DbErr::from))
156            .collect()
157    }
158
159    /// The status of every declared migration, in declaration order.
160    ///
161    /// # Errors
162    ///
163    /// Returns [`DbErr::Migration`] when two declared migrations share a
164    /// name; the errors of
165    /// [`get_applied_migrations`](Self::get_applied_migrations).
166    async fn status(db: &Database) -> Result<Vec<MigrationStatus>, DbErr> {
167        let migrations = Self::migrations();
168        ensure_unique(&migrations)?;
169        let applied = Self::get_applied_migrations(db).await?;
170        Ok(migrations
171            .iter()
172            .map(|m| MigrationStatus {
173                name: m.name().to_owned(),
174                applied: applied.iter().any(|a| a == m.name()),
175            })
176            .collect())
177    }
178
179    /// Applies the pending migrations, all of them or the first `steps`.
180    ///
181    /// Each migration and its version row are committed together; on the
182    /// first failure the transaction is dropped and rolled back, and the
183    /// error is returned without touching later migrations. Every
184    /// [`MigrationIssue`] is logged as a warning first, or fails the call
185    /// when [`strict`](Self::strict) is set.
186    ///
187    /// # Errors
188    ///
189    /// Returns [`DbErr::Migration`] when two declared migrations share a
190    /// name, or when [`strict`](Self::strict) is set and
191    /// [`check`](Self::check) finds an issue, in both cases before anything
192    /// runs; [`DbErr::Driver`] when a statement or the transaction fails;
193    /// any error the migration's `up` returns.
194    async fn up(db: &Database, steps: Option<u32>) -> Result<(), DbErr> {
195        let migrations = Self::migrations();
196        ensure_unique(&migrations)?;
197        Self::install(db).await?;
198        let applied = Self::get_applied_migrations(db).await?;
199        report_issues(&find_issues(&migrations, &applied), Self::strict())?;
200        let mut remaining = steps.map_or(usize::MAX, |s| s as usize);
201        for migration in migrations {
202            if remaining == 0 {
203                break;
204            }
205            if applied.iter().any(|a| a == migration.name()) {
206                continue;
207            }
208            tracing::info!(name = migration.name(), "applying migration");
209            // `IMMEDIATE` takes the write lock now rather than at the first
210            // write, so the migration cannot hit a busy error halfway.
211            let txn = db.begin_with_mode(TransactionMode::Immediate).await?;
212            // Another migrator may have applied it since the list was read;
213            // the write lock now held makes this check final.
214            if is_applied(&txn, Self::migration_table_name(), migration.name()).await? {
215                txn.rollback().await?;
216                continue;
217            }
218            {
219                let manager = SchemaManager::new(&txn);
220                migration.up(&manager).await?;
221            }
222            let now = std::time::SystemTime::now()
223                .duration_since(std::time::UNIX_EPOCH)
224                .map(|d| i64::try_from(d.as_secs()).unwrap_or(i64::MAX))
225                .unwrap_or_default();
226            let insert = Query::insert()
227                .into_table(Self::migration_table_name())
228                .columns(["version", "applied_at"])
229                .values([Expr::val(migration.name()), Expr::val(now)]);
230            txn.execute(turso_orm::Build::to_statement(&insert)).await?;
231            txn.commit().await?;
232            remaining -= 1;
233        }
234        Ok(())
235    }
236
237    /// Reverts the applied migrations, newest first, all of them or `steps`.
238    ///
239    /// Each migration's `down` and the deletion of its version row are
240    /// committed together, mirroring [`up`](Self::up).
241    ///
242    /// Unknown and out-of-order migrations are left alone: reverting does
243    /// not depend on them, and failing here would block the way back from
244    /// the state they describe.
245    ///
246    /// # Errors
247    ///
248    /// Returns [`DbErr::Migration`] when two declared migrations share a
249    /// name, before anything runs; [`DbErr::Driver`] when a statement or
250    /// the transaction fails; any error the migration's `down` returns,
251    /// including the default [`DbErr::Migration`] of an irreversible
252    /// migration.
253    async fn down(db: &Database, steps: Option<u32>) -> Result<(), DbErr> {
254        let migrations = Self::migrations();
255        ensure_unique(&migrations)?;
256        Self::install(db).await?;
257        let applied = Self::get_applied_migrations(db).await?;
258        let mut remaining = steps.map_or(usize::MAX, |s| s as usize);
259        for migration in migrations.into_iter().rev() {
260            if remaining == 0 {
261                break;
262            }
263            if !applied.iter().any(|a| a == migration.name()) {
264                continue;
265            }
266            tracing::info!(name = migration.name(), "reverting migration");
267            let txn = db.begin_with_mode(TransactionMode::Immediate).await?;
268            // Another migrator may have reverted it since the list was read;
269            // the write lock now held makes this check final.
270            if !is_applied(&txn, Self::migration_table_name(), migration.name()).await? {
271                txn.rollback().await?;
272                continue;
273            }
274            {
275                let manager = SchemaManager::new(&txn);
276                migration.down(&manager).await?;
277            }
278            let delete = Query::delete()
279                .from_table(Self::migration_table_name())
280                .and_where(Expr::col("version").eq(Expr::val(migration.name())));
281            txn.execute(turso_orm::Build::to_statement(&delete)).await?;
282            txn.commit().await?;
283            remaining -= 1;
284        }
285        Ok(())
286    }
287
288    /// Drops every user table, including the bookkeeping table, then applies all migrations.
289    ///
290    /// Internal `sqlite_*` and `__turso_*` tables are left alone.
291    ///
292    /// # Errors
293    ///
294    /// Returns [`DbErr::Migration`] when two declared migrations share a
295    /// name, before any table is dropped; [`DbErr::Driver`] when a query, a
296    /// drop or the transaction fails; the errors of [`up`](Self::up).
297    async fn fresh(db: &Database) -> Result<(), DbErr> {
298        ensure_unique(&Self::migrations())?;
299        let rows = db
300            .query_all(Statement::from_string(
301                "SELECT name FROM sqlite_schema WHERE type = 'table' AND name NOT LIKE 'sqlite_%' AND name NOT LIKE '__turso_%'",
302            ))
303            .await?;
304        let txn = db.begin_with_mode(TransactionMode::Immediate).await?;
305        // `PRAGMA foreign_keys` has no effect inside a transaction, so a
306        // table referenced by another one may refuse to drop first. Drop in
307        // rounds, retrying the tables that failed, until nothing is left or a
308        // round makes no progress — then the last error is the real one.
309        let mut pending: Vec<String> = rows
310            .iter()
311            .map(|row| row.get::<String>("name"))
312            .collect::<Result<_, _>>()?;
313        while !pending.is_empty() {
314            let before = pending.len();
315            let mut failed = Vec::new();
316            let mut last_error = None;
317            for name in pending {
318                let drop = Table::drop().table(name.clone()).if_exists();
319                if let Err(err) = txn.execute(turso_orm::Build::to_statement(&drop)).await {
320                    failed.push(name);
321                    last_error = Some(err);
322                }
323            }
324            if failed.len() == before
325                && let Some(err) = last_error
326            {
327                return Err(err.into());
328            }
329            pending = failed;
330        }
331        txn.commit().await?;
332        Self::up(db, None).await
333    }
334
335    /// Reverts every migration, then applies every migration.
336    ///
337    /// # Errors
338    ///
339    /// Returns [`DbErr::Migration`] when [`strict`](Self::strict) is set
340    /// and an applied migration is not declared, before anything is
341    /// reverted; the errors of [`down`](Self::down) and [`up`](Self::up).
342    async fn refresh(db: &Database) -> Result<(), DbErr> {
343        // An unknown migration survives `down`, so a strict `up` would only
344        // refuse it once everything had been reverted; check it first.
345        if Self::strict() {
346            let unknown: Vec<MigrationIssue> = Self::check(db)
347                .await?
348                .into_iter()
349                .filter(|issue| matches!(issue, MigrationIssue::Unknown(_)))
350                .collect();
351            report_issues(&unknown, true)?;
352        }
353        Self::down(db, None).await?;
354        Self::up(db, None).await
355    }
356
357    /// Reverts every migration.
358    ///
359    /// # Errors
360    ///
361    /// Returns the errors of [`down`](Self::down).
362    async fn reset(db: &Database) -> Result<(), DbErr> {
363        Self::down(db, None).await
364    }
365}
366
367/// Whether `version` has a row in the bookkeeping table `table`.
368///
369/// Read inside the migration's transaction, so that the answer reflects
370/// what other migrators committed before the write lock was taken.
371///
372/// # Errors
373///
374/// Returns [`DbErr::Driver`] when the query fails.
375async fn is_applied<C: ConnectionTrait>(
376    conn: &C,
377    table: &'static str,
378    version: &str,
379) -> Result<bool, DbErr> {
380    let stmt = Query::select()
381        .column("version")
382        .from(table)
383        .and_where(Expr::col("version").eq(Expr::val(version)));
384    Ok(conn
385        .query_one(turso_orm::Build::to_statement(&stmt))
386        .await?
387        .is_some())
388}
389
390/// Checks that no two declared migrations share a name.
391///
392/// # Errors
393///
394/// Returns [`DbErr::Migration`] naming the first duplicate.
395fn ensure_unique(migrations: &[Box<dyn MigrationTrait>]) -> Result<(), DbErr> {
396    let mut seen = HashSet::new();
397    match migrations.iter().find(|m| !seen.insert(m.name())) {
398        Some(duplicate) => Err(DbErr::Migration(format!(
399            "duplicate migration name `{}`",
400            duplicate.name()
401        ))),
402        None => Ok(()),
403    }
404}
405
406/// Compares the declared migrations with the applied ones.
407///
408/// A pending migration is out of order when any migration declared after
409/// it is applied.
410fn find_issues(migrations: &[Box<dyn MigrationTrait>], applied: &[String]) -> Vec<MigrationIssue> {
411    let declared: HashSet<&str> = migrations.iter().map(|m| m.name()).collect();
412    let is_applied = |name: &str| applied.iter().any(|a| a == name);
413    let last_applied = migrations.iter().rposition(|m| is_applied(m.name()));
414    let unknown = applied
415        .iter()
416        .filter(|a| !declared.contains(a.as_str()))
417        .map(|a| MigrationIssue::Unknown(a.clone()));
418    let out_of_order = migrations
419        .iter()
420        .take(last_applied.unwrap_or(0))
421        .filter(|m| !is_applied(m.name()))
422        .map(|m| MigrationIssue::OutOfOrder(m.name().to_owned()));
423    unknown.chain(out_of_order).collect()
424}
425
426/// Logs each issue as a warning, or turns them into one error when `strict`.
427///
428/// # Errors
429///
430/// Returns [`DbErr::Migration`] listing every issue when `strict` is set
431/// and there is at least one.
432fn report_issues(issues: &[MigrationIssue], strict: bool) -> Result<(), DbErr> {
433    if strict && !issues.is_empty() {
434        let list: Vec<String> = issues.iter().map(ToString::to_string).collect();
435        return Err(DbErr::Migration(list.join("; ")));
436    }
437    for issue in issues {
438        tracing::warn!("{issue}");
439    }
440    Ok(())
441}