Skip to main content

systemprompt_users/repository/user/
merge.rs

1//! Anonymous-to-identified user merge operations.
2//!
3//! Copyright (c) systemprompt.io — Business Source License 1.1.
4//! See <https://systemprompt.io> for licensing details.
5
6use 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    // Why: security artifacts (oauth tokens/codes, webauthn, API keys, device
21    // certs, bridge auth state, federated links) are deliberately NOT rekeyed
22    // — they are bound to the source identity and die with it via FK CASCADE
23    // when the source row is deleted. Only data rows transfer.
24    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}