Files
Oxicloud/src/infrastructure/repositories/pg/user_pg_repository.rs
T
Edouard Vanbelle d8b3f2e026 refactor(oidc): migrate provider into issuer
this make OIDC compliant with the invariant binding (issuer and subject)
admin can now rename their provider without breaking

clarifing federation_kind: report the kind of federation wired not the allowed login method
hybryd login method are still allowed
2026-08-08 16:37:45 +02:00

1593 lines
56 KiB
Rust

use futures::future::BoxFuture;
use sqlx::{PgPool, Row};
use std::sync::Arc;
use uuid::Uuid;
use crate::application::ports::auth_ports::UserStoragePort;
use crate::common::errors::DomainError;
use crate::domain::entities::user::{User, UserFlags, UserRole};
use crate::domain::repositories::user_repository::{
StorageStats, UserListEntry, UserRepository, UserRepositoryError, UserRepositoryResult,
};
use crate::infrastructure::repositories::pg::transaction_utils::with_transaction;
// Implement From<sqlx::Error> for UserRepositoryError to allow automatic conversions
impl From<sqlx::Error> for UserRepositoryError {
fn from(err: sqlx::Error) -> Self {
UserPgRepository::map_sqlx_error(err)
}
}
pub struct UserPgRepository {
pool: Arc<PgPool>,
}
impl UserPgRepository {
pub fn new(pool: Arc<PgPool>) -> Self {
Self { pool }
}
/// Borrowed access to the connection pool. Exposed so callers can
/// open transactions that span this repo and other repos / hooks
/// (e.g. `AuthApplicationService::delete_user_admin` opening a tx
/// that wraps the lifecycle dispatcher + the DELETE).
pub fn pool(&self) -> &PgPool {
&self.pool
}
// Helper method to map SQL errors to domain errors
pub fn map_sqlx_error(err: sqlx::Error) -> UserRepositoryError {
match err {
sqlx::Error::RowNotFound => UserRepositoryError::NotFound("User not found".to_string()),
sqlx::Error::Database(db_err) => {
if db_err.code().is_some_and(|code| code == "23505") {
// PostgreSQL uniqueness violation code
UserRepositoryError::AlreadyExists("User or email already exists".to_string())
} else {
UserRepositoryError::DatabaseError(format!("Database error: {}", db_err))
}
}
_ => UserRepositoryError::DatabaseError(format!("Database error: {}", err)),
}
}
/// Fetch only the authorization-relevant flags of a user. Not part of
/// the `UserRepository` trait — called directly from
/// `AuthApplicationService::get_user_flags`.
///
/// Deliberately selects three tiny columns instead of the full row:
/// the full-row SELECT includes `image` (a data URI of up to 512 KiB),
/// which per-request middleware guards were paying on every WebDAV /
/// CalDAV / CardDAV request just to read `is_external` or `role`.
pub async fn get_user_flags(&self, id: Uuid) -> UserRepositoryResult<UserFlags> {
let row = sqlx::query(
r#"
SELECT role::text as role_text,
is_external,
active,
force_password_change_at_next_login
FROM auth.users
WHERE id = $1
"#,
)
.bind(id)
.fetch_one(&*self.pool)
.await
.map_err(Self::map_sqlx_error)?;
let role_str: Option<String> = row.try_get("role_text").unwrap_or(None);
let role = match role_str.as_deref() {
Some("admin") => UserRole::Admin,
_ => UserRole::User,
};
Ok(UserFlags {
role,
is_external: row.get("is_external"),
active: row.get("active"),
force_password_change: row.get("force_password_change_at_next_login"),
})
}
/// Fetch only `(storage_used_bytes, storage_quota_bytes)`. Not part of
/// the `UserRepository` trait — called from `StorageUsageService`.
///
/// Same rationale as [`Self::get_user_flags`]: the full-row SELECT drags
/// `image` (a data URI of up to 512 KiB), `password_hash`,
/// `ui_preferences`, … across the wire, and the quota path runs on every
/// folder PROPFIND and every upload quota check just to read two i64s.
/// Measured in `benches/QUOTA-PATH.md`.
pub async fn get_storage_usage(&self, id: Uuid) -> UserRepositoryResult<(i64, i64)> {
let row = sqlx::query(
r#"
SELECT storage_used_bytes, storage_quota_bytes
FROM auth.users
WHERE id = $1
"#,
)
.bind(id)
.fetch_one(&*self.pool)
.await
.map_err(Self::map_sqlx_error)?;
Ok((
row.get("storage_used_bytes"),
row.get("storage_quota_bytes"),
))
}
/// Read `force_password_change_at_next_login`. Written TRUE by the
/// admin password-reset flow (via `OpaquePgRepository::clear_registration`,
/// which sets it alongside the envelope invalidation in one UPDATE)
/// and by admin-side `set_user_password`. Cleared on a successful
/// user-initiated `change_password`.
///
/// Reads via a single-column SELECT to avoid dragging the full row
/// (with its up-to-512 KiB `image`) on every login-response mint.
/// Returns `false` for missing users so the login path — which has
/// already resolved the user by id — treats a lost race the same
/// as "flag not set" rather than surfacing a 5xx.
pub async fn is_force_password_change(&self, id: Uuid) -> UserRepositoryResult<bool> {
let row: Option<(bool,)> = sqlx::query_as(
r#"
SELECT force_password_change_at_next_login
FROM auth.users
WHERE id = $1
"#,
)
.bind(id)
.fetch_optional(&*self.pool)
.await
.map_err(Self::map_sqlx_error)?;
Ok(row.map(|(v,)| v).unwrap_or(false))
}
/// Clear `force_password_change_at_next_login`. Called by the
/// change-password flow on success so a legitimate self-service
/// password rotation lifts the admin-set "temporary" marker in
/// one round-trip.
///
/// Deliberately does NOT gate on the current value — flipping FALSE
/// to FALSE is a no-op at the row level. That keeps the caller from
/// needing a read-modify-write.
pub async fn clear_force_password_change(&self, id: Uuid) -> UserRepositoryResult<()> {
sqlx::query(
r#"
UPDATE auth.users
SET force_password_change_at_next_login = FALSE
WHERE id = $1
"#,
)
.bind(id)
.execute(&*self.pool)
.await
.map_err(Self::map_sqlx_error)?;
Ok(())
}
/// Set `force_password_change_at_next_login = TRUE`. Used by
/// admin-initiated password reset when the OPAQUE substrate is NOT
/// wired. When it IS wired, callers should prefer
/// `OpaquePgRepository::clear_registration` which does the same
/// flag flip AND invalidates the OPAQUE envelope in one UPDATE
/// (see the port doc on `clear_registration` for the atomicity
/// contract). This method exists so OPAQUE-off deployments still
/// get the "admin's temp password prompts change on next login"
/// behaviour without having to depend on the OPAQUE code path.
pub async fn set_force_password_change(&self, id: Uuid) -> UserRepositoryResult<()> {
sqlx::query(
r#"
UPDATE auth.users
SET force_password_change_at_next_login = TRUE
WHERE id = $1
"#,
)
.bind(id)
.execute(&*self.pool)
.await
.map_err(Self::map_sqlx_error)?;
Ok(())
}
/// Updates a user's profile image (URL or data URI). Not part of the
/// `UserRepository` trait — called directly from `AuthApplicationService`.
pub async fn update_image(
&self,
user_id: Uuid,
image: Option<String>,
) -> UserRepositoryResult<()> {
sqlx::query(
r#"
UPDATE auth.users
SET image = $2, updated_at = NOW()
WHERE id = $1
"#,
)
.bind(user_id)
.bind(&image)
.execute(&*self.pool)
.await
.map_err(Self::map_sqlx_error)?;
Ok(())
}
/// Shallow-merge a partial UI-preferences patch into
/// `ui_preferences`. The Postgres `||` operator merges top-level
/// keys — `{"a":1,"b":2} || {"b":3,"c":4}` → `{"a":1,"b":3,"c":4}`,
/// which is exactly the semantic PATCH callers want: a partial
/// write only touches the keys it mentions, so a preference set on
/// one device isn't wiped by a partial write from another.
///
/// `jsonb_strip_nulls` removes any key whose incoming value is
/// null, giving callers a documented delete-a-key path (`PATCH
/// {"foo": null}` clears `foo`). Nested nulls inside a value
/// object survive — we only strip at the top level via the merge
/// result.
///
/// Not part of the `UserRepository` trait — called directly from
/// `AuthApplicationService::update_profile`. Bumps `updated_at`
/// so the standard "when did this row change" audits stay useful.
///
/// The CHECK constraints
/// (`users_ui_preferences_is_object` + `_size_cap`) enforce shape
/// and cap at the schema layer; a violating patch surfaces as an
/// sqlx error and returns to the handler as 400.
pub async fn update_ui_preferences(
&self,
user_id: Uuid,
patch: &serde_json::Value,
) -> UserRepositoryResult<()> {
sqlx::query(
r#"
UPDATE auth.users
SET ui_preferences = jsonb_strip_nulls(ui_preferences || $2::jsonb),
updated_at = NOW()
WHERE id = $1
"#,
)
.bind(user_id)
.bind(patch)
.execute(&*self.pool)
.await
.map_err(Self::map_sqlx_error)?;
Ok(())
}
}
impl UserRepository for UserPgRepository {
/// Creates a new user using a transaction
async fn create_user(&self, user: User) -> UserRepositoryResult<User> {
// Create a copy of the user for the closure
let user_clone = user.clone();
with_transaction(&self.pool, "create_user", |tx| {
// We need to move the closure into a BoxFuture to return inside
// the with_transaction call
Box::pin(async move {
// Use getters to extract the values
// Convert user.role() to string to pass it as plain text
let role_str = user_clone.role().to_string();
// Modify the SQL to do an explicit cast to the auth.userrole type
// `image` is included here (was missing pre-fix); without
// it a JIT-provisioned OIDC user landed in the row with
// a NULL profile picture even when the IdP's `picture`
// claim was non-empty. `update_user` already wrote the
// column so existing-user re-logins worked, but the
// first-time INSERT silently dropped it — surfaced by
// tests/oidc/oidc.hurl Step 6 asserting on `$.image`.
let _result = sqlx::query(
r#"
INSERT INTO auth.users (
id, username, email, password_hash, role,
storage_quota_bytes, storage_used_bytes,
created_at, updated_at, last_login_at, active,
federation_kind, federation_issuer, federation_subject,
image, is_external,
given_name, family_name, email_verified_at,
preferred_locale, notify_on_share, ui_preferences
) VALUES (
$1, $2, $3, $4, $5::auth.userrole, $6, $7, $8, $9, $10, $11,
$12, $13, $14, $15, $16, $17, $18, $19, $20, $21, $22
)
RETURNING *
"#,
)
.bind(user_clone.id())
.bind(user_clone.username())
.bind(user_clone.email())
.bind(user_clone.password_hash())
.bind(&role_str) // Convert to string but with explicit cast in SQL
.bind(user_clone.storage_quota_bytes())
.bind(user_clone.storage_used_bytes())
.bind(user_clone.created_at())
.bind(user_clone.updated_at())
.bind(user_clone.last_login_at())
.bind(user_clone.is_active())
.bind(user_clone.federation_kind().map(|k| k.as_str()))
.bind(user_clone.federation_issuer())
.bind(user_clone.federation_subject())
.bind(user_clone.image())
.bind(user_clone.is_external())
.bind(user_clone.given_name())
.bind(user_clone.family_name())
.bind(user_clone.email_verified_at())
.bind(user_clone.preferred_locale())
.bind(user_clone.notify_on_share())
// ui_preferences bind: always a JSON object. `User::new`
// initialises the bag to `{}`; ownership stays with the
// repo for shallow-merge writes via `update_ui_preferences`.
.bind(user_clone.ui_preferences())
.execute(&mut **tx)
.await
.map_err(Self::map_sqlx_error)?;
// We could perform additional operations here,
// such as configuring permissions, roles, etc.
Ok(user_clone)
}) as BoxFuture<'_, UserRepositoryResult<User>>
})
.await?;
Ok(user) // Return the original user for simplicity
}
/// Gets a user by ID
async fn get_user_by_id(&self, id: Uuid) -> UserRepositoryResult<User> {
let row = sqlx::query(
r#"
SELECT
id, username, email, password_hash, role::text as role_text,
storage_quota_bytes, storage_used_bytes,
created_at, updated_at, last_login_at, active,
federation_kind, federation_issuer, federation_subject, image, is_external,
given_name, family_name, email_verified_at, preferred_locale, notify_on_share,
ui_preferences
FROM auth.users
WHERE id = $1
"#,
)
.bind(id)
.fetch_one(&*self.pool)
.await
.map_err(Self::map_sqlx_error)?;
// Convert role string to UserRole enum
let role_str: Option<String> = row.try_get("role_text").unwrap_or(None);
let role = match role_str.as_deref() {
Some("admin") => UserRole::Admin,
_ => UserRole::User,
};
Ok(User::from_data_full(
row.get("id"),
row.get("username"),
row.get("email"),
row.get("password_hash"),
role,
row.get("storage_quota_bytes"),
row.get("storage_used_bytes"),
row.get("created_at"),
row.get("updated_at"),
row.get("last_login_at"),
row.get("active"),
row.get::<Option<String>, _>("federation_kind")
.as_deref()
.and_then(crate::domain::entities::user::FederationKind::parse),
row.get("federation_issuer"),
row.get("federation_subject"),
row.get("image"),
row.get("is_external"),
row.get("given_name"),
row.get("family_name"),
row.get("email_verified_at"),
row.get("preferred_locale"),
row.get("notify_on_share"),
row.get::<serde_json::Value, _>("ui_preferences"),
))
}
/// Gets a user by username
async fn get_user_by_username(&self, username: &str) -> UserRepositoryResult<User> {
let row = sqlx::query(
r#"
SELECT
id, username, email, password_hash, role::text as role_text,
storage_quota_bytes, storage_used_bytes,
created_at, updated_at, last_login_at, active,
federation_kind, federation_issuer, federation_subject, image, is_external,
given_name, family_name, email_verified_at, preferred_locale, notify_on_share,
ui_preferences
FROM auth.users
WHERE username = $1
"#,
)
.bind(username)
.fetch_one(&*self.pool)
.await
.map_err(Self::map_sqlx_error)?;
// Convert role string to UserRole enum
let role_str: Option<String> = row.try_get("role_text").unwrap_or(None);
let role = match role_str.as_deref() {
Some("admin") => UserRole::Admin,
_ => UserRole::User,
};
Ok(User::from_data_full(
row.get("id"),
row.get("username"),
row.get("email"),
row.get("password_hash"),
role,
row.get("storage_quota_bytes"),
row.get("storage_used_bytes"),
row.get("created_at"),
row.get("updated_at"),
row.get("last_login_at"),
row.get("active"),
row.get::<Option<String>, _>("federation_kind")
.as_deref()
.and_then(crate::domain::entities::user::FederationKind::parse),
row.get("federation_issuer"),
row.get("federation_subject"),
row.get("image"),
row.get("is_external"),
row.get("given_name"),
row.get("family_name"),
row.get("email_verified_at"),
row.get("preferred_locale"),
row.get("notify_on_share"),
row.get::<serde_json::Value, _>("ui_preferences"),
))
}
/// Gets a user by email
async fn get_user_by_email(&self, email: &str) -> UserRepositoryResult<User> {
let row = sqlx::query(
r#"
SELECT
id, username, email, password_hash, role::text as role_text,
storage_quota_bytes, storage_used_bytes,
created_at, updated_at, last_login_at, active,
federation_kind, federation_issuer, federation_subject, image, is_external,
given_name, family_name, email_verified_at, preferred_locale, notify_on_share,
ui_preferences
FROM auth.users
WHERE email = $1
"#,
)
.bind(email)
.fetch_one(&*self.pool)
.await
.map_err(Self::map_sqlx_error)?;
// Convert role string to UserRole enum
let role_str: Option<String> = row.try_get("role_text").unwrap_or(None);
let role = match role_str.as_deref() {
Some("admin") => UserRole::Admin,
_ => UserRole::User,
};
Ok(User::from_data_full(
row.get("id"),
row.get("username"),
row.get("email"),
row.get("password_hash"),
role,
row.get("storage_quota_bytes"),
row.get("storage_used_bytes"),
row.get("created_at"),
row.get("updated_at"),
row.get("last_login_at"),
row.get("active"),
row.get::<Option<String>, _>("federation_kind")
.as_deref()
.and_then(crate::domain::entities::user::FederationKind::parse),
row.get("federation_issuer"),
row.get("federation_subject"),
row.get("image"),
row.get("is_external"),
row.get("given_name"),
row.get("family_name"),
row.get("email_verified_at"),
row.get("preferred_locale"),
row.get("notify_on_share"),
row.get::<serde_json::Value, _>("ui_preferences"),
))
}
/// Batch loads users by id in one query (avoids N+1 for group-
/// recipient expansion). Missing ids are silently skipped — the
/// caller treats absent rows as "no such recipient", same as
/// `get_user_by_id` returning `NotFound` for a single lookup.
///
/// Notification-recipient projection: the up-to-512 KiB avatar `image`
/// and the `ui_preferences` JSONB are NOT hydrated (both come back as
/// `None`/`Null`) — the sole caller
/// (`RecipientNotificationService`) reads only the email/eligibility
/// fields, and a group fan-out of M members otherwise detoasted +
/// shipped + parsed M avatars purely to discard them (the ROUND12 §Q1
/// avatar-narrowing pattern; benches/ROUND13.md §Q1). If a future
/// caller needs the avatar, add a wide sibling rather than widening
/// this one back.
async fn get_users_by_ids(&self, ids: Vec<Uuid>) -> UserRepositoryResult<Vec<User>> {
if ids.is_empty() {
return Ok(Vec::new());
}
let rows = sqlx::query(
r#"
SELECT
id, username, email, password_hash, role::text as role_text,
storage_quota_bytes, storage_used_bytes,
created_at, updated_at, last_login_at, active,
federation_kind, federation_issuer, federation_subject, is_external,
given_name, family_name, email_verified_at, preferred_locale, notify_on_share
FROM auth.users
WHERE id = ANY($1)
"#,
)
.bind(&ids)
.fetch_all(&*self.pool)
.await
.map_err(Self::map_sqlx_error)?;
Ok(rows
.into_iter()
.map(|row| {
let role_str: Option<String> = row.try_get("role_text").unwrap_or(None);
let role = match role_str.as_deref() {
Some("admin") => UserRole::Admin,
_ => UserRole::User,
};
User::from_data_full(
row.get("id"),
row.get("username"),
row.get("email"),
row.get("password_hash"),
role,
row.get("storage_quota_bytes"),
row.get("storage_used_bytes"),
row.get("created_at"),
row.get("updated_at"),
row.get("last_login_at"),
row.get("active"),
row.get::<Option<String>, _>("federation_kind")
.as_deref()
.and_then(crate::domain::entities::user::FederationKind::parse),
row.get("federation_issuer"),
row.get("federation_subject"),
None, // image — not projected (notification-recipient path)
row.get("is_external"),
row.get("given_name"),
row.get("family_name"),
row.get("email_verified_at"),
row.get("preferred_locale"),
row.get("notify_on_share"),
serde_json::Value::Null, // ui_preferences — not projected
)
})
.collect())
}
/// Updates an existing user using a transaction
async fn update_user(&self, user: User) -> UserRepositoryResult<User> {
// Create a copy of the user for the closure
let user_clone = user.clone();
with_transaction(&self.pool, "update_user", |tx| {
Box::pin(async move {
// Update the user
sqlx::query(
r#"
UPDATE auth.users
SET
username = $2,
email = $3,
password_hash = $4,
role = $5::auth.userrole,
storage_quota_bytes = $6,
storage_used_bytes = $7,
updated_at = $8,
last_login_at = $9,
active = $10,
image = $11,
given_name = $12,
family_name = $13,
email_verified_at = $14,
preferred_locale = $15,
notify_on_share = $16,
-- Include `is_external` so the external →
-- internal upgrade path
-- (`AuthApplicationService::upgrade_to_internal`)
-- can flip this flag. Previously omitted
-- because no code path mutated it after
-- creation. The DB CHECK
-- `users_external_no_storage`
-- (`is_external=false OR quota=0`) is
-- satisfied by the upgrade because it
-- writes both fields in the same UPDATE:
-- `is_external=false, quota>0`.
is_external = $17
WHERE id = $1
"#,
)
.bind(user_clone.id())
.bind(user_clone.username())
.bind(user_clone.email())
.bind(user_clone.password_hash())
.bind(user_clone.role().to_string())
.bind(user_clone.storage_quota_bytes())
.bind(user_clone.storage_used_bytes())
.bind(user_clone.updated_at())
.bind(user_clone.last_login_at())
.bind(user_clone.is_active())
.bind(user_clone.image())
.bind(user_clone.given_name())
.bind(user_clone.family_name())
.bind(user_clone.email_verified_at())
.bind(user_clone.preferred_locale())
.bind(user_clone.notify_on_share())
.bind(user_clone.is_external())
.execute(&mut **tx)
.await
.map_err(Self::map_sqlx_error)?;
// We could perform additional operations here inside
// the same transaction, such as updating permissions, etc.
Ok(user_clone)
}) as BoxFuture<'_, UserRepositoryResult<User>>
})
.await?;
Ok(user)
}
/// Updates only the storage usage of a user.
///
/// The `IS DISTINCT FROM` guard makes this a no-op when the value is
/// unchanged — which is the common case for the periodic reconciliation
/// sweep — so it produces no dead tuple and no WAL when nothing changed.
async fn update_storage_usage(
&self,
user_id: Uuid,
usage_bytes: i64,
) -> UserRepositoryResult<()> {
sqlx::query(
r#"
UPDATE auth.users
SET
storage_used_bytes = $2,
updated_at = NOW()
WHERE id = $1 AND storage_used_bytes IS DISTINCT FROM $2
"#,
)
.bind(user_id)
.bind(usage_bytes)
.execute(&*self.pool)
.await
.map_err(Self::map_sqlx_error)?;
Ok(())
}
/// Updates the last login date
async fn update_last_login(&self, user_id: Uuid) -> UserRepositoryResult<()> {
sqlx::query(
r#"
UPDATE auth.users
SET
last_login_at = NOW(),
updated_at = NOW()
WHERE id = $1
"#,
)
.bind(user_id)
.execute(&*self.pool)
.await
.map_err(Self::map_sqlx_error)?;
Ok(())
}
/// Lists users with pagination
async fn list_users(
&self,
limit: i64,
offset: i64,
include_external: bool,
) -> UserRepositoryResult<Vec<User>> {
let rows = sqlx::query(
r#"
SELECT
id, username, email, password_hash, role::text as role_text,
storage_quota_bytes, storage_used_bytes,
created_at, updated_at, last_login_at, active,
federation_kind, federation_issuer, federation_subject, image, is_external,
given_name, family_name, email_verified_at, preferred_locale, notify_on_share,
ui_preferences
FROM auth.users
WHERE ($3 OR is_external = FALSE)
ORDER BY created_at DESC, id DESC
LIMIT $1 OFFSET $2
"#,
)
.bind(limit)
.bind(offset)
.bind(include_external)
.fetch_all(&*self.pool)
.await
.map_err(Self::map_sqlx_error)?;
let users = rows
.into_iter()
.map(|row| {
// Convert role string to UserRole enum for each row
let role_str: Option<String> = row.try_get("role_text").unwrap_or(None);
let role = match role_str.as_deref() {
Some("admin") => UserRole::Admin,
_ => UserRole::User,
};
User::from_data_full(
row.get("id"),
row.get("username"),
row.get("email"),
row.get("password_hash"),
role,
row.get("storage_quota_bytes"),
row.get("storage_used_bytes"),
row.get("created_at"),
row.get("updated_at"),
row.get("last_login_at"),
row.get("active"),
row.get::<Option<String>, _>("federation_kind")
.as_deref()
.and_then(crate::domain::entities::user::FederationKind::parse),
row.get("federation_issuer"),
row.get("federation_subject"),
row.get("image"),
row.get("is_external"),
row.get("given_name"),
row.get("family_name"),
row.get("email_verified_at"),
row.get("preferred_locale"),
row.get("notify_on_share"),
row.get::<serde_json::Value, _>("ui_preferences"),
)
})
.collect();
Ok(users)
}
async fn list_user_summaries(
&self,
limit: i64,
offset: i64,
include_external: bool,
) -> UserRepositoryResult<Vec<UserListEntry>> {
let rows = sqlx::query_as::<
_,
(
Uuid,
Option<String>,
String,
String,
i64,
i64,
Option<chrono::DateTime<chrono::Utc>>,
bool,
Option<String>,
Option<String>,
bool,
bool,
bool,
bool,
),
>(
// Auth-credential columns projected as booleans via `IS NOT
// NULL` rather than as timestamps / hashes so the row-mapping
// tuple stays small and the wire shape is exactly what the
// admin table needs. Per-row scalar tests — no cost beyond
// the full-table sequential scan the LIMIT/OFFSET already
// pays. `has_password` on the password_hash column tells
// the admin table whether a server-verifiable password is
// on file; combined with the two OPAQUE flags and
// federation_kind / federation_issuer, the SPA derives the
// full "capability set" per user (password / OPAQUE / SSO /
// passwordless).
r#"
SELECT
id, username, email, role::text,
storage_quota_bytes, storage_used_bytes,
last_login_at, active,
federation_kind, federation_issuer, is_external,
(password_hash IS NOT NULL) AS has_password,
(opaque_envelope IS NOT NULL) AS opaque_registered,
(opaque_migrated_at IS NOT NULL) AS opaque_migrated
FROM auth.users
WHERE ($3 OR is_external = FALSE)
ORDER BY created_at DESC, id DESC
LIMIT $1 OFFSET $2
"#,
)
.bind(limit)
.bind(offset)
.bind(include_external)
.fetch_all(self.pool.as_ref())
.await
.map_err(Self::map_sqlx_error)?;
Ok(rows
.into_iter()
.map(
|(
id,
username,
email,
role,
storage_quota_bytes,
storage_used_bytes,
last_login_at,
active,
federation_kind,
federation_issuer,
is_external,
has_password,
opaque_registered,
opaque_migrated,
)| UserListEntry {
id,
username,
email,
role: if role == "admin" {
UserRole::Admin
} else {
UserRole::User
},
storage_quota_bytes,
storage_used_bytes,
last_login_at,
active,
federation_kind,
federation_issuer,
is_external,
has_password,
opaque_registered,
opaque_migrated,
},
)
.collect())
}
async fn search_users(
&self,
query: &str,
limit: i64,
include_external: bool,
) -> UserRepositoryResult<Vec<User>> {
let pattern = format!("%{}%", query);
let rows = sqlx::query(
r#"
SELECT
id, username, email, password_hash, role::text as role_text,
storage_quota_bytes, storage_used_bytes,
created_at, updated_at, last_login_at, active,
federation_kind, federation_issuer, federation_subject, image, is_external,
given_name, family_name, email_verified_at, preferred_locale, notify_on_share,
ui_preferences
FROM auth.users
WHERE (username ILIKE $1 OR email ILIKE $1)
AND ($3 OR is_external = FALSE)
ORDER BY username
LIMIT $2
"#,
)
.bind(&pattern)
.bind(limit)
.bind(include_external)
.fetch_all(&*self.pool)
.await
.map_err(Self::map_sqlx_error)?;
let users = rows
.into_iter()
.map(|row| {
let role_str: Option<String> = row.try_get("role_text").unwrap_or(None);
let role = match role_str.as_deref() {
Some("admin") => UserRole::Admin,
_ => UserRole::User,
};
User::from_data_full(
row.get("id"),
row.get("username"),
row.get("email"),
row.get("password_hash"),
role,
row.get("storage_quota_bytes"),
row.get("storage_used_bytes"),
row.get("created_at"),
row.get("updated_at"),
row.get("last_login_at"),
row.get("active"),
row.get::<Option<String>, _>("federation_kind")
.as_deref()
.and_then(crate::domain::entities::user::FederationKind::parse),
row.get("federation_issuer"),
row.get("federation_subject"),
row.get("image"),
row.get("is_external"),
row.get("given_name"),
row.get("family_name"),
row.get("email_verified_at"),
row.get("preferred_locale"),
row.get("notify_on_share"),
row.get::<serde_json::Value, _>("ui_preferences"),
)
})
.collect();
Ok(users)
}
/// Activates or deactivates a user
async fn set_user_active_status(
&self,
user_id: Uuid,
active: bool,
) -> UserRepositoryResult<()> {
sqlx::query(
r#"
UPDATE auth.users
SET
active = $2,
updated_at = NOW()
WHERE id = $1
"#,
)
.bind(user_id)
.bind(active)
.execute(&*self.pool)
.await
.map_err(Self::map_sqlx_error)?;
Ok(())
}
/// Changes a user's password
async fn change_password(
&self,
user_id: Uuid,
password_hash: &str,
) -> UserRepositoryResult<()> {
sqlx::query(
r#"
UPDATE auth.users
SET
password_hash = $2,
updated_at = NOW()
WHERE id = $1
"#,
)
.bind(user_id)
.bind(password_hash)
.execute(&*self.pool)
.await
.map_err(Self::map_sqlx_error)?;
Ok(())
}
/// Changes a user's role
async fn change_role(&self, user_id: Uuid, role: UserRole) -> UserRepositoryResult<()> {
// Convert the role to string for the binding
let role_str = role.to_string();
sqlx::query(
r#"
UPDATE auth.users
SET
role = $2::auth.userrole,
updated_at = NOW()
WHERE id = $1
"#,
)
.bind(user_id)
.bind(&role_str)
.execute(&*self.pool)
.await
.map_err(Self::map_sqlx_error)?;
Ok(())
}
/// Counts users by role with a scalar `COUNT(*)` — no row hydration.
async fn count_users_by_role(&self, role: &str) -> UserRepositoryResult<i64> {
sqlx::query_scalar("SELECT COUNT(*) FROM auth.users WHERE role::text = $1")
.bind(role)
.fetch_one(&*self.pool)
.await
.map_err(Self::map_sqlx_error)
}
/// Lists users by role
async fn list_users_by_role(&self, role: &str) -> UserRepositoryResult<Vec<User>> {
let rows = sqlx::query(
r#"
SELECT
id, username, email, password_hash, role::text as role_text,
storage_quota_bytes, storage_used_bytes,
created_at, updated_at, last_login_at, active,
federation_kind, federation_issuer, federation_subject, image, is_external,
given_name, family_name, email_verified_at, preferred_locale, notify_on_share,
ui_preferences
FROM auth.users
WHERE role::text = $1
ORDER BY created_at DESC
"#,
)
.bind(role)
.fetch_all(&*self.pool)
.await
.map_err(Self::map_sqlx_error)?;
let users = rows
.into_iter()
.map(|row| {
// Convert role string to UserRole enum for each row
let role_str: Option<String> = row.try_get("role_text").unwrap_or(None);
let role = match role_str.as_deref() {
Some("admin") => UserRole::Admin,
_ => UserRole::User,
};
User::from_data_full(
row.get("id"),
row.get("username"),
row.get("email"),
row.get("password_hash"),
role,
row.get("storage_quota_bytes"),
row.get("storage_used_bytes"),
row.get("created_at"),
row.get("updated_at"),
row.get("last_login_at"),
row.get("active"),
row.get::<Option<String>, _>("federation_kind")
.as_deref()
.and_then(crate::domain::entities::user::FederationKind::parse),
row.get("federation_issuer"),
row.get("federation_subject"),
row.get("image"),
row.get("is_external"),
row.get("given_name"),
row.get("family_name"),
row.get("email_verified_at"),
row.get("preferred_locale"),
row.get("notify_on_share"),
row.get::<serde_json::Value, _>("ui_preferences"),
)
})
.collect();
Ok(users)
}
/// Deletes a user
async fn delete_user(&self, user_id: Uuid) -> UserRepositoryResult<()> {
sqlx::query(
r#"
DELETE FROM auth.users
WHERE id = $1
"#,
)
.bind(user_id)
.execute(&*self.pool)
.await
.map_err(Self::map_sqlx_error)?;
Ok(())
}
/// Finds a user by OIDC provider + subject pair
async fn get_user_by_federation_subject(
&self,
provider: &str,
subject: &str,
) -> UserRepositoryResult<User> {
let row = sqlx::query(
r#"
SELECT
id, username, email, password_hash, role::text as role_text,
storage_quota_bytes, storage_used_bytes,
created_at, updated_at, last_login_at, active,
federation_kind, federation_issuer, federation_subject, image, is_external,
given_name, family_name, email_verified_at, preferred_locale, notify_on_share,
ui_preferences
FROM auth.users
WHERE federation_issuer = $1 AND federation_subject = $2
"#,
)
.bind(provider)
.bind(subject)
.fetch_one(&*self.pool)
.await
.map_err(Self::map_sqlx_error)?;
let role_str: Option<String> = row.try_get("role_text").unwrap_or(None);
let role = match role_str.as_deref() {
Some("admin") => UserRole::Admin,
_ => UserRole::User,
};
Ok(User::from_data_full(
row.get("id"),
row.get("username"),
row.get("email"),
row.get("password_hash"),
role,
row.get("storage_quota_bytes"),
row.get("storage_used_bytes"),
row.get("created_at"),
row.get("updated_at"),
row.get("last_login_at"),
row.get("active"),
row.get::<Option<String>, _>("federation_kind")
.as_deref()
.and_then(crate::domain::entities::user::FederationKind::parse),
row.get("federation_issuer"),
row.get("federation_subject"),
row.get("image"),
row.get("is_external"),
row.get("given_name"),
row.get("family_name"),
row.get("email_verified_at"),
row.get("preferred_locale"),
row.get("notify_on_share"),
row.get::<serde_json::Value, _>("ui_preferences"),
))
}
/// Updates a user's storage quota
async fn update_storage_quota(
&self,
user_id: Uuid,
quota_bytes: i64,
) -> UserRepositoryResult<()> {
sqlx::query(
r#"
UPDATE auth.users
SET
storage_quota_bytes = $2,
updated_at = NOW()
WHERE id = $1
"#,
)
.bind(user_id)
.bind(quota_bytes)
.execute(&*self.pool)
.await
.map_err(Self::map_sqlx_error)?;
Ok(())
}
/// Counts the total number of users
async fn count_users(&self) -> UserRepositoryResult<i64> {
let row = sqlx::query("SELECT COUNT(*) as count FROM auth.users")
.fetch_one(&*self.pool)
.await
.map_err(Self::map_sqlx_error)?;
let count: i64 = row.get("count");
Ok(count)
}
/// Gets aggregated storage statistics
async fn get_storage_stats(&self) -> UserRepositoryResult<StorageStats> {
let row = sqlx::query(
r#"
SELECT
COUNT(*) as total_users,
COUNT(*) FILTER (WHERE active = true) as active_users,
COALESCE(SUM(storage_quota_bytes), 0) as total_quota_bytes,
COALESCE(SUM(storage_used_bytes), 0) as total_used_bytes,
COUNT(*) FILTER (WHERE storage_quota_bytes > 0 AND storage_used_bytes > storage_quota_bytes * 0.8) as users_over_80_percent,
COUNT(*) FILTER (WHERE storage_quota_bytes > 0 AND storage_used_bytes > storage_quota_bytes) as users_over_quota
FROM auth.users
"#
)
.fetch_one(&*self.pool)
.await
.map_err(Self::map_sqlx_error)?;
Ok(StorageStats {
total_users: row.get("total_users"),
active_users: row.get("active_users"),
total_quota_bytes: row.get("total_quota_bytes"),
total_used_bytes: row.get("total_used_bytes"),
users_over_80_percent: row.get("users_over_80_percent"),
users_over_quota: row.get("users_over_quota"),
})
}
}
// Storage port implementation for the application layer
impl UserStoragePort for UserPgRepository {
async fn create_user(&self, user: User) -> Result<User, DomainError> {
UserRepository::create_user(self, user)
.await
.map_err(DomainError::from)
}
async fn get_user_by_id(&self, id: Uuid) -> Result<User, DomainError> {
UserRepository::get_user_by_id(self, id)
.await
.map_err(DomainError::from)
}
async fn get_users_by_ids(&self, ids: Vec<Uuid>) -> Result<Vec<User>, DomainError> {
UserRepository::get_users_by_ids(self, ids)
.await
.map_err(DomainError::from)
}
async fn get_user_by_username(&self, username: &str) -> Result<User, DomainError> {
UserRepository::get_user_by_username(self, username)
.await
.map_err(DomainError::from)
}
async fn get_user_by_email(&self, email: &str) -> Result<User, DomainError> {
UserRepository::get_user_by_email(self, email)
.await
.map_err(DomainError::from)
}
async fn update_user(&self, user: User) -> Result<User, DomainError> {
UserRepository::update_user(self, user)
.await
.map_err(DomainError::from)
}
async fn update_storage_usage(
&self,
user_id: Uuid,
usage_bytes: i64,
) -> Result<(), DomainError> {
UserRepository::update_storage_usage(self, user_id, usage_bytes)
.await
.map_err(DomainError::from)
}
async fn list_users(
&self,
limit: i64,
offset: i64,
include_external: bool,
) -> Result<Vec<User>, DomainError> {
UserRepository::list_users(self, limit, offset, include_external)
.await
.map_err(DomainError::from)
}
async fn list_user_summaries(
&self,
limit: i64,
offset: i64,
include_external: bool,
) -> Result<Vec<UserListEntry>, DomainError> {
UserRepository::list_user_summaries(self, limit, offset, include_external)
.await
.map_err(DomainError::from)
}
async fn search_users(
&self,
query: &str,
limit: i64,
include_external: bool,
) -> Result<Vec<User>, DomainError> {
UserRepository::search_users(self, query, limit, include_external)
.await
.map_err(DomainError::from)
}
async fn search_usernames(
&self,
query: &str,
limit: i64,
include_external: bool,
) -> Result<Vec<Option<String>>, DomainError> {
// Same predicate / order / limit as `search_users`, username-only
// projection — the sharee autocomplete path reads nothing else, and
// the wide row drags the avatar `image` per matched user.
let pattern = format!("%{}%", query);
let rows = sqlx::query(
r#"
SELECT username
FROM auth.users
WHERE (username ILIKE $1 OR email ILIKE $1)
AND ($3 OR is_external = FALSE)
ORDER BY username
LIMIT $2
"#,
)
.bind(&pattern)
.bind(limit)
.bind(include_external)
.fetch_all(&*self.pool)
.await
.map_err(Self::map_sqlx_error)
.map_err(DomainError::from)?;
Ok(rows.into_iter().map(|row| row.get("username")).collect())
}
async fn mark_email_verified(&self, user_id: Uuid) -> Result<(), DomainError> {
// SQL twin of `User::mark_email_verified` — stamps once, keeps the
// first timestamp, and touches only the two columns involved.
sqlx::query(
r#"
UPDATE auth.users
SET email_verified_at = NOW(), updated_at = NOW()
WHERE id = $1 AND email_verified_at IS NULL
"#,
)
.bind(user_id)
.execute(&*self.pool)
.await
.map_err(Self::map_sqlx_error)
.map_err(DomainError::from)?;
Ok(())
}
async fn sync_oidc_login_profile(
&self,
user_id: Uuid,
image: Option<&str>,
) -> Result<(), DomainError> {
// `IS DISTINCT FROM` guard (the `update_storage_usage` pattern): the
// common repeat-login case — same IdP avatar, already verified —
// writes nothing at all (no dead tuple, no WAL).
sqlx::query(
r#"
UPDATE auth.users
SET image = $2,
email_verified_at = COALESCE(email_verified_at, NOW()),
updated_at = NOW()
WHERE id = $1
AND (image IS DISTINCT FROM $2 OR email_verified_at IS NULL)
"#,
)
.bind(user_id)
.bind(image)
.execute(&*self.pool)
.await
.map_err(Self::map_sqlx_error)
.map_err(DomainError::from)?;
Ok(())
}
async fn rebind_federation_issuer(
&self,
user_id: Uuid,
new_issuer: &str,
) -> Result<(), DomainError> {
// Same `IS DISTINCT FROM` guard as sync_oidc_login_profile: this
// fires on every OIDC login, so the common already-migrated case
// must be a zero-write no-op. Only actually flips the column
// when the stored value is stale (legacy display label vs the
// real issuer URL from the id_token's `iss` claim).
sqlx::query(
r#"
UPDATE auth.users
SET federation_issuer = $2,
updated_at = NOW()
WHERE id = $1
AND federation_issuer IS DISTINCT FROM $2
"#,
)
.bind(user_id)
.bind(new_issuer)
.execute(&*self.pool)
.await
.map_err(Self::map_sqlx_error)
.map_err(DomainError::from)?;
Ok(())
}
async fn list_users_by_role(&self, role: &str) -> Result<Vec<User>, DomainError> {
UserRepository::list_users_by_role(self, role)
.await
.map_err(DomainError::from)
}
async fn count_users_by_role(&self, role: &str) -> Result<i64, DomainError> {
UserRepository::count_users_by_role(self, role)
.await
.map_err(DomainError::from)
}
async fn delete_user(&self, user_id: Uuid) -> Result<(), DomainError> {
UserRepository::delete_user(self, user_id)
.await
.map_err(DomainError::from)
}
async fn change_password(&self, user_id: Uuid, password_hash: &str) -> Result<(), DomainError> {
UserRepository::change_password(self, user_id, password_hash)
.await
.map_err(DomainError::from)
}
async fn get_user_by_federation_subject(
&self,
provider: &str,
subject: &str,
) -> Result<User, DomainError> {
UserRepository::get_user_by_federation_subject(self, provider, subject)
.await
.map_err(DomainError::from)
}
async fn set_user_active_status(&self, user_id: Uuid, active: bool) -> Result<(), DomainError> {
UserRepository::set_user_active_status(self, user_id, active)
.await
.map_err(DomainError::from)
}
async fn change_role(&self, user_id: Uuid, role: &str) -> Result<(), DomainError> {
let user_role = match role {
"admin" => UserRole::Admin,
_ => UserRole::User,
};
UserRepository::change_role(self, user_id, user_role)
.await
.map_err(DomainError::from)
}
async fn update_storage_quota(
&self,
user_id: Uuid,
quota_bytes: i64,
) -> Result<(), DomainError> {
UserRepository::update_storage_quota(self, user_id, quota_bytes)
.await
.map_err(DomainError::from)
}
async fn count_users(&self) -> Result<i64, DomainError> {
UserRepository::count_users(self)
.await
.map_err(DomainError::from)
}
}
#[cfg(integration_tests)]
#[allow(dead_code)]
mod integration_tests {
use super::*;
use crate::integration_test_support::{ensure_clean_test_db, test_db_url};
use sqlx::postgres::PgPoolOptions;
async fn test_repo() -> UserPgRepository {
let pool = PgPoolOptions::new()
.max_connections(2)
.connect(&test_db_url())
.await
.expect("connect to integration-test PostgreSQL");
ensure_clean_test_db(&pool).await;
UserPgRepository::new(Arc::new(pool))
}
async fn insert_summary_fixture(
repo: &UserPgRepository,
id: Uuid,
username: Option<&str>,
email: &str,
role: &str,
is_external: bool,
) {
sqlx::query(
r#"
INSERT INTO auth.users (
id, username, email, password_hash, role,
storage_quota_bytes, storage_used_bytes,
created_at, updated_at, last_login_at, active,
federation_issuer, is_external
) VALUES (
$1, $2, $3, NULL, $4::auth.userrole,
$5, 0,
'9999-12-31 23:59:59+00', '9999-12-31 23:59:59+00', NULL, TRUE,
$6, $7
)
"#,
)
.bind(id)
.bind(username)
.bind(email)
.bind(role)
.bind(if is_external {
0_i64
} else {
10_737_418_240_i64
})
.bind(is_external.then_some("integration-idp"))
.bind(is_external)
.execute(repo.pool.as_ref())
.await
.expect("insert compact-list fixture");
}
#[tokio::test]
async fn compact_listing_maps_narrow_columns_and_stably_breaks_timestamp_ties() {
let repo = test_repo().await;
sqlx::query("DELETE FROM auth.users WHERE email LIKE 'perf-summary-%@example.invalid'")
.execute(repo.pool.as_ref())
.await
.expect("clean stale compact-list fixtures");
let mut ids = [Uuid::new_v4(), Uuid::new_v4(), Uuid::new_v4()];
ids.sort_unstable_by(|left, right| right.cmp(left));
let username_a = format!("perf-summary-a-{}", ids[0]);
let username_b = format!("perf-summary-b-{}", ids[2]);
insert_summary_fixture(
&repo,
ids[0],
Some(&username_a),
&format!("perf-summary-{}@example.invalid", ids[0]),
"admin",
false,
)
.await;
insert_summary_fixture(
&repo,
ids[1],
None,
&format!("perf-summary-{}@example.invalid", ids[1]),
"user",
true,
)
.await;
insert_summary_fixture(
&repo,
ids[2],
Some(&username_b),
&format!("perf-summary-{}@example.invalid", ids[2]),
"user",
false,
)
.await;
let page = UserRepository::list_user_summaries(&repo, 3, 0, true)
.await
.expect("compact projection query must decode");
assert_eq!(page.iter().map(|entry| entry.id).collect::<Vec<_>>(), ids);
assert_eq!(page[0].username.as_deref(), Some(username_a.as_str()));
assert_eq!(page[0].role, UserRole::Admin);
assert_eq!(page[0].storage_quota_bytes, 10_737_418_240);
assert_eq!(page[1].username, None);
assert!(page[1].is_external);
assert_eq!(page[1].federation_issuer.as_deref(), Some("integration-idp"));
let internal = UserRepository::list_user_summaries(&repo, 10, 0, false)
.await
.expect("internal compact projection query must decode");
assert!(internal.iter().any(|entry| entry.id == ids[0]));
assert!(internal.iter().any(|entry| entry.id == ids[2]));
assert!(!internal.iter().any(|entry| entry.id == ids[1]));
sqlx::query("DELETE FROM auth.users WHERE id = ANY($1)")
.bind(ids.as_slice())
.execute(repo.pool.as_ref())
.await
.expect("clean compact-list fixtures");
}
}