systemprompt_users/repository/user/
merge.rs1use 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
70async 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}