fraiseql_auth/account_linking/
postgres.rs1use 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
33pub 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
86pub struct PostgresAccountStore {
92 db: PgPool,
93}
94
95impl PostgresAccountStore {
96 #[must_use]
103 pub const fn new(db: PgPool) -> Self {
104 Self { db }
105 }
106
107 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
126fn 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#[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 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 let verified_email = email.map(normalize_email).filter(|e| !e.is_empty() && email_verified);
176
177 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 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
273async 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}