Skip to main content

platform_core/
migrations.rs

1use crate::error::{AppError, AppResult, ErrorCode};
2use sha2::{Digest as _, Sha256};
3use sqlx::PgPool;
4
5#[derive(Debug, Clone, Copy)]
6pub struct Migration {
7    pub name: &'static str,
8    pub sql: &'static str,
9}
10
11pub const PLATFORM_MIGRATIONS: &[Migration] = &[
12    Migration {
13        name: "platform/0001_create_platform_schema",
14        sql: include_str!("../migrations/0001_create_platform_schema.sql"),
15    },
16    Migration {
17        name: "platform/0002_create_outbox",
18        sql: include_str!("../migrations/0002_create_outbox.sql"),
19    },
20    Migration {
21        name: "platform/0003_extend_outbox_delivery_fields",
22        sql: include_str!("../migrations/0003_extend_outbox_delivery_fields.sql"),
23    },
24    Migration {
25        name: "platform/0004_add_outbox_summary_index",
26        sql: include_str!("../migrations/0004_add_outbox_summary_index.sql"),
27    },
28    Migration {
29        name: "platform/0005_create_execution_logs",
30        sql: include_str!("../migrations/0005_create_execution_logs.sql"),
31    },
32    Migration {
33        name: "platform/0006_create_story_events",
34        sql: include_str!("../migrations/0006_create_story_events.sql"),
35    },
36    Migration {
37        name: "platform/0007_create_config_schema",
38        sql: include_str!("../migrations/0007_create_config_schema.sql"),
39    },
40    Migration {
41        name: "platform/0008_create_provider_http_calls",
42        sql: include_str!("../migrations/0008_create_provider_http_calls.sql"),
43    },
44    Migration {
45        name: "platform/0009_add_story_query_indexes",
46        sql: include_str!("../migrations/0009_add_story_query_indexes.sql"),
47    },
48    Migration {
49        name: "platform/0010_create_idempotency_claims",
50        sql: include_str!("../migrations/0010_create_idempotency_claims.sql"),
51    },
52    Migration {
53        name: "platform/0011_create_extraction_artifacts",
54        sql: include_str!("../migrations/0011_create_extraction_artifacts.sql"),
55    },
56    Migration {
57        name: "platform/0012_create_delivery_artifacts",
58        sql: include_str!("../migrations/0012_create_delivery_artifacts.sql"),
59    },
60    Migration {
61        name: "platform/0013_create_provider_host_effect_commits",
62        sql: include_str!("../migrations/0013_create_provider_host_effect_commits.sql"),
63    },
64];
65
66pub async fn apply_migrations(pool: &PgPool, migrations: &[Migration]) -> AppResult<()> {
67    ensure_migration_table(pool).await?;
68
69    for migration in migrations {
70        apply_migration(pool, migration).await?;
71    }
72
73    Ok(())
74}
75
76pub async fn apply_module_migration(
77    pool: &PgPool,
78    name: &str,
79    artifact_digest: &str,
80    sql: &str,
81) -> AppResult<()> {
82    let observed_digest = {
83        use std::fmt::Write as _;
84        let mut value = String::from("sha256:");
85        for byte in Sha256::digest(sql.as_bytes()) {
86            write!(&mut value, "{byte:02x}").expect("writing to a String cannot fail");
87        }
88        value
89    };
90    if name.trim().is_empty() || observed_digest != artifact_digest {
91        return Err(AppError::new(
92            ErrorCode::Internal,
93            "Module migration identity or artifact digest is invalid",
94        ));
95    }
96    let mut tx = pool.begin().await.map_err(map_migration_error)?;
97    sqlx::raw_sql(
98        r#"
99        create schema if not exists platform;
100        create table if not exists platform.module_schema_migrations (
101            name text primary key,
102            artifact_digest text not null,
103            applied_at timestamptz not null default now()
104        );
105        "#,
106    )
107    .execute(&mut *tx)
108    .await
109    .map_err(map_migration_error)?;
110    let existing: Option<String> = sqlx::query_scalar(
111        "select artifact_digest from platform.module_schema_migrations where name = $1",
112    )
113    .bind(name)
114    .fetch_optional(&mut *tx)
115    .await
116    .map_err(map_migration_error)?;
117    if let Some(existing) = existing {
118        if existing != artifact_digest {
119            return Err(AppError::new(
120                ErrorCode::Conflict,
121                "Applied Module migration digest differs from the reviewed artifact",
122            ));
123        }
124        tx.commit().await.map_err(map_migration_error)?;
125        return Ok(());
126    }
127    sqlx::query(sqlx::AssertSqlSafe(sql.to_owned()))
128        .execute(&mut *tx)
129        .await
130        .map_err(map_migration_error)?;
131    sqlx::query(
132        "insert into platform.module_schema_migrations (name, artifact_digest) values ($1, $2)",
133    )
134    .bind(name)
135    .bind(artifact_digest)
136    .execute(&mut *tx)
137    .await
138    .map_err(map_migration_error)?;
139    tx.commit().await.map_err(map_migration_error)
140}
141
142async fn ensure_migration_table(pool: &PgPool) -> AppResult<()> {
143    sqlx::raw_sql(
144        r#"
145        create schema if not exists platform;
146
147        create table if not exists platform.schema_migrations (
148            name text primary key,
149            applied_at timestamptz not null default now()
150        );
151        "#,
152    )
153    .execute(pool)
154    .await
155    .map(|_| ())
156    .map_err(map_migration_error)
157}
158
159async fn apply_migration(pool: &PgPool, migration: &Migration) -> AppResult<()> {
160    let mut tx = pool.begin().await.map_err(map_migration_error)?;
161
162    let already_applied: Option<String> = sqlx::query_scalar(
163        r#"
164        select name
165        from platform.schema_migrations
166        where name = $1
167        "#,
168    )
169    .bind(migration.name)
170    .fetch_optional(&mut *tx)
171    .await
172    .map_err(map_migration_error)?;
173
174    if already_applied.is_some() {
175        tx.commit().await.map_err(map_migration_error)?;
176        return Ok(());
177    }
178
179    sqlx::raw_sql(migration.sql)
180        .execute(&mut *tx)
181        .await
182        .map_err(map_migration_error)?;
183
184    sqlx::query(
185        r#"
186        insert into platform.schema_migrations (name)
187        values ($1)
188        on conflict (name) do nothing
189        "#,
190    )
191    .bind(migration.name)
192    .execute(&mut *tx)
193    .await
194    .map_err(map_migration_error)?;
195
196    tx.commit().await.map_err(map_migration_error)
197}
198
199fn map_migration_error(source: sqlx::Error) -> AppError {
200    AppError::new(ErrorCode::Internal, "Database migration failed").with_source(source)
201}