Skip to main content

systemprompt_users/repository/user/
merge.rs

1//! The users-owned half of an account merge.
2//!
3//! Every other domain's rows are moved by its own
4//! [`OwnerReassignment`](systemprompt_traits::OwnerReassignment) before this
5//! runs; what is left is this crate's own: the source's sessions move to the
6//! target, the merge is recorded as a governance decision, and the source user
7//! is deleted — in one transaction, so a failure leaves the source in place
8//! for a rerun.
9//!
10//! Copyright (c) systemprompt.io — Business Source License 1.1.
11//! See <https://systemprompt.io> for licensing details.
12
13use sqlx::{Acquire, Postgres, Transaction};
14use systemprompt_identifiers::{Actor, ContextId, SessionId, UserId};
15use systemprompt_security::authz::types::DecisionTag;
16use systemprompt_security::authz::{GovernanceDecisionRecord, insert_governance_decision};
17
18use crate::error::Result;
19use crate::repository::UserRepository;
20
21const MERGE_TOOL_NAME: &str = "users.merge";
22const MERGE_POLICY: &str = "account_merge";
23
24#[derive(Debug, Clone, Copy)]
25pub struct MergeResult {
26    pub sessions: u64,
27    pub tasks: u64,
28    pub total_rows: u64,
29}
30
31pub const MERGE_EXCLUDED_SECURITY_TABLES: &[&str] = &[
32    "oauth_auth_codes",
33    "oauth_refresh_tokens",
34    "oauth_clients",
35    "webauthn_credentials",
36    "webauthn_challenges",
37    "webauthn_setup_tokens",
38    "user_api_keys",
39    "user_device_certs",
40    "bridge_sessions",
41    "bridge_exchange_codes",
42    "federated_identities",
43];
44
45impl UserRepository {
46    pub async fn complete_merge(&self, source_id: &UserId, target_id: &UserId) -> Result<u64> {
47        let mut conn = self.write_pool.acquire().await?;
48        let mut tx = conn.begin().await?;
49
50        let sessions = sqlx::query!(
51            "UPDATE user_sessions SET user_id = $1 WHERE user_id = $2",
52            target_id.as_str(),
53            source_id.as_str()
54        )
55        .execute(&mut *tx)
56        .await?
57        .rows_affected();
58
59        record_merge_attribution(&mut tx, source_id, target_id).await?;
60
61        sqlx::query!("DELETE FROM users WHERE id = $1", source_id.as_str())
62            .execute(&mut *tx)
63            .await?;
64
65        tx.commit().await?;
66        Ok(sessions)
67    }
68}
69
70// Why: governance_decisions is append-only — a decision is evidence of what was
71// authorised for whom at the time, so the merge is recorded as a new decision
72// rather than by re-attributing the source user's history to the target. A
73// reader following the target's trail finds this row and the source id in it.
74async fn record_merge_attribution(
75    tx: &mut Transaction<'_, Postgres>,
76    source_id: &UserId,
77    target_id: &UserId,
78) -> Result<()> {
79    let id = uuid::Uuid::new_v4().to_string();
80    let session_id = SessionId::new(id.clone());
81    let context_id = ContextId::derived_from_session(&session_id);
82    let actor = Actor::system(target_id.clone());
83    let reason = format!("account merge: {source_id} merged into {target_id}");
84    let evaluated_rules = serde_json::json!([]);
85    let record = GovernanceDecisionRecord {
86        id: &id,
87        actor: &actor,
88        session_id: Some(&session_id),
89        tool_name: MERGE_TOOL_NAME,
90        agent_id: None,
91        agent_scope: None,
92        decision: DecisionTag::Allow,
93        policy: MERGE_POLICY,
94        reason: &reason,
95        evaluated_rules: &evaluated_rules,
96        plugin_id: None,
97        act_chain: &[],
98        context_id: &context_id,
99        task_id: None,
100        trace_id: None,
101        client_id: None,
102        tool_use_id: None,
103    };
104    insert_governance_decision(&mut **tx, &record).await?;
105    Ok(())
106}