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!(
37 "systemprompt-core",
38 [
39 ("ai_requests", "user_id"),
40 ("mcp_artifacts", "user_id"),
41 ("mcp_tool_executions", "user_id"),
42 ("mcp_sessions", "user_id"),
43 ("mcp_external_sessions", "user_id"),
44 ("mcp_proxy_identities", "user_id"),
45 ("governance_decisions", "user_id"),
46 ("agent_tasks", "user_id"),
47 ("task_messages", "user_id"),
48 ("files", "user_id"),
49 ("event_outbox", "user_id", "channel <> 'reporting'"),
50 ("logs", "user_id"),
51 ("user_rate_limit_buckets", "user_id"),
52 ("oauth_jti_revocations", "user_id"),
53 ]
54);
55
56orphan_sweeps!(
57 "systemprompt-core",
58 [{
59 table: "artifact_payloads",
60 key: "sha256",
61 referenced_by: "mcp_artifacts",
62 via: "payload_sha256",
63 user_column: "user_id",
64 }]
65);
66
67#[derive(Debug, Clone, Copy)]
69pub struct PurgeCount {
70 pub owner: &'static str,
71 pub table: &'static str,
72 pub rows: i64,
73}
74
75fn identifier(kind: &str, raw: &str) -> Result<SafeIdentifier> {
76 SafeIdentifier::parse(raw)
77 .map_err(|e| UserError::Validation(format!("purge {kind} {raw}: {e}")))
78}
79
80fn predicate(entry: &UserPurgeTable) -> Result<String> {
84 let Some(raw) = entry.predicate else {
85 return Ok(String::new());
86 };
87 let allowed = |c: char| c.is_ascii_alphanumeric() || " _'<>=!.()".contains(c);
88 if raw.is_empty() || raw.contains("''") || !raw.chars().all(allowed) {
89 return Err(UserError::Validation(format!(
90 "purge predicate for {}: not a plain comparison: {raw}",
91 entry.table
92 )));
93 }
94 Ok(format!(" AND ({raw})"))
95}
96
97fn purge_where(entry: &UserPurgeTable) -> Result<String> {
98 let table = identifier("table", entry.table)?;
99 let column = identifier("column", entry.column)?;
100 Ok(format!(
103 "FROM {} WHERE {}::text = $1{}",
104 table.quoted(),
105 column.quoted(),
106 predicate(entry)?
107 ))
108}
109
110struct SweepSql {
111 delete: String,
112 preview: String,
113}
114
115fn sweep_sql(sweep: &OrphanSweep) -> Result<SweepSql> {
116 let table = identifier("sweep table", sweep.table)?.quoted();
117 let key = identifier("sweep key", sweep.key)?.quoted();
118 let referenced_by = identifier("sweep referrer", sweep.referenced_by)?.quoted();
119 let via = identifier("sweep reference", sweep.via)?.quoted();
120 let user = identifier("sweep user column", sweep.user_column)?.quoted();
121 let referrer = format!("SELECT 1 FROM {referenced_by} r WHERE r.{via} = t.{key}");
122 Ok(SweepSql {
123 delete: format!("DELETE FROM {table} t WHERE NOT EXISTS ({referrer})"),
124 preview: format!(
125 "SELECT COUNT(*) FROM {table} t \
126 WHERE EXISTS ({referrer} AND r.{user}::text = $1) \
127 AND NOT EXISTS ({referrer} AND r.{user}::text IS DISTINCT FROM $1)"
128 ),
129 })
130}
131
132impl UserRepository {
133 pub(super) async fn purge_user_rows(
134 tx: &mut Transaction<'_, Postgres>,
135 id: &UserId,
136 ) -> Result<Vec<PurgeCount>> {
137 let mut removed = Vec::new();
138 for entry in registered_user_purge_tables() {
139 let sql = format!("DELETE {}", purge_where(entry)?);
140 let result = sqlx::query(AssertSqlSafe(sql))
141 .bind(id.as_str())
142 .execute(&mut **tx)
143 .await?;
144 removed.push(PurgeCount {
145 owner: entry.owner,
146 table: entry.table,
147 rows: i64::try_from(result.rows_affected()).unwrap_or(i64::MAX),
148 });
149 }
150 for sweep in registered_orphan_sweeps() {
151 let result = sqlx::query(AssertSqlSafe(sweep_sql(sweep)?.delete))
152 .execute(&mut **tx)
153 .await?;
154 removed.push(PurgeCount {
155 owner: sweep.owner,
156 table: sweep.table,
157 rows: i64::try_from(result.rows_affected()).unwrap_or(i64::MAX),
158 });
159 }
160 Ok(removed)
161 }
162
163 pub async fn purge_preview(&self, id: &UserId) -> Result<Vec<PurgeCount>> {
164 let mut counts = Vec::new();
165 let sessions: i64 = sqlx::query_scalar!(
166 r#"SELECT COUNT(*) AS "count!" FROM user_sessions WHERE user_id = $1"#,
167 id.as_str()
168 )
169 .fetch_one(&*self.pool)
170 .await?;
171 counts.push(PurgeCount {
172 owner: "systemprompt-core",
173 table: "user_sessions",
174 rows: sessions,
175 });
176 for entry in registered_user_purge_tables() {
177 let sql = format!("SELECT COUNT(*) {}", purge_where(entry)?);
178 let rows: i64 = sqlx::query_scalar(AssertSqlSafe(sql))
179 .bind(id.as_str())
180 .fetch_one(&*self.pool)
181 .await?;
182 counts.push(PurgeCount {
183 owner: entry.owner,
184 table: entry.table,
185 rows,
186 });
187 }
188 for sweep in registered_orphan_sweeps() {
189 let rows: i64 = sqlx::query_scalar(AssertSqlSafe(sweep_sql(sweep)?.preview))
190 .bind(id.as_str())
191 .fetch_one(&*self.pool)
192 .await?;
193 counts.push(PurgeCount {
194 owner: sweep.owner,
195 table: sweep.table,
196 rows,
197 });
198 }
199 Ok(counts)
200 }
201}