Skip to main content

systemprompt_database/lifecycle/migrations/
down.rs

1//! Reverting applied migrations via their declared `down` SQL.
2//!
3//! Copyright (c) systemprompt.io — Business Source License 1.1.
4//! See <https://systemprompt.io> for licensing details.
5
6use 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}