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