systemprompt_users/repository/user/
purge.rs1use 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
31user_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#[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
82fn 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 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}