systemprompt_database/lifecycle/migrations/
down.rs1use super::exec::{TrackingWrite, execute_statements_transactional};
7use super::{MigrationResult, MigrationService};
8use crate::services::SqlExecutor;
9use systemprompt_extension::{Extension, LoaderError, Migration};
10use systemprompt_identifiers::{ExtensionId, ToDbValue};
11use tracing::info;
12
13impl MigrationService<'_> {
14 pub async fn run_down_migrations(
15 &self,
16 extension: &dyn Extension,
17 count: u32,
18 ) -> Result<MigrationResult, LoaderError> {
19 if count == 0 {
20 return Ok(MigrationResult::default());
21 }
22
23 let ext_id = &ExtensionId::new(extension.metadata().id);
24 self.ensure_migrations_table_exists().await?;
25
26 let result = self
27 .db
28 .query_raw_with(
29 &"SELECT version FROM extension_migrations WHERE extension_id = $1 ORDER BY \
30 version DESC LIMIT $2",
31 &[&ext_id, &count],
32 )
33 .await
34 .map_err(|e| LoaderError::MigrationStepFailed {
35 extension: ext_id.clone(),
36 context: "Failed to query applied migrations for revert".to_owned(),
37 source: Box::new(e),
38 })?;
39
40 let versions = result
41 .rows
42 .iter()
43 .map(|row| {
44 row.get("version")
45 .and_then(serde_json::Value::as_i64)
46 .and_then(|v| u32::try_from(v).ok())
47 .ok_or_else(|| LoaderError::MigrationFailed {
48 extension: ext_id.clone(),
49 message: "extension_migrations row has a malformed `version` column"
50 .to_owned(),
51 })
52 })
53 .collect::<Result<Vec<u32>, LoaderError>>()?;
54
55 if versions.is_empty() {
56 return Ok(MigrationResult::default());
57 }
58
59 let migrations = extension.migrations();
60 let mut migrations_run = 0;
61
62 for version in versions {
63 self.revert_version(ext_id, version, &migrations).await?;
64 migrations_run += 1;
65 }
66
67 Ok(MigrationResult {
68 migrations_run,
69 migrations_skipped: 0,
70 })
71 }
72
73 async fn revert_version(
74 &self,
75 ext_id: &ExtensionId,
76 version: u32,
77 migrations: &[Migration],
78 ) -> Result<(), LoaderError> {
79 let migration = migrations
80 .iter()
81 .find(|m| m.version == version)
82 .ok_or_else(|| LoaderError::MigrationFailed {
83 extension: ext_id.clone(),
84 message: format!(
85 "Cannot revert migration {version}: not declared in Extension::migrations()"
86 ),
87 })?;
88
89 if migration.tombstone {
90 return Err(LoaderError::MigrationFailed {
91 extension: ext_id.clone(),
92 message: format!(
93 "Cannot revert migration {version} ('{}'): the slot is tombstoned — its file \
94 was deleted, so there is no down SQL to run",
95 migration.name
96 ),
97 });
98 }
99
100 let down_sql = migration
101 .down
102 .ok_or_else(|| LoaderError::MigrationNotReversible {
103 extension: ext_id.clone(),
104 version,
105 })?;
106
107 info!(
108 extension = %ext_id,
109 version = migration.version,
110 name = %migration.name,
111 "Reverting migration"
112 );
113
114 let statements = SqlExecutor::parse_sql_statements(down_sql).map_err(|e| {
115 LoaderError::MigrationStepFailed {
116 extension: ext_id.clone(),
117 context: format!(
118 "Failed to parse down migration {} ({})",
119 migration.version, migration.name
120 ),
121 source: Box::new(e),
122 }
123 })?;
124 let delete_params: [&dyn ToDbValue; 2] = [&ext_id, &version];
125 execute_statements_transactional(
126 self.db,
127 &statements,
128 ext_id,
129 migration,
130 Some(TrackingWrite {
131 sql: "DELETE FROM extension_migrations WHERE extension_id = $1 AND version = $2",
132 params: &delete_params,
133 }),
134 )
135 .await
136 }
137}