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