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    Migration {
65        name: "platform/0014_capture_provider_http_body_evidence",
66        sql: include_str!("../migrations/0014_capture_provider_http_body_evidence.sql"),
67    },
68];
69
70pub async fn apply_migrations(pool: &PgPool, migrations: &[Migration]) -> AppResult<()> {
71    ensure_migration_table(pool).await?;
72
73    for migration in migrations {
74        apply_migration(pool, migration).await?;
75    }
76
77    Ok(())
78}
79
80pub async fn apply_module_migration(
81    pool: &PgPool,
82    name: &str,
83    artifact_digest: &str,
84    sql: &str,
85) -> AppResult<()> {
86    let observed_digest = {
87        use std::fmt::Write as _;
88        let mut value = String::from("sha256:");
89        for byte in Sha256::digest(sql.as_bytes()) {
90            write!(&mut value, "{byte:02x}").expect("writing to a String cannot fail");
91        }
92        value
93    };
94    if name.trim().is_empty() || observed_digest != artifact_digest {
95        return Err(AppError::new(
96            ErrorCode::Internal,
97            "Module migration identity or artifact digest is invalid",
98        ));
99    }
100    let mut tx = pool.begin().await.map_err(map_migration_error)?;
101    sqlx::raw_sql(
102        r#"
103        create schema if not exists platform;
104        create table if not exists platform.module_schema_migrations (
105            name text primary key,
106            artifact_digest text not null,
107            applied_at timestamptz not null default now()
108        );
109        "#,
110    )
111    .execute(&mut *tx)
112    .await
113    .map_err(map_migration_error)?;
114    let existing: Option<String> = sqlx::query_scalar(
115        "select artifact_digest from platform.module_schema_migrations where name = $1",
116    )
117    .bind(name)
118    .fetch_optional(&mut *tx)
119    .await
120    .map_err(map_migration_error)?;
121    if let Some(existing) = existing {
122        if existing != artifact_digest {
123            return Err(AppError::new(
124                ErrorCode::Conflict,
125                "Applied Module migration digest differs from the reviewed artifact",
126            ));
127        }
128        tx.commit().await.map_err(map_migration_error)?;
129        return Ok(());
130    }
131    sqlx::query(sqlx::AssertSqlSafe(sql.to_owned()))
132        .execute(&mut *tx)
133        .await
134        .map_err(map_migration_error)?;
135    sqlx::query(
136        "insert into platform.module_schema_migrations (name, artifact_digest) values ($1, $2)",
137    )
138    .bind(name)
139    .bind(artifact_digest)
140    .execute(&mut *tx)
141    .await
142    .map_err(map_migration_error)?;
143    tx.commit().await.map_err(map_migration_error)
144}
145
146async fn ensure_migration_table(pool: &PgPool) -> AppResult<()> {
147    sqlx::raw_sql(
148        r#"
149        create schema if not exists platform;
150
151        create table if not exists platform.schema_migrations (
152            name text primary key,
153            applied_at timestamptz not null default now()
154        );
155        "#,
156    )
157    .execute(pool)
158    .await
159    .map(|_| ())
160    .map_err(map_migration_error)
161}
162
163async fn apply_migration(pool: &PgPool, migration: &Migration) -> AppResult<()> {
164    let mut tx = pool.begin().await.map_err(map_migration_error)?;
165
166    let already_applied: Option<String> = sqlx::query_scalar(
167        r#"
168        select name
169        from platform.schema_migrations
170        where name = $1
171        "#,
172    )
173    .bind(migration.name)
174    .fetch_optional(&mut *tx)
175    .await
176    .map_err(map_migration_error)?;
177
178    if already_applied.is_some() {
179        tx.commit().await.map_err(map_migration_error)?;
180        return Ok(());
181    }
182
183    sqlx::raw_sql(migration.sql)
184        .execute(&mut *tx)
185        .await
186        .map_err(map_migration_error)?;
187
188    sqlx::query(
189        r#"
190        insert into platform.schema_migrations (name)
191        values ($1)
192        on conflict (name) do nothing
193        "#,
194    )
195    .bind(migration.name)
196    .execute(&mut *tx)
197    .await
198    .map_err(map_migration_error)?;
199
200    tx.commit().await.map_err(map_migration_error)
201}
202
203fn map_migration_error(source: sqlx::Error) -> AppError {
204    AppError::new(ErrorCode::Internal, "Database migration failed").with_source(source)
205}