feat: complete Phase 0 Enterprise IAM

This commit is contained in:
thakares committed 2026-07-21 15:26:16 +05:30
1 parent 3d2d291006
commit c2f5ba3f54
202 files changed
+21371 -1360

No files matched your search

+1 -60
View File
@@ -1,60 +1 @@
use sqlx::SqlitePool;
use crate::db::models::Application;
pub async fn create(
pool: &SqlitePool,
id: &str,
tenant_id: &str,
name: &str,
slug: &str,
) -> Result<Application, sqlx::Error> {
sqlx::query_as::<_, Application>(
r#"
INSERT INTO applications (id, tenant_id, name, slug)
VALUES (?, ?, ?, ?)
RETURNING *
"#,
)
.bind(id)
.bind(tenant_id)
.bind(name)
.bind(slug)
.fetch_one(pool)
.await
}
pub async fn find_by_slug(
pool: &SqlitePool,
slug: &str,
) -> Result<Option<Application>, sqlx::Error> {
sqlx::query_as::<_, Application>("SELECT * FROM applications WHERE slug = ?")
.bind(slug)
.fetch_optional(pool)
.await
}
pub async fn find_by_id(pool: &SqlitePool, id: &str) -> Result<Option<Application>, sqlx::Error> {
sqlx::query_as::<_, Application>("SELECT * FROM applications WHERE id = ?")
.bind(id)
.fetch_optional(pool)
.await
}
pub async fn list(pool: &SqlitePool, tenant_id: &str) -> Result<Vec<Application>, sqlx::Error> {
sqlx::query_as::<_, Application>("SELECT * FROM applications WHERE tenant_id = ? ORDER BY name")
.bind(tenant_id)
.fetch_all(pool)
.await
}
pub async fn set_enabled(pool: &SqlitePool, id: &str, enabled: bool) -> Result<(), sqlx::Error> {
sqlx::query(
"UPDATE applications SET enabled = ?, updated_at = strftime('%Y-%m-%dT%H:%M:%SZ', 'now') WHERE id = ?",
)
.bind(enabled)
.bind(id)
.execute(pool)
.await?;
Ok(())
}
pub use crate::db::repository::sqlite::applications::*;
+35 -30
View File
@@ -1,10 +1,14 @@
use sqlx::SqlitePool;
pub use crate::db::repository::sqlite::audit::*;
use crate::db::models::AuditLog;
use crate::db::provider::DatabaseProvider;
use std::sync::Arc;
// Removed direct import of AuditFilter to avoid conflict with traits version
/// Insert an audit log entry using the provided DatabaseProvider.
#[allow(clippy::too_many_arguments)]
pub async fn insert(
tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>,
provider: &Arc<dyn DatabaseProvider>,
id: &str,
actor_user_id: Option<&str>,
target_user_id: Option<&str>,
@@ -16,34 +20,35 @@ pub async fn insert(
user_agent: Option<&str>,
metadata_json: Option<&str>,
) -> Result<AuditLog, sqlx::Error> {
sqlx::query_as::<_, AuditLog>(
r#"
INSERT INTO audit_logs (
id, actor_user_id, target_user_id,
action, resource_type, resource_id,
severity, ip_address, user_agent, metadata_json
provider
.audit()
.insert(
id,
actor_user_id,
target_user_id,
action,
resource_type,
resource_id,
severity,
ip_address,
user_agent,
metadata_json,
)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
RETURNING *
"#,
)
.bind(id)
.bind(actor_user_id)
.bind(target_user_id)
.bind(action)
.bind(resource_type)
.bind(resource_id)
.bind(severity)
.bind(ip_address)
.bind(user_agent)
.bind(metadata_json)
.fetch_one(&mut **tx)
.await
}
pub async fn list_recent(pool: &SqlitePool, limit: i64) -> Result<Vec<AuditLog>, sqlx::Error> {
sqlx::query_as::<_, AuditLog>("SELECT * FROM audit_logs ORDER BY created_at DESC LIMIT ?")
.bind(limit)
.fetch_all(pool)
.await
}
/// Count filtered audit logs using the provided DatabaseProvider.
pub async fn count_filtered(
provider: &Arc<dyn DatabaseProvider>,
filter: &AuditFilter,
) -> Result<i64, sqlx::Error> {
provider.audit().count_filtered(filter).await
}
/// List filtered audit logs using the provided DatabaseProvider.
pub async fn list_filtered(
provider: &Arc<dyn DatabaseProvider>,
filter: &AuditFilter,
) -> Result<Vec<AuditLog>, sqlx::Error> {
provider.audit().list_filtered(filter).await
}
+13 -6
View File
@@ -1,8 +1,15 @@
pub mod applications;
pub mod traits;
pub use traits::*;
#[cfg(feature = "sqlite")]
pub mod sqlite;
#[cfg(feature = "sqlite")]
pub use sqlite::*;
#[cfg(feature = "postgres")]
pub mod postgres;
#[cfg(all(feature = "postgres", not(feature = "sqlite")))]
pub use postgres::*;
pub mod audit;
pub mod permissions;
pub mod roles;
pub mod service_accounts;
pub mod sessions;
pub mod tokens;
pub mod users;
+1 -41
View File
@@ -1,41 +1 @@
use sqlx::SqlitePool;
/// Return all permission names held by a user (via their roles).
pub async fn list_for_user(pool: &SqlitePool, user_id: &str) -> Result<Vec<String>, sqlx::Error> {
let rows: Vec<(String,)> = sqlx::query_as(
r#"
SELECT DISTINCT p.name
FROM permissions p
JOIN role_permissions rp ON rp.permission_id = p.id
JOIN user_roles ur ON ur.role_id = rp.role_id
WHERE ur.user_id = ?
ORDER BY p.name
"#,
)
.bind(user_id)
.fetch_all(pool)
.await?;
Ok(rows.into_iter().map(|(name,)| name).collect())
}
/// Check if a user holds a specific named permission.
pub async fn user_has_permission(
pool: &SqlitePool,
user_id: &str,
permission_name: &str,
) -> Result<bool, sqlx::Error> {
let row: (i64,) = sqlx::query_as(
r#"
SELECT COUNT(*)
FROM permissions p
JOIN role_permissions rp ON rp.permission_id = p.id
JOIN user_roles ur ON ur.role_id = rp.role_id
WHERE ur.user_id = ? AND p.name = ?
"#,
)
.bind(user_id)
.bind(permission_name)
.fetch_one(pool)
.await?;
Ok(row.0 > 0)
}
pub use crate::db::repository::sqlite::permissions::*;
+108
View File
@@ -0,0 +1,108 @@
use crate::db::repository::traits::ApplicationsRepository;
use async_trait::async_trait;
use sqlx::PgPool;
use crate::db::models::Application;
pub struct PostgresApplicationsRepository {
pub pool: PgPool,
}
#[async_trait]
impl ApplicationsRepository for PostgresApplicationsRepository {
async fn create(
&self,
id: &str,
tenant_id: &str,
name: &str,
slug: &str,
) -> Result<Application, sqlx::Error> {
sqlx::query_as::<_, Application>(
r#"
INSERT INTO applications (id, tenant_id, name, slug)
VALUES ($1, $2, $3, $4)
RETURNING *
"#,
)
.bind(id)
.bind(tenant_id)
.bind(name)
.bind(slug)
.fetch_one(&self.pool)
.await
}
async fn find_by_slug(&self, slug: &str) -> Result<Option<Application>, sqlx::Error> {
sqlx::query_as::<_, Application>("SELECT * FROM applications WHERE slug = $1")
.bind(slug)
.fetch_optional(&self.pool)
.await
}
async fn find_by_id(&self, id: &str) -> Result<Option<Application>, sqlx::Error> {
sqlx::query_as::<_, Application>("SELECT * FROM applications WHERE id = $1")
.bind(id)
.fetch_optional(&self.pool)
.await
}
async fn list(&self, tenant_id: &str) -> Result<Vec<Application>, sqlx::Error> {
sqlx::query_as::<_, Application>(
"SELECT * FROM applications WHERE tenant_id = $1 ORDER BY name",
)
.bind(tenant_id)
.fetch_all(&self.pool)
.await
}
async fn set_enabled(&self, id: &str, enabled: bool) -> Result<(), sqlx::Error> {
sqlx::query(
"UPDATE applications SET enabled = $1, updated_at = strftime('%Y-%m-%dT%H:%M:%SZ', 'now') WHERE id = $2",
)
.bind(enabled)
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
async fn update(
&self,
id: &str,
name: &str,
slug: &str,
enabled: bool,
) -> Result<(), sqlx::Error> {
sqlx::query(
r#"
UPDATE applications
SET name = $1, slug = $2, enabled = $3,
updated_at = strftime('%Y-%m-%dT%H:%M:%SZ', 'now')
WHERE id = $4
"#,
)
.bind(name)
.bind(slug)
.bind(enabled)
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
async fn delete(&self, id: &str) -> Result<(), sqlx::Error> {
sqlx::query("DELETE FROM applications WHERE id = $1")
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
async fn count(&self, tenant_id: &str) -> Result<i64, sqlx::Error> {
let row: (i64,) = sqlx::query_as("SELECT COUNT(*) FROM applications WHERE tenant_id = $1")
.bind(tenant_id)
.fetch_one(&self.pool)
.await?;
Ok(row.0)
}
}
+146
View File
@@ -0,0 +1,146 @@
use crate::db::repository::traits::AuditRepository;
use async_trait::async_trait;
use sqlx::PgPool;
use crate::db::models::AuditLog;
pub struct PostgresAuditRepository {
pub pool: PgPool,
}
use crate::db::repository::sqlite::audit::AuditFilter;
#[async_trait]
impl AuditRepository for PostgresAuditRepository {
async fn count(&self) -> Result<i64, sqlx::Error> {
let row: (i64,) = sqlx::query_as("SELECT COUNT(*) FROM audit_logs")
.fetch_one(&self.pool)
.await?;
Ok(row.0)
}
#[allow(clippy::too_many_arguments)]
async fn insert(
&self,
id: &str,
actor_user_id: Option<&str>,
target_user_id: Option<&str>,
action: &str,
resource_type: &str,
resource_id: Option<&str>,
severity: &str,
ip_address: Option<&str>,
user_agent: Option<&str>,
metadata_json: Option<&str>,
) -> Result<AuditLog, sqlx::Error> {
sqlx::query_as::<_, AuditLog>(
r#"
INSERT INTO audit_logs (
id, actor_user_id, target_user_id,
action, resource_type, resource_id,
severity, ip_address, user_agent, metadata_json
)
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10)
RETURNING *
"#,
)
.bind(id)
.bind(actor_user_id)
.bind(target_user_id)
.bind(action)
.bind(resource_type)
.bind(resource_id)
.bind(severity)
.bind(ip_address)
.bind(user_agent)
.bind(metadata_json)
.fetch_one(&self.pool)
.await
}
async fn list_recent(&self, limit: i64) -> Result<Vec<AuditLog>, sqlx::Error> {
sqlx::query_as::<_, AuditLog>("SELECT * FROM audit_logs ORDER BY created_at DESC LIMIT $1")
.bind(limit)
.fetch_all(&self.pool)
.await
}
async fn list_filtered(&self, filter: &AuditFilter) -> Result<Vec<AuditLog>, sqlx::Error> {
// Build a dynamic but simple filter using COALESCE-style optional matches.
// Empty optionals are treated as wildcards via OR IS NULL pattern with bind of None.
let search_like = filter
.search
.as_ref()
.map(|s| format!("%{}%", s.replace('%', "\\%")));
sqlx::query_as::<_, AuditLog>(
r#"
SELECT * FROM audit_logs
WHERE ($11 IS NULL OR actor_user_id = $21)
AND ($32 IS NULL OR action = $42)
AND ($53 IS NULL OR resource_type = $63)
AND ($74 IS NULL OR severity = $84)
AND ($95 IS NULL OR created_at >= $105)
AND ($116 IS NULL OR created_at <= $126)
AND (
$137 IS NULL
OR action LIKE $147 ESCAPE '\'
OR resource_type LIKE $157 ESCAPE '\'
OR resource_id LIKE $167 ESCAPE '\'
OR ip_address LIKE $177 ESCAPE '\'
OR metadata_json LIKE $187 ESCAPE '\'
)
ORDER BY created_at DESC
LIMIT $198 OFFSET $209
"#,
)
.bind(filter.actor_user_id.as_deref())
.bind(filter.action.as_deref())
.bind(filter.resource_type.as_deref())
.bind(filter.severity.as_deref())
.bind(filter.since.as_deref())
.bind(filter.until.as_deref())
.bind(search_like.as_deref())
.bind(filter.limit)
.bind(filter.offset)
.fetch_all(&self.pool)
.await
}
async fn count_filtered(&self, filter: &AuditFilter) -> Result<i64, sqlx::Error> {
let search_like = filter
.search
.as_ref()
.map(|s| format!("%{}%", s.replace('%', "\\%")));
let row: (i64,) = sqlx::query_as(
r#"
SELECT COUNT(*) FROM audit_logs
WHERE ($11 IS NULL OR actor_user_id = $21)
AND ($32 IS NULL OR action = $42)
AND ($53 IS NULL OR resource_type = $63)
AND ($74 IS NULL OR severity = $84)
AND ($95 IS NULL OR created_at >= $105)
AND ($116 IS NULL OR created_at <= $126)
AND (
$137 IS NULL
OR action LIKE $147 ESCAPE '\'
OR resource_type LIKE $157 ESCAPE '\'
OR resource_id LIKE $167 ESCAPE '\'
OR ip_address LIKE $177 ESCAPE '\'
OR metadata_json LIKE $187 ESCAPE '\'
)
"#,
)
.bind(filter.actor_user_id.as_deref())
.bind(filter.action.as_deref())
.bind(filter.resource_type.as_deref())
.bind(filter.severity.as_deref())
.bind(filter.since.as_deref())
.bind(filter.until.as_deref())
.bind(search_like.as_deref())
.fetch_one(&self.pool)
.await?;
Ok(row.0)
}
}
+62
View File
@@ -0,0 +1,62 @@
use crate::db::models::{Group, User};
use crate::db::repository::traits::GroupsRepository;
use async_trait::async_trait;
use sqlx::PgPool;
pub struct PostgresGroupsRepository {
pub pool: PgPool,
}
#[async_trait]
impl GroupsRepository for PostgresGroupsRepository {
async fn list(&self, _tenant_id: &str) -> Result<Vec<Group>, sqlx::Error> {
unimplemented!()
}
async fn find_by_id(&self, _id: &str) -> Result<Option<Group>, sqlx::Error> {
unimplemented!()
}
async fn create(
&self,
_id: &str,
_tenant_id: &str,
_name: &str,
_description: Option<&str>,
) -> Result<Group, sqlx::Error> {
unimplemented!()
}
async fn update(
&self,
_id: &str,
_name: &str,
_description: Option<&str>,
) -> Result<(), sqlx::Error> {
unimplemented!()
}
async fn delete(&self, _id: &str) -> Result<(), sqlx::Error> {
unimplemented!()
}
async fn count_members(&self, _group_id: &str) -> Result<i64, sqlx::Error> {
unimplemented!()
}
async fn list_members(&self, _group_id: &str) -> Result<Vec<User>, sqlx::Error> {
unimplemented!()
}
async fn add_member(&self, _group_id: &str, _user_id: &str) -> Result<(), sqlx::Error> {
unimplemented!()
}
async fn remove_member(&self, _group_id: &str, _user_id: &str) -> Result<(), sqlx::Error> {
unimplemented!()
}
async fn count(&self, _tenant_id: &str) -> Result<i64, sqlx::Error> {
unimplemented!()
}
}
+11
View File
@@ -0,0 +1,11 @@
pub mod applications;
pub mod audit;
pub mod groups;
pub mod permissions;
pub mod refresh_tokens;
pub mod roles;
pub mod service_accounts;
pub mod sessions;
pub mod tenants;
pub mod tokens;
pub mod users;
+125
View File
@@ -0,0 +1,125 @@
use crate::db::repository::traits::PermissionsRepository;
use async_trait::async_trait;
use sqlx::PgPool;
use crate::db::models::Permission;
pub struct PostgresPermissionsRepository {
pub pool: PgPool,
}
#[async_trait]
impl PermissionsRepository for PostgresPermissionsRepository {
/// List all permissions defined in the system.
async fn list_all(&self) -> Result<Vec<Permission>, sqlx::Error> {
sqlx::query_as::<_, Permission>("SELECT * FROM permissions ORDER BY name")
.fetch_all(&self.pool)
.await
}
/// List permissions assigned to a role.
async fn list_for_role(&self, role_id: &str) -> Result<Vec<Permission>, sqlx::Error> {
sqlx::query_as::<_, Permission>(
r#"
SELECT p.* FROM permissions p
JOIN role_permissions rp ON rp.permission_id = p.id
WHERE rp.role_id = $1
ORDER BY p.name
"#,
)
.bind(role_id)
.fetch_all(&self.pool)
.await
}
/// Assign a permission to a role (no-op if already assigned).
async fn assign_to_role(&self, role_id: &str, permission_id: &str) -> Result<(), sqlx::Error> {
sqlx::query(
"INSERT OR IGNORE INTO role_permissions (role_id, permission_id) VALUES ($1, $2)",
)
.bind(role_id)
.bind(permission_id)
.execute(&self.pool)
.await?;
Ok(())
}
/// Remove a permission from a role.
async fn remove_from_role(
&self,
role_id: &str,
permission_id: &str,
) -> Result<(), sqlx::Error> {
sqlx::query("DELETE FROM role_permissions WHERE role_id = $1 AND permission_id = $2")
.bind(role_id)
.bind(permission_id)
.execute(&self.pool)
.await?;
Ok(())
}
/// Clear all permissions for a role.
async fn clear_for_role(&self, role_id: &str) -> Result<(), sqlx::Error> {
sqlx::query("DELETE FROM role_permissions WHERE role_id = $1")
.bind(role_id)
.execute(&self.pool)
.await?;
Ok(())
}
/// Find a permission by name.
async fn find_by_name(&self, name: &str) -> Result<Option<Permission>, sqlx::Error> {
sqlx::query_as::<_, Permission>("SELECT * FROM permissions WHERE name = $1")
.bind(name)
.fetch_optional(&self.pool)
.await
}
/// Find a permission by id.
async fn find_by_id(&self, id: &str) -> Result<Option<Permission>, sqlx::Error> {
sqlx::query_as::<_, Permission>("SELECT * FROM permissions WHERE id = $1")
.bind(id)
.fetch_optional(&self.pool)
.await
}
/// Return all permission names held by a user (via their roles).
async fn list_for_user(&self, user_id: &str) -> Result<Vec<String>, sqlx::Error> {
let rows: Vec<(String,)> = sqlx::query_as(
r#"
SELECT DISTINCT p.name
FROM permissions p
JOIN role_permissions rp ON rp.permission_id = p.id
JOIN user_roles ur ON ur.role_id = rp.role_id
WHERE ur.user_id = $1
ORDER BY p.name
"#,
)
.bind(user_id)
.fetch_all(&self.pool)
.await?;
Ok(rows.into_iter().map(|(name,)| name).collect())
}
/// Check if a user holds a specific named permission.
async fn user_has_permission(
&self,
user_id: &str,
permission_name: &str,
) -> Result<bool, sqlx::Error> {
let row: (i64,) = sqlx::query_as(
r#"
SELECT COUNT(*)
FROM permissions p
JOIN role_permissions rp ON rp.permission_id = p.id
JOIN user_roles ur ON ur.role_id = rp.role_id
WHERE ur.user_id = $1 AND p.name = $2
"#,
)
.bind(user_id)
.bind(permission_name)
.fetch_one(&self.pool)
.await?;
Ok(row.0 > 0)
}
}
@@ -0,0 +1,59 @@
use crate::db::repository::traits::RefreshTokensRepository;
use async_trait::async_trait;
use sqlx::PgPool;
pub struct PostgresRefreshTokensRepository {
pub pool: PgPool,
}
use crate::db::repository::sqlite::refresh_tokens::RefreshToken;
#[async_trait]
impl RefreshTokensRepository for PostgresRefreshTokensRepository {
async fn create(
&self,
id: &str,
user_id: &str,
token_hash: &str,
expires_at: &str,
) -> Result<RefreshToken, sqlx::Error> {
sqlx::query_as::<_, RefreshToken>(
r#"
INSERT INTO refresh_tokens (id, user_id, token_hash, expires_at)
VALUES ($1, $2, $3, $4)
RETURNING *
"#,
)
.bind(id)
.bind(user_id)
.bind(token_hash)
.bind(expires_at)
.fetch_one(&self.pool)
.await
}
async fn find_by_hash(&self, token_hash: &str) -> Result<Option<RefreshToken>, sqlx::Error> {
sqlx::query_as::<_, RefreshToken>(
"SELECT * FROM refresh_tokens WHERE token_hash = $1 AND revoked = 0",
)
.bind(token_hash)
.fetch_optional(&self.pool)
.await
}
async fn revoke(&self, id: &str) -> Result<(), sqlx::Error> {
sqlx::query("UPDATE refresh_tokens SET revoked = 1 WHERE id = $1")
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
async fn revoke_all_for_user(&self, user_id: &str) -> Result<(), sqlx::Error> {
sqlx::query("UPDATE refresh_tokens SET revoked = 1 WHERE user_id = $1")
.bind(user_id)
.execute(&self.pool)
.await?;
Ok(())
}
}
+127
View File
@@ -0,0 +1,127 @@
use crate::db::repository::traits::RolesRepository;
use async_trait::async_trait;
use sqlx::PgPool;
use crate::db::models::Role;
pub struct PostgresRolesRepository {
pub pool: PgPool,
}
#[async_trait]
impl RolesRepository for PostgresRolesRepository {
async fn list_all(&self) -> Result<Vec<Role>, sqlx::Error> {
sqlx::query_as::<_, Role>("SELECT * FROM roles ORDER BY name")
.fetch_all(&self.pool)
.await
}
async fn find_by_name(&self, name: &str) -> Result<Option<Role>, sqlx::Error> {
sqlx::query_as::<_, Role>("SELECT * FROM roles WHERE name = $1")
.bind(name)
.fetch_optional(&self.pool)
.await
}
async fn find_by_id(&self, id: &str) -> Result<Option<Role>, sqlx::Error> {
sqlx::query_as::<_, Role>("SELECT * FROM roles WHERE id = $1")
.bind(id)
.fetch_optional(&self.pool)
.await
}
async fn list_for_user(&self, user_id: &str) -> Result<Vec<Role>, sqlx::Error> {
sqlx::query_as::<_, Role>(
r#"
SELECT r.* FROM roles r
JOIN user_roles ur ON ur.role_id = r.id
WHERE ur.user_id = $1
ORDER BY r.name
"#,
)
.bind(user_id)
.fetch_all(&self.pool)
.await
}
async fn assign_to_user(&self, user_id: &str, role_id: &str) -> Result<(), sqlx::Error> {
sqlx::query("INSERT OR IGNORE INTO user_roles (user_id, role_id) VALUES ($1, $2)")
.bind(user_id)
.bind(role_id)
.execute(&self.pool)
.await?;
Ok(())
}
async fn remove_from_user(&self, user_id: &str, role_id: &str) -> Result<(), sqlx::Error> {
sqlx::query("DELETE FROM user_roles WHERE user_id = $1 AND role_id = $2")
.bind(user_id)
.bind(role_id)
.execute(&self.pool)
.await?;
Ok(())
}
async fn admin_role_exists(&self) -> Result<bool, sqlx::Error> {
let row: (i64,) = sqlx::query_as("SELECT COUNT(*) FROM roles WHERE name = 'admin'")
.fetch_one(&self.pool)
.await?;
Ok(row.0 > 0)
}
/// Create a new role.
async fn create(
&self,
id: &str,
name: &str,
description: Option<&str>,
) -> Result<Role, sqlx::Error> {
sqlx::query_as::<_, Role>(
r#"
INSERT INTO roles (id, name, description)
VALUES ($1, $2, $3)
RETURNING *
"#,
)
.bind(id)
.bind(name)
.bind(description)
.fetch_one(&self.pool)
.await
}
/// Update role name/description.
async fn update(
&self,
id: &str,
name: &str,
description: Option<&str>,
) -> Result<(), sqlx::Error> {
sqlx::query("UPDATE roles SET name = $1, description = $2 WHERE id = $3")
.bind(name)
.bind(description)
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
/// Delete a role by id.
async fn delete(&self, id: &str) -> Result<(), sqlx::Error> {
sqlx::query("DELETE FROM roles WHERE id = $1")
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
/// List user ids that hold a given role.
async fn list_user_ids_for_role(&self, role_id: &str) -> Result<Vec<String>, sqlx::Error> {
let rows: Vec<(String,)> =
sqlx::query_as("SELECT user_id FROM user_roles WHERE role_id = $1 ORDER BY user_id")
.bind(role_id)
.fetch_all(&self.pool)
.await?;
Ok(rows.into_iter().map(|(id,)| id).collect())
}
}
@@ -0,0 +1,78 @@
use crate::db::repository::traits::ServiceAccountsRepository;
use async_trait::async_trait;
use sqlx::PgPool;
use crate::db::models::ServiceAccount;
pub struct PostgresServiceAccountsRepository {
pub pool: PgPool,
}
#[async_trait]
impl ServiceAccountsRepository for PostgresServiceAccountsRepository {
async fn create(
&self,
id: &str,
tenant_id: &str,
name: &str,
description: Option<&str>,
) -> Result<ServiceAccount, sqlx::Error> {
sqlx::query_as::<_, ServiceAccount>(
r#"
INSERT INTO service_accounts (id, tenant_id, name, description)
VALUES ($1, $2, $3, $4)
RETURNING *
"#,
)
.bind(id)
.bind(tenant_id)
.bind(name)
.bind(description)
.fetch_one(&self.pool)
.await
}
async fn find_by_id(&self, id: &str) -> Result<Option<ServiceAccount>, sqlx::Error> {
sqlx::query_as::<_, ServiceAccount>("SELECT * FROM service_accounts WHERE id = $1")
.bind(id)
.fetch_optional(&self.pool)
.await
}
async fn list(&self, tenant_id: &str) -> Result<Vec<ServiceAccount>, sqlx::Error> {
sqlx::query_as::<_, ServiceAccount>(
"SELECT * FROM service_accounts WHERE tenant_id = $1 ORDER BY name",
)
.bind(tenant_id)
.fetch_all(&self.pool)
.await
}
async fn set_enabled(&self, id: &str, enabled: bool) -> Result<(), sqlx::Error> {
sqlx::query(
"UPDATE service_accounts SET enabled = $1, updated_at = strftime('%Y-%m-%dT%H:%M:%SZ', 'now') WHERE id = $2",
)
.bind(enabled)
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
async fn delete(&self, id: &str) -> Result<(), sqlx::Error> {
sqlx::query("DELETE FROM service_accounts WHERE id = $1")
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
async fn count(&self, tenant_id: &str) -> Result<i64, sqlx::Error> {
let row: (i64,) =
sqlx::query_as("SELECT COUNT(*) FROM service_accounts WHERE tenant_id = $1")
.bind(tenant_id)
.fetch_one(&self.pool)
.await?;
Ok(row.0)
}
}
+123
View File
@@ -0,0 +1,123 @@
use crate::db::repository::traits::SessionsRepository;
use async_trait::async_trait;
use sqlx::PgPool;
use crate::db::models::Session;
pub struct PostgresSessionsRepository {
pub pool: PgPool,
}
#[async_trait]
impl SessionsRepository for PostgresSessionsRepository {
async fn create(
&self,
id: &str,
user_id: &str,
token_hash: &str,
ip_address: Option<&str>,
user_agent: Option<&str>,
expires_at: &str,
) -> Result<Session, sqlx::Error> {
sqlx::query_as::<_, Session>(
r#"
INSERT INTO sessions (id, user_id, token_hash, ip_address, user_agent, expires_at)
VALUES ($1, $2, $3, $4, $5, $6)
RETURNING *
"#,
)
.bind(id)
.bind(user_id)
.bind(token_hash)
.bind(ip_address)
.bind(user_agent)
.bind(expires_at)
.fetch_one(&self.pool)
.await
}
async fn find_by_token_hash(&self, token_hash: &str) -> Result<Option<Session>, sqlx::Error> {
sqlx::query_as::<_, Session>("SELECT * FROM sessions WHERE token_hash = $1 AND revoked = 0")
.bind(token_hash)
.fetch_optional(&self.pool)
.await
}
async fn revoke(&self, id: &str) -> Result<(), sqlx::Error> {
sqlx::query("UPDATE sessions SET revoked = 1 WHERE id = $1")
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
async fn revoke_all_for_user(&self, user_id: &str) -> Result<(), sqlx::Error> {
sqlx::query("UPDATE sessions SET revoked = 1 WHERE user_id = $1")
.bind(user_id)
.execute(&self.pool)
.await?;
Ok(())
}
async fn update_last_seen(&self, id: &str) -> Result<(), sqlx::Error> {
sqlx::query(
"UPDATE sessions SET last_seen_at = strftime('%Y-%m-%dT%H:%M:%SZ', 'now') WHERE id = $1",
)
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
/// List active (non-revoked, non-expired) sessions for a user.
async fn list_active_for_user(&self, user_id: &str) -> Result<Vec<Session>, sqlx::Error> {
sqlx::query_as::<_, Session>(
r#"
SELECT * FROM sessions
WHERE user_id = $1
AND revoked = 0
AND expires_at >= strftime('%Y-%m-%dT%H:%M:%SZ', 'now')
ORDER BY last_seen_at DESC
"#,
)
.bind(user_id)
.fetch_all(&self.pool)
.await
}
async fn list_all_active(&self) -> Result<Vec<Session>, sqlx::Error> {
unimplemented!()
}
/// Count active sessions system-wide.
async fn count_active(&self) -> Result<i64, sqlx::Error> {
let row: (i64,) = sqlx::query_as(
r#"
SELECT COUNT(*) FROM sessions
WHERE revoked = 0
AND expires_at >= strftime('%Y-%m-%dT%H:%M:%SZ', 'now')
"#,
)
.fetch_one(&self.pool)
.await?;
Ok(row.0)
}
/// Delete sessions that are expired or revoked. Called once at startup.
async fn cleanup_expired(&self) -> Result<u64, sqlx::Error> {
let result = sqlx::query(
r#"
DELETE FROM sessions
WHERE revoked = 1
OR expires_at < strftime('%Y-%m-%dT%H:%M:%SZ', 'now')
"#,
)
.execute(&self.pool)
.await?;
Ok(result.rows_affected())
}
async fn revoke_others(&self, _user_id: &str, _except_id: &str) -> Result<u64, sqlx::Error> {
unimplemented!()
}
}
+108
View File
@@ -0,0 +1,108 @@
use sqlx::PgPool;
use crate::db::models::Tenant;
use crate::db::repository::traits::TenantsRepository;
pub struct PostgresTenantsRepository {
pub pool: PgPool,
}
#[async_trait::async_trait]
impl TenantsRepository for PostgresTenantsRepository {
async fn find_by_id(&self, id: &str) -> Result<Option<Tenant>, sqlx::Error> {
let row = sqlx::query_as::<_, Tenant>(
"SELECT id, name, slug, enabled, created_at::text, updated_at::text FROM tenants WHERE id = $1",
)
.bind(id)
.fetch_optional(&self.pool)
.await?;
Ok(row)
}
async fn find_by_slug(&self, slug: &str) -> Result<Option<Tenant>, sqlx::Error> {
let row = sqlx::query_as::<_, Tenant>(
"SELECT id, name, slug, enabled, created_at::text, updated_at::text FROM tenants WHERE slug = $1",
)
.bind(slug)
.fetch_optional(&self.pool)
.await?;
Ok(row)
}
async fn list(&self) -> Result<Vec<Tenant>, sqlx::Error> {
let rows = sqlx::query_as::<_, Tenant>(
"SELECT id, name, slug, enabled, created_at::text, updated_at::text FROM tenants ORDER BY name ASC",
)
.fetch_all(&self.pool)
.await?;
Ok(rows)
}
async fn create(
&self,
id: &str,
name: &str,
slug: Option<&str>,
) -> Result<Tenant, sqlx::Error> {
let slug = slug.unwrap_or(id);
let row = sqlx::query_as::<_, Tenant>(
r#"
INSERT INTO tenants (id, name, slug, enabled)
VALUES ($1, $2, $3, true)
RETURNING id, name, slug, enabled, created_at::text, updated_at::text
"#,
)
.bind(id)
.bind(name)
.bind(slug)
.fetch_one(&self.pool)
.await?;
Ok(row)
}
async fn update(&self, id: &str, name: &str, slug: Option<&str>) -> Result<(), sqlx::Error> {
let slug = slug.unwrap_or(name);
sqlx::query(
r#"
UPDATE tenants
SET name = $1, slug = $2, updated_at = CURRENT_TIMESTAMP
WHERE id = $3
"#,
)
.bind(name)
.bind(slug)
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
async fn set_enabled(&self, id: &str, enabled: bool) -> Result<(), sqlx::Error> {
sqlx::query(
r#"
UPDATE tenants
SET enabled = $1, updated_at = CURRENT_TIMESTAMP
WHERE id = $2
"#,
)
.bind(enabled)
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
async fn delete(&self, id: &str) -> Result<(), sqlx::Error> {
sqlx::query("DELETE FROM tenants WHERE id = $1")
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
}
+79
View File
@@ -0,0 +1,79 @@
use crate::db::repository::traits::TokensRepository;
use async_trait::async_trait;
use sqlx::PgPool;
use crate::db::models::ApiToken;
pub struct PostgresTokensRepository {
pub pool: PgPool,
}
#[async_trait]
impl TokensRepository for PostgresTokensRepository {
async fn create(
&self,
id: &str,
user_id: &str,
name: &str,
token_hash: &str,
expires_at: Option<&str>,
) -> Result<ApiToken, sqlx::Error> {
sqlx::query_as::<_, ApiToken>(
r#"
INSERT INTO api_tokens (id, user_id, name, token_hash, expires_at)
VALUES ($1, $2, $3, $4, $5)
RETURNING *
"#,
)
.bind(id)
.bind(user_id)
.bind(name)
.bind(token_hash)
.bind(expires_at)
.fetch_one(&self.pool)
.await
}
async fn find_by_hash(&self, token_hash: &str) -> Result<Option<ApiToken>, sqlx::Error> {
sqlx::query_as::<_, ApiToken>(
"SELECT * FROM api_tokens WHERE token_hash = $1 AND revoked = 0",
)
.bind(token_hash)
.fetch_optional(&self.pool)
.await
}
async fn list_for_user(&self, user_id: &str) -> Result<Vec<ApiToken>, sqlx::Error> {
sqlx::query_as::<_, ApiToken>(
"SELECT * FROM api_tokens WHERE user_id = $1 ORDER BY created_at DESC",
)
.bind(user_id)
.fetch_all(&self.pool)
.await
}
async fn find_by_id(&self, id: &str) -> Result<Option<ApiToken>, sqlx::Error> {
sqlx::query_as::<_, ApiToken>("SELECT * FROM api_tokens WHERE id = $1")
.bind(id)
.fetch_optional(&self.pool)
.await
}
async fn revoke(&self, id: &str) -> Result<(), sqlx::Error> {
sqlx::query("UPDATE api_tokens SET revoked = 1 WHERE id = $1")
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
async fn update_last_used(&self, id: &str) -> Result<(), sqlx::Error> {
sqlx::query(
"UPDATE api_tokens SET last_used_at = strftime('%Y-%m-%dT%H:%M:%SZ', 'now') WHERE id = $1",
)
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
}
+163
View File
@@ -0,0 +1,163 @@
use crate::db::repository::traits::UsersRepository;
use async_trait::async_trait;
use sqlx::PgPool;
use crate::db::models::User;
pub struct PostgresUsersRepository {
pub pool: PgPool,
}
use crate::db::repository::sqlite::users::UserProfile;
#[async_trait]
impl UsersRepository for PostgresUsersRepository {
async fn count_admins(&self) -> Result<i64, sqlx::Error> {
let row: (i64,) = sqlx::query_as(
r#"
SELECT COUNT(DISTINCT ur.user_id)
FROM user_roles ur
JOIN roles r ON r.id = ur.role_id
WHERE r.name = 'admin'
"#,
)
.fetch_one(&self.pool)
.await?;
Ok(row.0)
}
async fn count(&self, tenant_id: &str) -> Result<i64, sqlx::Error> {
let row: (i64,) = sqlx::query_as("SELECT COUNT(*) FROM users WHERE tenant_id = $1")
.bind(tenant_id)
.fetch_one(&self.pool)
.await?;
Ok(row.0)
}
async fn count_by_status(&self, tenant_id: &str, status: i32) -> Result<i64, sqlx::Error> {
let row: (i64,) =
sqlx::query_as("SELECT COUNT(*) FROM users WHERE tenant_id = $1 AND status = $2")
.bind(tenant_id)
.bind(status)
.fetch_one(&self.pool)
.await?;
Ok(row.0)
}
async fn find_by_id(&self, id: &str) -> Result<Option<User>, sqlx::Error> {
sqlx::query_as::<_, User>("SELECT * FROM users WHERE id = $1")
.bind(id)
.fetch_optional(&self.pool)
.await
}
async fn find_by_username(&self, username: &str) -> Result<Option<User>, sqlx::Error> {
sqlx::query_as::<_, User>("SELECT * FROM users WHERE username = $1")
.bind(username)
.fetch_optional(&self.pool)
.await
}
async fn list(&self, tenant_id: &str) -> Result<Vec<User>, sqlx::Error> {
sqlx::query_as::<_, User>(
"SELECT * FROM users WHERE tenant_id = $1 ORDER BY created_at DESC",
)
.bind(tenant_id)
.fetch_all(&self.pool)
.await
}
async fn create(
&self,
id: &str,
tenant_id: &str,
username: &str,
password_hash: &str,
) -> Result<User, sqlx::Error> {
sqlx::query_as::<_, User>(
r#"
INSERT INTO users (id, tenant_id, username, password_hash, status)
VALUES ($1, $2, $3, $4, 1)
RETURNING *
"#,
)
.bind(id)
.bind(tenant_id)
.bind(username)
.bind(password_hash)
.fetch_one(&self.pool)
.await
}
async fn update_status(&self, id: &str, status: i32) -> Result<(), sqlx::Error> {
sqlx::query(
"UPDATE users SET status = $1, updated_at = strftime('%Y-%m-%dT%H:%M:%SZ', 'now') WHERE id = $2",
)
.bind(status)
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
async fn update_password_hash(&self, id: &str, password_hash: &str) -> Result<(), sqlx::Error> {
sqlx::query(
"UPDATE users SET password_hash = $1, updated_at = strftime('%Y-%m-%dT%H:%M:%SZ', 'now') WHERE id = $2",
)
.bind(password_hash)
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
async fn set_last_login(&self, id: &str) -> Result<(), sqlx::Error> {
sqlx::query(
"UPDATE users SET last_login_at = strftime('%Y-%m-%dT%H:%M:%SZ', 'now'), updated_at = strftime('%Y-%m-%dT%H:%M:%SZ', 'now') WHERE id = $1",
)
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
async fn username_exists(&self, tenant_id: &str, username: &str) -> Result<bool, sqlx::Error> {
let row: (i64,) =
sqlx::query_as("SELECT COUNT(*) FROM users WHERE tenant_id = $1 AND username = $2")
.bind(tenant_id)
.bind(username)
.fetch_one(&self.pool)
.await?;
Ok(row.0 > 0)
}
async fn get_profile(&self, user_id: &str) -> Result<Option<UserProfile>, sqlx::Error> {
sqlx::query_as::<_, UserProfile>("SELECT * FROM user_profiles WHERE user_id = $1")
.bind(user_id)
.fetch_optional(&self.pool)
.await
}
async fn upsert_profile(
&self,
user_id: &str,
email: Option<&str>,
full_name: Option<&str>,
) -> Result<UserProfile, sqlx::Error> {
sqlx::query_as::<_, UserProfile>(
r#"
INSERT INTO user_profiles (user_id, email, full_name)
VALUES ($1, $2, $3)
ON CONFLICT(user_id) DO UPDATE SET
email = excluded.email,
full_name = excluded.full_name
RETURNING *
"#,
)
.bind(user_id)
.bind(email)
.bind(full_name)
.fetch_one(&self.pool)
.await
}
}
+1 -70
View File
@@ -1,70 +1 @@
use sqlx::SqlitePool;
use crate::db::models::Role;
pub async fn list_all(pool: &SqlitePool) -> Result<Vec<Role>, sqlx::Error> {
sqlx::query_as::<_, Role>("SELECT * FROM roles ORDER BY name")
.fetch_all(pool)
.await
}
pub async fn find_by_name(pool: &SqlitePool, name: &str) -> Result<Option<Role>, sqlx::Error> {
sqlx::query_as::<_, Role>("SELECT * FROM roles WHERE name = ?")
.bind(name)
.fetch_optional(pool)
.await
}
pub async fn find_by_id(pool: &SqlitePool, id: &str) -> Result<Option<Role>, sqlx::Error> {
sqlx::query_as::<_, Role>("SELECT * FROM roles WHERE id = ?")
.bind(id)
.fetch_optional(pool)
.await
}
pub async fn list_for_user(pool: &SqlitePool, user_id: &str) -> Result<Vec<Role>, sqlx::Error> {
sqlx::query_as::<_, Role>(
r#"
SELECT r.* FROM roles r
JOIN user_roles ur ON ur.role_id = r.id
WHERE ur.user_id = ?
ORDER BY r.name
"#,
)
.bind(user_id)
.fetch_all(pool)
.await
}
pub async fn assign_to_user(
tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>,
user_id: &str,
role_id: &str,
) -> Result<(), sqlx::Error> {
sqlx::query("INSERT OR IGNORE INTO user_roles (user_id, role_id) VALUES (?, ?)")
.bind(user_id)
.bind(role_id)
.execute(&mut **tx)
.await?;
Ok(())
}
pub async fn remove_from_user(
tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>,
user_id: &str,
role_id: &str,
) -> Result<(), sqlx::Error> {
sqlx::query("DELETE FROM user_roles WHERE user_id = ? AND role_id = ?")
.bind(user_id)
.bind(role_id)
.execute(&mut **tx)
.await?;
Ok(())
}
pub async fn admin_role_exists(pool: &SqlitePool) -> Result<bool, sqlx::Error> {
let row: (i64,) = sqlx::query_as("SELECT COUNT(*) FROM roles WHERE name = 'admin'")
.fetch_one(pool)
.await?;
Ok(row.0 > 0)
}
pub use crate::db::repository::sqlite::roles::*;
-59
View File
@@ -1,59 +0,0 @@
use sqlx::SqlitePool;
use crate::db::models::ServiceAccount;
pub async fn create(
tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>,
id: &str,
tenant_id: &str,
name: &str,
description: Option<&str>,
) -> Result<ServiceAccount, sqlx::Error> {
sqlx::query_as::<_, ServiceAccount>(
r#"
INSERT INTO service_accounts (id, tenant_id, name, description)
VALUES (?, ?, ?, ?)
RETURNING *
"#,
)
.bind(id)
.bind(tenant_id)
.bind(name)
.bind(description)
.fetch_one(&mut **tx)
.await
}
pub async fn find_by_id(
pool: &SqlitePool,
id: &str,
) -> Result<Option<ServiceAccount>, sqlx::Error> {
sqlx::query_as::<_, ServiceAccount>("SELECT * FROM service_accounts WHERE id = ?")
.bind(id)
.fetch_optional(pool)
.await
}
pub async fn list(pool: &SqlitePool, tenant_id: &str) -> Result<Vec<ServiceAccount>, sqlx::Error> {
sqlx::query_as::<_, ServiceAccount>(
"SELECT * FROM service_accounts WHERE tenant_id = ? ORDER BY name",
)
.bind(tenant_id)
.fetch_all(pool)
.await
}
pub async fn set_enabled(
tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>,
id: &str,
enabled: bool,
) -> Result<(), sqlx::Error> {
sqlx::query(
"UPDATE service_accounts SET enabled = ?, updated_at = strftime('%Y-%m-%dT%H:%M:%SZ', 'now') WHERE id = ?",
)
.bind(enabled)
.bind(id)
.execute(&mut **tx)
.await?;
Ok(())
}
-79
View File
@@ -1,79 +0,0 @@
use sqlx::SqlitePool;
use crate::db::models::Session;
pub async fn create(
pool: &SqlitePool,
id: &str,
user_id: &str,
token_hash: &str,
ip_address: Option<&str>,
user_agent: Option<&str>,
expires_at: &str,
) -> Result<Session, sqlx::Error> {
sqlx::query_as::<_, Session>(
r#"
INSERT INTO sessions (id, user_id, token_hash, ip_address, user_agent, expires_at)
VALUES (?, ?, ?, ?, ?, ?)
RETURNING *
"#,
)
.bind(id)
.bind(user_id)
.bind(token_hash)
.bind(ip_address)
.bind(user_agent)
.bind(expires_at)
.fetch_one(pool)
.await
}
pub async fn find_by_token_hash(
pool: &SqlitePool,
token_hash: &str,
) -> Result<Option<Session>, sqlx::Error> {
sqlx::query_as::<_, Session>("SELECT * FROM sessions WHERE token_hash = ? AND revoked = 0")
.bind(token_hash)
.fetch_optional(pool)
.await
}
pub async fn revoke(pool: &SqlitePool, id: &str) -> Result<(), sqlx::Error> {
sqlx::query("UPDATE sessions SET revoked = 1 WHERE id = ?")
.bind(id)
.execute(pool)
.await?;
Ok(())
}
pub async fn revoke_all_for_user(pool: &SqlitePool, user_id: &str) -> Result<(), sqlx::Error> {
sqlx::query("UPDATE sessions SET revoked = 1 WHERE user_id = ?")
.bind(user_id)
.execute(pool)
.await?;
Ok(())
}
pub async fn update_last_seen(pool: &SqlitePool, id: &str) -> Result<(), sqlx::Error> {
sqlx::query(
"UPDATE sessions SET last_seen_at = strftime('%Y-%m-%dT%H:%M:%SZ', 'now') WHERE id = ?",
)
.bind(id)
.execute(pool)
.await?;
Ok(())
}
/// Delete sessions that are expired or revoked. Called once at startup.
pub async fn cleanup_expired(pool: &SqlitePool) -> Result<u64, sqlx::Error> {
let result = sqlx::query(
r#"
DELETE FROM sessions
WHERE revoked = 1
OR expires_at < strftime('%Y-%m-%dT%H:%M:%SZ', 'now')
"#,
)
.execute(pool)
.await?;
Ok(result.rows_affected())
}
+108
View File
@@ -0,0 +1,108 @@
use crate::db::repository::traits::ApplicationsRepository;
use async_trait::async_trait;
use sqlx::SqlitePool;
use crate::db::models::Application;
pub struct SqliteApplicationsRepository {
pub pool: SqlitePool,
}
#[async_trait]
impl ApplicationsRepository for SqliteApplicationsRepository {
async fn create(
&self,
id: &str,
tenant_id: &str,
name: &str,
slug: &str,
) -> Result<Application, sqlx::Error> {
sqlx::query_as::<_, Application>(
r#"
INSERT INTO applications (id, tenant_id, name, slug)
VALUES (?, ?, ?, ?)
RETURNING id, tenant_id, name, slug, enabled, created_at, updated_at, NULL as description, NULL as client_secret_hash, NULL as redirect_uris
"#,
)
.bind(id)
.bind(tenant_id)
.bind(name)
.bind(slug)
.fetch_one(&self.pool)
.await
}
async fn find_by_slug(&self, slug: &str) -> Result<Option<Application>, sqlx::Error> {
sqlx::query_as::<_, Application>("SELECT id, tenant_id, name, slug, enabled, created_at, updated_at, NULL as description, NULL as client_secret_hash, NULL as redirect_uris FROM applications WHERE slug = ?")
.bind(slug)
.fetch_optional(&self.pool)
.await
}
async fn find_by_id(&self, id: &str) -> Result<Option<Application>, sqlx::Error> {
sqlx::query_as::<_, Application>("SELECT id, tenant_id, name, slug, enabled, created_at, updated_at, NULL as description, NULL as client_secret_hash, NULL as redirect_uris FROM applications WHERE id = ?")
.bind(id)
.fetch_optional(&self.pool)
.await
}
async fn list(&self, tenant_id: &str) -> Result<Vec<Application>, sqlx::Error> {
sqlx::query_as::<_, Application>(
"SELECT id, tenant_id, name, slug, enabled, created_at, updated_at, NULL as description, NULL as client_secret_hash, NULL as redirect_uris FROM applications WHERE tenant_id = ? ORDER BY name",
)
.bind(tenant_id)
.fetch_all(&self.pool)
.await
}
async fn set_enabled(&self, id: &str, enabled: bool) -> Result<(), sqlx::Error> {
sqlx::query(
"UPDATE applications SET enabled = ?, updated_at = strftime('%Y-%m-%dT%H:%M:%SZ', 'now') WHERE id = ?",
)
.bind(enabled)
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
async fn update(
&self,
id: &str,
name: &str,
slug: &str,
enabled: bool,
) -> Result<(), sqlx::Error> {
sqlx::query(
r#"
UPDATE applications
SET name = ?, slug = ?, enabled = ?,
updated_at = strftime('%Y-%m-%dT%H:%M:%SZ', 'now')
WHERE id = ?
"#,
)
.bind(name)
.bind(slug)
.bind(enabled)
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
async fn delete(&self, id: &str) -> Result<(), sqlx::Error> {
sqlx::query("DELETE FROM applications WHERE id = ?")
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
async fn count(&self, tenant_id: &str) -> Result<i64, sqlx::Error> {
let row: (i64,) = sqlx::query_as("SELECT COUNT(*) FROM applications WHERE tenant_id = ?")
.bind(tenant_id)
.fetch_one(&self.pool)
.await?;
Ok(row.0)
}
}
+159
View File
@@ -0,0 +1,159 @@
use crate::db::repository::traits::AuditRepository;
use async_trait::async_trait;
use sqlx::SqlitePool;
use crate::db::models::AuditLog;
pub struct SqliteAuditRepository {
pub pool: SqlitePool,
}
/// Filtered audit log query. All filters are optional.
#[derive(Debug, Default)]
pub struct AuditFilter {
pub actor_user_id: Option<String>,
pub action: Option<String>,
pub resource_type: Option<String>,
pub severity: Option<String>,
pub since: Option<String>,
pub until: Option<String>,
pub search: Option<String>,
pub limit: i64,
pub offset: i64,
}
#[async_trait]
impl AuditRepository for SqliteAuditRepository {
/// Count all audit log entries.
async fn count(&self) -> Result<i64, sqlx::Error> {
let row: (i64,) = sqlx::query_as("SELECT COUNT(*) FROM audit_logs")
.fetch_one(&self.pool)
.await?;
Ok(row.0)
}
#[allow(clippy::too_many_arguments)]
#[allow(clippy::too_many_arguments)]
async fn insert(
&self,
id: &str,
actor_user_id: Option<&str>,
target_user_id: Option<&str>,
action: &str,
resource_type: &str,
resource_id: Option<&str>,
severity: &str,
ip_address: Option<&str>,
user_agent: Option<&str>,
metadata_json: Option<&str>,
) -> Result<AuditLog, sqlx::Error> {
sqlx::query_as::<_, AuditLog>(
r#"
INSERT INTO audit_logs (
id, actor_user_id, target_user_id,
action, resource_type, resource_id,
severity, ip_address, user_agent, metadata_json
)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
RETURNING *
"#,
)
.bind(id)
.bind(actor_user_id)
.bind(target_user_id)
.bind(action)
.bind(resource_type)
.bind(resource_id)
.bind(severity)
.bind(ip_address)
.bind(user_agent)
.bind(metadata_json)
.fetch_one(&self.pool)
.await
}
async fn list_recent(&self, limit: i64) -> Result<Vec<AuditLog>, sqlx::Error> {
sqlx::query_as::<_, AuditLog>("SELECT * FROM audit_logs ORDER BY created_at DESC LIMIT ?")
.bind(limit)
.fetch_all(&self.pool)
.await
}
async fn list_filtered(&self, filter: &AuditFilter) -> Result<Vec<AuditLog>, sqlx::Error> {
// Build a dynamic but simple filter using COALESCE-style optional matches.
// Empty optionals are treated as wildcards via OR IS NULL pattern with bind of None.
let search_like = filter
.search
.as_ref()
.map(|s| format!("%{}%", s.replace('%', "\\%")));
sqlx::query_as::<_, AuditLog>(
r#"
SELECT * FROM audit_logs
WHERE (?1 IS NULL OR actor_user_id = ?1)
AND (?2 IS NULL OR action = ?2)
AND (?3 IS NULL OR resource_type = ?3)
AND (?4 IS NULL OR severity = ?4)
AND (?5 IS NULL OR created_at >= ?5)
AND (?6 IS NULL OR created_at <= ?6)
AND (
?7 IS NULL
OR action LIKE ?7 ESCAPE '\'
OR resource_type LIKE ?7 ESCAPE '\'
OR resource_id LIKE ?7 ESCAPE '\'
OR ip_address LIKE ?7 ESCAPE '\'
OR metadata_json LIKE ?7 ESCAPE '\'
)
ORDER BY created_at DESC
LIMIT ?8 OFFSET ?9
"#,
)
.bind(filter.actor_user_id.as_deref())
.bind(filter.action.as_deref())
.bind(filter.resource_type.as_deref())
.bind(filter.severity.as_deref())
.bind(filter.since.as_deref())
.bind(filter.until.as_deref())
.bind(search_like.as_deref())
.bind(filter.limit)
.bind(filter.offset)
.fetch_all(&self.pool)
.await
}
async fn count_filtered(&self, filter: &AuditFilter) -> Result<i64, sqlx::Error> {
let search_like = filter
.search
.as_ref()
.map(|s| format!("%{}%", s.replace('%', "\\%")));
let row: (i64,) = sqlx::query_as(
r#"
SELECT COUNT(*) FROM audit_logs
WHERE (?1 IS NULL OR actor_user_id = ?1)
AND (?2 IS NULL OR action = ?2)
AND (?3 IS NULL OR resource_type = ?3)
AND (?4 IS NULL OR severity = ?4)
AND (?5 IS NULL OR created_at >= ?5)
AND (?6 IS NULL OR created_at <= ?6)
AND (
?7 IS NULL
OR action LIKE ?7 ESCAPE '\'
OR resource_type LIKE ?7 ESCAPE '\'
OR resource_id LIKE ?7 ESCAPE '\'
OR ip_address LIKE ?7 ESCAPE '\'
OR metadata_json LIKE ?7 ESCAPE '\'
)
"#,
)
.bind(filter.actor_user_id.as_deref())
.bind(filter.action.as_deref())
.bind(filter.resource_type.as_deref())
.bind(filter.severity.as_deref())
.bind(filter.since.as_deref())
.bind(filter.until.as_deref())
.bind(search_like.as_deref())
.fetch_one(&self.pool)
.await?;
Ok(row.0)
}
}
+139
View File
@@ -0,0 +1,139 @@
use crate::db::models::{Group, User};
use crate::db::repository::traits::GroupsRepository;
use async_trait::async_trait;
use sqlx::SqlitePool;
pub struct SqliteGroupsRepository {
pub pool: SqlitePool,
}
#[async_trait]
impl GroupsRepository for SqliteGroupsRepository {
async fn list(&self, tenant_id: &str) -> Result<Vec<Group>, sqlx::Error> {
let rows = sqlx::query_as::<_, Group>(
r#"
SELECT * FROM groups
WHERE tenant_id = ?
ORDER BY name ASC
"#,
)
.bind(tenant_id)
.fetch_all(&self.pool)
.await?;
Ok(rows)
}
async fn find_by_id(&self, id: &str) -> Result<Option<Group>, sqlx::Error> {
sqlx::query_as::<_, Group>("SELECT * FROM groups WHERE id = ?")
.bind(id)
.fetch_optional(&self.pool)
.await
}
async fn create(
&self,
id: &str,
tenant_id: &str,
name: &str,
description: Option<&str>,
) -> Result<Group, sqlx::Error> {
sqlx::query_as::<_, Group>(
r#"
INSERT INTO groups (id, tenant_id, name, description)
VALUES (?, ?, ?, ?)
RETURNING *
"#,
)
.bind(id)
.bind(tenant_id)
.bind(name)
.bind(description)
.fetch_one(&self.pool)
.await
}
async fn update(
&self,
id: &str,
name: &str,
description: Option<&str>,
) -> Result<(), sqlx::Error> {
sqlx::query(
r#"
UPDATE groups
SET name = ?, description = ?, updated_at = strftime('%Y-%m-%dT%H:%M:%SZ', 'now')
WHERE id = ?
"#,
)
.bind(name)
.bind(description)
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
async fn delete(&self, id: &str) -> Result<(), sqlx::Error> {
sqlx::query("DELETE FROM groups WHERE id = ?")
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
async fn count_members(&self, group_id: &str) -> Result<i64, sqlx::Error> {
let row: (i64,) = sqlx::query_as("SELECT COUNT(*) FROM user_groups WHERE group_id = ?")
.bind(group_id)
.fetch_one(&self.pool)
.await?;
Ok(row.0)
}
async fn list_members(&self, group_id: &str) -> Result<Vec<User>, sqlx::Error> {
sqlx::query_as::<_, User>(
r#"
SELECT u.*
FROM users u
JOIN user_groups ug ON u.id = ug.user_id
WHERE ug.group_id = ?
ORDER BY u.username ASC
"#,
)
.bind(group_id)
.fetch_all(&self.pool)
.await
}
async fn add_member(&self, group_id: &str, user_id: &str) -> Result<(), sqlx::Error> {
sqlx::query(
r#"
INSERT INTO user_groups (user_id, group_id)
VALUES (?, ?)
ON CONFLICT (user_id, group_id) DO NOTHING
"#,
)
.bind(user_id)
.bind(group_id)
.execute(&self.pool)
.await?;
Ok(())
}
async fn remove_member(&self, group_id: &str, user_id: &str) -> Result<(), sqlx::Error> {
sqlx::query("DELETE FROM user_groups WHERE user_id = ? AND group_id = ?")
.bind(user_id)
.bind(group_id)
.execute(&self.pool)
.await?;
Ok(())
}
async fn count(&self, tenant_id: &str) -> Result<i64, sqlx::Error> {
let row: (i64,) = sqlx::query_as("SELECT COUNT(*) FROM groups WHERE tenant_id = ?")
.bind(tenant_id)
.fetch_one(&self.pool)
.await?;
Ok(row.0)
}
}
+11
View File
@@ -0,0 +1,11 @@
pub mod applications;
pub mod audit;
pub mod groups;
pub mod permissions;
pub mod refresh_tokens;
pub mod roles;
pub mod service_accounts;
pub mod sessions;
pub mod tenants;
pub mod tokens;
pub mod users;
+122
View File
@@ -0,0 +1,122 @@
use crate::db::repository::traits::PermissionsRepository;
use async_trait::async_trait;
use sqlx::SqlitePool;
use crate::db::models::Permission;
pub struct SqlitePermissionsRepository {
pub pool: SqlitePool,
}
#[async_trait]
impl PermissionsRepository for SqlitePermissionsRepository {
/// List all permissions defined in the system.
async fn list_all(&self) -> Result<Vec<Permission>, sqlx::Error> {
sqlx::query_as::<_, Permission>("SELECT * FROM permissions ORDER BY name")
.fetch_all(&self.pool)
.await
}
/// List permissions assigned to a role.
async fn list_for_role(&self, role_id: &str) -> Result<Vec<Permission>, sqlx::Error> {
sqlx::query_as::<_, Permission>(
r#"
SELECT p.* FROM permissions p
JOIN role_permissions rp ON rp.permission_id = p.id
WHERE rp.role_id = ?
ORDER BY p.name
"#,
)
.bind(role_id)
.fetch_all(&self.pool)
.await
}
/// Assign a permission to a role (no-op if already assigned).
async fn assign_to_role(&self, role_id: &str, permission_id: &str) -> Result<(), sqlx::Error> {
sqlx::query(
"INSERT OR IGNORE INTO role_permissions (role_id, permission_id) VALUES (?, ?)",
)
.bind(role_id)
.bind(permission_id)
.execute(&self.pool)
.await?;
Ok(())
}
/// Remove a permission from a role.
async fn remove_from_role(
&self,
role_id: &str,
permission_id: &str,
) -> Result<(), sqlx::Error> {
sqlx::query("DELETE FROM role_permissions WHERE role_id = ? AND permission_id = ?")
.bind(role_id)
.bind(permission_id)
.execute(&self.pool)
.await?;
Ok(())
}
async fn clear_for_role(&self, role_id: &str) -> Result<(), sqlx::Error> {
sqlx::query("DELETE FROM role_permissions WHERE role_id = ?")
.bind(role_id)
.execute(&self.pool)
.await?;
Ok(())
}
async fn find_by_name(&self, name: &str) -> Result<Option<Permission>, sqlx::Error> {
sqlx::query_as::<_, Permission>("SELECT * FROM permissions WHERE name = ?")
.bind(name)
.fetch_optional(&self.pool)
.await
}
/// Find a permission by id.
async fn find_by_id(&self, id: &str) -> Result<Option<Permission>, sqlx::Error> {
sqlx::query_as::<_, Permission>("SELECT * FROM permissions WHERE id = ?")
.bind(id)
.fetch_optional(&self.pool)
.await
}
/// Return all permission names held by a user (via their roles).
async fn list_for_user(&self, user_id: &str) -> Result<Vec<String>, sqlx::Error> {
let rows: Vec<(String,)> = sqlx::query_as(
r#"
SELECT DISTINCT p.name
FROM permissions p
JOIN role_permissions rp ON rp.permission_id = p.id
JOIN user_roles ur ON ur.role_id = rp.role_id
WHERE ur.user_id = ?
ORDER BY p.name
"#,
)
.bind(user_id)
.fetch_all(&self.pool)
.await?;
Ok(rows.into_iter().map(|(name,)| name).collect())
}
/// Check if a user holds a specific named permission.
async fn user_has_permission(
&self,
user_id: &str,
permission_name: &str,
) -> Result<bool, sqlx::Error> {
let row: (i64,) = sqlx::query_as(
r#"
SELECT COUNT(*)
FROM permissions p
JOIN role_permissions rp ON rp.permission_id = p.id
JOIN user_roles ur ON ur.role_id = rp.role_id
WHERE ur.user_id = ? AND p.name = ?
"#,
)
.bind(user_id)
.bind(permission_name)
.fetch_one(&self.pool)
.await?;
Ok(row.0 > 0)
}
}
@@ -0,0 +1,67 @@
use crate::db::repository::traits::RefreshTokensRepository;
use async_trait::async_trait;
use sqlx::SqlitePool;
pub struct SqliteRefreshTokensRepository {
pub pool: SqlitePool,
}
#[derive(Debug, Clone, sqlx::FromRow)]
pub struct RefreshToken {
pub id: String,
pub user_id: String,
pub token_hash: String,
pub expires_at: String,
pub revoked: bool,
pub created_at: String,
}
#[async_trait]
impl RefreshTokensRepository for SqliteRefreshTokensRepository {
async fn create(
&self,
id: &str,
user_id: &str,
token_hash: &str,
expires_at: &str,
) -> Result<RefreshToken, sqlx::Error> {
sqlx::query_as::<_, RefreshToken>(
r#"
INSERT INTO refresh_tokens (id, user_id, token_hash, expires_at)
VALUES (?, ?, ?, ?)
RETURNING *
"#,
)
.bind(id)
.bind(user_id)
.bind(token_hash)
.bind(expires_at)
.fetch_one(&self.pool)
.await
}
async fn find_by_hash(&self, token_hash: &str) -> Result<Option<RefreshToken>, sqlx::Error> {
sqlx::query_as::<_, RefreshToken>(
"SELECT * FROM refresh_tokens WHERE token_hash = ? AND revoked = 0",
)
.bind(token_hash)
.fetch_optional(&self.pool)
.await
}
async fn revoke(&self, id: &str) -> Result<(), sqlx::Error> {
sqlx::query("UPDATE refresh_tokens SET revoked = 1 WHERE id = ?")
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
async fn revoke_all_for_user(&self, user_id: &str) -> Result<(), sqlx::Error> {
sqlx::query("UPDATE refresh_tokens SET revoked = 1 WHERE user_id = ?")
.bind(user_id)
.execute(&self.pool)
.await?;
Ok(())
}
}
+127
View File
@@ -0,0 +1,127 @@
use crate::db::repository::traits::RolesRepository;
use async_trait::async_trait;
use sqlx::SqlitePool;
use crate::db::models::Role;
pub struct SqliteRolesRepository {
pub pool: SqlitePool,
}
#[async_trait]
impl RolesRepository for SqliteRolesRepository {
async fn list_all(&self) -> Result<Vec<Role>, sqlx::Error> {
sqlx::query_as::<_, Role>("SELECT * FROM roles ORDER BY name")
.fetch_all(&self.pool)
.await
}
async fn find_by_name(&self, name: &str) -> Result<Option<Role>, sqlx::Error> {
sqlx::query_as::<_, Role>("SELECT * FROM roles WHERE name = ?")
.bind(name)
.fetch_optional(&self.pool)
.await
}
async fn find_by_id(&self, id: &str) -> Result<Option<Role>, sqlx::Error> {
sqlx::query_as::<_, Role>("SELECT * FROM roles WHERE id = ?")
.bind(id)
.fetch_optional(&self.pool)
.await
}
async fn list_for_user(&self, user_id: &str) -> Result<Vec<Role>, sqlx::Error> {
sqlx::query_as::<_, Role>(
r#"
SELECT r.* FROM roles r
JOIN user_roles ur ON ur.role_id = r.id
WHERE ur.user_id = ?
ORDER BY r.name
"#,
)
.bind(user_id)
.fetch_all(&self.pool)
.await
}
async fn assign_to_user(&self, user_id: &str, role_id: &str) -> Result<(), sqlx::Error> {
sqlx::query("INSERT OR IGNORE INTO user_roles (user_id, role_id) VALUES (?, ?)")
.bind(user_id)
.bind(role_id)
.execute(&self.pool)
.await?;
Ok(())
}
async fn remove_from_user(&self, user_id: &str, role_id: &str) -> Result<(), sqlx::Error> {
sqlx::query("DELETE FROM user_roles WHERE user_id = ? AND role_id = ?")
.bind(user_id)
.bind(role_id)
.execute(&self.pool)
.await?;
Ok(())
}
async fn admin_role_exists(&self) -> Result<bool, sqlx::Error> {
let row: (i64,) = sqlx::query_as("SELECT COUNT(*) FROM roles WHERE name = 'admin'")
.fetch_one(&self.pool)
.await?;
Ok(row.0 > 0)
}
/// Create a new role.
async fn create(
&self,
id: &str,
name: &str,
description: Option<&str>,
) -> Result<Role, sqlx::Error> {
sqlx::query_as::<_, Role>(
r#"
INSERT INTO roles (id, name, description)
VALUES (?, ?, ?)
RETURNING *
"#,
)
.bind(id)
.bind(name)
.bind(description)
.fetch_one(&self.pool)
.await
}
/// Update role name/description.
async fn update(
&self,
id: &str,
name: &str,
description: Option<&str>,
) -> Result<(), sqlx::Error> {
sqlx::query("UPDATE roles SET name = ?, description = ? WHERE id = ?")
.bind(name)
.bind(description)
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
/// Delete a role by id.
async fn delete(&self, id: &str) -> Result<(), sqlx::Error> {
sqlx::query("DELETE FROM roles WHERE id = ?")
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
/// List user ids that hold a given role.
async fn list_user_ids_for_role(&self, role_id: &str) -> Result<Vec<String>, sqlx::Error> {
let rows: Vec<(String,)> =
sqlx::query_as("SELECT user_id FROM user_roles WHERE role_id = ? ORDER BY user_id")
.bind(role_id)
.fetch_all(&self.pool)
.await?;
Ok(rows.into_iter().map(|(id,)| id).collect())
}
}
@@ -0,0 +1,78 @@
use crate::db::repository::traits::ServiceAccountsRepository;
use async_trait::async_trait;
use sqlx::SqlitePool;
use crate::db::models::ServiceAccount;
pub struct SqliteServiceAccountsRepository {
pub pool: SqlitePool,
}
#[async_trait]
impl ServiceAccountsRepository for SqliteServiceAccountsRepository {
async fn create(
&self,
id: &str,
tenant_id: &str,
name: &str,
description: Option<&str>,
) -> Result<ServiceAccount, sqlx::Error> {
sqlx::query_as::<_, ServiceAccount>(
r#"
INSERT INTO service_accounts (id, tenant_id, name, description)
VALUES (?, ?, ?, ?)
RETURNING *
"#,
)
.bind(id)
.bind(tenant_id)
.bind(name)
.bind(description)
.fetch_one(&self.pool)
.await
}
async fn find_by_id(&self, id: &str) -> Result<Option<ServiceAccount>, sqlx::Error> {
sqlx::query_as::<_, ServiceAccount>("SELECT * FROM service_accounts WHERE id = ?")
.bind(id)
.fetch_optional(&self.pool)
.await
}
async fn list(&self, tenant_id: &str) -> Result<Vec<ServiceAccount>, sqlx::Error> {
sqlx::query_as::<_, ServiceAccount>(
"SELECT * FROM service_accounts WHERE tenant_id = ? ORDER BY name",
)
.bind(tenant_id)
.fetch_all(&self.pool)
.await
}
async fn set_enabled(&self, id: &str, enabled: bool) -> Result<(), sqlx::Error> {
sqlx::query(
"UPDATE service_accounts SET enabled = ?, updated_at = strftime('%Y-%m-%dT%H:%M:%SZ', 'now') WHERE id = ?",
)
.bind(enabled)
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
async fn delete(&self, id: &str) -> Result<(), sqlx::Error> {
sqlx::query("DELETE FROM service_accounts WHERE id = ?")
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
async fn count(&self, tenant_id: &str) -> Result<i64, sqlx::Error> {
let row: (i64,) =
sqlx::query_as("SELECT COUNT(*) FROM service_accounts WHERE tenant_id = ?")
.bind(tenant_id)
.fetch_one(&self.pool)
.await?;
Ok(row.0)
}
}
+141
View File
@@ -0,0 +1,141 @@
use crate::db::repository::traits::SessionsRepository;
use async_trait::async_trait;
use sqlx::SqlitePool;
use crate::db::models::Session;
pub struct SqliteSessionsRepository {
pub pool: SqlitePool,
}
#[async_trait]
impl SessionsRepository for SqliteSessionsRepository {
async fn create(
&self,
id: &str,
user_id: &str,
token_hash: &str,
ip_address: Option<&str>,
user_agent: Option<&str>,
expires_at: &str,
) -> Result<Session, sqlx::Error> {
sqlx::query_as::<_, Session>(
r#"
INSERT INTO sessions (id, user_id, token_hash, ip_address, user_agent, expires_at)
VALUES (?, ?, ?, ?, ?, ?)
RETURNING *
"#,
)
.bind(id)
.bind(user_id)
.bind(token_hash)
.bind(ip_address)
.bind(user_agent)
.bind(expires_at)
.fetch_one(&self.pool)
.await
}
async fn find_by_token_hash(&self, token_hash: &str) -> Result<Option<Session>, sqlx::Error> {
sqlx::query_as::<_, Session>("SELECT * FROM sessions WHERE token_hash = ? AND revoked = 0")
.bind(token_hash)
.fetch_optional(&self.pool)
.await
}
async fn revoke(&self, id: &str) -> Result<(), sqlx::Error> {
sqlx::query("UPDATE sessions SET revoked = 1 WHERE id = ?")
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
async fn revoke_all_for_user(&self, user_id: &str) -> Result<(), sqlx::Error> {
sqlx::query("UPDATE sessions SET revoked = 1 WHERE user_id = ?")
.bind(user_id)
.execute(&self.pool)
.await?;
Ok(())
}
async fn update_last_seen(&self, id: &str) -> Result<(), sqlx::Error> {
sqlx::query(
"UPDATE sessions SET last_seen_at = strftime('%Y-%m-%dT%H:%M:%SZ', 'now') WHERE id = ?",
)
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
/// List active (non-revoked, non-expired) sessions for a user.
async fn list_active_for_user(&self, user_id: &str) -> Result<Vec<Session>, sqlx::Error> {
sqlx::query_as::<_, Session>(
r#"
SELECT * FROM sessions
WHERE user_id = ?
AND revoked = 0
AND expires_at >= strftime('%Y-%m-%dT%H:%M:%SZ', 'now')
ORDER BY last_seen_at DESC
"#,
)
.bind(user_id)
.fetch_all(&self.pool)
.await
}
/// List ALL active sessions system-wide (admin only)
async fn list_all_active(&self) -> Result<Vec<Session>, sqlx::Error> {
sqlx::query_as::<_, Session>(
r#"
SELECT * FROM sessions
WHERE revoked = 0
AND expires_at >= strftime('%Y-%m-%dT%H:%M:%SZ', 'now')
ORDER BY last_seen_at DESC
"#,
)
.fetch_all(&self.pool)
.await
}
/// Count active sessions system-wide.
async fn count_active(&self) -> Result<i64, sqlx::Error> {
let row: (i64,) = sqlx::query_as(
r#"
SELECT COUNT(*) FROM sessions
WHERE revoked = 0
AND expires_at >= strftime('%Y-%m-%dT%H:%M:%SZ', 'now')
"#,
)
.fetch_one(&self.pool)
.await?;
Ok(row.0)
}
/// Delete sessions that are expired or revoked. Called once at startup.
async fn cleanup_expired(&self) -> Result<u64, sqlx::Error> {
let result = sqlx::query(
r#"
DELETE FROM sessions
WHERE revoked = 1
OR expires_at < strftime('%Y-%m-%dT%H:%M:%SZ', 'now')
"#,
)
.execute(&self.pool)
.await?;
Ok(result.rows_affected())
}
/// Revoke all sessions for a user EXCEPT the given session_id
async fn revoke_others(&self, user_id: &str, except_id: &str) -> Result<u64, sqlx::Error> {
let result = sqlx::query(
"UPDATE sessions SET revoked = 1 WHERE user_id = ? AND id != ? AND revoked = 0",
)
.bind(user_id)
.bind(except_id)
.execute(&self.pool)
.await?;
Ok(result.rows_affected())
}
}
+108
View File
@@ -0,0 +1,108 @@
use sqlx::SqlitePool;
use crate::db::models::Tenant;
use crate::db::repository::traits::TenantsRepository;
pub struct SqliteTenantsRepository {
pub pool: SqlitePool,
}
#[async_trait::async_trait]
impl TenantsRepository for SqliteTenantsRepository {
async fn find_by_id(&self, id: &str) -> Result<Option<Tenant>, sqlx::Error> {
let row = sqlx::query_as::<_, Tenant>(
"SELECT id, name, slug, enabled, created_at, updated_at FROM tenants WHERE id = ?",
)
.bind(id)
.fetch_optional(&self.pool)
.await?;
Ok(row)
}
async fn find_by_slug(&self, slug: &str) -> Result<Option<Tenant>, sqlx::Error> {
let row = sqlx::query_as::<_, Tenant>(
"SELECT id, name, slug, enabled, created_at, updated_at FROM tenants WHERE slug = ?",
)
.bind(slug)
.fetch_optional(&self.pool)
.await?;
Ok(row)
}
async fn list(&self) -> Result<Vec<Tenant>, sqlx::Error> {
let rows = sqlx::query_as::<_, Tenant>(
"SELECT id, name, slug, enabled, created_at, updated_at FROM tenants ORDER BY name ASC",
)
.fetch_all(&self.pool)
.await?;
Ok(rows)
}
async fn create(
&self,
id: &str,
name: &str,
slug: Option<&str>,
) -> Result<Tenant, sqlx::Error> {
let slug = slug.unwrap_or(id);
let row = sqlx::query_as::<_, Tenant>(
r#"
INSERT INTO tenants (id, name, slug, enabled)
VALUES (?, ?, ?, 1)
RETURNING id, name, slug, enabled, created_at, updated_at
"#,
)
.bind(id)
.bind(name)
.bind(slug)
.fetch_one(&self.pool)
.await?;
Ok(row)
}
async fn update(&self, id: &str, name: &str, slug: Option<&str>) -> Result<(), sqlx::Error> {
let slug = slug.unwrap_or(name);
sqlx::query(
r#"
UPDATE tenants
SET name = ?, slug = ?, updated_at = strftime('%Y-%m-%dT%H:%M:%SZ', 'now')
WHERE id = ?
"#,
)
.bind(name)
.bind(slug)
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
async fn set_enabled(&self, id: &str, enabled: bool) -> Result<(), sqlx::Error> {
sqlx::query(
r#"
UPDATE tenants
SET enabled = ?, updated_at = strftime('%Y-%m-%dT%H:%M:%SZ', 'now')
WHERE id = ?
"#,
)
.bind(enabled as i32)
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
async fn delete(&self, id: &str) -> Result<(), sqlx::Error> {
sqlx::query("DELETE FROM tenants WHERE id = ?")
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
}
+79
View File
@@ -0,0 +1,79 @@
use crate::db::repository::traits::TokensRepository;
use async_trait::async_trait;
use sqlx::SqlitePool;
use crate::db::models::ApiToken;
pub struct SqliteTokensRepository {
pub pool: SqlitePool,
}
#[async_trait]
impl TokensRepository for SqliteTokensRepository {
async fn create(
&self,
id: &str,
user_id: &str,
name: &str,
token_hash: &str,
expires_at: Option<&str>,
) -> Result<ApiToken, sqlx::Error> {
sqlx::query_as::<_, ApiToken>(
r#"
INSERT INTO api_tokens (id, user_id, name, token_hash, expires_at)
VALUES (?, ?, ?, ?, ?)
RETURNING *
"#,
)
.bind(id)
.bind(user_id)
.bind(name)
.bind(token_hash)
.bind(expires_at)
.fetch_one(&self.pool)
.await
}
async fn find_by_hash(&self, token_hash: &str) -> Result<Option<ApiToken>, sqlx::Error> {
sqlx::query_as::<_, ApiToken>(
"SELECT * FROM api_tokens WHERE token_hash = ? AND revoked = 0",
)
.bind(token_hash)
.fetch_optional(&self.pool)
.await
}
async fn list_for_user(&self, user_id: &str) -> Result<Vec<ApiToken>, sqlx::Error> {
sqlx::query_as::<_, ApiToken>(
"SELECT * FROM api_tokens WHERE user_id = ? ORDER BY created_at DESC",
)
.bind(user_id)
.fetch_all(&self.pool)
.await
}
async fn find_by_id(&self, id: &str) -> Result<Option<ApiToken>, sqlx::Error> {
sqlx::query_as::<_, ApiToken>("SELECT * FROM api_tokens WHERE id = ?")
.bind(id)
.fetch_optional(&self.pool)
.await
}
async fn revoke(&self, id: &str) -> Result<(), sqlx::Error> {
sqlx::query("UPDATE api_tokens SET revoked = 1 WHERE id = ?")
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
async fn update_last_used(&self, id: &str) -> Result<(), sqlx::Error> {
sqlx::query(
"UPDATE api_tokens SET last_used_at = strftime('%Y-%m-%dT%H:%M:%SZ', 'now') WHERE id = ?",
)
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
}
+173
View File
@@ -0,0 +1,173 @@
use crate::db::repository::traits::UsersRepository;
use async_trait::async_trait;
use sqlx::SqlitePool;
use crate::db::models::User;
pub struct SqliteUsersRepository {
pub pool: SqlitePool,
}
/// User profile fields from `user_profiles`.
#[derive(Debug, Clone, sqlx::FromRow, serde::Serialize, serde::Deserialize)]
pub struct UserProfile {
pub user_id: String,
pub email: Option<String>,
pub full_name: Option<String>,
pub avatar_url: Option<String>,
pub metadata_json: Option<String>,
}
#[async_trait]
impl UsersRepository for SqliteUsersRepository {
/// Count users with a given status in a tenant.
async fn count_by_status(&self, tenant_id: &str, status: i32) -> Result<i64, sqlx::Error> {
let row: (i64,) =
sqlx::query_as("SELECT COUNT(*) FROM users WHERE tenant_id = ? AND status = ?")
.bind(tenant_id)
.bind(status)
.fetch_one(&self.pool)
.await?;
Ok(row.0)
}
/// Count all users in a tenant.
async fn count(&self, tenant_id: &str) -> Result<i64, sqlx::Error> {
let row: (i64,) = sqlx::query_as("SELECT COUNT(*) FROM users WHERE tenant_id = ?")
.bind(tenant_id)
.fetch_one(&self.pool)
.await?;
Ok(row.0)
}
/// Count users that have the admin role.
async fn count_admins(&self) -> Result<i64, sqlx::Error> {
let row: (i64,) = sqlx::query_as(
r#"
SELECT COUNT(DISTINCT ur.user_id)
FROM user_roles ur
JOIN roles r ON r.id = ur.role_id
WHERE r.name = 'admin'
"#,
)
.fetch_one(&self.pool)
.await?;
Ok(row.0)
}
async fn find_by_id(&self, id: &str) -> Result<Option<User>, sqlx::Error> {
sqlx::query_as::<_, User>("SELECT * FROM users WHERE id = ?")
.bind(id)
.fetch_optional(&self.pool)
.await
}
async fn find_by_username(&self, username: &str) -> Result<Option<User>, sqlx::Error> {
sqlx::query_as::<_, User>("SELECT * FROM users WHERE username = ?")
.bind(username)
.fetch_optional(&self.pool)
.await
}
async fn list(&self, tenant_id: &str) -> Result<Vec<User>, sqlx::Error> {
sqlx::query_as::<_, User>(
"SELECT * FROM users WHERE tenant_id = ? ORDER BY created_at DESC",
)
.bind(tenant_id)
.fetch_all(&self.pool)
.await
}
async fn create(
&self,
id: &str,
tenant_id: &str,
username: &str,
password_hash: &str,
) -> Result<User, sqlx::Error> {
sqlx::query_as::<_, User>(
r#"
INSERT INTO users (id, tenant_id, username, password_hash, status)
VALUES (?, ?, ?, ?, 1)
RETURNING *
"#,
)
.bind(id)
.bind(tenant_id)
.bind(username)
.bind(password_hash)
.fetch_one(&self.pool)
.await
}
async fn update_status(&self, id: &str, status: i32) -> Result<(), sqlx::Error> {
sqlx::query(
"UPDATE users SET status = ?, updated_at = strftime('%Y-%m-%dT%H:%M:%SZ', 'now') WHERE id = ?",
)
.bind(status)
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
async fn update_password_hash(&self, id: &str, password_hash: &str) -> Result<(), sqlx::Error> {
sqlx::query(
"UPDATE users SET password_hash = ?, updated_at = strftime('%Y-%m-%dT%H:%M:%SZ', 'now') WHERE id = ?",
)
.bind(password_hash)
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
async fn set_last_login(&self, id: &str) -> Result<(), sqlx::Error> {
sqlx::query(
"UPDATE users SET last_login_at = strftime('%Y-%m-%dT%H:%M:%SZ', 'now'), updated_at = strftime('%Y-%m-%dT%H:%M:%SZ', 'now') WHERE id = ?",
)
.bind(id)
.execute(&self.pool)
.await?;
Ok(())
}
async fn username_exists(&self, tenant_id: &str, username: &str) -> Result<bool, sqlx::Error> {
let row: (i64,) =
sqlx::query_as("SELECT COUNT(*) FROM users WHERE tenant_id = ? AND username = ?")
.bind(tenant_id)
.bind(username)
.fetch_one(&self.pool)
.await?;
Ok(row.0 > 0)
}
async fn get_profile(&self, user_id: &str) -> Result<Option<UserProfile>, sqlx::Error> {
sqlx::query_as::<_, UserProfile>("SELECT * FROM user_profiles WHERE user_id = ?")
.bind(user_id)
.fetch_optional(&self.pool)
.await
}
async fn upsert_profile(
&self,
user_id: &str,
email: Option<&str>,
full_name: Option<&str>,
) -> Result<UserProfile, sqlx::Error> {
sqlx::query_as::<_, UserProfile>(
r#"
INSERT INTO user_profiles (user_id, email, full_name)
VALUES (?, ?, ?)
ON CONFLICT(user_id) DO UPDATE SET
email = excluded.email,
full_name = excluded.full_name
RETURNING *
"#,
)
.bind(user_id)
.bind(email)
.bind(full_name)
.fetch_one(&self.pool)
.await
}
}
+13 -66
View File
@@ -1,74 +1,21 @@
use sqlx::SqlitePool;
pub use crate::db::repository::sqlite::tokens::*;
use crate::db::models::ApiToken;
use crate::db::provider::DatabaseProvider;
use std::sync::Arc;
pub async fn create(
tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>,
id: &str,
/// List tokens for a user using the provided DatabaseProvider.
pub async fn list_for_user(
provider: &Arc<dyn DatabaseProvider>,
user_id: &str,
name: &str,
token_hash: &str,
expires_at: Option<&str>,
) -> Result<ApiToken, sqlx::Error> {
sqlx::query_as::<_, ApiToken>(
r#"
INSERT INTO api_tokens (id, user_id, name, token_hash, expires_at)
VALUES (?, ?, ?, ?, ?)
RETURNING *
"#,
)
.bind(id)
.bind(user_id)
.bind(name)
.bind(token_hash)
.bind(expires_at)
.fetch_one(&mut **tx)
.await
) -> Result<Vec<ApiToken>, sqlx::Error> {
provider.tokens().list_for_user(user_id).await
}
pub async fn find_by_hash(
pool: &SqlitePool,
token_hash: &str,
) -> Result<Option<ApiToken>, sqlx::Error> {
sqlx::query_as::<_, ApiToken>("SELECT * FROM api_tokens WHERE token_hash = ? AND revoked = 0")
.bind(token_hash)
.fetch_optional(pool)
.await
}
pub async fn list_for_user(pool: &SqlitePool, user_id: &str) -> Result<Vec<ApiToken>, sqlx::Error> {
sqlx::query_as::<_, ApiToken>(
"SELECT * FROM api_tokens WHERE user_id = ? ORDER BY created_at DESC",
)
.bind(user_id)
.fetch_all(pool)
.await
}
pub async fn find_by_id(pool: &SqlitePool, id: &str) -> Result<Option<ApiToken>, sqlx::Error> {
sqlx::query_as::<_, ApiToken>("SELECT * FROM api_tokens WHERE id = ?")
.bind(id)
.fetch_optional(pool)
.await
}
pub async fn revoke(
tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>,
/// Find a token by its ID using the provided DatabaseProvider.
pub async fn find_by_id(
provider: &Arc<dyn DatabaseProvider>,
id: &str,
) -> Result<(), sqlx::Error> {
sqlx::query("UPDATE api_tokens SET revoked = 1 WHERE id = ?")
.bind(id)
.execute(&mut **tx)
.await?;
Ok(())
}
pub async fn update_last_used(pool: &SqlitePool, id: &str) -> Result<(), sqlx::Error> {
sqlx::query(
"UPDATE api_tokens SET last_used_at = strftime('%Y-%m-%dT%H:%M:%SZ', 'now') WHERE id = ?",
)
.bind(id)
.execute(pool)
.await?;
Ok(())
) -> Result<Option<ApiToken>, sqlx::Error> {
provider.tokens().find_by_id(id).await
}
+256
View File
@@ -0,0 +1,256 @@
use crate::db::models::{
ApiToken, Application, AuditLog, Group, Permission, Role, ServiceAccount, Session, Tenant, User,
};
use crate::db::repository::sqlite::audit::AuditFilter;
use crate::db::repository::sqlite::refresh_tokens::RefreshToken;
use crate::db::repository::sqlite::users::UserProfile;
#[async_trait::async_trait]
pub trait UsersRepository: Send + Sync {
async fn find_by_id(&self, id: &str) -> Result<Option<User>, sqlx::Error>;
async fn find_by_username(&self, username: &str) -> Result<Option<User>, sqlx::Error>;
async fn list(&self, tenant_id: &str) -> Result<Vec<User>, sqlx::Error>;
async fn create(
&self,
id: &str,
tenant_id: &str,
username: &str,
password_hash: &str,
) -> Result<User, sqlx::Error>;
async fn update_status(&self, id: &str, status: i32) -> Result<(), sqlx::Error>;
async fn update_password_hash(&self, id: &str, password_hash: &str) -> Result<(), sqlx::Error>;
async fn set_last_login(&self, id: &str) -> Result<(), sqlx::Error>;
async fn username_exists(&self, tenant_id: &str, username: &str) -> Result<bool, sqlx::Error>;
async fn count_admins(&self) -> Result<i64, sqlx::Error>;
async fn count(&self, tenant_id: &str) -> Result<i64, sqlx::Error>;
async fn count_by_status(&self, tenant_id: &str, status: i32) -> Result<i64, sqlx::Error>;
async fn get_profile(&self, user_id: &str) -> Result<Option<UserProfile>, sqlx::Error>;
async fn upsert_profile(
&self,
user_id: &str,
email: Option<&str>,
full_name: Option<&str>,
) -> Result<UserProfile, sqlx::Error>;
}
#[async_trait::async_trait]
pub trait ServiceAccountsRepository: Send + Sync {
async fn create(
&self,
id: &str,
tenant_id: &str,
name: &str,
description: Option<&str>,
) -> Result<ServiceAccount, sqlx::Error>;
async fn find_by_id(&self, id: &str) -> Result<Option<ServiceAccount>, sqlx::Error>;
async fn list(&self, tenant_id: &str) -> Result<Vec<ServiceAccount>, sqlx::Error>;
async fn set_enabled(&self, id: &str, enabled: bool) -> Result<(), sqlx::Error>;
async fn delete(&self, id: &str) -> Result<(), sqlx::Error>;
async fn count(&self, tenant_id: &str) -> Result<i64, sqlx::Error>;
}
#[async_trait::async_trait]
pub trait RolesRepository: Send + Sync {
async fn list_all(&self) -> Result<Vec<Role>, sqlx::Error>;
async fn find_by_name(&self, name: &str) -> Result<Option<Role>, sqlx::Error>;
async fn find_by_id(&self, id: &str) -> Result<Option<Role>, sqlx::Error>;
async fn list_for_user(&self, user_id: &str) -> Result<Vec<Role>, sqlx::Error>;
async fn assign_to_user(&self, user_id: &str, role_id: &str) -> Result<(), sqlx::Error>;
async fn remove_from_user(&self, user_id: &str, role_id: &str) -> Result<(), sqlx::Error>;
async fn admin_role_exists(&self) -> Result<bool, sqlx::Error>;
async fn create(
&self,
id: &str,
name: &str,
description: Option<&str>,
) -> Result<Role, sqlx::Error>;
async fn update(
&self,
id: &str,
name: &str,
description: Option<&str>,
) -> Result<(), sqlx::Error>;
async fn delete(&self, id: &str) -> Result<(), sqlx::Error>;
async fn list_user_ids_for_role(&self, role_id: &str) -> Result<Vec<String>, sqlx::Error>;
}
#[async_trait::async_trait]
pub trait SessionsRepository: Send + Sync {
async fn create(
&self,
id: &str,
user_id: &str,
token_hash: &str,
ip_address: Option<&str>,
user_agent: Option<&str>,
expires_at: &str,
) -> Result<Session, sqlx::Error>;
async fn find_by_token_hash(&self, token_hash: &str) -> Result<Option<Session>, sqlx::Error>;
async fn revoke(&self, id: &str) -> Result<(), sqlx::Error>;
async fn revoke_all_for_user(&self, user_id: &str) -> Result<(), sqlx::Error>;
async fn update_last_seen(&self, id: &str) -> Result<(), sqlx::Error>;
async fn list_active_for_user(&self, user_id: &str) -> Result<Vec<Session>, sqlx::Error>;
async fn list_all_active(&self) -> Result<Vec<Session>, sqlx::Error>;
async fn count_active(&self) -> Result<i64, sqlx::Error>;
async fn cleanup_expired(&self) -> Result<u64, sqlx::Error>;
async fn revoke_others(&self, user_id: &str, except_id: &str) -> Result<u64, sqlx::Error>;
}
#[async_trait::async_trait]
pub trait ApplicationsRepository: Send + Sync {
async fn create(
&self,
id: &str,
tenant_id: &str,
name: &str,
slug: &str,
) -> Result<Application, sqlx::Error>;
async fn find_by_slug(&self, slug: &str) -> Result<Option<Application>, sqlx::Error>;
async fn find_by_id(&self, id: &str) -> Result<Option<Application>, sqlx::Error>;
async fn list(&self, tenant_id: &str) -> Result<Vec<Application>, sqlx::Error>;
async fn set_enabled(&self, id: &str, enabled: bool) -> Result<(), sqlx::Error>;
async fn update(
&self,
id: &str,
name: &str,
slug: &str,
enabled: bool,
) -> Result<(), sqlx::Error>;
async fn delete(&self, id: &str) -> Result<(), sqlx::Error>;
async fn count(&self, tenant_id: &str) -> Result<i64, sqlx::Error>;
}
#[async_trait::async_trait]
pub trait RefreshTokensRepository: Send + Sync {
async fn create(
&self,
id: &str,
user_id: &str,
token_hash: &str,
expires_at: &str,
) -> Result<RefreshToken, sqlx::Error>;
async fn find_by_hash(&self, token_hash: &str) -> Result<Option<RefreshToken>, sqlx::Error>;
async fn revoke(&self, id: &str) -> Result<(), sqlx::Error>;
async fn revoke_all_for_user(&self, user_id: &str) -> Result<(), sqlx::Error>;
}
#[async_trait::async_trait]
pub trait AuditRepository: Send + Sync {
#[allow(clippy::too_many_arguments)]
async fn insert(
&self,
id: &str,
actor_user_id: Option<&str>,
target_user_id: Option<&str>,
action: &str,
resource_type: &str,
resource_id: Option<&str>,
severity: &str,
ip_address: Option<&str>,
user_agent: Option<&str>,
metadata_json: Option<&str>,
) -> Result<AuditLog, sqlx::Error>;
async fn list_recent(&self, limit: i64) -> Result<Vec<AuditLog>, sqlx::Error>;
async fn count(&self) -> Result<i64, sqlx::Error>;
async fn list_filtered(&self, filter: &AuditFilter) -> Result<Vec<AuditLog>, sqlx::Error>;
async fn count_filtered(&self, filter: &AuditFilter) -> Result<i64, sqlx::Error>;
}
#[async_trait::async_trait]
pub trait AuditRepositoryExt: Send + Sync {
async fn log(&self, event: crate::audit::AuditEvent<'_>) -> Result<AuditLog, sqlx::Error>;
}
#[async_trait::async_trait]
impl<T: ?Sized + AuditRepository> AuditRepositoryExt for T {
async fn log(&self, event: crate::audit::AuditEvent<'_>) -> Result<AuditLog, sqlx::Error> {
self.insert(
&uuid::Uuid::new_v4().to_string(),
event.actor_id,
event.target_id,
event.action,
event.resource_type,
event.resource_id,
event.severity.as_str(),
event.ip,
event.ua,
event.metadata,
)
.await
}
}
#[async_trait::async_trait]
pub trait TokensRepository: Send + Sync {
async fn create(
&self,
id: &str,
user_id: &str,
name: &str,
token_hash: &str,
expires_at: Option<&str>,
) -> Result<ApiToken, sqlx::Error>;
async fn find_by_hash(&self, token_hash: &str) -> Result<Option<ApiToken>, sqlx::Error>;
async fn list_for_user(&self, user_id: &str) -> Result<Vec<ApiToken>, sqlx::Error>;
async fn find_by_id(&self, id: &str) -> Result<Option<ApiToken>, sqlx::Error>;
async fn revoke(&self, id: &str) -> Result<(), sqlx::Error>;
async fn update_last_used(&self, id: &str) -> Result<(), sqlx::Error>;
}
#[async_trait::async_trait]
pub trait PermissionsRepository: Send + Sync {
async fn list_all(&self) -> Result<Vec<Permission>, sqlx::Error>;
async fn list_for_role(&self, role_id: &str) -> Result<Vec<Permission>, sqlx::Error>;
async fn assign_to_role(&self, role_id: &str, permission_id: &str) -> Result<(), sqlx::Error>;
async fn remove_from_role(&self, role_id: &str, permission_id: &str)
-> Result<(), sqlx::Error>;
async fn clear_for_role(&self, role_id: &str) -> Result<(), sqlx::Error>;
async fn find_by_name(&self, name: &str) -> Result<Option<Permission>, sqlx::Error>;
async fn find_by_id(&self, id: &str) -> Result<Option<Permission>, sqlx::Error>;
async fn list_for_user(&self, user_id: &str) -> Result<Vec<String>, sqlx::Error>;
async fn user_has_permission(
&self,
user_id: &str,
permission_name: &str,
) -> Result<bool, sqlx::Error>;
}
#[async_trait::async_trait]
pub trait TenantsRepository: Send + Sync {
async fn find_by_id(&self, id: &str) -> Result<Option<Tenant>, sqlx::Error>;
async fn find_by_slug(&self, slug: &str) -> Result<Option<Tenant>, sqlx::Error>;
async fn list(&self) -> Result<Vec<Tenant>, sqlx::Error>;
async fn create(&self, id: &str, name: &str, slug: Option<&str>)
-> Result<Tenant, sqlx::Error>;
async fn update(&self, id: &str, name: &str, slug: Option<&str>) -> Result<(), sqlx::Error>;
async fn set_enabled(&self, id: &str, enabled: bool) -> Result<(), sqlx::Error>;
async fn delete(&self, id: &str) -> Result<(), sqlx::Error>;
}
#[async_trait::async_trait]
pub trait GroupsRepository: Send + Sync {
async fn list(&self, tenant_id: &str) -> Result<Vec<Group>, sqlx::Error>;
async fn find_by_id(&self, id: &str) -> Result<Option<Group>, sqlx::Error>;
async fn create(
&self,
id: &str,
tenant_id: &str,
name: &str,
description: Option<&str>,
) -> Result<Group, sqlx::Error>;
async fn update(
&self,
id: &str,
name: &str,
description: Option<&str>,
) -> Result<(), sqlx::Error>;
async fn delete(&self, id: &str) -> Result<(), sqlx::Error>;
async fn count_members(&self, group_id: &str) -> Result<i64, sqlx::Error>;
async fn list_members(
&self,
group_id: &str,
) -> Result<Vec<crate::db::models::User>, sqlx::Error>;
async fn add_member(&self, group_id: &str, user_id: &str) -> Result<(), sqlx::Error>;
async fn remove_member(&self, group_id: &str, user_id: &str) -> Result<(), sqlx::Error>;
async fn count(&self, tenant_id: &str) -> Result<i64, sqlx::Error>;
}
+1 -121
View File
@@ -1,121 +1 @@
use sqlx::SqlitePool;
use crate::db::models::User;
pub async fn find_by_id(pool: &SqlitePool, id: &str) -> Result<Option<User>, sqlx::Error> {
sqlx::query_as::<_, User>("SELECT * FROM users WHERE id = ?")
.bind(id)
.fetch_optional(pool)
.await
}
pub async fn find_by_username(
pool: &SqlitePool,
username: &str,
) -> Result<Option<User>, sqlx::Error> {
sqlx::query_as::<_, User>("SELECT * FROM users WHERE username = ?")
.bind(username)
.fetch_optional(pool)
.await
}
pub async fn list(pool: &SqlitePool, tenant_id: &str) -> Result<Vec<User>, sqlx::Error> {
sqlx::query_as::<_, User>("SELECT * FROM users WHERE tenant_id = ? ORDER BY created_at DESC")
.bind(tenant_id)
.fetch_all(pool)
.await
}
pub async fn create(
tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>,
id: &str,
tenant_id: &str,
username: &str,
password_hash: &str,
) -> Result<User, sqlx::Error> {
sqlx::query_as::<_, User>(
r#"
INSERT INTO users (id, tenant_id, username, password_hash, status)
VALUES (?, ?, ?, ?, 1)
RETURNING *
"#,
)
.bind(id)
.bind(tenant_id)
.bind(username)
.bind(password_hash)
.fetch_one(&mut **tx)
.await
}
pub async fn update_status(
tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>,
id: &str,
status: i32,
) -> Result<(), sqlx::Error> {
sqlx::query(
"UPDATE users SET status = ?, updated_at = strftime('%Y-%m-%dT%H:%M:%SZ', 'now') WHERE id = ?",
)
.bind(status)
.bind(id)
.execute(&mut **tx)
.await?;
Ok(())
}
pub async fn update_password_hash(
tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>,
id: &str,
password_hash: &str,
) -> Result<(), sqlx::Error> {
sqlx::query(
"UPDATE users SET password_hash = ?, updated_at = strftime('%Y-%m-%dT%H:%M:%SZ', 'now') WHERE id = ?",
)
.bind(password_hash)
.bind(id)
.execute(&mut **tx)
.await?;
Ok(())
}
pub async fn set_last_login(
tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>,
id: &str,
) -> Result<(), sqlx::Error> {
sqlx::query(
"UPDATE users SET last_login_at = strftime('%Y-%m-%dT%H:%M:%SZ', 'now'), updated_at = strftime('%Y-%m-%dT%H:%M:%SZ', 'now') WHERE id = ?",
)
.bind(id)
.execute(&mut **tx)
.await?;
Ok(())
}
pub async fn username_exists(
pool: &SqlitePool,
tenant_id: &str,
username: &str,
) -> Result<bool, sqlx::Error> {
let row: (i64,) =
sqlx::query_as("SELECT COUNT(*) FROM users WHERE tenant_id = ? AND username = ?")
.bind(tenant_id)
.bind(username)
.fetch_one(pool)
.await?;
Ok(row.0 > 0)
}
/// Count users that have the admin role.
pub async fn count_admins(pool: &SqlitePool) -> Result<i64, sqlx::Error> {
let row: (i64,) = sqlx::query_as(
r#"
SELECT COUNT(DISTINCT ur.user_id)
FROM user_roles ur
JOIN roles r ON r.id = ur.role_id
WHERE r.name = 'admin'
"#,
)
.fetch_one(pool)
.await?;
Ok(row.0)
}
pub use crate::db::repository::sqlite::users::*;