Skip to main content

systemprompt_users/repository/user/
purge.rs

1//! Deleting a user's rows from every table that keys on them.
2//!
3//! The schema cannot cascade a user delete: activity tables carry `user_id`
4//! without a foreign key, because the column also holds sentinel principals
5//! and the identifiers a session column means differ table to table. Each
6//! owning crate therefore declares its tables in the
7//! [`systemprompt_extension::purge`] registry and the delete runs them here,
8//! inside the same transaction that deletes the `users` row. Content
9//! the user's rows referenced but did not own — an artifact body shared by
10//! digest — is declared as an orphan sweep: the keys the user's rows reference
11//! are captured first, and once the user-keyed tables are gone only those keys
12//! that nothing references any more are cleared, in the same transaction. The
13//! same lists, counted instead of deleted, are the dry-run report: what a
14//! deletion would remove, table by table — and, run against a user who no
15//! longer exists, an orphan report.
16//!
17//! Copyright (c) systemprompt.io — Business Source License 1.1.
18//! See <https://systemprompt.io> for licensing details.
19
20use sqlx::{AssertSqlSafe, Postgres, Transaction};
21use systemprompt_database::admin::SafeIdentifier;
22use systemprompt_extension::purge::{
23    OrphanSweep, UserPurgeTable, registered_orphan_sweeps, registered_user_purge_tables,
24};
25use systemprompt_extension::{orphan_sweeps, user_purge_tables};
26use systemprompt_identifiers::UserId;
27
28use crate::error::{Result, UserError};
29use crate::repository::UserRepository;
30
31// Why: core's own user-keyed tables without a foreign key to `users`, one
32// entry per owning crate's table. `ai_requests` cascades to its children;
33// `user_sessions` is deleted by the caller first because `ai_requests` and
34// `user_contexts` point at it with SET NULL and must go before the session.
35user_purge_tables!(
36    "systemprompt-core",
37    [
38        ("ai_requests", "user_id"),
39        ("mcp_artifacts", "user_id"),
40        ("mcp_tool_executions", "user_id"),
41        ("mcp_sessions", "user_id"),
42        ("mcp_external_sessions", "user_id"),
43        ("mcp_proxy_identities", "user_id"),
44        ("governance_decisions", "user_id"),
45        ("agent_tasks", "user_id"),
46        ("task_messages", "user_id"),
47        ("files", "user_id"),
48        ("event_outbox", "user_id"),
49        ("logs", "user_id"),
50        ("user_rate_limit_buckets", "user_id"),
51        ("oauth_jti_revocations", "user_id"),
52    ]
53);
54
55orphan_sweeps!(
56    "systemprompt-core",
57    [{
58        table: "artifact_payloads",
59        key: "sha256",
60        referenced_by: "mcp_artifacts",
61        via: "payload_sha256",
62        user_column: "user_id",
63    }]
64);
65
66/// How many rows one purge table holds for a user.
67#[derive(Debug, Clone, Copy)]
68pub struct PurgeCount {
69    pub owner: &'static str,
70    pub table: &'static str,
71    pub rows: i64,
72}
73
74fn identifier(kind: &'static str, raw: &str) -> Result<SafeIdentifier> {
75    SafeIdentifier::parse(raw).map_err(|source| UserError::PurgeIdentifier {
76        kind,
77        name: raw.to_owned(),
78        source,
79    })
80}
81
82// Why: the predicate is a fixed fragment from a crate's registration, never
83// user input, but it is interpolated into SQL — so it is held to a character
84// set that cannot terminate the statement or open a comment.
85fn predicate(entry: &UserPurgeTable) -> Result<String> {
86    let Some(raw) = entry.predicate else {
87        return Ok(String::new());
88    };
89    let allowed = |c: char| c.is_ascii_alphanumeric() || " _'<>=!.()".contains(c);
90    if raw.is_empty() || raw.contains("''") || !raw.chars().all(allowed) {
91        return Err(UserError::Validation(format!(
92            "purge predicate for {}: not a plain comparison: {raw}",
93            entry.table
94        )));
95    }
96    Ok(format!(" AND ({raw})"))
97}
98
99fn purge_where(entry: &UserPurgeTable) -> Result<String> {
100    let table = identifier("table", entry.table)?;
101    let column = identifier("column", entry.column)?;
102    // Why: `::text` comparison — purge tables key on TEXT and VARCHAR
103    // columns, and the cast keeps one text bind valid for every one of them.
104    Ok(format!(
105        "FROM {} WHERE {}::text = $1{}",
106        table.quoted(),
107        column.quoted(),
108        predicate(entry)?
109    ))
110}
111
112struct SweepSql {
113    candidates: String,
114    delete: String,
115    preview: String,
116}
117
118fn sweep_sql(sweep: &OrphanSweep) -> Result<SweepSql> {
119    let table = identifier("sweep table", sweep.table)?.quoted();
120    let key = identifier("sweep key", sweep.key)?.quoted();
121    let referenced_by = identifier("sweep referrer", sweep.referenced_by)?.quoted();
122    let via = identifier("sweep reference", sweep.via)?.quoted();
123    let user = identifier("sweep user column", sweep.user_column)?.quoted();
124    let referrer = format!("SELECT 1 FROM {referenced_by} r WHERE r.{via} = t.{key}");
125    Ok(SweepSql {
126        candidates: format!(
127            "SELECT DISTINCT r.{via}::text FROM {referenced_by} r \
128             WHERE r.{user}::text = $1 AND r.{via} IS NOT NULL"
129        ),
130        delete: format!(
131            "DELETE FROM {table} t WHERE t.{key}::text = ANY($1) AND NOT EXISTS ({referrer})"
132        ),
133        preview: format!(
134            "SELECT COUNT(*) FROM {table} t \
135             WHERE EXISTS ({referrer} AND r.{user}::text = $1) \
136             AND NOT EXISTS ({referrer} AND r.{user}::text IS DISTINCT FROM $1)"
137        ),
138    })
139}
140
141impl UserRepository {
142    pub(super) async fn purge_user_rows(
143        tx: &mut Transaction<'_, Postgres>,
144        id: &UserId,
145    ) -> Result<Vec<PurgeCount>> {
146        let mut sweeps = Vec::new();
147        for sweep in registered_orphan_sweeps() {
148            let sql = sweep_sql(sweep)?;
149            let keys: Vec<String> = sqlx::query_scalar(AssertSqlSafe(sql.candidates))
150                .bind(id.as_str())
151                .fetch_all(&mut **tx)
152                .await?;
153            sweeps.push((sweep, sql.delete, keys));
154        }
155        let mut removed = Vec::new();
156        for entry in registered_user_purge_tables() {
157            let sql = format!("DELETE {}", purge_where(entry)?);
158            let result = sqlx::query(AssertSqlSafe(sql))
159                .bind(id.as_str())
160                .execute(&mut **tx)
161                .await?;
162            removed.push(PurgeCount {
163                owner: entry.owner,
164                table: entry.table,
165                rows: i64::try_from(result.rows_affected()).unwrap_or(i64::MAX),
166            });
167        }
168        for (sweep, delete, keys) in sweeps {
169            let result = sqlx::query(AssertSqlSafe(delete))
170                .bind(keys)
171                .execute(&mut **tx)
172                .await?;
173            removed.push(PurgeCount {
174                owner: sweep.owner,
175                table: sweep.table,
176                rows: i64::try_from(result.rows_affected()).unwrap_or(i64::MAX),
177            });
178        }
179        Ok(removed)
180    }
181
182    pub async fn purge_preview(&self, id: &UserId) -> Result<Vec<PurgeCount>> {
183        let mut counts = Vec::new();
184        let sessions: i64 = sqlx::query_scalar!(
185            r#"SELECT COUNT(*) AS "count!" FROM user_sessions WHERE user_id = $1"#,
186            id.as_str()
187        )
188        .fetch_one(&*self.pool)
189        .await?;
190        counts.push(PurgeCount {
191            owner: "systemprompt-core",
192            table: "user_sessions",
193            rows: sessions,
194        });
195        for entry in registered_user_purge_tables() {
196            let sql = format!("SELECT COUNT(*) {}", purge_where(entry)?);
197            let rows: i64 = sqlx::query_scalar(AssertSqlSafe(sql))
198                .bind(id.as_str())
199                .fetch_one(&*self.pool)
200                .await?;
201            counts.push(PurgeCount {
202                owner: entry.owner,
203                table: entry.table,
204                rows,
205            });
206        }
207        for sweep in registered_orphan_sweeps() {
208            let rows: i64 = sqlx::query_scalar(AssertSqlSafe(sweep_sql(sweep)?.preview))
209                .bind(id.as_str())
210                .fetch_one(&*self.pool)
211                .await?;
212            counts.push(PurgeCount {
213                owner: sweep.owner,
214                table: sweep.table,
215                rows,
216            });
217        }
218        Ok(counts)
219    }
220}