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}