systemprompt_users/repository/
federated_identity.rs1use chrono::Utc;
8use sqlx::Acquire;
9use systemprompt_identifiers::UserId;
10use systemprompt_traits::FederatedIdentityClaims;
11
12use crate::error::Result;
13use crate::models::{User, UserRole, UserStatus, normalise_email};
14use crate::repository::UserRepository;
15
16impl UserRepository {
17 pub async fn find_federated(&self, issuer: &str, external_sub: &str) -> Result<Option<UserId>> {
18 let row = sqlx::query!(
19 "SELECT user_id FROM federated_identities WHERE issuer = $1 AND external_sub = $2",
20 issuer,
21 external_sub
22 )
23 .fetch_optional(&*self.pool)
24 .await?;
25
26 Ok(row.map(|r| UserId::new(r.user_id)))
27 }
28
29 pub async fn find_or_create_federated(
30 &self,
31 issuer: &str,
32 external_sub: &str,
33 claims: &FederatedIdentityClaims,
34 ) -> Result<User> {
35 let mut conn = self.write_pool.acquire().await?;
36 let mut tx = conn.begin().await?;
37
38 if let Some(existing) = sqlx::query!(
39 "UPDATE federated_identities SET last_seen_at = CURRENT_TIMESTAMP WHERE issuer = $1 \
40 AND external_sub = $2 RETURNING user_id",
41 issuer,
42 external_sub
43 )
44 .fetch_optional(&mut *tx)
45 .await?
46 {
47 let user = sqlx::query_as!(
48 User,
49 r#"
50 SELECT id, name, email, full_name, display_name, status,
51 email_verified, roles, avatar_url, is_bot, is_scanner,
52 created_at, updated_at
53 FROM users WHERE id = $1
54 "#,
55 existing.user_id
56 )
57 .fetch_one(&mut *tx)
58 .await?;
59 tx.commit().await?;
60 return Ok(user);
61 }
62
63 if let Some(existing) =
64 link_by_verified_email(&mut tx, issuer, external_sub, claims).await?
65 {
66 tx.commit().await?;
67 return Ok(existing);
68 }
69
70 let fields = NewFederatedUser::derive(issuer, external_sub, claims);
71
72 let user = sqlx::query_as!(
73 User,
74 r#"
75 INSERT INTO users (
76 id, name, email, full_name, display_name,
77 status, email_verified, roles, is_bot,
78 created_at, updated_at
79 )
80 VALUES ($1, $2, $3, $4, $5, $6, false, $7::TEXT[], false, $8, $8)
81 RETURNING id, name, email, full_name, display_name, status, email_verified,
82 roles, avatar_url, is_bot, is_scanner, created_at, updated_at
83 "#,
84 fields.id.as_str(),
85 fields.name,
86 fields.email,
87 fields.display_name.as_deref(),
88 fields.display_name.as_deref(),
89 fields.status,
90 &fields.roles,
91 fields.now,
92 )
93 .fetch_one(&mut *tx)
94 .await?;
95
96 sqlx::query!(
97 "INSERT INTO federated_identities (issuer, external_sub, user_id) VALUES ($1, $2, $3)",
98 issuer,
99 external_sub,
100 user.id.as_str()
101 )
102 .execute(&mut *tx)
103 .await?;
104
105 tx.commit().await?;
106 Ok(user)
107 }
108}
109
110async fn link_by_verified_email(
111 tx: &mut sqlx::Transaction<'_, sqlx::Postgres>,
112 issuer: &str,
113 external_sub: &str,
114 claims: &FederatedIdentityClaims,
115) -> Result<Option<User>> {
116 if !claims.email_verified {
117 return Ok(None);
118 }
119 let Some(addr) = claims.email.as_deref() else {
120 return Ok(None);
121 };
122 let email = normalise_email(addr);
123 let deleted_status = UserStatus::Deleted.as_str();
124 let Some(existing) = sqlx::query_as!(
125 User,
126 r#"
127 SELECT id, name, email, full_name, display_name, status,
128 email_verified, roles, avatar_url, is_bot, is_scanner,
129 created_at, updated_at
130 FROM users WHERE email = $1 AND status != $2
131 "#,
132 email,
133 deleted_status
134 )
135 .fetch_optional(&mut **tx)
136 .await?
137 else {
138 return Ok(None);
139 };
140
141 sqlx::query!(
142 "INSERT INTO federated_identities (issuer, external_sub, user_id) VALUES ($1, $2, $3)",
143 issuer,
144 external_sub,
145 existing.id.as_str()
146 )
147 .execute(&mut **tx)
148 .await?;
149 Ok(Some(existing))
150}
151
152struct NewFederatedUser {
153 id: UserId,
154 name: String,
155 email: String,
156 display_name: Option<String>,
157 status: &'static str,
158 roles: Vec<String>,
159 now: chrono::DateTime<Utc>,
160}
161
162impl NewFederatedUser {
163 fn derive(issuer: &str, external_sub: &str, claims: &FederatedIdentityClaims) -> Self {
164 let name = claims
165 .preferred_username
166 .clone()
167 .or_else(|| claims.name.clone())
168 .unwrap_or_else(|| format!("fed_{}_{}", short_hash(issuer), short_hash(external_sub)));
169 let synthetic_email = || {
170 format!(
171 "{}@{}.federated.local",
172 short_hash(external_sub),
173 short_host(issuer)
174 )
175 };
176 let email = match (claims.email.as_deref(), claims.email_verified) {
177 (Some(addr), true) => normalise_email(addr),
178 (Some(addr), false) => {
179 tracing::warn!(
180 issuer,
181 external_sub,
182 upstream_email = addr,
183 "upstream IdP did not assert email_verified; using synthetic local email to \
184 prevent account-claim attacks"
185 );
186 synthetic_email()
187 },
188 (None, _) => synthetic_email(),
189 };
190
191 Self {
192 id: UserId::new(uuid::Uuid::new_v4().to_string()),
193 name,
194 email,
195 display_name: claims.name.clone(),
196 status: UserStatus::Active.as_str(),
197 roles: normalised_roles(&claims.roles),
198 now: Utc::now(),
199 }
200 }
201}
202
203fn normalised_roles(claim_roles: &[String]) -> Vec<String> {
204 if claim_roles.is_empty() {
205 vec![UserRole::User.as_str().to_owned()]
206 } else {
207 claim_roles.to_vec()
208 }
209}
210
211fn short_hash(s: &str) -> String {
212 use sha2::{Digest, Sha256};
213 let digest = Sha256::digest(s.as_bytes());
214 hex::encode(&digest[..6])
215}
216
217fn short_host(issuer: &str) -> String {
218 issuer
219 .trim_start_matches("https://")
220 .trim_start_matches("http://")
221 .split('/')
222 .next()
223 .unwrap_or("issuer")
224 .replace(['.', ':'], "-")
225}