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 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}