|
1 | | -use anyhow::Result; |
| 1 | +use anyhow::{Context, Result}; |
2 | 2 | use chrono::{DateTime, Utc}; |
3 | 3 | use serde::{Deserialize, Serialize}; |
4 | 4 | use sqlx::{PgPool, Row}; |
| 5 | +use tracing::info; |
5 | 6 | use uuid::Uuid; |
6 | 7 |
|
7 | 8 | // ── Public data types ───────────────────────────────────────────────────────── |
@@ -202,8 +203,174 @@ impl Db { |
202 | 203 | Ok(db) |
203 | 204 | } |
204 | 205 |
|
| 206 | + /// Run all pending versioned migrations in order, inside a single |
| 207 | + /// transaction per migration. Idempotent — migrations whose version is |
| 208 | + /// already recorded in `schema_migrations` are skipped. |
| 209 | + /// |
| 210 | + /// Concurrency: the whole routine is guarded by a Postgres advisory lock so |
| 211 | + /// two node instances pointed at the same database (e.g. during a |
| 212 | + /// blue/green or rolling deploy) cannot race to apply the same migration |
| 213 | + /// and trip the `schema_migrations` primary key. |
| 214 | + /// |
| 215 | + /// Legacy installs: v1 bundles the entire pre-versioning schema, and every |
| 216 | + /// statement in it is idempotent (`CREATE TABLE IF NOT EXISTS`, |
| 217 | + /// `CREATE INDEX IF NOT EXISTS`, `ADD COLUMN IF NOT EXISTS`). So an existing |
| 218 | + /// node that predates this system just runs v1 once: existing objects are |
| 219 | + /// no-ops, and any objects it was missing are created. We deliberately do |
| 220 | + /// *not* short-circuit on the presence of a single canonical table — a node |
| 221 | + /// that was behind on schema would then be marked complete while still |
| 222 | + /// missing newer objects. |
205 | 223 | async fn migrate(&self) -> Result<()> { |
206 | | - let stmts = [ |
| 224 | + // Bootstrap: ensure the `schema_migrations` table itself exists. |
| 225 | + sqlx::query( |
| 226 | + r#"CREATE TABLE IF NOT EXISTS schema_migrations ( |
| 227 | + version BIGINT NOT NULL PRIMARY KEY, |
| 228 | + name TEXT NOT NULL, |
| 229 | + applied_at TEXT NOT NULL |
| 230 | + )"#, |
| 231 | + ) |
| 232 | + .execute(&self.pool) |
| 233 | + .await |
| 234 | + .context("creating schema_migrations table")?; |
| 235 | + |
| 236 | + // Serialize migrations across processes: hold a session-level advisory |
| 237 | + // lock on a dedicated connection for the whole run. Another instance |
| 238 | + // starting up blocks here until we finish. The lock is released when we |
| 239 | + // explicitly unlock below, or automatically if the connection is |
| 240 | + // dropped (e.g. on panic), so a crash can't wedge future restarts. |
| 241 | + let mut lock_conn = self |
| 242 | + .pool |
| 243 | + .acquire() |
| 244 | + .await |
| 245 | + .context("acquiring connection for migration advisory lock")?; |
| 246 | + sqlx::query("SELECT pg_advisory_lock($1)") |
| 247 | + .bind(MIGRATION_ADVISORY_LOCK) |
| 248 | + .execute(&mut *lock_conn) |
| 249 | + .await |
| 250 | + .context("acquiring migration advisory lock")?; |
| 251 | + |
| 252 | + let result = self.run_pending_migrations().await; |
| 253 | + |
| 254 | + let _ = sqlx::query("SELECT pg_advisory_unlock($1)") |
| 255 | + .bind(MIGRATION_ADVISORY_LOCK) |
| 256 | + .execute(&mut *lock_conn) |
| 257 | + .await; |
| 258 | + |
| 259 | + result |
| 260 | + } |
| 261 | + |
| 262 | + /// Apply every migration whose version isn't yet recorded, in order. |
| 263 | + /// Must be called while holding the migration advisory lock. |
| 264 | + async fn run_pending_migrations(&self) -> Result<()> { |
| 265 | + for m in MIGRATIONS { |
| 266 | + let already: bool = sqlx::query( |
| 267 | + "SELECT EXISTS(SELECT 1 FROM schema_migrations WHERE version = $1) AS applied", |
| 268 | + ) |
| 269 | + .bind(m.version) |
| 270 | + .fetch_one(&self.pool) |
| 271 | + .await? |
| 272 | + .get::<bool, _>("applied"); |
| 273 | + |
| 274 | + if already { |
| 275 | + continue; |
| 276 | + } |
| 277 | + |
| 278 | + let started = std::time::Instant::now(); |
| 279 | + info!( |
| 280 | + version = m.version, |
| 281 | + name = m.name, |
| 282 | + statements = m.stmts.len(), |
| 283 | + "applying migration" |
| 284 | + ); |
| 285 | + |
| 286 | + // Run the migration body in a single transaction so a failure |
| 287 | + // mid-way leaves the database in its prior state rather than |
| 288 | + // partially mutated. |
| 289 | + let mut tx = self.pool.begin().await?; |
| 290 | + for stmt in m.stmts { |
| 291 | + sqlx::query(stmt).execute(&mut *tx).await.with_context(|| { |
| 292 | + format!( |
| 293 | + "migration v{} ({}) failed on statement: {}", |
| 294 | + m.version, m.name, stmt |
| 295 | + ) |
| 296 | + })?; |
| 297 | + } |
| 298 | + sqlx::query( |
| 299 | + "INSERT INTO schema_migrations (version, name, applied_at) |
| 300 | + VALUES ($1, $2, $3)", |
| 301 | + ) |
| 302 | + .bind(m.version) |
| 303 | + .bind(m.name) |
| 304 | + .bind(Utc::now().to_rfc3339()) |
| 305 | + .execute(&mut *tx) |
| 306 | + .await |
| 307 | + .context("recording migration as applied")?; |
| 308 | + tx.commit() |
| 309 | + .await |
| 310 | + .with_context(|| format!("committing migration v{}", m.version))?; |
| 311 | + |
| 312 | + info!( |
| 313 | + version = m.version, |
| 314 | + name = m.name, |
| 315 | + elapsed_ms = started.elapsed().as_millis() as u64, |
| 316 | + "migration applied" |
| 317 | + ); |
| 318 | + } |
| 319 | + |
| 320 | + Ok(()) |
| 321 | + } |
| 322 | + |
| 323 | + /// Returns `(version, name, applied_at)` for every applied migration, |
| 324 | + /// oldest first. Useful for ops/observability — surface via `gl status` |
| 325 | + /// or `/api/v1/stats` in a follow-up. |
| 326 | + #[allow(dead_code)] |
| 327 | + pub async fn migration_status(&self) -> Result<Vec<(i64, String, String)>> { |
| 328 | + let rows = sqlx::query( |
| 329 | + "SELECT version, name, applied_at FROM schema_migrations ORDER BY version ASC", |
| 330 | + ) |
| 331 | + .fetch_all(&self.pool) |
| 332 | + .await?; |
| 333 | + Ok(rows |
| 334 | + .into_iter() |
| 335 | + .map(|r| { |
| 336 | + ( |
| 337 | + r.get::<i64, _>("version"), |
| 338 | + r.get("name"), |
| 339 | + r.get("applied_at"), |
| 340 | + ) |
| 341 | + }) |
| 342 | + .collect()) |
| 343 | + } |
| 344 | +} |
| 345 | + |
| 346 | +// ── Migration catalogue ────────────────────────────────────────────────────── |
| 347 | +// |
| 348 | +// All schema statements are bundled into a single v1 migration so we can ship |
| 349 | +// versioned migrations on a live network without breaking the existing |
| 350 | +// install base. Future schema changes MUST be added as v2, v3, … — never |
| 351 | +// appended to v1. Operators can read `schema_migrations` to confirm a node |
| 352 | +// is at the expected version. |
| 353 | +// |
| 354 | +// Each migration runs in a single transaction, so statements that Postgres |
| 355 | +// forbids inside a transaction (notably `CREATE INDEX CONCURRENTLY`) cannot be |
| 356 | +// used here. Build such indexes the ordinary, transaction-safe way, or stage |
| 357 | +// them as a dedicated out-of-band operational step. |
| 358 | + |
| 359 | +// Arbitrary but stable key for the migration advisory lock ("gitlawb_" bytes). |
| 360 | +const MIGRATION_ADVISORY_LOCK: i64 = 0x6769_746C_6177_625F; |
| 361 | + |
| 362 | +const MIGRATION_V1_NAME: &str = "initial_schema"; |
| 363 | + |
| 364 | +struct Migration { |
| 365 | + version: i64, |
| 366 | + name: &'static str, |
| 367 | + stmts: &'static [&'static str], |
| 368 | +} |
| 369 | + |
| 370 | +const MIGRATIONS: &[Migration] = &[Migration { |
| 371 | + version: 1, |
| 372 | + name: MIGRATION_V1_NAME, |
| 373 | + stmts: &[ |
207 | 374 | r#"CREATE TABLE IF NOT EXISTS repos ( |
208 | 375 | id TEXT NOT NULL PRIMARY KEY, |
209 | 376 | name TEXT NOT NULL, |
@@ -458,14 +625,8 @@ impl Db { |
458 | 625 | "CREATE INDEX IF NOT EXISTS idx_bounties_status ON bounties(status)", |
459 | 626 | "CREATE INDEX IF NOT EXISTS idx_bounties_repo ON bounties(repo_owner, repo_name)", |
460 | 627 | "CREATE INDEX IF NOT EXISTS idx_bounties_claimant ON bounties(claimant_did)", |
461 | | - ]; |
462 | | - |
463 | | - for stmt in &stmts { |
464 | | - sqlx::query(stmt).execute(&self.pool).await?; |
465 | | - } |
466 | | - Ok(()) |
467 | | - } |
468 | | -} |
| 628 | + ], |
| 629 | +}]; |
469 | 630 |
|
470 | 631 | // ── Repos ───────────────────────────────────────────────────────────────────── |
471 | 632 |
|
@@ -2194,3 +2355,93 @@ impl Db { |
2194 | 2355 | } |
2195 | 2356 | } |
2196 | 2357 | } |
| 2358 | + |
| 2359 | +// ── Tests ───────────────────────────────────────────────────────────────────── |
| 2360 | +// |
| 2361 | +// These tests don't require a live Postgres connection. They validate the |
| 2362 | +// static migration catalogue is well-formed so a future maintainer can't |
| 2363 | +// ship a regression like duplicate versions, negative versions, or empty |
| 2364 | +// migration bodies. The actual SQL execution is exercised by integration |
| 2365 | +// tests / first-run on a real node. |
| 2366 | + |
| 2367 | +#[cfg(test)] |
| 2368 | +mod migration_tests { |
| 2369 | + use super::{MIGRATIONS, MIGRATION_V1_NAME}; |
| 2370 | + |
| 2371 | + #[test] |
| 2372 | + fn migrations_are_non_empty() { |
| 2373 | + assert!( |
| 2374 | + !MIGRATIONS.is_empty(), |
| 2375 | + "MIGRATIONS must contain at least the initial v1 schema" |
| 2376 | + ); |
| 2377 | + } |
| 2378 | + |
| 2379 | + #[test] |
| 2380 | + fn migration_versions_are_strictly_increasing() { |
| 2381 | + let mut last = i64::MIN; |
| 2382 | + for m in MIGRATIONS { |
| 2383 | + assert!( |
| 2384 | + m.version > last, |
| 2385 | + "migration versions must be strictly increasing; \ |
| 2386 | + found {} after {}", |
| 2387 | + m.version, |
| 2388 | + last |
| 2389 | + ); |
| 2390 | + last = m.version; |
| 2391 | + } |
| 2392 | + } |
| 2393 | + |
| 2394 | + #[test] |
| 2395 | + fn migration_versions_start_at_one() { |
| 2396 | + // A version of 0 (or negative) would be a footgun: any future |
| 2397 | + // `WHERE version > current_max` style query would skip it. |
| 2398 | + assert_eq!( |
| 2399 | + MIGRATIONS.first().map(|m| m.version), |
| 2400 | + Some(1), |
| 2401 | + "the first migration must have version 1" |
| 2402 | + ); |
| 2403 | + } |
| 2404 | + |
| 2405 | + #[test] |
| 2406 | + fn migration_names_are_non_empty_and_distinct() { |
| 2407 | + let mut seen = std::collections::HashSet::new(); |
| 2408 | + for m in MIGRATIONS { |
| 2409 | + assert!( |
| 2410 | + !m.name.is_empty(), |
| 2411 | + "migration v{} has empty name", |
| 2412 | + m.version |
| 2413 | + ); |
| 2414 | + assert!( |
| 2415 | + !m.name.contains(char::is_whitespace), |
| 2416 | + "migration v{} name {:?} contains whitespace", |
| 2417 | + m.version, |
| 2418 | + m.name |
| 2419 | + ); |
| 2420 | + assert!( |
| 2421 | + seen.insert(m.name), |
| 2422 | + "duplicate migration name: {:?}", |
| 2423 | + m.name |
| 2424 | + ); |
| 2425 | + } |
| 2426 | + } |
| 2427 | + |
| 2428 | + #[test] |
| 2429 | + fn migration_bodies_are_non_empty() { |
| 2430 | + for m in MIGRATIONS { |
| 2431 | + assert!( |
| 2432 | + !m.stmts.is_empty(), |
| 2433 | + "migration v{} ({}) has no SQL statements", |
| 2434 | + m.version, |
| 2435 | + m.name |
| 2436 | + ); |
| 2437 | + } |
| 2438 | + } |
| 2439 | + |
| 2440 | + #[test] |
| 2441 | + fn v1_name_is_the_initial_schema() { |
| 2442 | + // This is what the legacy-install backfill writes to |
| 2443 | + // `schema_migrations` when an existing node upgrades. If you rename |
| 2444 | + // it, you must also update the backfill. |
| 2445 | + assert_eq!(MIGRATIONS[0].name, MIGRATION_V1_NAME); |
| 2446 | + } |
| 2447 | +} |
0 commit comments