Skip to main content

systemprompt_database/lifecycle/migrations/
repair.rs

1//! Migration checksum-drift repair.
2//!
3//! When an already-applied migration file is edited in place, its stored
4//! checksum stops matching the file and the runner refuses to proceed.
5//! [`MigrationService::repair_drift`] re-executes each drifted migration and
6//! rewrites its stored checksum in the same transaction; the tracking row is
7//! never deleted, so a failed re-apply rolls back to "drifted but tracked"
8//! instead of leaving the migration untracked and crash-looping the next
9//! boot. Re-applying requires the migration SQL to be re-executable against
10//! the current schema — a later migration may have invalidated that, in which
11//! case [`MigrationService::reconcile_drift`] rewrites the stored checksum
12//! without executing any SQL. `no_transaction` migrations cannot be repaired
13//! atomically: a mid-SQL failure leaves the row tracked with the old
14//! checksum, which still reports as drift rather than crash-looping.
15//!
16//! Copyright (c) systemprompt.io — Business Source License 1.1.
17//! See <https://systemprompt.io> for licensing details.
18
19use super::exec::{TrackingWrite, check_cross_extension_alters, execute_statements_transactional};
20use super::{ChecksumDrift, MigrationService};
21use crate::lifecycle::installation::BootstrapLockGuard;
22use crate::services::SqlExecutor;
23use systemprompt_extension::{Extension, LoaderError, Migration};
24use systemprompt_identifiers::ToDbValue;
25
26const UPDATE_CHECKSUM_SQL: &str =
27    "UPDATE extension_migrations SET checksum = $3 WHERE extension_id = $1 AND version = $2";
28
29#[derive(Debug, Default, Clone)]
30pub struct RepairResult {
31    pub repaired: Vec<ChecksumDrift>,
32    pub migrations_run: usize,
33}
34
35impl MigrationService<'_> {
36    pub async fn repair_drift(
37        &self,
38        extension: &dyn Extension,
39    ) -> Result<RepairResult, LoaderError> {
40        let status = self.status(extension).await?;
41
42        if status.drift.is_empty() {
43            return Ok(RepairResult::default());
44        }
45
46        let guard = BootstrapLockGuard::acquire(self.db).await?;
47        let outcome = self.reapply_drifted(extension, &status.drift).await;
48        let pending = match outcome {
49            Ok(()) => self.run_pending_migrations(extension).await,
50            Err(e) => Err(e),
51        };
52        guard.release().await;
53        let result = pending?;
54
55        Ok(RepairResult {
56            repaired: status.drift,
57            migrations_run: result.migrations_run,
58        })
59    }
60
61    pub async fn reconcile_drift(
62        &self,
63        extension: &dyn Extension,
64    ) -> Result<RepairResult, LoaderError> {
65        let status = self.status(extension).await?;
66
67        if status.drift.is_empty() {
68            return Ok(RepairResult::default());
69        }
70
71        let guard = BootstrapLockGuard::acquire(self.db).await?;
72        let mut outcome = Ok(());
73        for drift in &status.drift {
74            if let Err(e) = self.rewrite_checksum(drift).await {
75                outcome = Err(e);
76                break;
77            }
78        }
79        guard.release().await;
80        outcome?;
81
82        Ok(RepairResult {
83            repaired: status.drift,
84            migrations_run: 0,
85        })
86    }
87
88    async fn rewrite_checksum(&self, drift: &ChecksumDrift) -> Result<(), LoaderError> {
89        self.db
90            .execute(
91                &UPDATE_CHECKSUM_SQL,
92                &[&drift.extension_id, &drift.version, &drift.current_checksum],
93            )
94            .await
95            .map_err(|e| LoaderError::MigrationFailed {
96                extension: drift.extension_id.clone(),
97                message: format!(
98                    "Failed to rewrite checksum for migration {} ('{}'): {e}",
99                    drift.version, drift.name
100                ),
101            })?;
102        Ok(())
103    }
104
105    async fn reapply_drifted(
106        &self,
107        extension: &dyn Extension,
108        drift: &[ChecksumDrift],
109    ) -> Result<(), LoaderError> {
110        let ext_id = extension.metadata().id;
111        let migrations = extension.migrations();
112
113        for d in drift {
114            let migration = migrations
115                .iter()
116                .find(|m| m.version == d.version)
117                .ok_or_else(|| LoaderError::MigrationFailed {
118                    extension: ext_id.to_owned(),
119                    message: format!(
120                        "Drifted migration {} ('{}') is no longer declared by extension \
121                         '{ext_id}'",
122                        d.version, d.name
123                    ),
124                })?;
125            self.reapply_one(extension, migration, d).await?;
126        }
127
128        Ok(())
129    }
130
131    async fn reapply_one(
132        &self,
133        extension: &dyn Extension,
134        migration: &Migration,
135        drift: &ChecksumDrift,
136    ) -> Result<(), LoaderError> {
137        let ext_id = extension.metadata().id;
138
139        check_cross_extension_alters(extension, migration)?;
140
141        tracing::info!(
142            extension = %ext_id,
143            version = migration.version,
144            name = %migration.name,
145            no_transaction = migration.no_transaction,
146            "Re-applying drifted migration"
147        );
148
149        let update_params: [&dyn ToDbValue; 3] =
150            [&drift.extension_id, &drift.version, &drift.current_checksum];
151
152        if migration.no_transaction {
153            SqlExecutor::execute_statements_parsed(self.db, migration.sql)
154                .await
155                .map_err(|e| LoaderError::MigrationFailed {
156                    extension: ext_id.to_owned(),
157                    message: format!(
158                        "Failed to re-apply drifted migration {} ({}): {e}",
159                        migration.version, migration.name
160                    ),
161                })?;
162            self.rewrite_checksum(drift).await
163        } else {
164            let statements = SqlExecutor::parse_sql_statements(migration.sql).map_err(|e| {
165                LoaderError::MigrationFailed {
166                    extension: ext_id.to_owned(),
167                    message: format!(
168                        "Failed to parse migration {} ({}): {e}",
169                        migration.version, migration.name
170                    ),
171                }
172            })?;
173            execute_statements_transactional(
174                self.db,
175                &statements,
176                ext_id,
177                migration,
178                Some(TrackingWrite {
179                    sql: UPDATE_CHECKSUM_SQL,
180                    params: &update_params,
181                }),
182            )
183            .await
184        }
185    }
186}