Skip to main content

systemprompt_users/repository/
federated_identity.rs

1//! Repository for `federated_identities` — the `{issuer, external_sub} ->
2//! users.id` mapping used by RFC 8693 token-exchange first-touch.
3//!
4//! Copyright (c) systemprompt.io — Business Source License 1.1.
5//! See <https://systemprompt.io> for licensing details.
6
7use 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}