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