Skip to main content

fraiseql_auth/account_linking/
postgres.rs

1//! PostgreSQL-backed [`AccountStore`] — durable user / identity persistence.
2//!
3//! This is the production backend for [`AccountStore`](super::AccountStore); the
4//! [`InMemoryAccountStore`](super::InMemoryAccountStore) loses all linkage on
5//! restart. It is a drop-in replacement (same trait, same `"user_<uuid>"`
6//! identifier format, so it joins the existing `_system.sessions.user_id`), so
7//! `multi_provider` / `phone_otp` need no change beyond which `Arc<dyn AccountStore>`
8//! they are handed.
9//!
10//! # Schema
11//!
12//! - `core.tb_user` — one row per stable account (`user_id`, optional verified email).
13//! - `core.tb_auth_identity` — one row per linked `(provider, provider_id)`, FK to a user.
14//!
15//! Both carry a `tenant_id` and RLS deny-by-default (mirroring the change-log RLS in
16//! observers migration `12`). RLS is `ENABLE`, not `FORCE`: this store runs as the
17//! table owner and bypasses the policies — exactly like the executor/poller for the
18//! change-log — while any other (non-`BYPASSRLS`) role reads zero rows unless it sets
19//! the `fraiseql.tenant_id` GUC. v1 operates single-tenant (`tenant_id` NULL, since the
20//! [`AccountStore`](super::AccountStore) trait carries no tenant parameter); per-tenant
21//! scoping is a forward-compatible extension.
22
23use async_trait::async_trait;
24use sqlx::{Row, postgres::PgPool};
25use uuid::Uuid;
26
27use super::{AccountLinkResult, AccountRecord, AccountStore, ProviderLink, normalize_email};
28use crate::{
29    audit::logger::{AuditEventType, SecretType, get_audit_logger},
30    error::{AuthError, Result},
31};
32
33/// Idempotent DDL for the user / identity store. Exposed so a migration runner can
34/// apply it explicitly; [`PostgresAccountStore::init`] runs the same statements.
35pub const SCHEMA_SQL: &str = r"
36CREATE SCHEMA IF NOT EXISTS core;
37
38CREATE TABLE IF NOT EXISTS core.tb_user (
39    pk_user    BIGINT GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
40    id         UUID NOT NULL DEFAULT gen_random_uuid(),
41    user_id    TEXT NOT NULL UNIQUE,
42    email      TEXT,
43    tenant_id  UUID,
44    created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
45    updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
46);
47CREATE UNIQUE INDEX IF NOT EXISTS uq_user_email ON core.tb_user (email) WHERE email IS NOT NULL;
48
49CREATE TABLE IF NOT EXISTS core.tb_auth_identity (
50    pk_auth_identity BIGINT GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
51    id          UUID NOT NULL DEFAULT gen_random_uuid(),
52    fk_user     BIGINT NOT NULL REFERENCES core.tb_user (pk_user) ON DELETE CASCADE,
53    user_id     TEXT NOT NULL,
54    provider    TEXT NOT NULL,
55    provider_id TEXT NOT NULL,
56    tenant_id   UUID,
57    created_at  TIMESTAMPTZ NOT NULL DEFAULT now(),
58    UNIQUE (provider, provider_id)
59);
60CREATE INDEX IF NOT EXISTS idx_auth_identity_user    ON core.tb_auth_identity (fk_user);
61CREATE INDEX IF NOT EXISTS idx_auth_identity_user_id ON core.tb_auth_identity (user_id);
62
63-- RLS deny-by-default (mirrors observers migration 12). ENABLE not FORCE so the
64-- owner (this store) and BYPASSRLS roles operate freely; a non-owner role reads a
65-- row only once it has set fraiseql.tenant_id to that row's tenant (fail-closed).
66ALTER TABLE core.tb_user          ENABLE ROW LEVEL SECURITY;
67ALTER TABLE core.tb_auth_identity ENABLE ROW LEVEL SECURITY;
68
69DROP POLICY IF EXISTS p_user_tenant_read ON core.tb_user;
70CREATE POLICY p_user_tenant_read ON core.tb_user
71    FOR SELECT USING (tenant_id = NULLIF(current_setting('fraiseql.tenant_id', true), '')::uuid);
72DROP POLICY IF EXISTS p_user_insert ON core.tb_user;
73CREATE POLICY p_user_insert ON core.tb_user FOR INSERT WITH CHECK (true);
74
75DROP POLICY IF EXISTS p_auth_identity_tenant_read ON core.tb_auth_identity;
76CREATE POLICY p_auth_identity_tenant_read ON core.tb_auth_identity
77    FOR SELECT USING (tenant_id = NULLIF(current_setting('fraiseql.tenant_id', true), '')::uuid);
78DROP POLICY IF EXISTS p_auth_identity_insert ON core.tb_auth_identity;
79CREATE POLICY p_auth_identity_insert ON core.tb_auth_identity FOR INSERT WITH CHECK (true);
80
81-- Least-privilege baseline: never world-readable. RLS is defence-in-depth on top.
82REVOKE ALL ON core.tb_user          FROM PUBLIC;
83REVOKE ALL ON core.tb_auth_identity FROM PUBLIC;
84";
85
86/// PostgreSQL-backed account store.
87///
88/// Persists user accounts and their linked provider identities, so account linking
89/// survives a process restart. See the module-level documentation for the schema and
90/// RLS posture.
91pub struct PostgresAccountStore {
92    db: PgPool,
93}
94
95impl PostgresAccountStore {
96    /// Create a new store over an existing pool.
97    ///
98    /// The pool's role must own (or `BYPASSRLS`) the `core.tb_user` /
99    /// `core.tb_auth_identity` tables — it runs the trusted login path and must not be
100    /// constrained by the deny-by-default RLS. Calling [`init`](Self::init) once on
101    /// startup creates the tables (so the connecting role owns them by construction).
102    #[must_use]
103    pub const fn new(db: PgPool) -> Self {
104        Self { db }
105    }
106
107    /// Create the `core.tb_user` / `core.tb_auth_identity` schema (idempotent).
108    ///
109    /// Call once on startup. Safe to re-run and safe on a database that predates this
110    /// store (the `CREATE … IF NOT EXISTS` form is the back-compat path for existing
111    /// deployments that have no user table).
112    ///
113    /// # Errors
114    ///
115    /// Returns [`AuthError::DatabaseError`] if the DDL fails.
116    pub async fn init(&self) -> Result<()> {
117        sqlx::raw_sql(SCHEMA_SQL).execute(&self.db).await.map_err(|e| {
118            AuthError::DatabaseError {
119                message: format!("Failed to initialize identity store: {e}"),
120            }
121        })?;
122        Ok(())
123    }
124}
125
126/// Generate a fresh stable user identifier, matching the in-memory store's format so
127/// the two backends are interchangeable and the value joins `_system.sessions.user_id`.
128fn new_user_id() -> String {
129    format!("user_{}", Uuid::new_v4().as_simple())
130}
131
132fn db_error(context: &str, e: &sqlx::Error) -> AuthError {
133    AuthError::DatabaseError {
134        message: format!("{context}: {e}"),
135    }
136}
137
138// Reason: AccountStore is defined with #[async_trait]; the impl must match its
139// transformed signatures. async_trait: dyn-dispatch required; remove when RTN + Send
140// is stable (RFC 3425).
141#[async_trait]
142impl AccountStore for PostgresAccountStore {
143    async fn link_or_create_user(
144        &self,
145        email: Option<&str>,
146        email_verified: bool,
147        provider: &str,
148        provider_id: &str,
149    ) -> Result<AccountLinkResult> {
150        let mut tx = self.db.begin().await.map_err(|e| db_error("begin tx", &e))?;
151
152        // 1. A known (provider, provider_id) is an idempotent re-login: same user, no new link. The
153        //    UNIQUE(provider, provider_id) constraint makes this the authoritative lookup.
154        if let Some(row) = sqlx::query(
155            "SELECT user_id FROM core.tb_auth_identity WHERE provider = $1 AND provider_id = $2",
156        )
157        .bind(provider)
158        .bind(provider_id)
159        .fetch_optional(&mut *tx)
160        .await
161        .map_err(|e| db_error("lookup identity", &e))?
162        {
163            let user_id: String = row.get("user_id");
164            tx.commit().await.map_err(|e| db_error("commit", &e))?;
165            return Ok(AccountLinkResult {
166                user_id,
167                is_new: false,
168                linked: false,
169            });
170        }
171
172        // 2. Resolve the linking key. A verified, non-empty email links across providers; anything
173        //    else is keyed on (provider, provider_id) so distinct identities can never collapse
174        //    (H26).
175        let verified_email = email.map(normalize_email).filter(|e| !e.is_empty() && email_verified);
176
177        // 3. Find the email-keyed user, or create a fresh account.
178        let (user_id, pk_user, is_new, linked) = if let Some(em) = verified_email.as_deref() {
179            if let Some(row) =
180                sqlx::query("SELECT pk_user, user_id FROM core.tb_user WHERE email = $1")
181                    .bind(em)
182                    .fetch_optional(&mut *tx)
183                    .await
184                    .map_err(|e| db_error("lookup user by email", &e))?
185            {
186                let pk_user: i64 = row.get("pk_user");
187                let user_id: String = row.get("user_id");
188                (user_id, pk_user, false, true)
189            } else {
190                let user_id = new_user_id();
191                let pk_user = insert_user(&mut tx, &user_id, Some(em)).await?;
192                (user_id, pk_user, true, false)
193            }
194        } else {
195            let user_id = new_user_id();
196            let pk_user = insert_user(&mut tx, &user_id, None).await?;
197            (user_id, pk_user, true, false)
198        };
199
200        // 4. Link the provider identity (new for this account by construction — step 1 ruled out an
201        //    existing one).
202        sqlx::query(
203            "INSERT INTO core.tb_auth_identity (fk_user, user_id, provider, provider_id) \
204             VALUES ($1, $2, $3, $4)",
205        )
206        .bind(pk_user)
207        .bind(&user_id)
208        .bind(provider)
209        .bind(provider_id)
210        .execute(&mut *tx)
211        .await
212        .map_err(|e| db_error("insert identity", &e))?;
213
214        tx.commit().await.map_err(|e| db_error("commit", &e))?;
215
216        let logger = get_audit_logger();
217        if is_new {
218            logger.log_success(
219                AuditEventType::SessionTokenCreated,
220                SecretType::SessionToken,
221                Some(user_id.clone()),
222                &format!("account_created:{provider}"),
223            );
224        } else {
225            logger.log_success(
226                AuditEventType::AuthSuccess,
227                SecretType::SessionToken,
228                Some(user_id.clone()),
229                &format!("account_linked:{provider}"),
230            );
231        }
232
233        Ok(AccountLinkResult {
234            user_id,
235            is_new,
236            linked,
237        })
238    }
239
240    async fn get_account(&self, user_id: &str) -> Result<AccountRecord> {
241        let email: Option<String> =
242            sqlx::query("SELECT email FROM core.tb_user WHERE user_id = $1")
243                .bind(user_id)
244                .fetch_optional(&self.db)
245                .await
246                .map_err(|e| db_error("lookup user", &e))?
247                .ok_or(AuthError::TokenNotFound)?
248                .get("email");
249
250        let providers = sqlx::query(
251            "SELECT provider, provider_id FROM core.tb_auth_identity \
252             WHERE user_id = $1 ORDER BY pk_auth_identity",
253        )
254        .bind(user_id)
255        .fetch_all(&self.db)
256        .await
257        .map_err(|e| db_error("lookup identities", &e))?
258        .into_iter()
259        .map(|row| ProviderLink {
260            provider:    row.get("provider"),
261            provider_id: row.get("provider_id"),
262        })
263        .collect();
264
265        Ok(AccountRecord {
266            user_id: user_id.to_string(),
267            email,
268            providers,
269        })
270    }
271}
272
273/// Insert a new user row and return its `pk_user`.
274async fn insert_user(
275    tx: &mut sqlx::Transaction<'_, sqlx::Postgres>,
276    user_id: &str,
277    email: Option<&str>,
278) -> Result<i64> {
279    let row =
280        sqlx::query("INSERT INTO core.tb_user (user_id, email) VALUES ($1, $2) RETURNING pk_user")
281            .bind(user_id)
282            .bind(email)
283            .fetch_one(&mut **tx)
284            .await
285            .map_err(|e| db_error("insert user", &e))?;
286    Ok(row.get("pk_user"))
287}