Skip to main content

platform_core/
migrations.rs

1use crate::error::{AppError, AppResult, ErrorCode};
2use sqlx::PgPool;
3
4#[derive(Debug, Clone, Copy)]
5pub struct Migration {
6    pub name: &'static str,
7    pub sql: &'static str,
8}
9
10pub const PLATFORM_MIGRATIONS: &[Migration] = &[
11    Migration {
12        name: "platform/0001_create_platform_schema",
13        sql: include_str!("../migrations/0001_create_platform_schema.sql"),
14    },
15    Migration {
16        name: "platform/0002_create_outbox",
17        sql: include_str!("../migrations/0002_create_outbox.sql"),
18    },
19    Migration {
20        name: "platform/0003_extend_outbox_delivery_fields",
21        sql: include_str!("../migrations/0003_extend_outbox_delivery_fields.sql"),
22    },
23    Migration {
24        name: "platform/0004_add_outbox_summary_index",
25        sql: include_str!("../migrations/0004_add_outbox_summary_index.sql"),
26    },
27    Migration {
28        name: "platform/0005_create_execution_logs",
29        sql: include_str!("../migrations/0005_create_execution_logs.sql"),
30    },
31    Migration {
32        name: "platform/0006_create_story_events",
33        sql: include_str!("../migrations/0006_create_story_events.sql"),
34    },
35    Migration {
36        name: "platform/0007_create_config_schema",
37        sql: include_str!("../migrations/0007_create_config_schema.sql"),
38    },
39    Migration {
40        name: "platform/0008_create_remote_http_proxy_calls",
41        sql: include_str!("../migrations/0008_create_remote_http_proxy_calls.sql"),
42    },
43    Migration {
44        name: "platform/0009_add_story_query_indexes",
45        sql: include_str!("../migrations/0009_add_story_query_indexes.sql"),
46    },
47    Migration {
48        name: "platform/0010_create_idempotency_claims",
49        sql: include_str!("../migrations/0010_create_idempotency_claims.sql"),
50    },
51    Migration {
52        name: "platform/0011_create_extraction_artifacts",
53        sql: include_str!("../migrations/0011_create_extraction_artifacts.sql"),
54    },
55    Migration {
56        name: "platform/0012_create_delivery_artifacts",
57        sql: include_str!("../migrations/0012_create_delivery_artifacts.sql"),
58    },
59];
60
61pub async fn apply_migrations(pool: &PgPool, migrations: &[Migration]) -> AppResult<()> {
62    ensure_migration_table(pool).await?;
63
64    for migration in migrations {
65        apply_migration(pool, migration).await?;
66    }
67
68    Ok(())
69}
70
71async fn ensure_migration_table(pool: &PgPool) -> AppResult<()> {
72    sqlx::raw_sql(
73        r#"
74        create schema if not exists platform;
75
76        create table if not exists platform.schema_migrations (
77            name text primary key,
78            applied_at timestamptz not null default now()
79        );
80        "#,
81    )
82    .execute(pool)
83    .await
84    .map(|_| ())
85    .map_err(map_migration_error)
86}
87
88async fn apply_migration(pool: &PgPool, migration: &Migration) -> AppResult<()> {
89    let mut tx = pool.begin().await.map_err(map_migration_error)?;
90
91    let already_applied: Option<String> = sqlx::query_scalar(
92        r#"
93        select name
94        from platform.schema_migrations
95        where name = $1
96        "#,
97    )
98    .bind(migration.name)
99    .fetch_optional(&mut *tx)
100    .await
101    .map_err(map_migration_error)?;
102
103    if already_applied.is_some() {
104        tx.commit().await.map_err(map_migration_error)?;
105        return Ok(());
106    }
107
108    sqlx::raw_sql(migration.sql)
109        .execute(&mut *tx)
110        .await
111        .map_err(map_migration_error)?;
112
113    sqlx::query(
114        r#"
115        insert into platform.schema_migrations (name)
116        values ($1)
117        on conflict (name) do nothing
118        "#,
119    )
120    .bind(migration.name)
121    .execute(&mut *tx)
122    .await
123    .map_err(map_migration_error)?;
124
125    tx.commit().await.map_err(map_migration_error)
126}
127
128fn map_migration_error(source: sqlx::Error) -> AppError {
129    AppError::new(ErrorCode::Internal, "Database migration failed").with_source(source)
130}