systemprompt_database/lifecycle/migrations/
repair.rs1use 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}