fix: resolve all clippy warnings and enforce cargo fmt
- display_helpers: convert module doc-comments to regular comments, merge identical text/markdown + text/ branches - search_service: replace needless range loops with slice-based pagination - folder_repository, folder_db_repository: collapse nested if statements - favorites_pg_repository: remove unnecessary borrow on generic arg - file_blob_read_repository: collapse 6 nested if-let blocks - file_blob_write_repository: collapse nested if for dedup ref decrement - chunked_upload_service: use div_ceil(), collapse 2 nested if blocks - folder_handler: collapse nested if-let for owner check - webdav_handler: replace 7x io::Error::new(ErrorKind::Other, ..) with io::Error::other(..) - cargo fmt applied to all files Passes: cargo clippy --all-targets --all-features -- -D warnings
This commit is contained in:
+14
-16
@@ -110,24 +110,22 @@ async fn create_pool_with_retries(
|
||||
.connect(connection_string)
|
||||
.await
|
||||
{
|
||||
Ok(pool) => {
|
||||
match sqlx::query("SELECT 1").execute(&pool).await {
|
||||
Ok(_) => {
|
||||
tracing::info!("PostgreSQL {} pool established successfully", label);
|
||||
return Ok(pool);
|
||||
}
|
||||
Err(e) => {
|
||||
tracing::error!("Error verifying {} pool connection: {}", label, e);
|
||||
if attempt >= MAX_ATTEMPTS {
|
||||
return Err(anyhow::anyhow!(
|
||||
"Error verifying PostgreSQL {} pool connection: {}",
|
||||
label,
|
||||
e
|
||||
));
|
||||
}
|
||||
Ok(pool) => match sqlx::query("SELECT 1").execute(&pool).await {
|
||||
Ok(_) => {
|
||||
tracing::info!("PostgreSQL {} pool established successfully", label);
|
||||
return Ok(pool);
|
||||
}
|
||||
Err(e) => {
|
||||
tracing::error!("Error verifying {} pool connection: {}", label, e);
|
||||
if attempt >= MAX_ATTEMPTS {
|
||||
return Err(anyhow::anyhow!(
|
||||
"Error verifying PostgreSQL {} pool connection: {}",
|
||||
label,
|
||||
e
|
||||
));
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
Err(e) => {
|
||||
tracing::error!(
|
||||
"Error connecting to PostgreSQL {} pool (attempt {}/{}): {}",
|
||||
|
||||
@@ -222,7 +222,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_uuid).bind(item_id).bind(item_type);
|
||||
}
|
||||
|
||||
let result = q.execute(&mut *tx).await.map_err(|e| {
|
||||
|
||||
@@ -441,10 +441,7 @@ impl FileReadPort for FileBlobReadRepository {
|
||||
///
|
||||
/// Single GiST-indexed query via ltree `<@`.
|
||||
/// Ordered by `(fo.path, fi.name)` so callers iterate in directory order.
|
||||
async fn list_files_in_subtree(
|
||||
&self,
|
||||
folder_id: &str,
|
||||
) -> Result<Vec<File>, DomainError> {
|
||||
async fn list_files_in_subtree(&self, folder_id: &str) -> Result<Vec<File>, DomainError> {
|
||||
let rows: Vec<(
|
||||
String,
|
||||
String,
|
||||
@@ -472,9 +469,7 @@ impl FileReadPort for FileBlobReadRepository {
|
||||
.bind(folder_id)
|
||||
.fetch_all(self.pool.as_ref())
|
||||
.await
|
||||
.map_err(|e| {
|
||||
DomainError::internal_error("FileBlobRead", format!("subtree files: {e}"))
|
||||
})?;
|
||||
.map_err(|e| DomainError::internal_error("FileBlobRead", format!("subtree files: {e}")))?;
|
||||
|
||||
rows.into_iter()
|
||||
.map(|(id, name, fid, fpath, size, mime, ca, ma, uid)| {
|
||||
@@ -520,11 +515,11 @@ impl FileReadPort for FileBlobReadRepository {
|
||||
conditions.push(format!("fi.folder_id = ${bind_idx}::uuid"));
|
||||
}
|
||||
|
||||
if let Some(name) = &criteria.name_contains {
|
||||
if !name.is_empty() {
|
||||
bind_idx += 1;
|
||||
conditions.push(format!("LOWER(fi.name) LIKE ${bind_idx}"));
|
||||
}
|
||||
if let Some(name) = &criteria.name_contains
|
||||
&& !name.is_empty()
|
||||
{
|
||||
bind_idx += 1;
|
||||
conditions.push(format!("LOWER(fi.name) LIKE ${bind_idx}"));
|
||||
}
|
||||
|
||||
let where_clause = conditions.join(" AND ");
|
||||
@@ -546,19 +541,30 @@ impl FileReadPort for FileBlobReadRepository {
|
||||
);
|
||||
|
||||
// ── Bind parameters dynamically ──────────────────────────────────
|
||||
let mut query = sqlx::query_as::<_, (
|
||||
String, String, Option<String>, Option<String>,
|
||||
i64, String, i64, i64, Option<String>, i64,
|
||||
)>(&sql)
|
||||
.bind(user_id);
|
||||
let mut query = sqlx::query_as::<
|
||||
_,
|
||||
(
|
||||
String,
|
||||
String,
|
||||
Option<String>,
|
||||
Option<String>,
|
||||
i64,
|
||||
String,
|
||||
i64,
|
||||
i64,
|
||||
Option<String>,
|
||||
i64,
|
||||
),
|
||||
>(&sql)
|
||||
.bind(user_id);
|
||||
|
||||
if let Some(fid) = folder_id {
|
||||
query = query.bind(fid);
|
||||
}
|
||||
if let Some(name) = &criteria.name_contains {
|
||||
if !name.is_empty() {
|
||||
query = query.bind(format!("%{}%", name.to_lowercase()));
|
||||
}
|
||||
if let Some(name) = &criteria.name_contains
|
||||
&& !name.is_empty()
|
||||
{
|
||||
query = query.bind(format!("%{}%", name.to_lowercase()));
|
||||
}
|
||||
query = query.bind(limit).bind(offset);
|
||||
|
||||
@@ -600,9 +606,7 @@ impl FileReadPort for FileBlobReadRepository {
|
||||
// When no root folder specified, delegate to existing paginated search
|
||||
let root_id = match root_folder_id {
|
||||
None => {
|
||||
return self
|
||||
.search_files_paginated(None, criteria, user_id)
|
||||
.await;
|
||||
return self.search_files_paginated(None, criteria, user_id).await;
|
||||
}
|
||||
Some(id) => id,
|
||||
};
|
||||
@@ -631,19 +635,19 @@ impl FileReadPort for FileBlobReadRepository {
|
||||
"fo.lpath <@ (SELECT lpath FROM storage.folders WHERE id = $2::uuid)".to_string(),
|
||||
);
|
||||
|
||||
if let Some(name) = &criteria.name_contains {
|
||||
if !name.is_empty() {
|
||||
bind_idx += 1;
|
||||
conditions.push(format!("LOWER(fi.name) LIKE ${bind_idx}"));
|
||||
}
|
||||
if let Some(name) = &criteria.name_contains
|
||||
&& !name.is_empty()
|
||||
{
|
||||
bind_idx += 1;
|
||||
conditions.push(format!("LOWER(fi.name) LIKE ${bind_idx}"));
|
||||
}
|
||||
if let Some(types) = &criteria.file_types {
|
||||
if !types.is_empty() {
|
||||
bind_idx += 1;
|
||||
conditions.push(format!(
|
||||
"LOWER(SUBSTRING(fi.name FROM '\\.([^.]+)$')) = ANY(${bind_idx})"
|
||||
));
|
||||
}
|
||||
if let Some(types) = &criteria.file_types
|
||||
&& !types.is_empty()
|
||||
{
|
||||
bind_idx += 1;
|
||||
conditions.push(format!(
|
||||
"LOWER(SUBSTRING(fi.name FROM '\\.([^.]+)$')) = ANY(${bind_idx})"
|
||||
));
|
||||
}
|
||||
if criteria.created_after.is_some() {
|
||||
bind_idx += 1;
|
||||
@@ -698,24 +702,34 @@ impl FileReadPort for FileBlobReadRepository {
|
||||
);
|
||||
|
||||
// ── Bind parameters dynamically ──
|
||||
let mut query = sqlx::query_as::<_, (
|
||||
String, String, Option<String>, Option<String>,
|
||||
i64, String, i64, i64, Option<String>, i64,
|
||||
)>(&sql)
|
||||
.bind(user_id)
|
||||
.bind(root_id);
|
||||
let mut query = sqlx::query_as::<
|
||||
_,
|
||||
(
|
||||
String,
|
||||
String,
|
||||
Option<String>,
|
||||
Option<String>,
|
||||
i64,
|
||||
String,
|
||||
i64,
|
||||
i64,
|
||||
Option<String>,
|
||||
i64,
|
||||
),
|
||||
>(&sql)
|
||||
.bind(user_id)
|
||||
.bind(root_id);
|
||||
|
||||
if let Some(name) = &criteria.name_contains {
|
||||
if !name.is_empty() {
|
||||
query = query.bind(format!("%{}%", name.to_lowercase()));
|
||||
}
|
||||
if let Some(name) = &criteria.name_contains
|
||||
&& !name.is_empty()
|
||||
{
|
||||
query = query.bind(format!("%{}%", name.to_lowercase()));
|
||||
}
|
||||
if let Some(types) = &criteria.file_types {
|
||||
if !types.is_empty() {
|
||||
let lower_types: Vec<String> =
|
||||
types.iter().map(|t| t.to_lowercase()).collect();
|
||||
query = query.bind(lower_types);
|
||||
}
|
||||
if let Some(types) = &criteria.file_types
|
||||
&& !types.is_empty()
|
||||
{
|
||||
let lower_types: Vec<String> = types.iter().map(|t| t.to_lowercase()).collect();
|
||||
query = query.bind(lower_types);
|
||||
}
|
||||
if let Some(v) = criteria.created_after {
|
||||
query = query.bind(v as i64);
|
||||
@@ -739,12 +753,9 @@ impl FileReadPort for FileBlobReadRepository {
|
||||
query = query.bind(limit).bind(offset);
|
||||
|
||||
// ── Execute single query ──
|
||||
let rows = query
|
||||
.fetch_all(self.pool.as_ref())
|
||||
.await
|
||||
.map_err(|e| {
|
||||
DomainError::internal_error("FileBlobRead", format!("subtree search: {e}"))
|
||||
})?;
|
||||
let rows = query.fetch_all(self.pool.as_ref()).await.map_err(|e| {
|
||||
DomainError::internal_error("FileBlobRead", format!("subtree search: {e}"))
|
||||
})?;
|
||||
|
||||
let total_count = rows.first().map_or(0, |r| r.9) as usize;
|
||||
|
||||
@@ -973,12 +984,9 @@ mod tests {
|
||||
.build(),
|
||||
};
|
||||
|
||||
repo.hash_cache
|
||||
.insert("a".to_string(), "ha".to_string());
|
||||
repo.hash_cache
|
||||
.insert("b".to_string(), "hb".to_string());
|
||||
repo.hash_cache
|
||||
.insert("c".to_string(), "hc".to_string());
|
||||
repo.hash_cache.insert("a".to_string(), "ha".to_string());
|
||||
repo.hash_cache.insert("b".to_string(), "hb".to_string());
|
||||
repo.hash_cache.insert("c".to_string(), "hc".to_string());
|
||||
|
||||
// Force moka to run pending eviction tasks
|
||||
repo.hash_cache.run_pending_tasks();
|
||||
@@ -1016,7 +1024,8 @@ mod tests {
|
||||
}
|
||||
|
||||
for h in handles {
|
||||
h.join().expect("Thread must not panic — no poison possible with moka");
|
||||
h.join()
|
||||
.expect("Thread must not panic — no poison possible with moka");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -166,14 +166,14 @@ impl FileBlobWriteRepository {
|
||||
};
|
||||
|
||||
// Decrement old blob ref (only if hash changed, best-effort)
|
||||
if old_hash != new_hash {
|
||||
if let Err(e) = self.dedup.remove_reference(&old_hash).await {
|
||||
tracing::warn!(
|
||||
"Failed to decrement old blob ref {}: {}",
|
||||
&old_hash[..12],
|
||||
e
|
||||
);
|
||||
}
|
||||
if old_hash != new_hash
|
||||
&& let Err(e) = self.dedup.remove_reference(&old_hash).await
|
||||
{
|
||||
tracing::warn!(
|
||||
"Failed to decrement old blob ref {}: {}",
|
||||
&old_hash[..12],
|
||||
e
|
||||
);
|
||||
}
|
||||
|
||||
Ok(())
|
||||
@@ -701,10 +701,7 @@ impl FileWritePort for FileBlobWriteRepository {
|
||||
);
|
||||
}
|
||||
}
|
||||
DomainError::internal_error(
|
||||
"FileBlobWrite",
|
||||
format!("copy_folder_tree: {e}"),
|
||||
)
|
||||
DomainError::internal_error("FileBlobWrite", format!("copy_folder_tree: {e}"))
|
||||
})?;
|
||||
|
||||
tracing::info!(
|
||||
|
||||
@@ -685,10 +685,7 @@ impl FolderRepository for FolderDbRepository {
|
||||
///
|
||||
/// Single GiST-indexed query: `fo.lpath <@ (root's lpath)`.
|
||||
/// Ordered by `fo.path` so callers can iterate in directory order.
|
||||
async fn list_subtree_folders(
|
||||
&self,
|
||||
folder_id: &str,
|
||||
) -> Result<Vec<Folder>, DomainError> {
|
||||
async fn list_subtree_folders(&self, folder_id: &str) -> Result<Vec<Folder>, DomainError> {
|
||||
let sql = "SELECT fo.id::text, fo.name, fo.path, fo.parent_id::text, \
|
||||
fo.user_id::text, \
|
||||
EXTRACT(EPOCH FROM fo.created_at)::bigint, \
|
||||
@@ -698,14 +695,21 @@ impl FolderRepository for FolderDbRepository {
|
||||
AND fo.lpath <@ (SELECT lpath FROM storage.folders WHERE id = $1::uuid) \
|
||||
ORDER BY fo.path";
|
||||
|
||||
let rows: Vec<(String, String, String, Option<String>, Option<String>, i64, i64)> =
|
||||
sqlx::query_as(sql)
|
||||
.bind(folder_id)
|
||||
.fetch_all(self.pool())
|
||||
.await
|
||||
.map_err(|e| {
|
||||
DomainError::internal_error("FolderDb", format!("subtree folders: {e}"))
|
||||
})?;
|
||||
let rows: Vec<(
|
||||
String,
|
||||
String,
|
||||
String,
|
||||
Option<String>,
|
||||
Option<String>,
|
||||
i64,
|
||||
i64,
|
||||
)> = sqlx::query_as(sql)
|
||||
.bind(folder_id)
|
||||
.fetch_all(self.pool())
|
||||
.await
|
||||
.map_err(|e| {
|
||||
DomainError::internal_error("FolderDb", format!("subtree folders: {e}"))
|
||||
})?;
|
||||
|
||||
rows.into_iter()
|
||||
.map(|(id, name, path, pid, uid, ca, ma)| {
|
||||
@@ -728,10 +732,10 @@ impl FolderRepository for FolderDbRepository {
|
||||
recursive: bool,
|
||||
) -> Result<Vec<Folder>, DomainError> {
|
||||
// Recursive with folder scope → existing optimised ltree scan
|
||||
if recursive {
|
||||
if let Some(fid) = parent_id {
|
||||
return self.list_descendant_folders(fid, name_contains, user_id).await;
|
||||
}
|
||||
if recursive && let Some(fid) = parent_id {
|
||||
return self
|
||||
.list_descendant_folders(fid, name_contains, user_id)
|
||||
.await;
|
||||
}
|
||||
|
||||
// Build optional name filter
|
||||
@@ -761,22 +765,27 @@ impl FolderRepository for FolderDbRepository {
|
||||
ORDER BY fo.name"
|
||||
);
|
||||
|
||||
let rows: Vec<(String, String, String, Option<String>, Option<String>, i64, i64)> =
|
||||
if let Some(ref pattern) = name_pattern {
|
||||
sqlx::query_as(&sql)
|
||||
.bind(user_id)
|
||||
.bind(pattern)
|
||||
.fetch_all(self.pool())
|
||||
.await
|
||||
} else {
|
||||
sqlx::query_as(&sql)
|
||||
.bind(user_id)
|
||||
.fetch_all(self.pool())
|
||||
.await
|
||||
}
|
||||
.map_err(|e| {
|
||||
DomainError::internal_error("FolderDb", format!("search_folders: {e}"))
|
||||
})?;
|
||||
let rows: Vec<(
|
||||
String,
|
||||
String,
|
||||
String,
|
||||
Option<String>,
|
||||
Option<String>,
|
||||
i64,
|
||||
i64,
|
||||
)> = if let Some(ref pattern) = name_pattern {
|
||||
sqlx::query_as(&sql)
|
||||
.bind(user_id)
|
||||
.bind(pattern)
|
||||
.fetch_all(self.pool())
|
||||
.await
|
||||
} else {
|
||||
sqlx::query_as(&sql)
|
||||
.bind(user_id)
|
||||
.fetch_all(self.pool())
|
||||
.await
|
||||
}
|
||||
.map_err(|e| DomainError::internal_error("FolderDb", format!("search_folders: {e}")))?;
|
||||
|
||||
return rows
|
||||
.into_iter()
|
||||
@@ -820,37 +829,42 @@ impl FolderRepository for FolderDbRepository {
|
||||
)
|
||||
};
|
||||
|
||||
let rows: Vec<(String, String, String, Option<String>, Option<String>, i64, i64)> =
|
||||
if let Some(pid) = parent_id {
|
||||
if let Some(ref pattern) = name_pattern {
|
||||
sqlx::query_as(&sql)
|
||||
.bind(pid)
|
||||
.bind(user_id)
|
||||
.bind(pattern)
|
||||
.fetch_all(self.pool())
|
||||
.await
|
||||
} else {
|
||||
sqlx::query_as(&sql)
|
||||
.bind(pid)
|
||||
.bind(user_id)
|
||||
.fetch_all(self.pool())
|
||||
.await
|
||||
}
|
||||
} else if let Some(ref pattern) = name_pattern {
|
||||
let rows: Vec<(
|
||||
String,
|
||||
String,
|
||||
String,
|
||||
Option<String>,
|
||||
Option<String>,
|
||||
i64,
|
||||
i64,
|
||||
)> = if let Some(pid) = parent_id {
|
||||
if let Some(ref pattern) = name_pattern {
|
||||
sqlx::query_as(&sql)
|
||||
.bind(pid)
|
||||
.bind(user_id)
|
||||
.bind(pattern)
|
||||
.fetch_all(self.pool())
|
||||
.await
|
||||
} else {
|
||||
sqlx::query_as(&sql)
|
||||
.bind(pid)
|
||||
.bind(user_id)
|
||||
.fetch_all(self.pool())
|
||||
.await
|
||||
}
|
||||
.map_err(|e| {
|
||||
DomainError::internal_error("FolderDb", format!("search_folders: {e}"))
|
||||
})?;
|
||||
} else if let Some(ref pattern) = name_pattern {
|
||||
sqlx::query_as(&sql)
|
||||
.bind(user_id)
|
||||
.bind(pattern)
|
||||
.fetch_all(self.pool())
|
||||
.await
|
||||
} else {
|
||||
sqlx::query_as(&sql)
|
||||
.bind(user_id)
|
||||
.fetch_all(self.pool())
|
||||
.await
|
||||
}
|
||||
.map_err(|e| DomainError::internal_error("FolderDb", format!("search_folders: {e}")))?;
|
||||
|
||||
rows.into_iter()
|
||||
.map(|(id, name, path, pid, uid, ca, ma)| {
|
||||
@@ -870,9 +884,10 @@ impl FolderRepository for FolderDbRepository {
|
||||
user_id: &str,
|
||||
) -> Result<Vec<Folder>, DomainError> {
|
||||
let (where_extra, name_pattern) = match name_contains {
|
||||
Some(name) if !name.is_empty() => {
|
||||
(" AND LOWER(fo.name) LIKE $3", Some(format!("%{}%", name.to_lowercase())))
|
||||
}
|
||||
Some(name) if !name.is_empty() => (
|
||||
" AND LOWER(fo.name) LIKE $3",
|
||||
Some(format!("%{}%", name.to_lowercase())),
|
||||
),
|
||||
_ => ("", None),
|
||||
};
|
||||
|
||||
@@ -890,24 +905,29 @@ impl FolderRepository for FolderDbRepository {
|
||||
ORDER BY fo.name"
|
||||
);
|
||||
|
||||
let rows: Vec<(String, String, String, Option<String>, Option<String>, i64, i64)> =
|
||||
if let Some(ref pattern) = name_pattern {
|
||||
sqlx::query_as(&sql)
|
||||
.bind(user_id)
|
||||
.bind(folder_id)
|
||||
.bind(pattern)
|
||||
.fetch_all(self.pool())
|
||||
.await
|
||||
} else {
|
||||
sqlx::query_as(&sql)
|
||||
.bind(user_id)
|
||||
.bind(folder_id)
|
||||
.fetch_all(self.pool())
|
||||
.await
|
||||
}
|
||||
.map_err(|e| {
|
||||
DomainError::internal_error("FolderDb", format!("descendant search: {e}"))
|
||||
})?;
|
||||
let rows: Vec<(
|
||||
String,
|
||||
String,
|
||||
String,
|
||||
Option<String>,
|
||||
Option<String>,
|
||||
i64,
|
||||
i64,
|
||||
)> = if let Some(ref pattern) = name_pattern {
|
||||
sqlx::query_as(&sql)
|
||||
.bind(user_id)
|
||||
.bind(folder_id)
|
||||
.bind(pattern)
|
||||
.fetch_all(self.pool())
|
||||
.await
|
||||
} else {
|
||||
sqlx::query_as(&sql)
|
||||
.bind(user_id)
|
||||
.bind(folder_id)
|
||||
.fetch_all(self.pool())
|
||||
.await
|
||||
}
|
||||
.map_err(|e| DomainError::internal_error("FolderDb", format!("descendant search: {e}")))?;
|
||||
|
||||
rows.into_iter()
|
||||
.map(|(id, name, path, pid, uid, ca, ma)| {
|
||||
|
||||
@@ -1,308 +1,307 @@
|
||||
use async_trait::async_trait;
|
||||
use sqlx::{PgPool, Row};
|
||||
use std::sync::Arc;
|
||||
|
||||
use crate::{
|
||||
application::ports::share_ports::ShareStoragePort,
|
||||
common::errors::DomainError,
|
||||
domain::entities::share::{Share, ShareItemType, SharePermissions},
|
||||
};
|
||||
|
||||
/// PostgreSQL implementation of [`ShareStoragePort`].
|
||||
///
|
||||
/// Replaces the legacy file-based `ShareFsRepository` that read/wrote the
|
||||
/// entire `shares.json` on every operation. Each method now issues a single
|
||||
/// indexed SQL statement — O(1) lookups, ACID transactions, and no data-race
|
||||
/// risk.
|
||||
pub struct SharePgRepository {
|
||||
db_pool: Arc<PgPool>,
|
||||
}
|
||||
|
||||
impl SharePgRepository {
|
||||
pub fn new(db_pool: Arc<PgPool>) -> Self {
|
||||
Self { db_pool }
|
||||
}
|
||||
|
||||
/// 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.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| {
|
||||
DomainError::internal_error("Share", format!("Failed to read item_id: {e}"))
|
||||
})?;
|
||||
let item_name: Option<String> = row.try_get("item_name").unwrap_or(None);
|
||||
let item_type_str: String = row.try_get("item_type").map_err(|e| {
|
||||
DomainError::internal_error("Share", format!("Failed to read item_type: {e}"))
|
||||
})?;
|
||||
let token: String = row.try_get("token").map_err(|e| {
|
||||
DomainError::internal_error("Share", format!("Failed to read token: {e}"))
|
||||
})?;
|
||||
let password_hash: Option<String> = row.try_get("password_hash").unwrap_or(None);
|
||||
let expires_at: Option<i64> = row.try_get("expires_at").unwrap_or(None);
|
||||
let permissions_read: bool = row.try_get("permissions_read").unwrap_or(true);
|
||||
let permissions_write: bool = row.try_get("permissions_write").unwrap_or(false);
|
||||
let permissions_reshare: bool = row.try_get("permissions_reshare").unwrap_or(false);
|
||||
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| {
|
||||
DomainError::internal_error("Share", format!("Failed to read created_by: {e}"))
|
||||
})?;
|
||||
let access_count: i64 = row.try_get("access_count").unwrap_or(0);
|
||||
|
||||
let item_type =
|
||||
ShareItemType::try_from(item_type_str.as_str()).unwrap_or(ShareItemType::File);
|
||||
let permissions = SharePermissions::new(permissions_read, permissions_write, permissions_reshare);
|
||||
|
||||
Ok(Share::from_raw(
|
||||
id,
|
||||
item_id,
|
||||
item_name,
|
||||
item_type,
|
||||
token,
|
||||
password_hash,
|
||||
expires_at.map(|v| v as u64),
|
||||
permissions,
|
||||
created_at as u64,
|
||||
created_by,
|
||||
access_count as u64,
|
||||
))
|
||||
}
|
||||
}
|
||||
|
||||
#[async_trait]
|
||||
impl ShareStoragePort for SharePgRepository {
|
||||
async fn save_share(&self, share: &Share) -> Result<Share, DomainError> {
|
||||
let row = sqlx::query(
|
||||
r#"
|
||||
INSERT INTO storage.shares
|
||||
(id, item_id, item_name, item_type, token, password_hash,
|
||||
expires_at, permissions_read, permissions_write, permissions_reshare,
|
||||
created_at, created_by, access_count)
|
||||
VALUES
|
||||
($1::UUID, $2, $3, $4, $5, $6,
|
||||
$7, $8, $9, $10,
|
||||
$11, $12, $13)
|
||||
ON CONFLICT (id) DO UPDATE SET
|
||||
item_name = EXCLUDED.item_name,
|
||||
password_hash = EXCLUDED.password_hash,
|
||||
expires_at = EXCLUDED.expires_at,
|
||||
permissions_read = EXCLUDED.permissions_read,
|
||||
permissions_write = EXCLUDED.permissions_write,
|
||||
permissions_reshare = EXCLUDED.permissions_reshare,
|
||||
access_count = EXCLUDED.access_count
|
||||
RETURNING
|
||||
id::TEXT, item_id, item_name, item_type, token, password_hash,
|
||||
expires_at, permissions_read, permissions_write, permissions_reshare,
|
||||
created_at, created_by, access_count
|
||||
"#,
|
||||
)
|
||||
.bind(share.id())
|
||||
.bind(share.item_id())
|
||||
.bind(share.item_name())
|
||||
.bind(share.item_type().to_string())
|
||||
.bind(share.token())
|
||||
.bind(share.password_hash())
|
||||
.bind(share.expires_at().map(|v| v as i64))
|
||||
.bind(share.permissions().read())
|
||||
.bind(share.permissions().write())
|
||||
.bind(share.permissions().reshare())
|
||||
.bind(share.created_at() as i64)
|
||||
.bind(share.created_by())
|
||||
.bind(share.access_count() as i64)
|
||||
.fetch_one(&*self.db_pool)
|
||||
.await
|
||||
.map_err(|e| {
|
||||
tracing::error!("Database error saving share: {}", e);
|
||||
DomainError::internal_error("Share", format!("Failed to save share: {e}"))
|
||||
})?;
|
||||
|
||||
Self::row_to_entity(&row)
|
||||
}
|
||||
|
||||
async fn find_share_by_id(&self, id: &str) -> Result<Share, DomainError> {
|
||||
let row = sqlx::query(
|
||||
r#"
|
||||
SELECT id::TEXT, 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
|
||||
"#,
|
||||
)
|
||||
.bind(id)
|
||||
.fetch_optional(&*self.db_pool)
|
||||
.await
|
||||
.map_err(|e| {
|
||||
tracing::error!("Database error finding share by id: {}", e);
|
||||
DomainError::internal_error("Share", format!("Failed to find share: {e}"))
|
||||
})?;
|
||||
|
||||
match row {
|
||||
Some(r) => Self::row_to_entity(&r),
|
||||
None => Err(DomainError::not_found(
|
||||
"Share",
|
||||
format!("Share with ID {id} not found"),
|
||||
)),
|
||||
}
|
||||
}
|
||||
|
||||
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,
|
||||
expires_at, permissions_read, permissions_write, permissions_reshare,
|
||||
created_at, created_by, access_count
|
||||
FROM storage.shares
|
||||
WHERE token = $1
|
||||
"#,
|
||||
)
|
||||
.bind(token)
|
||||
.fetch_optional(&*self.db_pool)
|
||||
.await
|
||||
.map_err(|e| {
|
||||
tracing::error!("Database error finding share by token: {}", e);
|
||||
DomainError::internal_error("Share", format!("Failed to find share by token: {e}"))
|
||||
})?;
|
||||
|
||||
match row {
|
||||
Some(r) => Self::row_to_entity(&r),
|
||||
None => Err(DomainError::not_found(
|
||||
"Share",
|
||||
format!("Share with token {token} not found"),
|
||||
)),
|
||||
}
|
||||
}
|
||||
|
||||
async fn find_shares_by_item(
|
||||
&self,
|
||||
item_id: &str,
|
||||
item_type: &ShareItemType,
|
||||
) -> Result<Vec<Share>, DomainError> {
|
||||
let rows = sqlx::query(
|
||||
r#"
|
||||
SELECT id::TEXT, 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 item_id = $1 AND item_type = $2
|
||||
ORDER BY created_at DESC
|
||||
"#,
|
||||
)
|
||||
.bind(item_id)
|
||||
.bind(item_type.to_string())
|
||||
.fetch_all(&*self.db_pool)
|
||||
.await
|
||||
.map_err(|e| {
|
||||
tracing::error!("Database error finding shares by item: {}", e);
|
||||
DomainError::internal_error("Share", format!("Failed to find shares by item: {e}"))
|
||||
})?;
|
||||
|
||||
rows.iter().map(Self::row_to_entity).collect()
|
||||
}
|
||||
|
||||
async fn update_share(&self, share: &Share) -> Result<Share, DomainError> {
|
||||
let row = sqlx::query(
|
||||
r#"
|
||||
UPDATE storage.shares SET
|
||||
item_name = $2,
|
||||
password_hash = $3,
|
||||
expires_at = $4,
|
||||
permissions_read = $5,
|
||||
permissions_write = $6,
|
||||
permissions_reshare = $7,
|
||||
access_count = $8
|
||||
WHERE id = $1::UUID
|
||||
RETURNING
|
||||
id::TEXT, item_id, item_name, item_type, token, password_hash,
|
||||
expires_at, permissions_read, permissions_write, permissions_reshare,
|
||||
created_at, created_by, access_count
|
||||
"#,
|
||||
)
|
||||
.bind(share.id())
|
||||
.bind(share.item_name())
|
||||
.bind(share.password_hash())
|
||||
.bind(share.expires_at().map(|v| v as i64))
|
||||
.bind(share.permissions().read())
|
||||
.bind(share.permissions().write())
|
||||
.bind(share.permissions().reshare())
|
||||
.bind(share.access_count() as i64)
|
||||
.fetch_optional(&*self.db_pool)
|
||||
.await
|
||||
.map_err(|e| {
|
||||
tracing::error!("Database error updating share: {}", e);
|
||||
DomainError::internal_error("Share", format!("Failed to update share: {e}"))
|
||||
})?;
|
||||
|
||||
match row {
|
||||
Some(r) => Self::row_to_entity(&r),
|
||||
None => Err(DomainError::not_found(
|
||||
"Share",
|
||||
format!("Share with ID {} not found for update", share.id()),
|
||||
)),
|
||||
}
|
||||
}
|
||||
|
||||
async fn delete_share(&self, id: &str) -> Result<(), DomainError> {
|
||||
let result = sqlx::query(
|
||||
"DELETE FROM storage.shares WHERE id = $1::UUID",
|
||||
)
|
||||
.bind(id)
|
||||
.execute(&*self.db_pool)
|
||||
.await
|
||||
.map_err(|e| {
|
||||
tracing::error!("Database error deleting share: {}", e);
|
||||
DomainError::internal_error("Share", format!("Failed to delete share: {e}"))
|
||||
})?;
|
||||
|
||||
if result.rows_affected() == 0 {
|
||||
return Err(DomainError::not_found(
|
||||
"Share",
|
||||
format!("Share with ID {id} not found for deletion"),
|
||||
));
|
||||
}
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn find_shares_by_user(
|
||||
&self,
|
||||
user_id: &str,
|
||||
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,
|
||||
expires_at, permissions_read, permissions_write, permissions_reshare,
|
||||
created_at, created_by, access_count,
|
||||
COUNT(*) OVER() AS total_count
|
||||
FROM storage.shares
|
||||
WHERE created_by = $1
|
||||
ORDER BY created_at DESC
|
||||
LIMIT $2 OFFSET $3
|
||||
"#,
|
||||
)
|
||||
.bind(user_id)
|
||||
.bind(limit as i64)
|
||||
.bind(offset as i64)
|
||||
.fetch_all(&*self.db_pool)
|
||||
.await
|
||||
.map_err(|e| {
|
||||
tracing::error!("Database error finding shares by user: {}", e);
|
||||
DomainError::internal_error("Share", format!("Failed to find shares by user: {e}"))
|
||||
})?;
|
||||
|
||||
let total: usize = rows
|
||||
.first()
|
||||
.and_then(|r| r.try_get::<i64, _>("total_count").ok())
|
||||
.unwrap_or(0) as usize;
|
||||
|
||||
let shares: Result<Vec<Share>, DomainError> =
|
||||
rows.iter().map(Self::row_to_entity).collect();
|
||||
|
||||
Ok((shares?, total))
|
||||
}
|
||||
}
|
||||
use async_trait::async_trait;
|
||||
use sqlx::{PgPool, Row};
|
||||
use std::sync::Arc;
|
||||
|
||||
use crate::{
|
||||
application::ports::share_ports::ShareStoragePort,
|
||||
common::errors::DomainError,
|
||||
domain::entities::share::{Share, ShareItemType, SharePermissions},
|
||||
};
|
||||
|
||||
/// PostgreSQL implementation of [`ShareStoragePort`].
|
||||
///
|
||||
/// Replaces the legacy file-based `ShareFsRepository` that read/wrote the
|
||||
/// entire `shares.json` on every operation. Each method now issues a single
|
||||
/// indexed SQL statement — O(1) lookups, ACID transactions, and no data-race
|
||||
/// risk.
|
||||
pub struct SharePgRepository {
|
||||
db_pool: Arc<PgPool>,
|
||||
}
|
||||
|
||||
impl SharePgRepository {
|
||||
pub fn new(db_pool: Arc<PgPool>) -> Self {
|
||||
Self { db_pool }
|
||||
}
|
||||
|
||||
/// 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
|
||||
.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| {
|
||||
DomainError::internal_error("Share", format!("Failed to read item_id: {e}"))
|
||||
})?;
|
||||
let item_name: Option<String> = row.try_get("item_name").unwrap_or(None);
|
||||
let item_type_str: String = row.try_get("item_type").map_err(|e| {
|
||||
DomainError::internal_error("Share", format!("Failed to read item_type: {e}"))
|
||||
})?;
|
||||
let token: String = row.try_get("token").map_err(|e| {
|
||||
DomainError::internal_error("Share", format!("Failed to read token: {e}"))
|
||||
})?;
|
||||
let password_hash: Option<String> = row.try_get("password_hash").unwrap_or(None);
|
||||
let expires_at: Option<i64> = row.try_get("expires_at").unwrap_or(None);
|
||||
let permissions_read: bool = row.try_get("permissions_read").unwrap_or(true);
|
||||
let permissions_write: bool = row.try_get("permissions_write").unwrap_or(false);
|
||||
let permissions_reshare: bool = row.try_get("permissions_reshare").unwrap_or(false);
|
||||
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| {
|
||||
DomainError::internal_error("Share", format!("Failed to read created_by: {e}"))
|
||||
})?;
|
||||
let access_count: i64 = row.try_get("access_count").unwrap_or(0);
|
||||
|
||||
let item_type =
|
||||
ShareItemType::try_from(item_type_str.as_str()).unwrap_or(ShareItemType::File);
|
||||
let permissions =
|
||||
SharePermissions::new(permissions_read, permissions_write, permissions_reshare);
|
||||
|
||||
Ok(Share::from_raw(
|
||||
id,
|
||||
item_id,
|
||||
item_name,
|
||||
item_type,
|
||||
token,
|
||||
password_hash,
|
||||
expires_at.map(|v| v as u64),
|
||||
permissions,
|
||||
created_at as u64,
|
||||
created_by,
|
||||
access_count as u64,
|
||||
))
|
||||
}
|
||||
}
|
||||
|
||||
#[async_trait]
|
||||
impl ShareStoragePort for SharePgRepository {
|
||||
async fn save_share(&self, share: &Share) -> Result<Share, DomainError> {
|
||||
let row = sqlx::query(
|
||||
r#"
|
||||
INSERT INTO storage.shares
|
||||
(id, item_id, item_name, item_type, token, password_hash,
|
||||
expires_at, permissions_read, permissions_write, permissions_reshare,
|
||||
created_at, created_by, access_count)
|
||||
VALUES
|
||||
($1::UUID, $2, $3, $4, $5, $6,
|
||||
$7, $8, $9, $10,
|
||||
$11, $12, $13)
|
||||
ON CONFLICT (id) DO UPDATE SET
|
||||
item_name = EXCLUDED.item_name,
|
||||
password_hash = EXCLUDED.password_hash,
|
||||
expires_at = EXCLUDED.expires_at,
|
||||
permissions_read = EXCLUDED.permissions_read,
|
||||
permissions_write = EXCLUDED.permissions_write,
|
||||
permissions_reshare = EXCLUDED.permissions_reshare,
|
||||
access_count = EXCLUDED.access_count
|
||||
RETURNING
|
||||
id::TEXT, item_id, item_name, item_type, token, password_hash,
|
||||
expires_at, permissions_read, permissions_write, permissions_reshare,
|
||||
created_at, created_by, access_count
|
||||
"#,
|
||||
)
|
||||
.bind(share.id())
|
||||
.bind(share.item_id())
|
||||
.bind(share.item_name())
|
||||
.bind(share.item_type().to_string())
|
||||
.bind(share.token())
|
||||
.bind(share.password_hash())
|
||||
.bind(share.expires_at().map(|v| v as i64))
|
||||
.bind(share.permissions().read())
|
||||
.bind(share.permissions().write())
|
||||
.bind(share.permissions().reshare())
|
||||
.bind(share.created_at() as i64)
|
||||
.bind(share.created_by())
|
||||
.bind(share.access_count() as i64)
|
||||
.fetch_one(&*self.db_pool)
|
||||
.await
|
||||
.map_err(|e| {
|
||||
tracing::error!("Database error saving share: {}", e);
|
||||
DomainError::internal_error("Share", format!("Failed to save share: {e}"))
|
||||
})?;
|
||||
|
||||
Self::row_to_entity(&row)
|
||||
}
|
||||
|
||||
async fn find_share_by_id(&self, id: &str) -> Result<Share, DomainError> {
|
||||
let row = sqlx::query(
|
||||
r#"
|
||||
SELECT id::TEXT, 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
|
||||
"#,
|
||||
)
|
||||
.bind(id)
|
||||
.fetch_optional(&*self.db_pool)
|
||||
.await
|
||||
.map_err(|e| {
|
||||
tracing::error!("Database error finding share by id: {}", e);
|
||||
DomainError::internal_error("Share", format!("Failed to find share: {e}"))
|
||||
})?;
|
||||
|
||||
match row {
|
||||
Some(r) => Self::row_to_entity(&r),
|
||||
None => Err(DomainError::not_found(
|
||||
"Share",
|
||||
format!("Share with ID {id} not found"),
|
||||
)),
|
||||
}
|
||||
}
|
||||
|
||||
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,
|
||||
expires_at, permissions_read, permissions_write, permissions_reshare,
|
||||
created_at, created_by, access_count
|
||||
FROM storage.shares
|
||||
WHERE token = $1
|
||||
"#,
|
||||
)
|
||||
.bind(token)
|
||||
.fetch_optional(&*self.db_pool)
|
||||
.await
|
||||
.map_err(|e| {
|
||||
tracing::error!("Database error finding share by token: {}", e);
|
||||
DomainError::internal_error("Share", format!("Failed to find share by token: {e}"))
|
||||
})?;
|
||||
|
||||
match row {
|
||||
Some(r) => Self::row_to_entity(&r),
|
||||
None => Err(DomainError::not_found(
|
||||
"Share",
|
||||
format!("Share with token {token} not found"),
|
||||
)),
|
||||
}
|
||||
}
|
||||
|
||||
async fn find_shares_by_item(
|
||||
&self,
|
||||
item_id: &str,
|
||||
item_type: &ShareItemType,
|
||||
) -> Result<Vec<Share>, DomainError> {
|
||||
let rows = sqlx::query(
|
||||
r#"
|
||||
SELECT id::TEXT, 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 item_id = $1 AND item_type = $2
|
||||
ORDER BY created_at DESC
|
||||
"#,
|
||||
)
|
||||
.bind(item_id)
|
||||
.bind(item_type.to_string())
|
||||
.fetch_all(&*self.db_pool)
|
||||
.await
|
||||
.map_err(|e| {
|
||||
tracing::error!("Database error finding shares by item: {}", e);
|
||||
DomainError::internal_error("Share", format!("Failed to find shares by item: {e}"))
|
||||
})?;
|
||||
|
||||
rows.iter().map(Self::row_to_entity).collect()
|
||||
}
|
||||
|
||||
async fn update_share(&self, share: &Share) -> Result<Share, DomainError> {
|
||||
let row = sqlx::query(
|
||||
r#"
|
||||
UPDATE storage.shares SET
|
||||
item_name = $2,
|
||||
password_hash = $3,
|
||||
expires_at = $4,
|
||||
permissions_read = $5,
|
||||
permissions_write = $6,
|
||||
permissions_reshare = $7,
|
||||
access_count = $8
|
||||
WHERE id = $1::UUID
|
||||
RETURNING
|
||||
id::TEXT, item_id, item_name, item_type, token, password_hash,
|
||||
expires_at, permissions_read, permissions_write, permissions_reshare,
|
||||
created_at, created_by, access_count
|
||||
"#,
|
||||
)
|
||||
.bind(share.id())
|
||||
.bind(share.item_name())
|
||||
.bind(share.password_hash())
|
||||
.bind(share.expires_at().map(|v| v as i64))
|
||||
.bind(share.permissions().read())
|
||||
.bind(share.permissions().write())
|
||||
.bind(share.permissions().reshare())
|
||||
.bind(share.access_count() as i64)
|
||||
.fetch_optional(&*self.db_pool)
|
||||
.await
|
||||
.map_err(|e| {
|
||||
tracing::error!("Database error updating share: {}", e);
|
||||
DomainError::internal_error("Share", format!("Failed to update share: {e}"))
|
||||
})?;
|
||||
|
||||
match row {
|
||||
Some(r) => Self::row_to_entity(&r),
|
||||
None => Err(DomainError::not_found(
|
||||
"Share",
|
||||
format!("Share with ID {} not found for update", share.id()),
|
||||
)),
|
||||
}
|
||||
}
|
||||
|
||||
async fn delete_share(&self, id: &str) -> Result<(), DomainError> {
|
||||
let result = sqlx::query("DELETE FROM storage.shares WHERE id = $1::UUID")
|
||||
.bind(id)
|
||||
.execute(&*self.db_pool)
|
||||
.await
|
||||
.map_err(|e| {
|
||||
tracing::error!("Database error deleting share: {}", e);
|
||||
DomainError::internal_error("Share", format!("Failed to delete share: {e}"))
|
||||
})?;
|
||||
|
||||
if result.rows_affected() == 0 {
|
||||
return Err(DomainError::not_found(
|
||||
"Share",
|
||||
format!("Share with ID {id} not found for deletion"),
|
||||
));
|
||||
}
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn find_shares_by_user(
|
||||
&self,
|
||||
user_id: &str,
|
||||
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,
|
||||
expires_at, permissions_read, permissions_write, permissions_reshare,
|
||||
created_at, created_by, access_count,
|
||||
COUNT(*) OVER() AS total_count
|
||||
FROM storage.shares
|
||||
WHERE created_by = $1
|
||||
ORDER BY created_at DESC
|
||||
LIMIT $2 OFFSET $3
|
||||
"#,
|
||||
)
|
||||
.bind(user_id)
|
||||
.bind(limit as i64)
|
||||
.bind(offset as i64)
|
||||
.fetch_all(&*self.db_pool)
|
||||
.await
|
||||
.map_err(|e| {
|
||||
tracing::error!("Database error finding shares by user: {}", e);
|
||||
DomainError::internal_error("Share", format!("Failed to find shares by user: {e}"))
|
||||
})?;
|
||||
|
||||
let total: usize = rows
|
||||
.first()
|
||||
.and_then(|r| r.try_get::<i64, _>("total_count").ok())
|
||||
.unwrap_or(0) as usize;
|
||||
|
||||
let shares: Result<Vec<Share>, DomainError> =
|
||||
rows.iter().map(Self::row_to_entity).collect();
|
||||
|
||||
Ok((shares?, total))
|
||||
}
|
||||
}
|
||||
|
||||
@@ -130,7 +130,7 @@ impl UploadSession {
|
||||
|
||||
/// Build the completed-chunks bitmask (1 bit per chunk).
|
||||
fn build_progress_bitmask(&self) -> Vec<u8> {
|
||||
let len = (self.chunks.len() + 7) / 8;
|
||||
let len = self.chunks.len().div_ceil(8);
|
||||
let mut bitmask = vec![0u8; len];
|
||||
for chunk in &self.chunks {
|
||||
if chunk.status == ChunkStatus::Complete {
|
||||
@@ -157,8 +157,8 @@ impl UploadSession {
|
||||
/// Persist the full session metadata once (on create).
|
||||
async fn persist_metadata(&self) -> Result<(), String> {
|
||||
let path = self.temp_dir.join(SESSION_META_FILE);
|
||||
let json = serde_json::to_vec(self)
|
||||
.map_err(|e| format!("Failed to serialise session: {e}"))?;
|
||||
let json =
|
||||
serde_json::to_vec(self).map_err(|e| format!("Failed to serialise session: {e}"))?;
|
||||
// Atomic write: write to .tmp then rename
|
||||
let tmp = self.temp_dir.join("session.json.tmp");
|
||||
fs::write(&tmp, &json)
|
||||
@@ -208,9 +208,7 @@ impl ChunkedUploadService {
|
||||
};
|
||||
|
||||
if recovered_count > 0 {
|
||||
tracing::info!(
|
||||
"♻️ Recovered {recovered_count} chunked-upload session(s) from disk"
|
||||
);
|
||||
tracing::info!("♻️ Recovered {recovered_count} chunked-upload session(s) from disk");
|
||||
}
|
||||
|
||||
// Start cleanup task
|
||||
@@ -357,14 +355,13 @@ impl ChunkedUploadService {
|
||||
let dir_name = path.file_name().and_then(|n| n.to_str()).unwrap_or("");
|
||||
|
||||
let sessions = sessions.read().await;
|
||||
if !sessions.contains_key(dir_name) {
|
||||
if let Ok(metadata) = fs::metadata(&path).await
|
||||
&& let Ok(modified) = metadata.modified()
|
||||
&& modified.elapsed().unwrap_or_default() > SESSION_EXPIRATION
|
||||
{
|
||||
let _ = fs::remove_dir_all(&path).await;
|
||||
tracing::info!("🧹 Cleaned orphaned upload dir: {:?}", path);
|
||||
}
|
||||
if !sessions.contains_key(dir_name)
|
||||
&& let Ok(metadata) = fs::metadata(&path).await
|
||||
&& let Ok(modified) = metadata.modified()
|
||||
&& modified.elapsed().unwrap_or_default() > SESSION_EXPIRATION
|
||||
{
|
||||
let _ = fs::remove_dir_all(&path).await;
|
||||
tracing::info!("🧹 Cleaned orphaned upload dir: {:?}", path);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -504,11 +501,10 @@ impl ChunkedUploadService {
|
||||
// worker free for other connections.
|
||||
if let Some(ref expected_checksum) = checksum {
|
||||
let data_clone = data.clone(); // Bytes::clone is O(1) — just an Arc increment
|
||||
let actual_checksum = tokio::task::spawn_blocking(move || {
|
||||
format!("{:x}", md5::compute(&data_clone))
|
||||
})
|
||||
.await
|
||||
.map_err(|e| format!("MD5 checksum task failed: {e}"))?;
|
||||
let actual_checksum =
|
||||
tokio::task::spawn_blocking(move || format!("{:x}", md5::compute(&data_clone)))
|
||||
.await
|
||||
.map_err(|e| format!("MD5 checksum task failed: {e}"))?;
|
||||
|
||||
if actual_checksum != *expected_checksum {
|
||||
return Err(format!(
|
||||
@@ -576,10 +572,7 @@ impl ChunkedUploadService {
|
||||
}
|
||||
|
||||
/// Get upload status
|
||||
async fn get_status_inner(
|
||||
&self,
|
||||
upload_id: &str,
|
||||
) -> Result<UploadStatusResponseDto, String> {
|
||||
async fn get_status_inner(&self, upload_id: &str) -> Result<UploadStatusResponseDto, String> {
|
||||
let sessions = self.sessions.read().await;
|
||||
let session = sessions
|
||||
.get(upload_id)
|
||||
@@ -645,12 +638,17 @@ impl ChunkedUploadService {
|
||||
let chunks_meta: Vec<(usize, PathBuf)> = session
|
||||
.chunks
|
||||
.iter()
|
||||
.map(|c| (c.index, session.temp_dir.join(format!("chunk_{:06}", c.index))))
|
||||
.map(|c| {
|
||||
(
|
||||
c.index,
|
||||
session.temp_dir.join(format!("chunk_{:06}", c.index)),
|
||||
)
|
||||
})
|
||||
.collect();
|
||||
let total_size = session.total_size;
|
||||
|
||||
let hash = tokio::task::spawn_blocking(move || -> Result<String, String> {
|
||||
use std::io::{Read, Write, BufWriter as StdBufWriter};
|
||||
use std::io::{BufWriter as StdBufWriter, Read, Write};
|
||||
|
||||
let raw_output = std::fs::OpenOptions::new()
|
||||
.create(true)
|
||||
@@ -721,10 +719,10 @@ impl ChunkedUploadService {
|
||||
/// Finalize upload: remove session from RAM and clean up temp directory
|
||||
async fn finalize_upload_inner(&self, upload_id: &str) -> Result<(), String> {
|
||||
let mut sessions = self.sessions.write().await;
|
||||
if let Some(session) = sessions.remove(upload_id) {
|
||||
if let Err(e) = fs::remove_dir_all(&session.temp_dir).await {
|
||||
tracing::warn!("Failed to cleanup upload {}: {}", upload_id, e);
|
||||
}
|
||||
if let Some(session) = sessions.remove(upload_id)
|
||||
&& let Err(e) = fs::remove_dir_all(&session.temp_dir).await
|
||||
{
|
||||
tracing::warn!("Failed to cleanup upload {}: {}", upload_id, e);
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
@@ -919,8 +917,7 @@ mod tests {
|
||||
};
|
||||
|
||||
let json = serde_json::to_vec(&session).expect("serialise");
|
||||
let restored: UploadSession =
|
||||
serde_json::from_slice(&json).expect("deserialise");
|
||||
let restored: UploadSession = serde_json::from_slice(&json).expect("deserialise");
|
||||
|
||||
assert_eq!(restored.id, session.id);
|
||||
assert_eq!(restored.filename, session.filename);
|
||||
@@ -995,7 +992,9 @@ mod tests {
|
||||
|
||||
let recovered = ChunkedUploadService::recover_sessions(&base).await;
|
||||
assert_eq!(recovered.len(), 1);
|
||||
let session = recovered.get(&upload_id).expect("session must be recovered");
|
||||
let session = recovered
|
||||
.get(&upload_id)
|
||||
.expect("session must be recovered");
|
||||
assert_eq!(session.filename, "bigfile.bin");
|
||||
assert_eq!(session.folder_id, Some("folder-x".into()));
|
||||
assert_eq!(session.chunks[0].status, ChunkStatus::Complete);
|
||||
@@ -1050,10 +1049,8 @@ mod tests {
|
||||
assert!(status.pending_chunks.is_empty());
|
||||
|
||||
// 4. Complete (assemble)
|
||||
let (path, filename, _folder, _ct, size, hash) = service
|
||||
.complete_upload_inner(&id)
|
||||
.await
|
||||
.expect("complete");
|
||||
let (path, filename, _folder, _ct, size, hash) =
|
||||
service.complete_upload_inner(&id).await.expect("complete");
|
||||
assert_eq!(filename, "test.txt");
|
||||
assert_eq!(size, 1024);
|
||||
assert!(!hash.is_empty());
|
||||
@@ -1077,7 +1074,13 @@ mod tests {
|
||||
let service = ChunkedUploadService::new(base.clone()).await;
|
||||
|
||||
let resp = service
|
||||
.create_session_inner("x.bin".into(), None, "application/octet-stream".into(), 512, Some(512))
|
||||
.create_session_inner(
|
||||
"x.bin".into(),
|
||||
None,
|
||||
"application/octet-stream".into(),
|
||||
512,
|
||||
Some(512),
|
||||
)
|
||||
.await
|
||||
.expect("create");
|
||||
|
||||
@@ -1155,12 +1158,18 @@ mod tests {
|
||||
chunk_size: 512,
|
||||
chunks: vec![
|
||||
ChunkInfo {
|
||||
index: 0, offset: 0, size: 512,
|
||||
status: ChunkStatus::Pending, checksum: None,
|
||||
index: 0,
|
||||
offset: 0,
|
||||
size: 512,
|
||||
status: ChunkStatus::Pending,
|
||||
checksum: None,
|
||||
},
|
||||
ChunkInfo {
|
||||
index: 1, offset: 512, size: 512,
|
||||
status: ChunkStatus::Pending, checksum: None,
|
||||
index: 1,
|
||||
offset: 512,
|
||||
size: 512,
|
||||
status: ChunkStatus::Pending,
|
||||
checksum: None,
|
||||
},
|
||||
],
|
||||
created_at: Utc::now(),
|
||||
@@ -1171,14 +1180,20 @@ mod tests {
|
||||
|
||||
// Write metadata
|
||||
let json = serde_json::to_vec(&session).unwrap();
|
||||
fs::write(session_dir.join(SESSION_META_FILE), &json).await.unwrap();
|
||||
fs::write(session_dir.join(SESSION_META_FILE), &json)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
// Write progress marking both chunks complete
|
||||
let bitmask = vec![0b00000011u8]; // bits 0 and 1
|
||||
fs::write(session_dir.join(PROGRESS_FILE), &bitmask).await.unwrap();
|
||||
fs::write(session_dir.join(PROGRESS_FILE), &bitmask)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
// But only create chunk_000000 on disk — chunk_000001 is "missing"
|
||||
fs::write(session_dir.join("chunk_000000"), &[0u8; 512]).await.unwrap();
|
||||
fs::write(session_dir.join("chunk_000000"), &[0u8; 512])
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let recovered = ChunkedUploadService::recover_sessions(&base).await;
|
||||
let s = recovered.get("partial-session").expect("must be recovered");
|
||||
|
||||
@@ -712,8 +712,7 @@ impl DedupService {
|
||||
let blob_root = blob_root.clone();
|
||||
async move {
|
||||
let prefix = &hash[0..2];
|
||||
let blob_path =
|
||||
blob_root.join(prefix).join(format!("{}.blob", hash));
|
||||
let blob_path = blob_root.join(prefix).join(format!("{}.blob", hash));
|
||||
|
||||
let mut issues = Vec::new();
|
||||
|
||||
|
||||
@@ -124,7 +124,13 @@ impl FileContentCache {
|
||||
/// Put file content into cache
|
||||
///
|
||||
/// Moka handles eviction automatically based on weight (content size).
|
||||
pub async fn put(&self, file_id: String, content: Bytes, etag: Arc<str>, content_type: Arc<str>) {
|
||||
pub async fn put(
|
||||
&self,
|
||||
file_id: String,
|
||||
content: Bytes,
|
||||
etag: Arc<str>,
|
||||
content_type: Arc<str>,
|
||||
) {
|
||||
let size = content.len();
|
||||
|
||||
// Don't cache if too large
|
||||
@@ -258,24 +264,14 @@ mod tests {
|
||||
// A file within the limit should be cached
|
||||
let small = Bytes::from(vec![0u8; 50]);
|
||||
cache
|
||||
.put(
|
||||
"small".to_string(),
|
||||
small,
|
||||
"e1".into(),
|
||||
"app/bin".into(),
|
||||
)
|
||||
.put("small".to_string(), small, "e1".into(), "app/bin".into())
|
||||
.await;
|
||||
assert!(cache.get("small").await.is_some());
|
||||
|
||||
// A file exceeding max_file_size is rejected by our own logic
|
||||
let big = Bytes::from(vec![1u8; 51]);
|
||||
cache
|
||||
.put(
|
||||
"big".to_string(),
|
||||
big,
|
||||
"e2".into(),
|
||||
"app/bin".into(),
|
||||
)
|
||||
.put("big".to_string(), big, "e2".into(), "app/bin".into())
|
||||
.await;
|
||||
assert!(
|
||||
cache.get("big").await.is_none(),
|
||||
@@ -296,12 +292,7 @@ mod tests {
|
||||
|
||||
let content = Bytes::from("test");
|
||||
cache
|
||||
.put(
|
||||
"file1".to_string(),
|
||||
content,
|
||||
"e".into(),
|
||||
"t".into(),
|
||||
)
|
||||
.put("file1".to_string(), content, "e".into(), "t".into())
|
||||
.await;
|
||||
|
||||
assert!(cache.get("file1").await.is_some());
|
||||
|
||||
@@ -258,11 +258,7 @@ impl ImageTranscodeService {
|
||||
original_size,
|
||||
transcoded_size
|
||||
);
|
||||
return Ok((
|
||||
original_content,
|
||||
original_mime.to_string(),
|
||||
false,
|
||||
));
|
||||
return Ok((original_content, original_mime.to_string(), false));
|
||||
}
|
||||
|
||||
let saved = original_size - transcoded_size;
|
||||
|
||||
@@ -39,13 +39,12 @@ impl Argon2PasswordHasher {
|
||||
///
|
||||
/// Panics at startup if the parameters are invalid (caught immediately).
|
||||
pub fn new(memory_cost: u32, time_cost: u32, parallelism: u32) -> Self {
|
||||
let params = Params::new(memory_cost, time_cost, parallelism, None)
|
||||
.unwrap_or_else(|e| {
|
||||
panic!(
|
||||
"Invalid Argon2 parameters (m={}, t={}, p={}): {}",
|
||||
memory_cost, time_cost, parallelism, e
|
||||
)
|
||||
});
|
||||
let params = Params::new(memory_cost, time_cost, parallelism, None).unwrap_or_else(|e| {
|
||||
panic!(
|
||||
"Invalid Argon2 parameters (m={}, t={}, p={}): {}",
|
||||
memory_cost, time_cost, parallelism, e
|
||||
)
|
||||
});
|
||||
|
||||
tracing::info!(
|
||||
"Argon2PasswordHasher initialized: m_cost={} KiB, t_cost={}, p_cost={}",
|
||||
|
||||
@@ -192,11 +192,7 @@ impl ThumbnailService {
|
||||
}
|
||||
|
||||
// 2. Generate thumbnail (CPU-bound, runs in spawn_blocking)
|
||||
tracing::info!(
|
||||
"🎨 Generating thumbnail: {} {:?}",
|
||||
file_id_owned,
|
||||
size
|
||||
);
|
||||
tracing::info!("🎨 Generating thumbnail: {} {:?}", file_id_owned, size);
|
||||
match self.generate_thumbnail(&original_owned, size).await {
|
||||
Ok(bytes) => {
|
||||
// Save to disk (best-effort — don't fail the request)
|
||||
@@ -245,7 +241,10 @@ impl ThumbnailService {
|
||||
let max_dim = size.max_dimension();
|
||||
|
||||
// Acquire semaphore permit — bounds peak RAM from concurrent decodes
|
||||
let _permit = self.decode_semaphore.acquire().await
|
||||
let _permit = self
|
||||
.decode_semaphore
|
||||
.acquire()
|
||||
.await
|
||||
.map_err(|_| ThumbnailError::TaskError("Decode semaphore closed".into()))?;
|
||||
|
||||
// Run image processing in blocking thread pool
|
||||
@@ -278,9 +277,9 @@ impl ThumbnailService {
|
||||
// Adaptive filter: faster filters for smaller sizes where
|
||||
// quality difference vs Lanczos3 is imperceptible
|
||||
let filter = match size {
|
||||
ThumbnailSize::Icon => FilterType::Triangle, // 150px — max speed
|
||||
ThumbnailSize::Icon => FilterType::Triangle, // 150px — max speed
|
||||
ThumbnailSize::Preview => FilterType::CatmullRom, // 400px — good balance
|
||||
ThumbnailSize::Large => FilterType::CatmullRom, // 800px — sufficient quality
|
||||
ThumbnailSize::Large => FilterType::CatmullRom, // 800px — sufficient quality
|
||||
};
|
||||
let thumbnail = img.resize(new_width, new_height, filter);
|
||||
|
||||
@@ -312,7 +311,10 @@ impl ThumbnailService {
|
||||
let _permit = match self.decode_semaphore.acquire().await {
|
||||
Ok(p) => p,
|
||||
Err(_) => {
|
||||
tracing::warn!("Decode semaphore closed, skipping thumbnails for {}", file_id);
|
||||
tracing::warn!(
|
||||
"Decode semaphore closed, skipping thumbnails for {}",
|
||||
file_id
|
||||
);
|
||||
return;
|
||||
}
|
||||
};
|
||||
@@ -333,8 +335,8 @@ impl ThumbnailService {
|
||||
)));
|
||||
}
|
||||
|
||||
let img = image::open(&path)
|
||||
.map_err(|e| ThumbnailError::ImageError(e.to_string()))?;
|
||||
let img =
|
||||
image::open(&path).map_err(|e| ThumbnailError::ImageError(e.to_string()))?;
|
||||
|
||||
let (orig_w, orig_h) = (img.width(), img.height());
|
||||
|
||||
@@ -352,18 +354,15 @@ impl ThumbnailService {
|
||||
};
|
||||
|
||||
let filter = match size {
|
||||
ThumbnailSize::Icon => FilterType::Triangle,
|
||||
ThumbnailSize::Icon => FilterType::Triangle,
|
||||
ThumbnailSize::Preview => FilterType::CatmullRom,
|
||||
ThumbnailSize::Large => FilterType::CatmullRom,
|
||||
ThumbnailSize::Large => FilterType::CatmullRom,
|
||||
};
|
||||
let thumb = img.resize(new_w, new_h, filter);
|
||||
|
||||
let mut buf = Vec::new();
|
||||
thumb
|
||||
.write_to(
|
||||
&mut std::io::Cursor::new(&mut buf),
|
||||
ImageFormat::WebP,
|
||||
)
|
||||
.write_to(&mut std::io::Cursor::new(&mut buf), ImageFormat::WebP)
|
||||
.map_err(|e| ThumbnailError::ImageError(e.to_string()))?;
|
||||
|
||||
Ok((size, Bytes::from(buf)))
|
||||
@@ -376,15 +375,11 @@ impl ThumbnailService {
|
||||
let thumbnails = match results {
|
||||
Ok(Ok(t)) => t,
|
||||
Ok(Err(e)) => {
|
||||
tracing::warn!(
|
||||
"Thumbnail generation failed for {}: {}", file_id, e
|
||||
);
|
||||
tracing::warn!("Thumbnail generation failed for {}: {}", file_id, e);
|
||||
return;
|
||||
}
|
||||
Err(e) => {
|
||||
tracing::warn!(
|
||||
"Thumbnail task panicked for {}: {}", file_id, e
|
||||
);
|
||||
tracing::warn!("Thumbnail task panicked for {}: {}", file_id, e);
|
||||
return;
|
||||
}
|
||||
};
|
||||
|
||||
@@ -8,8 +8,8 @@ use crate::{
|
||||
use async_trait::async_trait;
|
||||
use async_zip::base::write::ZipFileWriter;
|
||||
use async_zip::{Compression, ZipEntryBuilder};
|
||||
use futures::io::AsyncWriteExt as FuturesWriteExt;
|
||||
use futures::StreamExt;
|
||||
use futures::io::AsyncWriteExt as FuturesWriteExt;
|
||||
use std::collections::HashMap;
|
||||
use std::sync::Arc;
|
||||
use tempfile::NamedTempFile;
|
||||
@@ -99,17 +99,13 @@ impl ZipService {
|
||||
.folder_service
|
||||
.list_subtree_folders(folder_id)
|
||||
.await
|
||||
.map_err(|e| {
|
||||
ZipError::FolderContentsError(format!("subtree folders: {}", e))
|
||||
})?;
|
||||
.map_err(|e| ZipError::FolderContentsError(format!("subtree folders: {}", e)))?;
|
||||
|
||||
let all_files = self
|
||||
.file_service
|
||||
.list_files_in_subtree(folder_id)
|
||||
.await
|
||||
.map_err(|e| {
|
||||
ZipError::FolderContentsError(format!("subtree files: {}", e))
|
||||
})?;
|
||||
.map_err(|e| ZipError::FolderContentsError(format!("subtree files: {}", e)))?;
|
||||
|
||||
info!(
|
||||
"ZIP subtree: {} folders, {} files",
|
||||
@@ -157,8 +153,7 @@ impl ZipService {
|
||||
let zip_dir = format!("{}/", folder_zip_path(&folder.path));
|
||||
|
||||
// Directory entry (Stored, zero-length body)
|
||||
let dir_entry =
|
||||
ZipEntryBuilder::new(zip_dir.clone().into(), Compression::Stored);
|
||||
let dir_entry = ZipEntryBuilder::new(zip_dir.clone().into(), Compression::Stored);
|
||||
match zip.write_entry_whole(dir_entry, &[]).await {
|
||||
Ok(()) => debug!("Folder added to ZIP: {}", zip_dir),
|
||||
Err(e) => {
|
||||
|
||||
Reference in New Issue
Block a user