systemprompt_users/repository/user/
merge.rs1use sqlx::{Acquire, Postgres, Transaction};
7use systemprompt_identifiers::UserId;
8
9use crate::error::Result;
10use crate::repository::UserRepository;
11
12#[derive(Debug, Clone, Copy)]
13pub struct MergeResult {
14 pub sessions: u64,
15 pub tasks: u64,
16 pub total_rows: u64,
17}
18
19impl UserRepository {
20 pub async fn merge_users(&self, source_id: &UserId, target_id: &UserId) -> Result<MergeResult> {
25 let mut conn = self.write_pool.acquire().await?;
26 let mut tx = conn.begin().await?;
27
28 let source = source_id.as_str();
29 let target = target_id.as_str();
30
31 let sessions = transfer_sessions(&mut tx, source, target).await?;
32 let tasks = transfer_tasks(&mut tx, source, target).await?;
33 let mut total_rows = sessions + tasks;
34 total_rows += transfer_audit_rows(&mut tx, source, target).await?;
35 total_rows += transfer_content_rows(&mut tx, source, target).await?;
36
37 sqlx::query!(
38 "UPDATE fingerprint_reputation SET associated_user_ids = \
39 array_replace(associated_user_ids, $2, $1) WHERE $2 = ANY(associated_user_ids)",
40 target,
41 source
42 )
43 .execute(&mut *tx)
44 .await?;
45
46 sqlx::query!(
47 "DELETE FROM ai_quota_buckets WHERE subject_kind = 'user' AND subject_id = $1",
48 source
49 )
50 .execute(&mut *tx)
51 .await?;
52
53 sqlx::query!("DELETE FROM users WHERE id = $1", source)
54 .execute(&mut *tx)
55 .await?;
56
57 tx.commit().await?;
58 Ok(MergeResult {
59 sessions,
60 tasks,
61 total_rows,
62 })
63 }
64}
65
66async fn transfer_sessions(
67 tx: &mut Transaction<'_, Postgres>,
68 source: &str,
69 target: &str,
70) -> Result<u64> {
71 let result = sqlx::query!(
72 "UPDATE user_sessions SET user_id = $1 WHERE user_id = $2",
73 target,
74 source
75 )
76 .execute(&mut **tx)
77 .await?;
78 Ok(result.rows_affected())
79}
80
81async fn transfer_tasks(
82 tx: &mut Transaction<'_, Postgres>,
83 source: &str,
84 target: &str,
85) -> Result<u64> {
86 let result = sqlx::query!(
87 "UPDATE agent_tasks SET user_id = $1 WHERE user_id = $2",
88 target,
89 source
90 )
91 .execute(&mut **tx)
92 .await?;
93 Ok(result.rows_affected())
94}
95
96async fn transfer_audit_rows(
97 tx: &mut Transaction<'_, Postgres>,
98 source: &str,
99 target: &str,
100) -> Result<u64> {
101 let mut moved = 0;
102 moved += sqlx::query!(
103 "UPDATE task_messages SET user_id = $1 WHERE user_id = $2",
104 target,
105 source
106 )
107 .execute(&mut **tx)
108 .await?
109 .rows_affected();
110 moved += sqlx::query!(
111 "UPDATE user_contexts SET user_id = $1 WHERE user_id = $2",
112 target,
113 source
114 )
115 .execute(&mut **tx)
116 .await?
117 .rows_affected();
118 moved += sqlx::query!(
119 "UPDATE mcp_tool_executions SET user_id = $1 WHERE user_id = $2",
120 target,
121 source
122 )
123 .execute(&mut **tx)
124 .await?
125 .rows_affected();
126 moved += sqlx::query!(
127 "UPDATE mcp_artifacts SET user_id = $1 WHERE user_id = $2",
128 target,
129 source
130 )
131 .execute(&mut **tx)
132 .await?
133 .rows_affected();
134 moved += sqlx::query!(
135 "UPDATE mcp_sessions SET user_id = $1 WHERE user_id = $2",
136 target,
137 source
138 )
139 .execute(&mut **tx)
140 .await?
141 .rows_affected();
142 moved += sqlx::query!(
143 "UPDATE governance_decisions SET user_id = $1 WHERE user_id = $2",
144 target,
145 source
146 )
147 .execute(&mut **tx)
148 .await?
149 .rows_affected();
150 moved += sqlx::query!(
151 "UPDATE logs SET user_id = $1 WHERE user_id = $2",
152 target,
153 source
154 )
155 .execute(&mut **tx)
156 .await?
157 .rows_affected();
158 Ok(moved)
159}
160
161async fn transfer_content_rows(
162 tx: &mut Transaction<'_, Postgres>,
163 source: &str,
164 target: &str,
165) -> Result<u64> {
166 let mut moved = 0;
167 moved += sqlx::query!(
168 "UPDATE ai_requests SET user_id = $1 WHERE user_id = $2",
169 target,
170 source
171 )
172 .execute(&mut **tx)
173 .await?
174 .rows_affected();
175 moved += sqlx::query!(
176 "UPDATE engagement_events SET user_id = $1 WHERE user_id = $2",
177 target,
178 source
179 )
180 .execute(&mut **tx)
181 .await?
182 .rows_affected();
183 moved += sqlx::query!(
184 "UPDATE analytics_events SET user_id = $1 WHERE user_id = $2",
185 target,
186 source
187 )
188 .execute(&mut **tx)
189 .await?
190 .rows_affected();
191 moved += sqlx::query!(
192 "UPDATE event_outbox SET user_id = $1 WHERE user_id = $2",
193 target,
194 source
195 )
196 .execute(&mut **tx)
197 .await?
198 .rows_affected();
199 moved += sqlx::query!(
200 "UPDATE files SET user_id = $1 WHERE user_id = $2",
201 target,
202 source
203 )
204 .execute(&mut **tx)
205 .await?
206 .rows_affected();
207 moved += sqlx::query!(
208 "UPDATE link_clicks SET user_id = $1 WHERE user_id = $2",
209 target,
210 source
211 )
212 .execute(&mut **tx)
213 .await?
214 .rows_affected();
215 Ok(moved)
216}