Skip to main content

systemprompt_database/lifecycle/migrations/
stamp.rs

1//! Fresh-install baseline stamping.
2//!
3//! The declarative schema (`schema/*.sql`) is the baseline: a fresh database
4//! reaches target shape from the structural/dependent DDL alone, so its
5//! migrations carry no information and must not execute. [`MigrationService::
6//! assess_freshness`] decides, before any DDL has run, whether an extension is
7//! landing on a fresh database; [`MigrationService::baseline_stamp_rows`] then
8//! yields the `extension_migrations` rows recording every defined migration as
9//! applied, which the installer commits alongside the structural DDL rather
10//! than executing their SQL. Established databases (any tracking history, or
11//! any owned table already present) take the normal incremental path.
12//!
13//! One class of migration is stamped **and** executed: a retirement, whose
14//! every statement is a `DROP … IF EXISTS` or a `DELETE FROM
15//! extension_migrations`. Such a migration retires relations that another,
16//! since-deleted extension left behind, and an extension whose own tables are
17//! all absent says nothing about theirs — a production database kept nineteen
18//! `eval_*` tables and three orphaned ledger rows because the migration that
19//! dropped them belonged to an extension the database was meeting for the
20//! first time. Every statement of a retirement is idempotent, so running it
21//! on a truly fresh database is a no-op.
22//!
23//! Copyright (c) systemprompt.io — Business Source License 1.1.
24//! See <https://systemprompt.io> for licensing details.
25
26use super::MigrationService;
27use super::exec::execute_statements_transactional;
28use crate::services::SqlExecutor;
29use pg_query::NodeEnum;
30use systemprompt_extension::{Extension, LoaderError, Migration};
31use systemprompt_identifiers::ExtensionId;
32use tracing::{info, warn};
33
34/// One `extension_migrations` row recording a migration as applied without
35/// having executed it.
36#[derive(Debug, Clone)]
37pub struct BaselineStamp {
38    pub id: String,
39    pub version: u32,
40    pub name: String,
41    pub checksum: String,
42}
43
44#[derive(Debug, Clone, Copy)]
45pub struct FreshnessCheck {
46    pub no_history: bool,
47    pub tables_present: usize,
48    pub tables_total: usize,
49}
50
51impl FreshnessCheck {
52    #[must_use]
53    pub const fn is_fresh(&self) -> bool {
54        self.no_history && self.tables_present == 0
55    }
56}
57
58impl MigrationService<'_> {
59    pub async fn assess_freshness(
60        &self,
61        extension_id: &ExtensionId,
62        owned_tables: &[String],
63    ) -> Result<FreshnessCheck, LoaderError> {
64        self.ensure_migrations_table_exists().await?;
65
66        let no_history = self.get_applied_migrations(extension_id).await?.is_empty();
67
68        let mut tables_present = 0usize;
69        for table in owned_tables {
70            let (schema, name) = table.split_once('.').unwrap_or(("public", table.as_str()));
71            let result = self
72                .db
73                .query_raw_with(
74                    &"SELECT 1 AS present FROM information_schema.tables WHERE table_schema = $1 \
75                      AND table_name = $2",
76                    &[&schema, &name],
77                )
78                .await
79                .map_err(|e| LoaderError::MigrationStepFailed {
80                    extension: extension_id.clone(),
81                    context: format!("Failed to check for existing table '{table}'"),
82                    source: Box::new(e),
83                })?;
84            if !result.rows.is_empty() {
85                tables_present += 1;
86            }
87        }
88
89        let check = FreshnessCheck {
90            no_history,
91            tables_present,
92            tables_total: owned_tables.len(),
93        };
94
95        if check.no_history && check.tables_present > 0 && check.tables_present < check.tables_total
96        {
97            warn!(
98                extension = %extension_id,
99                tables_present = check.tables_present,
100                tables_total = check.tables_total,
101                "Extension has no migration history but some owned tables already exist; \
102                 treating as an established database and executing migrations normally"
103            );
104        }
105
106        Ok(check)
107    }
108
109    pub async fn run_stamped_retirements(
110        &self,
111        extension: &dyn Extension,
112    ) -> Result<usize, LoaderError> {
113        let ext_id = &ExtensionId::new(extension.metadata().id);
114        let mut ran = 0usize;
115        for migration in extension
116            .migrations()
117            .iter()
118            .filter(|migration| !migration.tombstone && is_retirement(migration))
119        {
120            let statements = SqlExecutor::parse_sql_statements(migration.sql).map_err(|e| {
121                LoaderError::MigrationStepFailed {
122                    extension: ext_id.clone(),
123                    context: format!(
124                        "Failed to parse retirement migration {} ({})",
125                        migration.version, migration.name
126                    ),
127                    source: Box::new(e),
128                }
129            })?;
130            info!(
131                extension = %ext_id,
132                version = migration.version,
133                name = %migration.name,
134                "Fresh install: executing stamped retirement migration"
135            );
136            execute_statements_transactional(self.db, &statements, ext_id, migration, None).await?;
137            ran += 1;
138        }
139        Ok(ran)
140    }
141
142    #[must_use]
143    pub fn baseline_stamp_rows(extension: &dyn Extension) -> Vec<BaselineStamp> {
144        let ext_id = &ExtensionId::new(extension.metadata().id);
145        extension
146            .migrations()
147            .iter()
148            .filter(|migration| !migration.tombstone)
149            .map(|migration| BaselineStamp {
150                id: format!("{}_{:03}", ext_id, migration.version),
151                version: migration.version,
152                name: migration.name.clone(),
153                checksum: migration.checksum(),
154            })
155            .collect()
156    }
157}
158
159#[must_use]
160pub fn is_retirement(migration: &Migration) -> bool {
161    let Ok(parsed) = pg_query::parse(migration.sql) else {
162        return false;
163    };
164    let mut statements = 0usize;
165    for raw in parsed.protobuf.stmts {
166        let Some(node) = raw.stmt.and_then(|s| s.node) else {
167            continue;
168        };
169        statements += 1;
170        let retires = match &node {
171            NodeEnum::DropStmt(drop) => drop.missing_ok,
172            NodeEnum::DeleteStmt(delete) => delete
173                .relation
174                .as_ref()
175                .is_some_and(|relation| relation.relname == "extension_migrations"),
176            _ => false,
177        };
178        if !retires {
179            return false;
180        }
181    }
182    statements > 0
183}