Skip to main content

systemprompt_database/lifecycle/installation/
migration_cost.rs

1//! Finds the statements in a migration that rewrite a hot table, and pairs
2//! them with the cost its author measured.
3//!
4//! A migration runs inside the boot, before the HTTP listener is bound, in
5//! one transaction, holding its locks, and every per-row trigger on the table
6//! fires for every row it touches. That is how a 3,644-row `UPDATE
7//! ai_requests` took 27 minutes on a production instance: one trigger
8//! re-enqueued the whole client session per row, 113 ms a row. Suspending it
9//! took the same statement to 2.0 s. The same shape emptied `logs` once and
10//! wrote 77,797 outbox tombstones and as many `pg_notify` calls in a single
11//! transaction.
12//!
13//! What this module reports is deliberately blunt: any write to a table on
14//! the hot list, plus the `ALTER`/`CREATE INDEX` forms that take a full scan
15//! or a blocking lock. Judging whether a `WHERE` clause is selective is not
16//! something a parser can do — the author is the one who can measure it, and
17//! [`CostDirective`] is where they say so.
18//!
19//! [`HOT_TABLES`] is the core-shipped list: row counts are from the 2026-09-22
20//! production analysis, and every one of them grows with traffic while none is
21//! ever pruned to a bounded size. It is passed in rather than read directly so
22//! an installation can add the tables it owns — core cannot know about a
23//! downstream repo's hottest table.
24//!
25//! Two consumers, one detector, and they read it over different populations.
26//!
27//! Each repo's test suite runs it over the whole catalogue, where the findings
28//! are a hard failure against a baseline of migrations that predate the gate.
29//! That is where an unmeasured backfill is supposed to be caught — in the pull
30//! request, not on a customer's server.
31//!
32//! The boot runs it over the migrations that are about to execute against a
33//! table that already holds rows, and only warns: a missing comment must never
34//! brick an upgrade, and the statement timeout derived from `measured` is what
35//! actually bounds the damage. It is deliberately not the whole catalogue. A
36//! migration that has already run cannot be made cheaper by a comment, so
37//! reporting it says nothing the operator can act on, and the grandfathered
38//! set is large enough that doing so buried the one finding that mattered.
39//!
40//! Copyright (c) systemprompt.io — Business Source License 1.1.
41//! See <https://systemprompt.io> for licensing details.
42
43use std::sync::Arc;
44
45use pg_query::NodeEnum;
46use pg_query::protobuf::AlterTableType;
47use systemprompt_extension::Extension;
48use systemprompt_extension::cost::{self, CostDirective};
49use systemprompt_identifiers::ExtensionId;
50
51pub const HOT_TABLES: &[&str] = &[
52    "ai_requests",
53    "ai_request_messages",
54    "ai_request_payloads",
55    "ai_request_client_evidence",
56    "ai_request_tool_calls",
57    "analytics_events",
58    "event_outbox",
59    "logs",
60    "user_sessions",
61];
62
63/// One statement that rewrites or rescans a hot table. `position` is 1-based,
64/// matching how the runner numbers statements when one fails.
65#[derive(Debug, Clone, PartialEq, Eq)]
66pub struct ExpensiveStatement {
67    pub position: usize,
68    pub table: String,
69    pub form: &'static str,
70}
71
72/// One migration's expensive statements and what it declared about them.
73/// `malformed` is set when the body carries a `@cost` line that does not parse.
74#[derive(Debug, Clone)]
75pub struct MigrationCost {
76    pub extension: ExtensionId,
77    pub migration: String,
78    pub statements: Vec<ExpensiveStatement>,
79    pub declared: Option<CostDirective>,
80    pub malformed: Option<String>,
81}
82
83impl MigrationCost {
84    #[must_use]
85    pub const fn is_undeclared(&self) -> bool {
86        !self.statements.is_empty() && self.declared.is_none()
87    }
88
89    #[must_use]
90    pub fn label(&self) -> String {
91        format!("{}/{}", self.extension, self.migration)
92    }
93
94    #[must_use]
95    pub fn statement_summary(&self) -> String {
96        self.statements
97            .iter()
98            .map(|s| format!("statement {} {} {}", s.position, s.form, s.table))
99            .collect::<Vec<_>>()
100            .join("; ")
101    }
102}
103
104#[must_use]
105pub fn audit_migration_cost(extensions: &[Arc<dyn Extension>], hot: &[&str]) -> Vec<MigrationCost> {
106    let mut out = Vec::new();
107    for ext in extensions {
108        let extension = ExtensionId::new(ext.id());
109        for migration in ext.migrations().into_iter().filter(|m| !m.tombstone) {
110            let label = format!("{:03}_{}", migration.version, migration.name);
111            if let Some(cost) = audit_one(&extension, &label, migration.sql, hot) {
112                out.push(cost);
113            }
114        }
115    }
116    out
117}
118
119#[must_use]
120pub fn audit_one(
121    extension: &ExtensionId,
122    migration: &str,
123    sql: &str,
124    hot: &[&str],
125) -> Option<MigrationCost> {
126    let (declared, malformed) = match cost::parse(sql) {
127        Ok(found) => (found, None),
128        Err(e) => (None, Some(e.to_string())),
129    };
130    let statements = expensive_statements(sql, hot);
131    if statements.is_empty() && declared.is_none() && malformed.is_none() {
132        return None;
133    }
134    Some(MigrationCost {
135        extension: extension.clone(),
136        migration: migration.to_owned(),
137        statements,
138        declared,
139        malformed,
140    })
141}
142
143// Why: an unparseable body is not this check's business — `migration_refs`
144// already refuses it with the parse error, and reporting it twice would only
145// bury that message.
146fn expensive_statements(sql: &str, hot: &[&str]) -> Vec<ExpensiveStatement> {
147    let Ok(parsed) = pg_query::parse(sql) else {
148        return Vec::new();
149    };
150    let mut out = Vec::new();
151    for (index, node) in parsed
152        .protobuf
153        .stmts
154        .iter()
155        .filter_map(|raw| raw.stmt.as_ref().and_then(|s| s.node.as_ref()))
156        .enumerate()
157    {
158        let position = index + 1;
159        if let Some((table, form)) = classify(node)
160            && hot.contains(&table.as_str())
161        {
162            out.push(ExpensiveStatement {
163                position,
164                table,
165                form,
166            });
167        }
168    }
169    out
170}
171
172fn is_select_driven(select: Option<&pg_query::protobuf::Node>) -> bool {
173    let Some(NodeEnum::SelectStmt(select)) = select.and_then(|n| n.node.as_ref()) else {
174        return false;
175    };
176    select.values_lists.is_empty()
177}
178
179fn classify(node: &NodeEnum) -> Option<(String, &'static str)> {
180    match node {
181        NodeEnum::UpdateStmt(stmt) => Some((stmt.relation.as_ref()?.relname.clone(), "UPDATE on")),
182        NodeEnum::DeleteStmt(stmt) => {
183            Some((stmt.relation.as_ref()?.relname.clone(), "DELETE from"))
184        },
185        // Why: `INSERT … VALUES` writes what the author typed; only a
186        // select-driven insert scales with the table it reads. The parser
187        // models both as a SelectStmt hanging off the insert, so the two are
188        // told apart by that node carrying rows of its own rather than a
189        // FROM — without this, every literal insert reads as a backfill.
190        NodeEnum::InsertStmt(stmt) if is_select_driven(stmt.select_stmt.as_deref()) => Some((
191            stmt.relation.as_ref()?.relname.clone(),
192            "INSERT … SELECT into",
193        )),
194        // Why: a non-concurrent index build holds a write lock for the whole
195        // build; concurrently is the form that does not stop traffic.
196        NodeEnum::IndexStmt(stmt) if !stmt.concurrent => Some((
197            stmt.relation.as_ref()?.relname.clone(),
198            "CREATE INDEX (not CONCURRENTLY) on",
199        )),
200        NodeEnum::AlterTableStmt(stmt) => {
201            let table = stmt.relation.as_ref()?.relname.clone();
202            let form = stmt.cmds.iter().find_map(|cmd| match cmd.node.as_ref() {
203                Some(NodeEnum::AlterTableCmd(c)) => scanning_alter(c.subtype),
204                _ => None,
205            })?;
206            Some((table, form))
207        },
208        _ => None,
209    }
210}
211
212// Why: both forms read every existing row before they can be recorded, and
213// both take an ACCESS EXCLUSIVE or SHARE UPDATE EXCLUSIVE lock while doing it.
214const fn scanning_alter(subtype: i32) -> Option<&'static str> {
215    if subtype == AlterTableType::AtValidateConstraint as i32 {
216        return Some("ALTER TABLE … VALIDATE CONSTRAINT on");
217    }
218    if subtype == AlterTableType::AtSetNotNull as i32 {
219        return Some("ALTER TABLE … SET NOT NULL on");
220    }
221    None
222}