platform_core/
migrations.rs1use 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}