perf: migrate all user/session/auth IDs from VARCHAR(36) to native UUID

- Schema: all ~15 VARCHAR(36) columns → UUID with DEFAULT gen_random_uuid()
- Domain entities: User, Session, DeviceCode, AppPassword, Share → id: Uuid
- DTOs: CurrentUser.id → Uuid (API boundary DTOs keep String for JSON)
- Auth middleware: parse JWT claims.sub (String) → Uuid at boundary
- All repository traits, port traits, service impls updated end-to-end
- Handlers: pass Uuid by value (Copy, 16 bytes) instead of String refs
- Settings chain: updated_by column → Uuid (was text, caused setup crash)
- Removed ~650 lines of String↔Uuid conversion boilerplate
- Eliminates per-request heap allocations for ID cloning
- 16-byte binary comparison vs 36-byte string comparison in all queries
- Native UUID indexing in PostgreSQL (btree on 16 bytes vs 36-char text)

85 files changed, 1090 insertions(+), 1739 deletions(-)
This commit is contained in:
Diocrafts
2026-03-07 14:59:32 +01:00
parent 9f08460027
commit 06ed0455ce
85 changed files with 1090 additions and 1739 deletions
@@ -47,7 +47,7 @@ impl CalendarStoragePort for CalendarStorageAdapter {
async fn create_calendar(
&self,
dto: CreateCalendarDto,
owner_id: &str,
owner_id: Uuid,
) -> Result<CalendarDto, DomainError> {
let calendar = Calendar::new(dto.name, owner_id.to_string(), dto.description, dto.color)?;
@@ -117,7 +117,7 @@ impl CalendarStoragePort for CalendarStorageAdapter {
async fn list_calendars_by_owner(
&self,
owner_id: &str,
owner_id: Uuid,
) -> Result<Vec<CalendarDto>, DomainError> {
let calendars = self
.calendar_repository
@@ -128,7 +128,7 @@ impl CalendarStoragePort for CalendarStorageAdapter {
async fn list_calendars_shared_with_user(
&self,
user_id: &str,
user_id: Uuid,
) -> Result<Vec<CalendarDto>, DomainError> {
let calendars = self
.calendar_repository
@@ -152,7 +152,7 @@ impl CalendarStoragePort for CalendarStorageAdapter {
async fn check_calendar_access(
&self,
calendar_id: &str,
user_id: &str,
user_id: Uuid,
) -> Result<bool, DomainError> {
let uuid = Uuid::parse_str(calendar_id).map_err(|_| {
DomainError::new(
@@ -172,7 +172,7 @@ impl CalendarStoragePort for CalendarStorageAdapter {
async fn share_calendar(
&self,
calendar_id: &str,
user_id: &str,
user_id: Uuid,
access_level: &str,
) -> Result<(), DomainError> {
let uuid = Uuid::parse_str(calendar_id).map_err(|_| {
@@ -191,7 +191,7 @@ impl CalendarStoragePort for CalendarStorageAdapter {
async fn remove_calendar_sharing(
&self,
calendar_id: &str,
user_id: &str,
user_id: Uuid,
) -> Result<(), DomainError> {
let uuid = Uuid::parse_str(calendar_id).map_err(|_| {
DomainError::new(
@@ -61,7 +61,7 @@ impl ContactStorageAdapter {
async fn check_address_book_access(
&self,
address_book_id: &Uuid,
user_id: &str,
user_id: Uuid,
) -> Result<AddressBook, DomainError> {
let address_book = self
.address_book_repository
@@ -72,7 +72,7 @@ impl ContactStorageAdapter {
})?;
// Check if user is owner
if address_book.owner_id() == user_id {
if address_book.owner_id() == user_id.to_string() {
return Ok(address_book);
}
@@ -86,7 +86,7 @@ impl ContactStorageAdapter {
.address_book_repository
.get_address_book_shares(address_book_id)
.await?;
if shares.iter().any(|(shared_user, _)| shared_user == user_id) {
if shares.iter().any(|(shared_user, _)| shared_user == &user_id.to_string()) {
return Ok(address_book);
}
@@ -101,7 +101,7 @@ impl ContactStorageAdapter {
async fn check_write_access(
&self,
address_book_id: &Uuid,
user_id: &str,
user_id: Uuid,
) -> Result<AddressBook, DomainError> {
let address_book = self
.address_book_repository
@@ -112,7 +112,7 @@ impl ContactStorageAdapter {
})?;
// Owner always has write access
if address_book.owner_id() == user_id {
if address_book.owner_id() == user_id.to_string() {
return Ok(address_book);
}
@@ -123,7 +123,7 @@ impl ContactStorageAdapter {
.await?;
if shares
.iter()
.any(|(shared_user, can_write)| shared_user == user_id && *can_write)
.any(|(shared_user, can_write)| shared_user == &user_id.to_string() && *can_write)
{
return Ok(address_book);
}
@@ -247,7 +247,10 @@ impl AddressBookUseCase for ContactStorageAdapter {
let uuid = Self::parse_uuid(address_book_id, "AddressBook")?;
// Check write access
let mut address_book = self.check_write_access(&uuid, &update.user_id).await?;
let user_id = Uuid::parse_str(&update.user_id).map_err(|_| {
DomainError::new(ErrorKind::InvalidInput, "AddressBook", "Invalid user ID format")
})?;
let mut address_book = self.check_write_access(&uuid, user_id).await?;
if let Some(name) = update.name {
address_book.set_name(name);
@@ -273,7 +276,7 @@ impl AddressBookUseCase for ContactStorageAdapter {
async fn delete_address_book(
&self,
address_book_id: &str,
user_id: &str,
user_id: Uuid,
) -> Result<(), DomainError> {
let uuid = Self::parse_uuid(address_book_id, "AddressBook")?;
@@ -286,7 +289,7 @@ impl AddressBookUseCase for ContactStorageAdapter {
DomainError::new(ErrorKind::NotFound, "AddressBook", "Address book not found")
})?;
if address_book.owner_id() != user_id {
if address_book.owner_id() != user_id.to_string() {
return Err(DomainError::new(
ErrorKind::AccessDenied,
"AddressBook",
@@ -302,7 +305,7 @@ impl AddressBookUseCase for ContactStorageAdapter {
async fn get_address_book(
&self,
address_book_id: &str,
user_id: &str,
user_id: Uuid,
) -> Result<AddressBookDto, DomainError> {
let uuid = Self::parse_uuid(address_book_id, "AddressBook")?;
let address_book = self.check_address_book_access(&uuid, user_id).await?;
@@ -311,7 +314,7 @@ impl AddressBookUseCase for ContactStorageAdapter {
async fn list_user_address_books(
&self,
user_id: &str,
user_id: Uuid,
) -> Result<Vec<AddressBookDto>, DomainError> {
let owned = self
.address_book_repository
@@ -339,7 +342,7 @@ impl AddressBookUseCase for ContactStorageAdapter {
async fn share_address_book(
&self,
dto: ShareAddressBookDto,
user_id: &str,
user_id: Uuid,
) -> Result<(), DomainError> {
let uuid = Self::parse_uuid(&dto.address_book_id, "AddressBook")?;
@@ -352,7 +355,7 @@ impl AddressBookUseCase for ContactStorageAdapter {
DomainError::new(ErrorKind::NotFound, "AddressBook", "Address book not found")
})?;
if address_book.owner_id() != user_id {
if address_book.owner_id() != user_id.to_string() {
return Err(DomainError::new(
ErrorKind::AccessDenied,
"AddressBook",
@@ -360,15 +363,19 @@ impl AddressBookUseCase for ContactStorageAdapter {
));
}
let target_user_id = Uuid::parse_str(&dto.user_id).map_err(|_| {
DomainError::new(ErrorKind::InvalidInput, "AddressBook", "Invalid target user ID format")
})?;
self.address_book_repository
.share_address_book(&uuid, &dto.user_id, dto.can_write)
.share_address_book(&uuid, target_user_id, dto.can_write)
.await
}
async fn unshare_address_book(
&self,
dto: UnshareAddressBookDto,
user_id: &str,
user_id: Uuid,
) -> Result<(), DomainError> {
let uuid = Self::parse_uuid(&dto.address_book_id, "AddressBook")?;
@@ -381,7 +388,7 @@ impl AddressBookUseCase for ContactStorageAdapter {
DomainError::new(ErrorKind::NotFound, "AddressBook", "Address book not found")
})?;
if address_book.owner_id() != user_id {
if address_book.owner_id() != user_id.to_string() {
return Err(DomainError::new(
ErrorKind::AccessDenied,
"AddressBook",
@@ -389,15 +396,19 @@ impl AddressBookUseCase for ContactStorageAdapter {
));
}
let target_user_id = Uuid::parse_str(&dto.user_id).map_err(|_| {
DomainError::new(ErrorKind::InvalidInput, "AddressBook", "Invalid target user ID format")
})?;
self.address_book_repository
.unshare_address_book(&uuid, &dto.user_id)
.unshare_address_book(&uuid, target_user_id)
.await
}
async fn get_address_book_shares(
&self,
address_book_id: &str,
user_id: &str,
user_id: Uuid,
) -> Result<Vec<(String, bool)>, DomainError> {
let uuid = Self::parse_uuid(address_book_id, "AddressBook")?;
@@ -410,7 +421,7 @@ impl AddressBookUseCase for ContactStorageAdapter {
DomainError::new(ErrorKind::NotFound, "AddressBook", "Address book not found")
})?;
if address_book.owner_id() != user_id {
if address_book.owner_id() != user_id.to_string() {
return Err(DomainError::new(
ErrorKind::AccessDenied,
"AddressBook",
@@ -429,7 +440,10 @@ impl ContactUseCase for ContactStorageAdapter {
let address_book_id = Self::parse_uuid(&dto.address_book_id, "AddressBook")?;
// Check write access
self.check_write_access(&address_book_id, &dto.user_id)
let user_id = Uuid::parse_str(&dto.user_id).map_err(|_| {
DomainError::new(ErrorKind::InvalidInput, "Contact", "Invalid user ID format")
})?;
self.check_write_access(&address_book_id, user_id)
.await?;
let now = chrono::Utc::now();
@@ -471,7 +485,10 @@ impl ContactUseCase for ContactStorageAdapter {
let address_book_id = Self::parse_uuid(&dto.address_book_id, "AddressBook")?;
// Check write access
self.check_write_access(&address_book_id, &dto.user_id)
let user_id = Uuid::parse_str(&dto.user_id).map_err(|_| {
DomainError::new(ErrorKind::InvalidInput, "Contact", "Invalid user ID format")
})?;
self.check_write_access(&address_book_id, user_id)
.await?;
// Parse vCard fields
@@ -591,7 +608,10 @@ impl ContactUseCase for ContactStorageAdapter {
.ok_or_else(|| DomainError::new(ErrorKind::NotFound, "Contact", "Contact not found"))?;
// Check write access to the address book
self.check_write_access(contact.address_book_id(), &update.user_id)
let user_id = Uuid::parse_str(&update.user_id).map_err(|_| {
DomainError::new(ErrorKind::InvalidInput, "Contact", "Invalid user ID format")
})?;
self.check_write_access(contact.address_book_id(), user_id)
.await?;
if let Some(full_name) = update.full_name {
@@ -643,7 +663,7 @@ impl ContactUseCase for ContactStorageAdapter {
Ok(ContactDto::from(updated))
}
async fn delete_contact(&self, contact_id: &str, user_id: &str) -> Result<(), DomainError> {
async fn delete_contact(&self, contact_id: &str, user_id: Uuid) -> Result<(), DomainError> {
let uuid = Self::parse_uuid(contact_id, "Contact")?;
let contact = self
@@ -662,7 +682,7 @@ impl ContactUseCase for ContactStorageAdapter {
async fn get_contact(
&self,
contact_id: &str,
user_id: &str,
user_id: Uuid,
) -> Result<ContactDto, DomainError> {
let uuid = Self::parse_uuid(contact_id, "Contact")?;
@@ -682,7 +702,7 @@ impl ContactUseCase for ContactStorageAdapter {
async fn list_contacts(
&self,
address_book_id: &str,
user_id: &str,
user_id: Uuid,
) -> Result<Vec<ContactDto>, DomainError> {
let uuid = Self::parse_uuid(address_book_id, "AddressBook")?;
@@ -700,7 +720,7 @@ impl ContactUseCase for ContactStorageAdapter {
&self,
address_book_id: &str,
query: &str,
user_id: &str,
user_id: Uuid,
) -> Result<Vec<ContactDto>, DomainError> {
let uuid = Self::parse_uuid(address_book_id, "AddressBook")?;
@@ -721,7 +741,10 @@ impl ContactUseCase for ContactStorageAdapter {
let address_book_id = Self::parse_uuid(&dto.address_book_id, "AddressBook")?;
// Check write access
self.check_write_access(&address_book_id, &dto.user_id)
let user_id = Uuid::parse_str(&dto.user_id).map_err(|_| {
DomainError::new(ErrorKind::InvalidInput, "ContactGroup", "Invalid user ID format")
})?;
self.check_write_access(&address_book_id, user_id)
.await?;
let group = ContactGroup::new(address_book_id, dto.name);
@@ -746,7 +769,10 @@ impl ContactUseCase for ContactStorageAdapter {
})?;
// Check write access
self.check_write_access(group.address_book_id(), &update.user_id)
let user_id = Uuid::parse_str(&update.user_id).map_err(|_| {
DomainError::new(ErrorKind::InvalidInput, "ContactGroup", "Invalid user ID format")
})?;
self.check_write_access(group.address_book_id(), user_id)
.await?;
group.set_name(update.name);
@@ -756,7 +782,7 @@ impl ContactUseCase for ContactStorageAdapter {
Ok(ContactGroupDto::from(updated))
}
async fn delete_group(&self, group_id: &str, user_id: &str) -> Result<(), DomainError> {
async fn delete_group(&self, group_id: &str, user_id: Uuid) -> Result<(), DomainError> {
let uuid = Self::parse_uuid(group_id, "ContactGroup")?;
let group = self
@@ -777,7 +803,7 @@ impl ContactUseCase for ContactStorageAdapter {
async fn get_group(
&self,
group_id: &str,
user_id: &str,
user_id: Uuid,
) -> Result<ContactGroupDto, DomainError> {
let uuid = Self::parse_uuid(group_id, "ContactGroup")?;
@@ -799,7 +825,7 @@ impl ContactUseCase for ContactStorageAdapter {
async fn list_groups(
&self,
address_book_id: &str,
user_id: &str,
user_id: Uuid,
) -> Result<Vec<ContactGroupDto>, DomainError> {
let uuid = Self::parse_uuid(address_book_id, "AddressBook")?;
@@ -816,7 +842,7 @@ impl ContactUseCase for ContactStorageAdapter {
async fn add_contact_to_group(
&self,
dto: GroupMembershipDto,
user_id: &str,
user_id: Uuid,
) -> Result<(), DomainError> {
let group_id = Self::parse_uuid(&dto.group_id, "ContactGroup")?;
let contact_id = Self::parse_uuid(&dto.contact_id, "Contact")?;
@@ -841,7 +867,7 @@ impl ContactUseCase for ContactStorageAdapter {
async fn remove_contact_from_group(
&self,
dto: GroupMembershipDto,
user_id: &str,
user_id: Uuid,
) -> Result<(), DomainError> {
let group_id = Self::parse_uuid(&dto.group_id, "ContactGroup")?;
let contact_id = Self::parse_uuid(&dto.contact_id, "Contact")?;
@@ -866,7 +892,7 @@ impl ContactUseCase for ContactStorageAdapter {
async fn list_contacts_in_group(
&self,
group_id: &str,
user_id: &str,
user_id: Uuid,
) -> Result<Vec<ContactDto>, DomainError> {
let uuid = Self::parse_uuid(group_id, "ContactGroup")?;
@@ -889,7 +915,7 @@ impl ContactUseCase for ContactStorageAdapter {
async fn list_groups_for_contact(
&self,
contact_id: &str,
user_id: &str,
user_id: Uuid,
) -> Result<Vec<ContactGroupDto>, DomainError> {
let uuid = Self::parse_uuid(contact_id, "Contact")?;
@@ -910,7 +936,7 @@ impl ContactUseCase for ContactStorageAdapter {
async fn get_contact_vcard(
&self,
contact_id: &str,
user_id: &str,
user_id: Uuid,
) -> Result<String, DomainError> {
let uuid = Self::parse_uuid(contact_id, "Contact")?;
@@ -930,7 +956,7 @@ impl ContactUseCase for ContactStorageAdapter {
async fn get_contacts_as_vcards(
&self,
address_book_id: &str,
user_id: &str,
user_id: Uuid,
) -> Result<Vec<(String, String)>, DomainError> {
let uuid = Self::parse_uuid(address_book_id, "AddressBook")?;
@@ -144,7 +144,7 @@ impl AddressBookRepository for AddressBookPgRepository {
async fn get_address_books_by_owner(
&self,
owner_id: &str,
owner_id: Uuid,
) -> AddressBookRepositoryResult<Vec<AddressBook>> {
let rows = sqlx::query(
r#"
@@ -182,7 +182,7 @@ impl AddressBookRepository for AddressBookPgRepository {
async fn get_shared_address_books(
&self,
user_id: &str,
user_id: Uuid,
) -> AddressBookRepositoryResult<Vec<AddressBook>> {
let rows = sqlx::query(
r#"
@@ -254,7 +254,7 @@ impl AddressBookRepository for AddressBookPgRepository {
async fn share_address_book(
&self,
address_book_id: &Uuid,
user_id: &str,
user_id: Uuid,
can_write: bool,
) -> AddressBookRepositoryResult<()> {
sqlx::query(
@@ -277,7 +277,7 @@ impl AddressBookRepository for AddressBookPgRepository {
async fn unshare_address_book(
&self,
address_book_id: &Uuid,
user_id: &str,
user_id: Uuid,
) -> AddressBookRepositoryResult<()> {
sqlx::query(
r#"
@@ -6,6 +6,7 @@ use crate::domain::entities::app_password::AppPassword;
use chrono::{DateTime, Utc};
use sqlx::PgPool;
use std::sync::Arc;
use uuid::Uuid;
pub struct AppPasswordPgRepository {
pool: Arc<PgPool>,
@@ -48,7 +49,7 @@ impl AppPasswordStoragePort for AppPasswordPgRepository {
Ok(ap)
}
async fn list_by_user(&self, user_id: &str) -> Result<Vec<AppPassword>, DomainError> {
async fn list_by_user(&self, user_id: Uuid) -> Result<Vec<AppPassword>, DomainError> {
let rows = sqlx::query_as::<_, AppPasswordRow>(
r#"
SELECT id, user_id, label, password_hash, prefix, scopes,
@@ -66,7 +67,7 @@ impl AppPasswordStoragePort for AppPasswordPgRepository {
Ok(rows.into_iter().map(|r| r.into()).collect())
}
async fn get_by_id(&self, id: &str) -> Result<AppPassword, DomainError> {
async fn get_by_id(&self, id: Uuid) -> Result<AppPassword, DomainError> {
let row = sqlx::query_as::<_, AppPasswordRow>(
r#"
SELECT id, user_id, label, password_hash, prefix, scopes,
@@ -79,12 +80,12 @@ impl AppPasswordStoragePort for AppPasswordPgRepository {
.fetch_optional(self.pool())
.await
.map_err(|e| DomainError::internal_error("AppPasswordPg", format!("get_by_id: {e}")))?
.ok_or_else(|| DomainError::not_found("AppPassword", id))?;
.ok_or_else(|| DomainError::not_found("AppPassword", id.to_string()))?;
Ok(row.into())
}
async fn get_active_by_user_id(&self, user_id: &str) -> Result<Vec<AppPassword>, DomainError> {
async fn get_active_by_user_id(&self, user_id: Uuid) -> Result<Vec<AppPassword>, DomainError> {
let rows = sqlx::query_as::<_, AppPasswordRow>(
r#"
SELECT id, user_id, label, password_hash, prefix, scopes,
@@ -105,7 +106,7 @@ impl AppPasswordStoragePort for AppPasswordPgRepository {
async fn get_active_by_user_prefix(
&self,
user_id: &str,
user_id: Uuid,
prefix: &str,
) -> Result<Vec<AppPassword>, DomainError> {
let rows = sqlx::query_as::<_, AppPasswordRow>(
@@ -131,7 +132,7 @@ impl AppPasswordStoragePort for AppPasswordPgRepository {
Ok(rows.into_iter().map(|r| r.into()).collect())
}
async fn touch_last_used(&self, id: &str) -> Result<(), DomainError> {
async fn touch_last_used(&self, id: Uuid) -> Result<(), DomainError> {
sqlx::query("UPDATE auth.app_passwords SET last_used_at = NOW() WHERE id = $1")
.bind(id)
.execute(self.pool())
@@ -140,7 +141,7 @@ impl AppPasswordStoragePort for AppPasswordPgRepository {
Ok(())
}
async fn revoke(&self, id: &str, user_id: &str) -> Result<(), DomainError> {
async fn revoke(&self, id: Uuid, user_id: Uuid) -> Result<(), DomainError> {
let result = sqlx::query(
"UPDATE auth.app_passwords SET active = FALSE WHERE id = $1 AND user_id = $2",
)
@@ -151,12 +152,12 @@ impl AppPasswordStoragePort for AppPasswordPgRepository {
.map_err(|e| DomainError::internal_error("AppPasswordPg", format!("revoke: {e}")))?;
if result.rows_affected() == 0 {
return Err(DomainError::not_found("AppPassword", id));
return Err(DomainError::not_found("AppPassword", id.to_string()));
}
Ok(())
}
async fn delete_by_user_and_id(&self, id: &str, user_id: &str) -> Result<bool, DomainError> {
async fn delete_by_user_and_id(&self, id: Uuid, user_id: Uuid) -> Result<bool, DomainError> {
let result = sqlx::query("DELETE FROM auth.app_passwords WHERE id = $1 AND user_id = $2")
.bind(id)
.bind(user_id)
@@ -190,8 +191,8 @@ impl AppPasswordStoragePort for AppPasswordPgRepository {
/// Internal row struct for sqlx mapping.
#[derive(sqlx::FromRow)]
struct AppPasswordRow {
id: String,
user_id: String,
id: Uuid,
user_id: Uuid,
label: String,
password_hash: String,
prefix: String,
@@ -140,7 +140,7 @@ impl CalendarRepository for CalendarPgRepository {
async fn list_calendars_by_owner(
&self,
owner_id: &str,
owner_id: Uuid,
) -> CalendarRepositoryResult<Vec<Calendar>> {
let rows = sqlx::query(
r#"
@@ -180,7 +180,7 @@ impl CalendarRepository for CalendarPgRepository {
async fn find_calendar_by_name_and_owner(
&self,
name: &str,
owner_id: &str,
owner_id: Uuid,
) -> CalendarRepositoryResult<Calendar> {
let row = sqlx::query(
r#"
@@ -218,7 +218,7 @@ impl CalendarRepository for CalendarPgRepository {
async fn list_calendars_shared_with_user(
&self,
user_id: &str,
user_id: Uuid,
) -> CalendarRepositoryResult<Vec<Calendar>> {
let rows = sqlx::query(
r#"
@@ -299,7 +299,7 @@ impl CalendarRepository for CalendarPgRepository {
async fn user_has_calendar_access(
&self,
calendar_id: &Uuid,
user_id: &str,
user_id: Uuid,
) -> CalendarRepositoryResult<bool> {
// Check if the user is the owner of the calendar or has a share
let row = sqlx::query(
@@ -327,7 +327,7 @@ impl CalendarRepository for CalendarPgRepository {
async fn share_calendar(
&self,
calendar_id: &Uuid,
user_id: &str,
user_id: Uuid,
access_level: &str,
) -> CalendarRepositoryResult<()> {
// Validate access level
@@ -358,7 +358,7 @@ impl CalendarRepository for CalendarPgRepository {
async fn remove_calendar_sharing(
&self,
calendar_id: &Uuid,
user_id: &str,
user_id: Uuid,
) -> CalendarRepositoryResult<()> {
sqlx::query(
r#"
@@ -2,6 +2,7 @@
use sqlx::{PgPool, Row};
use std::sync::Arc;
use uuid::Uuid;
use crate::application::ports::auth_ports::DeviceCodeStoragePort;
use crate::common::errors::{DomainError, ErrorKind};
@@ -28,7 +29,7 @@ impl DeviceCodePgRepository {
let status = DeviceCodeStatus::parse(&status_str).unwrap_or(DeviceCodeStatus::Expired);
Ok(DeviceCode::from_raw(
row.try_get("id").unwrap_or_default(),
row.try_get("id").unwrap(),
row.try_get("device_code").unwrap_or_default(),
row.try_get("user_code").unwrap_or_default(),
row.try_get("client_name").unwrap_or_default(),
@@ -212,7 +213,7 @@ impl DeviceCodeStoragePort for DeviceCodePgRepository {
Ok(result.rows_affected())
}
async fn list_by_user(&self, user_id: &str) -> Result<Vec<DeviceCode>, DomainError> {
async fn list_by_user(&self, user_id: Uuid) -> Result<Vec<DeviceCode>, DomainError> {
let rows = sqlx::query(
r#"
SELECT id, device_code, user_code, client_name, scopes,
@@ -239,7 +240,7 @@ impl DeviceCodeStoragePort for DeviceCodePgRepository {
rows.iter().map(Self::map_row).collect()
}
async fn delete_by_id(&self, id: &str) -> Result<(), DomainError> {
async fn delete_by_id(&self, id: Uuid) -> Result<(), DomainError> {
sqlx::query("DELETE FROM auth.device_codes WHERE id = $1")
.bind(id)
.execute(self.pool.as_ref())
@@ -20,9 +20,7 @@ impl FavoritesPgRepository {
}
impl FavoritesRepositoryPort for FavoritesPgRepository {
async fn get_favorites(&self, user_id: &str) -> Result<Vec<FavoriteItemDto>> {
let user_uuid = Uuid::parse_str(user_id)?;
async fn get_favorites(&self, user_id: Uuid) -> Result<Vec<FavoriteItemDto>> {
let rows = sqlx::query(
r#"
SELECT
@@ -41,12 +39,12 @@ impl FavoritesRepositoryPort for FavoritesPgRepository {
AND f.id = uf.item_id::UUID
LEFT JOIN storage.folders fld ON uf.item_type = 'folder'
AND fld.id = uf.item_id::UUID
WHERE uf.user_id = $1::TEXT
WHERE uf.user_id = $1
ORDER BY uf.created_at DESC
LIMIT 500
"#,
)
.bind(user_uuid)
.bind(user_id)
.fetch_all(&*self.db_pool)
.await
.map_err(|e| {
@@ -85,17 +83,15 @@ impl FavoritesRepositoryPort for FavoritesPgRepository {
Ok(favorites)
}
async fn add_favorite(&self, user_id: &str, item_id: &str, item_type: &str) -> Result<()> {
let user_uuid = Uuid::parse_str(user_id)?;
async fn add_favorite(&self, user_id: Uuid, item_id: &str, item_type: &str) -> Result<()> {
sqlx::query(
r#"
INSERT INTO auth.user_favorites (user_id, item_id, item_type)
VALUES ($1::TEXT, $2, $3)
VALUES ($1, $2, $3)
ON CONFLICT (user_id, item_id, item_type) DO NOTHING
"#,
)
.bind(user_uuid)
.bind(user_id)
.bind(item_id)
.bind(item_type)
.execute(&*self.db_pool)
@@ -112,16 +108,14 @@ impl FavoritesRepositoryPort for FavoritesPgRepository {
Ok(())
}
async fn remove_favorite(&self, user_id: &str, item_id: &str, item_type: &str) -> Result<bool> {
let user_uuid = Uuid::parse_str(user_id)?;
async fn remove_favorite(&self, user_id: Uuid, item_id: &str, item_type: &str) -> Result<bool> {
let result = sqlx::query(
r#"
DELETE FROM auth.user_favorites
WHERE user_id = $1::TEXT AND item_id = $2 AND item_type = $3
WHERE user_id = $1 AND item_id = $2 AND item_type = $3
"#,
)
.bind(user_uuid)
.bind(user_id)
.bind(item_id)
.bind(item_type)
.execute(&*self.db_pool)
@@ -138,18 +132,16 @@ impl FavoritesRepositoryPort for FavoritesPgRepository {
Ok(result.rows_affected() > 0)
}
async fn is_favorite(&self, user_id: &str, item_id: &str, item_type: &str) -> Result<bool> {
let user_uuid = Uuid::parse_str(user_id)?;
async fn is_favorite(&self, user_id: Uuid, item_id: &str, item_type: &str) -> Result<bool> {
let row = sqlx::query(
r#"
SELECT EXISTS (
SELECT 1 FROM auth.user_favorites
WHERE user_id = $1::TEXT AND item_id = $2 AND item_type = $3
WHERE user_id = $1 AND item_id = $2 AND item_type = $3
) AS "is_favorite"
"#,
)
.bind(user_uuid)
.bind(user_id)
.bind(item_id)
.bind(item_type)
.fetch_one(&*self.db_pool)
@@ -166,13 +158,11 @@ impl FavoritesRepositoryPort for FavoritesPgRepository {
Ok(row.try_get("is_favorite").unwrap_or(false))
}
async fn add_favorites_batch(&self, user_id: &str, items: &[(String, String)]) -> Result<u64> {
async fn add_favorites_batch(&self, user_id: Uuid, items: &[(String, String)]) -> Result<u64> {
if items.is_empty() {
return Ok(0);
}
let user_uuid = Uuid::parse_str(user_id)?;
// Validate all item_types upfront
for (_, item_type) in items {
if item_type != "file" && item_type != "folder" {
@@ -210,7 +200,7 @@ impl FavoritesRepositoryPort for FavoritesPgRepository {
query.push_str(", ");
}
query.push_str(&format!(
"(${}::TEXT, ${}, ${})",
"(${}, ${}, ${})",
param_idx,
param_idx + 1,
param_idx + 2
@@ -222,7 +212,7 @@ impl FavoritesRepositoryPort for FavoritesPgRepository {
let mut q = sqlx::query(&query);
for (item_id, item_type) in chunk {
q = q.bind(user_uuid).bind(item_id).bind(item_type);
q = q.bind(user_id).bind(item_id).bind(item_type);
}
let result = q.execute(&mut *tx).await.map_err(|e| {
@@ -251,22 +241,20 @@ impl FavoritesRepositoryPort for FavoritesPgRepository {
async fn batch_check_favorites(
&self,
user_id: &str,
user_id: Uuid,
item_ids: &[(&str, &str)],
) -> Result<HashSet<String>> {
if item_ids.is_empty() {
return Ok(HashSet::new());
}
let user_uuid = Uuid::parse_str(user_id)?;
// Collect just the IDs for the IN clause
let ids: Vec<String> = item_ids.iter().map(|(id, _)| id.to_string()).collect();
let rows = sqlx::query(
"SELECT item_id FROM auth.user_favorites WHERE user_id = $1::TEXT AND item_id = ANY($2)",
"SELECT item_id FROM auth.user_favorites WHERE user_id = $1 AND item_id = ANY($2)",
)
.bind(user_uuid)
.bind(user_id)
.bind(&ids)
.fetch_all(&*self.db_pool)
.await
@@ -35,6 +35,7 @@ use crate::common::errors::DomainError;
use crate::domain::entities::file::File;
use crate::domain::services::path_service::StoragePath;
use crate::infrastructure::services::dedup_service::DedupService;
use uuid::Uuid;
/// Type alias for file metadata rows from SQL queries.
type FileRow = (
@@ -166,7 +167,7 @@ impl FileBlobReadRepository {
/// `sort_date` epoch for each file (used as pagination cursor).
pub async fn list_media_files(
&self,
owner_id: &str,
owner_id: Uuid,
before: Option<i64>,
limit: i64,
) -> Result<(Vec<File>, Vec<i64>), DomainError> {
@@ -255,7 +256,7 @@ impl FileReadPort for FileBlobReadRepository {
)
}
async fn get_file_for_owner(&self, id: &str, owner_id: &str) -> Result<File, DomainError> {
async fn get_file_for_owner(&self, id: &str, owner_id: Uuid) -> Result<File, DomainError> {
let row = sqlx::query_as::<
_,
(
@@ -286,7 +287,7 @@ impl FileReadPort for FileBlobReadRepository {
"#,
)
.bind(id)
.bind(owner_id)
.bind(owner_id.to_string())
.fetch_optional(self.pool.as_ref())
.await
.map_err(|e| DomainError::internal_error("FileBlobRead", format!("get_for_owner: {e}")))?
@@ -350,8 +351,9 @@ impl FileReadPort for FileBlobReadRepository {
async fn list_files_for_owner(
&self,
folder_id: Option<&str>,
owner_id: &str,
owner_id: Uuid,
) -> Result<Vec<File>, DomainError> {
let owner_str = owner_id.to_string();
let rows: Vec<FileRow> = if let Some(fid) = folder_id {
sqlx::query_as(
r#"
@@ -368,7 +370,7 @@ impl FileReadPort for FileBlobReadRepository {
"#,
)
.bind(fid)
.bind(owner_id)
.bind(&owner_str)
.fetch_all(self.pool.as_ref())
.await
} else {
@@ -386,7 +388,7 @@ impl FileReadPort for FileBlobReadRepository {
ORDER BY fi.name
"#,
)
.bind(owner_id)
.bind(&owner_str)
.fetch_all(self.pool.as_ref())
.await
}
@@ -468,10 +470,11 @@ impl FileReadPort for FileBlobReadRepository {
async fn list_files_batch_for_owner(
&self,
folder_id: Option<&str>,
owner_id: &str,
owner_id: Uuid,
offset: i64,
limit: i64,
) -> Result<Vec<File>, DomainError> {
let owner_str = owner_id.to_string();
let rows: Vec<FileRow> = if let Some(fid) = folder_id {
sqlx::query_as(
r#"
@@ -491,7 +494,7 @@ impl FileReadPort for FileBlobReadRepository {
.bind(fid)
.bind(limit)
.bind(offset)
.bind(owner_id)
.bind(&owner_str)
.fetch_all(self.pool.as_ref())
.await
} else {
@@ -512,7 +515,7 @@ impl FileReadPort for FileBlobReadRepository {
)
.bind(limit)
.bind(offset)
.bind(owner_id)
.bind(&owner_str)
.fetch_all(self.pool.as_ref())
.await
}
@@ -753,7 +756,7 @@ impl FileReadPort for FileBlobReadRepository {
&self,
folder_id: Option<&str>,
criteria: &SearchCriteriaDto,
user_id: &str,
user_id: Uuid,
) -> Result<(Vec<File>, usize), DomainError> {
let offset = criteria.offset as i64;
let limit = criteria.limit as i64;
@@ -822,7 +825,7 @@ impl FileReadPort for FileBlobReadRepository {
i64,
),
>(&sql)
.bind(user_id);
.bind(user_id.to_string());
if let Some(fid) = folder_id {
query = query.bind(fid);
@@ -867,7 +870,7 @@ impl FileReadPort for FileBlobReadRepository {
&self,
root_folder_id: Option<&str>,
criteria: &SearchCriteriaDto,
user_id: &str,
user_id: Uuid,
) -> Result<(Vec<File>, usize), DomainError> {
// When no root folder specified, delegate to existing paginated search
let root_id = match root_folder_id {
@@ -983,7 +986,7 @@ impl FileReadPort for FileBlobReadRepository {
i64,
),
>(&sql)
.bind(user_id)
.bind(user_id.to_string())
.bind(root_id);
if let Some(name) = &criteria.name_contains
@@ -1043,7 +1046,7 @@ impl FileReadPort for FileBlobReadRepository {
&self,
folder_id: Option<&str>,
criteria: &SearchCriteriaDto,
user_id: &str,
user_id: Uuid,
) -> Result<usize, DomainError> {
let (_, count) = self
.search_files_paginated(folder_id, criteria, user_id)
@@ -19,9 +19,7 @@ impl RecentItemsPgRepository {
}
impl RecentItemsRepositoryPort for RecentItemsPgRepository {
async fn get_recent_items(&self, user_id: &str, limit: i32) -> Result<Vec<RecentItemDto>> {
let user_uuid = Uuid::parse_str(user_id)?;
async fn get_recent_items(&self, user_id: Uuid, limit: i32) -> Result<Vec<RecentItemDto>> {
let rows = sqlx::query(
r#"
SELECT
@@ -39,12 +37,12 @@ impl RecentItemsRepositoryPort for RecentItemsPgRepository {
AND f.id = ur.item_id::UUID
LEFT JOIN storage.folders fld ON ur.item_type = 'folder'
AND fld.id = ur.item_id::UUID
WHERE ur.user_id = $1::TEXT
WHERE ur.user_id = $1
ORDER BY ur.accessed_at DESC
LIMIT $2
"#,
)
.bind(user_uuid)
.bind(user_id)
.bind(limit)
.fetch_all(&*self.db_pool)
.await
@@ -83,18 +81,16 @@ impl RecentItemsRepositoryPort for RecentItemsPgRepository {
Ok(items)
}
async fn upsert_access(&self, user_id: &str, item_id: &str, item_type: &str) -> Result<()> {
let user_uuid = Uuid::parse_str(user_id)?;
async fn upsert_access(&self, user_id: Uuid, item_id: &str, item_type: &str) -> Result<()> {
sqlx::query(
r#"
INSERT INTO auth.user_recent_files (user_id, item_id, item_type, accessed_at)
VALUES ($1::TEXT, $2, $3, CURRENT_TIMESTAMP)
VALUES ($1, $2, $3, CURRENT_TIMESTAMP)
ON CONFLICT (user_id, item_id, item_type)
DO UPDATE SET accessed_at = CURRENT_TIMESTAMP
"#,
)
.bind(user_uuid)
.bind(user_id)
.bind(item_id)
.bind(item_type)
.execute(&*self.db_pool)
@@ -111,16 +107,14 @@ impl RecentItemsRepositoryPort for RecentItemsPgRepository {
Ok(())
}
async fn remove_item(&self, user_id: &str, item_id: &str, item_type: &str) -> Result<bool> {
let user_uuid = Uuid::parse_str(user_id)?;
async fn remove_item(&self, user_id: Uuid, item_id: &str, item_type: &str) -> Result<bool> {
let result = sqlx::query(
r#"
DELETE FROM auth.user_recent_files
WHERE user_id = $1::TEXT AND item_id = $2 AND item_type = $3
WHERE user_id = $1 AND item_id = $2 AND item_type = $3
"#,
)
.bind(user_uuid)
.bind(user_id)
.bind(item_id)
.bind(item_type)
.execute(&*self.db_pool)
@@ -137,16 +131,14 @@ impl RecentItemsRepositoryPort for RecentItemsPgRepository {
Ok(result.rows_affected() > 0)
}
async fn clear_all(&self, user_id: &str) -> Result<()> {
let user_uuid = Uuid::parse_str(user_id)?;
async fn clear_all(&self, user_id: Uuid) -> Result<()> {
sqlx::query(
r#"
DELETE FROM auth.user_recent_files
WHERE user_id = $1::TEXT
WHERE user_id = $1
"#,
)
.bind(user_uuid)
.bind(user_id)
.execute(&*self.db_pool)
.await
.map_err(|e| {
@@ -161,21 +153,19 @@ impl RecentItemsRepositoryPort for RecentItemsPgRepository {
Ok(())
}
async fn prune(&self, user_id: &str, max_items: i32) -> Result<()> {
let user_uuid = Uuid::parse_str(user_id)?;
async fn prune(&self, user_id: Uuid, max_items: i32) -> Result<()> {
sqlx::query(
r#"
DELETE FROM auth.user_recent_files
WHERE id IN (
SELECT id FROM auth.user_recent_files
WHERE user_id = $1::TEXT
WHERE user_id = $1
ORDER BY accessed_at DESC
OFFSET $2
)
"#,
)
.bind(user_uuid)
.bind(user_id)
.bind(max_items)
.execute(&*self.db_pool)
.await
@@ -2,6 +2,7 @@ use chrono::Utc;
use futures::future::BoxFuture;
use sqlx::{PgPool, Row};
use std::sync::Arc;
use uuid::Uuid;
use crate::application::ports::auth_ports::SessionStoragePort;
use crate::common::errors::DomainError;
@@ -104,7 +105,7 @@ impl SessionRepository for SessionPgRepository {
}
/// Gets a session by ID
async fn get_session_by_id(&self, id: &str) -> SessionRepositoryResult<Session> {
async fn get_session_by_id(&self, id: Uuid) -> SessionRepositoryResult<Session> {
let row = sqlx::query(
r#"
SELECT
@@ -165,7 +166,7 @@ impl SessionRepository for SessionPgRepository {
/// Gets all sessions for a user
async fn get_sessions_by_user_id(
&self,
user_id: &str,
user_id: Uuid,
) -> SessionRepositoryResult<Vec<Session>> {
let rows = sqlx::query(
r#"
@@ -202,8 +203,8 @@ impl SessionRepository for SessionPgRepository {
}
/// Revokes a specific session using a transaction
async fn revoke_session(&self, session_id: &str) -> SessionRepositoryResult<()> {
let id = session_id.to_string(); // Clone for use in closure
async fn revoke_session(&self, session_id: Uuid) -> SessionRepositoryResult<()> {
let id = session_id; // Copy for use in closure
with_transaction(&self.pool, "revoke_session", |tx| {
Box::pin(async move {
@@ -216,14 +217,14 @@ impl SessionRepository for SessionPgRepository {
RETURNING user_id
"#,
)
.bind(&id)
.bind(id)
.fetch_optional(&mut **tx)
.await
.map_err(Self::map_sqlx_error)?;
// If we found the session, we can log a security event
if let Some(row) = result {
let user_id: String = row.try_get("user_id").unwrap_or_default();
let user_id: Uuid = row.try_get("user_id").unwrap_or_default();
// Log security event (in a security table)
// This is optional but shows how additional operations
@@ -238,8 +239,8 @@ impl SessionRepository for SessionPgRepository {
}
/// Revokes all sessions for a user using a transaction
async fn revoke_all_user_sessions(&self, user_id: &str) -> SessionRepositoryResult<u64> {
let user_id_clone = user_id.to_string(); // Clone for use in closure
async fn revoke_all_user_sessions(&self, user_id: Uuid) -> SessionRepositoryResult<u64> {
let user_id_copy = user_id; // Copy for use in closure
with_transaction(&self.pool, "revoke_all_user_sessions", |tx| {
Box::pin(async move {
@@ -251,7 +252,7 @@ impl SessionRepository for SessionPgRepository {
WHERE user_id = $1 AND revoked = false
"#,
)
.bind(&user_id_clone)
.bind(user_id_copy)
.execute(&mut **tx)
.await
.map_err(Self::map_sqlx_error)?;
@@ -260,7 +261,7 @@ impl SessionRepository for SessionPgRepository {
// Log security event
if affected > 0 {
tracing::info!("Revoked {} sessions for user {}", affected, user_id_clone);
tracing::info!("Revoked {} sessions for user {}", affected, user_id_copy);
}
Ok(affected)
@@ -305,13 +306,13 @@ impl SessionStoragePort for SessionPgRepository {
.map_err(DomainError::from)
}
async fn revoke_session(&self, session_id: &str) -> Result<(), DomainError> {
async fn revoke_session(&self, session_id: Uuid) -> Result<(), DomainError> {
SessionRepository::revoke_session(self, session_id)
.await
.map_err(DomainError::from)
}
async fn revoke_all_user_sessions(&self, user_id: &str) -> Result<u64, DomainError> {
async fn revoke_all_user_sessions(&self, user_id: Uuid) -> Result<u64, DomainError> {
SessionRepository::revoke_all_user_sessions(self, user_id)
.await
.map_err(DomainError::from)
@@ -1,6 +1,7 @@
use sqlx::PgPool;
use std::collections::HashMap;
use std::sync::Arc;
use uuid::Uuid;
use crate::common::errors::{DomainError, ErrorKind};
use crate::domain::repositories::settings_repository::SettingsRepository;
@@ -60,7 +61,7 @@ impl SettingsRepository for SettingsPgRepository {
value: &str,
category: &str,
is_secret: bool,
updated_by: Option<&str>,
updated_by: Option<Uuid>,
) -> Result<(), DomainError> {
sqlx::query(
"INSERT INTO auth.admin_settings (key, value, category, is_secret, updated_by, updated_at)
@@ -102,7 +103,7 @@ impl SettingsRepository for SettingsPgRepository {
///
/// Only the first caller that inserts the row gets `rows_affected == 1`;
/// concurrent callers see 0 rows affected and receive `false`.
async fn try_claim_initialization(&self, admin_user_id: &str) -> Result<bool, DomainError> {
async fn try_claim_initialization(&self, admin_user_id: Uuid) -> Result<bool, DomainError> {
let result = sqlx::query(
"INSERT INTO auth.admin_settings (key, value, category, is_secret, updated_by, updated_at)
VALUES ('system_initialized', 'true', 'system', false, $1, NOW())
@@ -1,5 +1,6 @@
use sqlx::{PgPool, Row};
use std::sync::Arc;
use uuid::Uuid;
use crate::{
application::ports::share_ports::ShareStoragePort,
@@ -37,7 +38,7 @@ impl SharePgRepository {
/// Maps a [`sqlx::postgres::PgRow`] to the domain [`Share`] entity.
fn row_to_entity(row: &sqlx::postgres::PgRow) -> Result<Share, DomainError> {
let id: String = row
let id: Uuid = row
.try_get("id")
.map_err(|e| DomainError::internal_error("Share", format!("Failed to read id: {e}")))?;
let item_id: String = row.try_get("item_id").map_err(|e| {
@@ -58,7 +59,7 @@ impl SharePgRepository {
let created_at: i64 = row.try_get("created_at").map_err(|e| {
DomainError::internal_error("Share", format!("Failed to read created_at: {e}"))
})?;
let created_by: String = row.try_get("created_by").map_err(|e| {
let created_by: Uuid = row.try_get("created_by").map_err(|e| {
DomainError::internal_error("Share", format!("Failed to read created_by: {e}"))
})?;
let access_count: i64 = row.try_get("access_count").unwrap_or(0);
@@ -93,7 +94,7 @@ impl ShareStoragePort for SharePgRepository {
expires_at, permissions_read, permissions_write, permissions_reshare,
created_at, created_by, access_count)
VALUES
($1::UUID, $2, $3, $4, $5, $6,
($1, $2, $3, $4, $5, $6,
$7, $8, $9, $10,
$11, $12, $13)
ON CONFLICT (id) DO UPDATE SET
@@ -105,7 +106,7 @@ impl ShareStoragePort for SharePgRepository {
permissions_reshare = EXCLUDED.permissions_reshare,
access_count = EXCLUDED.access_count
RETURNING
id::TEXT, item_id, item_name, item_type, token, password_hash,
id, item_id, item_name, item_type, token, password_hash,
expires_at, permissions_read, permissions_write, permissions_reshare,
created_at, created_by, access_count
"#,
@@ -136,7 +137,7 @@ impl ShareStoragePort for SharePgRepository {
async fn find_share_by_token(&self, token: &str) -> Result<Share, DomainError> {
let row = sqlx::query(
r#"
SELECT id::TEXT, item_id, item_name, item_type, token, password_hash,
SELECT id, item_id, item_name, item_type, token, password_hash,
expires_at, permissions_read, permissions_write, permissions_reshare,
created_at, created_by, access_count
FROM storage.shares
@@ -162,16 +163,16 @@ impl ShareStoragePort for SharePgRepository {
async fn find_share_by_id_for_user(
&self,
id: &str,
user_id: &str,
id: Uuid,
user_id: Uuid,
) -> Result<Share, DomainError> {
let row = sqlx::query(
r#"
SELECT id::TEXT, item_id, item_name, item_type, token, password_hash,
SELECT id, item_id, item_name, item_type, token, password_hash,
expires_at, permissions_read, permissions_write, permissions_reshare,
created_at, created_by, access_count
FROM storage.shares
WHERE id = $1::UUID AND created_by = $2
WHERE id = $1 AND created_by = $2
"#,
)
.bind(id)
@@ -193,9 +194,9 @@ impl ShareStoragePort for SharePgRepository {
}
}
async fn delete_share_for_user(&self, id: &str, user_id: &str) -> Result<(), DomainError> {
async fn delete_share_for_user(&self, id: Uuid, user_id: Uuid) -> Result<(), DomainError> {
let result =
sqlx::query("DELETE FROM storage.shares WHERE id = $1::UUID AND created_by = $2")
sqlx::query("DELETE FROM storage.shares WHERE id = $1 AND created_by = $2")
.bind(id)
.bind(user_id)
.execute(&*self.db_pool)
@@ -220,11 +221,11 @@ impl ShareStoragePort for SharePgRepository {
&self,
item_id: &str,
item_type: &ShareItemType,
user_id: &str,
user_id: Uuid,
) -> Result<Vec<Share>, DomainError> {
let rows = sqlx::query(
r#"
SELECT id::TEXT, item_id, item_name, item_type, token, password_hash,
SELECT id, item_id, item_name, item_type, token, password_hash,
expires_at, permissions_read, permissions_write, permissions_reshare,
created_at, created_by, access_count
FROM storage.shares
@@ -256,9 +257,9 @@ impl ShareStoragePort for SharePgRepository {
permissions_write = $6,
permissions_reshare = $7,
access_count = $8
WHERE id = $1::UUID
WHERE id = $1
RETURNING
id::TEXT, item_id, item_name, item_type, token, password_hash,
id, item_id, item_name, item_type, token, password_hash,
expires_at, permissions_read, permissions_write, permissions_reshare,
created_at, created_by, access_count
"#,
@@ -289,14 +290,14 @@ impl ShareStoragePort for SharePgRepository {
async fn find_shares_by_user(
&self,
user_id: &str,
user_id: Uuid,
offset: usize,
limit: usize,
) -> Result<(Vec<Share>, usize), DomainError> {
// Single query with window function — count + rows in one roundtrip
let rows = sqlx::query(
r#"
SELECT id::TEXT, item_id, item_name, item_type, token, password_hash,
SELECT id, item_id, item_name, item_type, token, password_hash,
expires_at, permissions_read, permissions_write, permissions_reshare,
created_at, created_by, access_count,
COUNT(*) OVER() AS total_count
@@ -50,7 +50,7 @@ impl TrashDbRepository {
id: Uuid,
name: String,
item_type: String,
user_id: String,
user_id: Uuid,
trashed_at: Option<DateTime<Utc>>,
) -> TrashedItem {
let trashed_at = trashed_at.unwrap_or_else(Utc::now);
@@ -61,14 +61,12 @@ impl TrashDbRepository {
_ => TrashedItemType::File,
};
let user_uuid = Uuid::parse_str(&user_id).unwrap_or_else(|_| Uuid::nil());
// In the soft-delete model, the trash entry ID is the same as the
// original item ID since there is no separate trash table.
TrashedItem::from_raw(
id, // trash entry id (same as original)
id, // original item id
user_uuid, // owner
id, // trash entry id (same as original)
id, // original item id
user_id, // owner
item_type_enum,
name.clone(),
String::new(), // original_path — not stored separately in soft-delete model
@@ -87,7 +85,7 @@ impl TrashRepository for TrashDbRepository {
}
async fn get_trash_items(&self, user_id: &Uuid) -> Result<Vec<TrashedItem>> {
let rows = sqlx::query_as::<_, (Uuid, String, String, String, Option<DateTime<Utc>>)>(
let rows = sqlx::query_as::<_, (Uuid, String, String, Uuid, Option<DateTime<Utc>>)>(
r#"
SELECT id, name, item_type, user_id, trashed_at
FROM storage.trash_items
@@ -95,7 +93,7 @@ impl TrashRepository for TrashDbRepository {
ORDER BY trashed_at DESC
"#,
)
.bind(user_id.to_string())
.bind(user_id)
.fetch_all(self.pool.as_ref())
.await
.map_err(|e| DomainError::internal_error("TrashDb", format!("list: {e}")))?;
@@ -109,7 +107,7 @@ impl TrashRepository for TrashDbRepository {
}
async fn get_trash_item(&self, id: &Uuid, user_id: &Uuid) -> Result<Option<TrashedItem>> {
let row = sqlx::query_as::<_, (Uuid, String, String, String, Option<DateTime<Utc>>)>(
let row = sqlx::query_as::<_, (Uuid, String, String, Uuid, Option<DateTime<Utc>>)>(
r#"
SELECT id, name, item_type, user_id, trashed_at
FROM storage.trash_items
@@ -117,7 +115,7 @@ impl TrashRepository for TrashDbRepository {
"#,
)
.bind(id)
.bind(user_id.to_string())
.bind(user_id)
.fetch_optional(self.pool.as_ref())
.await
.map_err(|e| DomainError::internal_error("TrashDb", format!("get: {e}")))?;
@@ -144,14 +142,14 @@ impl TrashRepository for TrashDbRepository {
async fn clear_trash(&self, user_id: &Uuid) -> Result<()> {
// Delete all trashed files for this user
sqlx::query("DELETE FROM storage.files WHERE user_id = $1 AND is_trashed = TRUE")
.bind(user_id.to_string())
.bind(user_id)
.execute(self.pool.as_ref())
.await
.map_err(|e| DomainError::internal_error("TrashDb", format!("clear files: {e}")))?;
// Delete all trashed folders for this user
sqlx::query("DELETE FROM storage.folders WHERE user_id = $1 AND is_trashed = TRUE")
.bind(user_id.to_string())
.bind(user_id)
.execute(self.pool.as_ref())
.await
.map_err(|e| DomainError::internal_error("TrashDb", format!("clear folders: {e}")))?;
@@ -1,6 +1,7 @@
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;
@@ -101,7 +102,7 @@ impl UserRepository for UserPgRepository {
}
/// Gets a user by ID
async fn get_user_by_id(&self, id: &str) -> UserRepositoryResult<User> {
async fn get_user_by_id(&self, id: Uuid) -> UserRepositoryResult<User> {
let row = sqlx::query(
r#"
SELECT
@@ -278,7 +279,7 @@ impl UserRepository for UserPgRepository {
/// Updates only the storage usage of a user
async fn update_storage_usage(
&self,
user_id: &str,
user_id: Uuid,
usage_bytes: i64,
) -> UserRepositoryResult<()> {
sqlx::query(
@@ -300,7 +301,7 @@ impl UserRepository for UserPgRepository {
}
/// Updates the last login date
async fn update_last_login(&self, user_id: &str) -> UserRepositoryResult<()> {
async fn update_last_login(&self, user_id: Uuid) -> UserRepositoryResult<()> {
sqlx::query(
r#"
UPDATE auth.users
@@ -423,7 +424,7 @@ impl UserRepository for UserPgRepository {
/// Activates or deactivates a user
async fn set_user_active_status(
&self,
user_id: &str,
user_id: Uuid,
active: bool,
) -> UserRepositoryResult<()> {
sqlx::query(
@@ -447,7 +448,7 @@ impl UserRepository for UserPgRepository {
/// Changes a user's password
async fn change_password(
&self,
user_id: &str,
user_id: Uuid,
password_hash: &str,
) -> UserRepositoryResult<()> {
sqlx::query(
@@ -469,7 +470,7 @@ impl UserRepository for UserPgRepository {
}
/// Changes a user's role
async fn change_role(&self, user_id: &str, role: UserRole) -> UserRepositoryResult<()> {
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();
@@ -542,7 +543,7 @@ impl UserRepository for UserPgRepository {
}
/// Deletes a user
async fn delete_user(&self, user_id: &str) -> UserRepositoryResult<()> {
async fn delete_user(&self, user_id: Uuid) -> UserRepositoryResult<()> {
sqlx::query(
r#"
DELETE FROM auth.users
@@ -606,7 +607,7 @@ impl UserRepository for UserPgRepository {
/// Updates a user's storage quota
async fn update_storage_quota(
&self,
user_id: &str,
user_id: Uuid,
quota_bytes: i64,
) -> UserRepositoryResult<()> {
sqlx::query(
@@ -675,7 +676,7 @@ impl UserStoragePort for UserPgRepository {
.map_err(DomainError::from)
}
async fn get_user_by_id(&self, id: &str) -> Result<User, DomainError> {
async fn get_user_by_id(&self, id: Uuid) -> Result<User, DomainError> {
UserRepository::get_user_by_id(self, id)
.await
.map_err(DomainError::from)
@@ -701,7 +702,7 @@ impl UserStoragePort for UserPgRepository {
async fn update_storage_usage(
&self,
user_id: &str,
user_id: Uuid,
usage_bytes: i64,
) -> Result<(), DomainError> {
UserRepository::update_storage_usage(self, user_id, usage_bytes)
@@ -727,13 +728,13 @@ impl UserStoragePort for UserPgRepository {
.map_err(DomainError::from)
}
async fn delete_user(&self, user_id: &str) -> Result<(), DomainError> {
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: &str, password_hash: &str) -> Result<(), DomainError> {
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)
@@ -749,13 +750,13 @@ impl UserStoragePort for UserPgRepository {
.map_err(DomainError::from)
}
async fn set_user_active_status(&self, user_id: &str, active: bool) -> Result<(), DomainError> {
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: &str, role: &str) -> Result<(), DomainError> {
async fn change_role(&self, user_id: Uuid, role: &str) -> Result<(), DomainError> {
let user_role = match role {
"admin" => UserRole::Admin,
_ => UserRole::User,
@@ -767,7 +768,7 @@ impl UserStoragePort for UserPgRepository {
async fn update_storage_quota(
&self,
user_id: &str,
user_id: Uuid,
quota_bytes: i64,
) -> Result<(), DomainError> {
UserRepository::update_storage_quota(self, user_id, quota_bytes)
@@ -801,7 +801,7 @@ impl ChunkedUploadService {
impl ChunkedUploadPort for ChunkedUploadService {
async fn create_session(
&self,
user_id: &str,
user_id: Uuid,
filename: String,
folder_id: Option<String>,
content_type: String,
@@ -809,7 +809,7 @@ impl ChunkedUploadPort for ChunkedUploadService {
chunk_size: Option<usize>,
) -> Result<CreateUploadResponseDto, DomainError> {
self.create_session_inner(
user_id.to_owned(),
user_id.to_string(),
filename,
folder_id,
content_type,
@@ -823,12 +823,12 @@ impl ChunkedUploadPort for ChunkedUploadService {
async fn upload_chunk(
&self,
upload_id: &str,
user_id: &str,
user_id: Uuid,
chunk_index: usize,
data: bytes::Bytes,
checksum: Option<String>,
) -> Result<ChunkUploadResponseDto, DomainError> {
self.upload_chunk_inner(upload_id, user_id, chunk_index, data, checksum)
self.upload_chunk_inner(upload_id, &user_id.to_string(), chunk_index, data, checksum)
.await
.map_err(|e| DomainError::new(ErrorKind::InternalError, "ChunkedUpload", e))
}
@@ -836,9 +836,9 @@ impl ChunkedUploadPort for ChunkedUploadService {
async fn get_status(
&self,
upload_id: &str,
user_id: &str,
user_id: Uuid,
) -> Result<UploadStatusResponseDto, DomainError> {
self.get_status_inner(upload_id, user_id)
self.get_status_inner(upload_id, &user_id.to_string())
.await
.map_err(|e| DomainError::new(ErrorKind::NotFound, "ChunkedUpload", e))
}
@@ -846,21 +846,21 @@ impl ChunkedUploadPort for ChunkedUploadService {
async fn complete_upload(
&self,
upload_id: &str,
user_id: &str,
user_id: Uuid,
) -> Result<(PathBuf, String, Option<String>, String, u64, String), DomainError> {
self.complete_upload_inner(upload_id, user_id)
self.complete_upload_inner(upload_id, &user_id.to_string())
.await
.map_err(|e| DomainError::new(ErrorKind::InternalError, "ChunkedUpload", e))
}
async fn finalize_upload(&self, upload_id: &str, user_id: &str) -> Result<(), DomainError> {
self.finalize_upload_inner(upload_id, user_id)
async fn finalize_upload(&self, upload_id: &str, user_id: Uuid) -> Result<(), DomainError> {
self.finalize_upload_inner(upload_id, &user_id.to_string())
.await
.map_err(|e| DomainError::new(ErrorKind::InternalError, "ChunkedUpload", e))
}
async fn cancel_upload(&self, upload_id: &str, user_id: &str) -> Result<(), DomainError> {
self.cancel_upload_inner(upload_id, user_id)
async fn cancel_upload(&self, upload_id: &str, user_id: Uuid) -> Result<(), DomainError> {
self.cancel_upload_inner(upload_id, &user_id.to_string())
.await
.map_err(|e| DomainError::new(ErrorKind::InternalError, "ChunkedUpload", e))
}
@@ -7,6 +7,7 @@
use sqlx::PgPool;
use std::sync::Arc;
use uuid::Uuid;
use crate::application::dtos::display_helpers::{
category_for, format_file_size, icon_class_for, icon_special_class_for,
@@ -39,7 +40,7 @@ impl PathResolverService {
pub async fn resolve_path_for_user(
&self,
path: &str,
user_id: &str,
user_id: Uuid,
) -> Result<ResolvedResource, DomainError> {
let path = path.trim_start_matches('/').trim_end_matches('/');
if path.is_empty() {
@@ -180,7 +181,7 @@ impl PathResolverService {
}
/// Returns `true` if the resource at `path` belongs to `user_id`.
pub async fn exists_for_user(&self, path: &str, user_id: &str) -> Result<bool, DomainError> {
pub async fn exists_for_user(&self, path: &str, user_id: Uuid) -> Result<bool, DomainError> {
let path = path.trim_start_matches('/').trim_end_matches('/');
if path.is_empty() {
return Ok(false);