This commit is contained in:
DioCrafts
2026-06-10 22:03:49 +02:00
parent 45fae4b27e
commit f678ff414e
12 changed files with 789 additions and 200 deletions
@@ -20,6 +20,7 @@ use crate::domain::entities::file::File;
use crate::domain::services::path_service::StoragePath;
use super::folder_db_repository::FolderDbRepository;
use super::transaction_utils::retry_on_deadlock;
use crate::infrastructure::services::dedup_service::DedupService;
/// File write repository backed by PostgreSQL metadata + blob storage.
@@ -157,24 +158,28 @@ impl FileBlobWriteRepository {
modified_at: Option<i64>,
) -> Result<(String, i64), DomainError> {
// Atomic CTE: capture old hash then update in one round-trip, no TOCTOU.
let (old_hash, updated_at) = match sqlx::query_as::<_, (String, i64)>(
r#"
WITH old AS (
SELECT id, blob_hash FROM storage.files WHERE id = $3::uuid FOR UPDATE
// Deadlock victims (40P01) retry before the compensation below runs —
// a successful retry must keep the new blob reference alive.
let (old_hash, updated_at) = match retry_on_deadlock("files.swap_blob_hash", || {
sqlx::query_as::<_, (String, i64)>(
r#"
WITH old AS (
SELECT id, blob_hash FROM storage.files WHERE id = $3::uuid FOR UPDATE
)
UPDATE storage.files f
SET blob_hash = $1, size = $2,
updated_at = COALESCE(to_timestamp($4), NOW())
FROM old
WHERE f.id = old.id
RETURNING old.blob_hash, EXTRACT(EPOCH FROM f.updated_at)::bigint
"#,
)
UPDATE storage.files f
SET blob_hash = $1, size = $2,
updated_at = COALESCE(to_timestamp($4), NOW())
FROM old
WHERE f.id = old.id
RETURNING old.blob_hash, EXTRACT(EPOCH FROM f.updated_at)::bigint
"#,
)
.bind(new_hash)
.bind(new_size)
.bind(file_id)
.bind(modified_at.map(|t| t as f64))
.fetch_optional(self.pool.as_ref())
.bind(new_hash)
.bind(new_size)
.bind(file_id)
.bind(modified_at.map(|t| t as f64))
.fetch_optional(self.pool.as_ref())
})
.await
{
Ok(Some(row)) => row,
@@ -236,23 +241,30 @@ impl FileBlobWriteRepository {
let is_new_blob = !dedup_result.was_deduplicated();
let blob_hash = dedup_result.hash().to_string();
let row = match sqlx::query_as::<_, (String, i64, i64)>(
r#"
INSERT INTO storage.files (name, folder_id, user_id, blob_hash, size, mime_type, category_order)
VALUES ($1, $2::uuid, $3, $4, $5, $6, $7)
RETURNING id::text,
EXTRACT(EPOCH FROM created_at)::bigint,
EXTRACT(EPOCH FROM updated_at)::bigint
"#,
)
.bind(&name)
.bind(&folder_id)
.bind(user_id)
.bind(&blob_hash)
.bind(size as i64)
.bind(&content_type)
.bind(category_order_for(&name, &content_type))
.fetch_one(self.pool.as_ref())
// Deadlock victims (40P01) retry before the compensation below runs —
// a successful retry must keep the blob reference alive. The final
// attempt's error falls through untouched so the 23505 mapping holds
// (a retried INSERT can legitimately lose to a concurrent identical
// upload).
let row = match retry_on_deadlock("files.insert", || {
sqlx::query_as::<_, (String, i64, i64)>(
r#"
INSERT INTO storage.files (name, folder_id, user_id, blob_hash, size, mime_type, category_order)
VALUES ($1, $2::uuid, $3, $4, $5, $6, $7)
RETURNING id::text,
EXTRACT(EPOCH FROM created_at)::bigint,
EXTRACT(EPOCH FROM updated_at)::bigint
"#,
)
.bind(&name)
.bind(&folder_id)
.bind(user_id)
.bind(&blob_hash)
.bind(size as i64)
.bind(&content_type)
.bind(category_order_for(&name, &content_type))
.fetch_one(self.pool.as_ref())
})
.await
{
Ok(row) => row,
@@ -372,46 +384,48 @@ impl FileWritePort for FileBlobWriteRepository {
// Single round-trip; blob content is NOT copied (dedup makes this zero-copy).
let target_fid = target_folder_id.clone();
let row = sqlx::query_as::<
_,
(
String,
String,
Option<String>,
i64,
String,
i64,
i64,
String,
),
>(
r#"
WITH src AS (
SELECT name, folder_id, user_id, blob_hash, size, mime_type, category_order
FROM storage.files
WHERE id = $1::uuid AND NOT is_trashed
),
new_file AS (
INSERT INTO storage.files (name, folder_id, user_id, blob_hash, size, mime_type, category_order)
SELECT name,
COALESCE($2::uuid, folder_id),
user_id,
blob_hash,
size,
mime_type,
category_order
FROM src
RETURNING id::text, name, folder_id::text, size, mime_type,
EXTRACT(EPOCH FROM created_at)::bigint,
EXTRACT(EPOCH FROM updated_at)::bigint,
blob_hash
let row = retry_on_deadlock("files.copy", || {
sqlx::query_as::<
_,
(
String,
String,
Option<String>,
i64,
String,
i64,
i64,
String,
),
>(
r#"
WITH src AS (
SELECT name, folder_id, user_id, blob_hash, size, mime_type, category_order
FROM storage.files
WHERE id = $1::uuid AND NOT is_trashed
),
new_file AS (
INSERT INTO storage.files (name, folder_id, user_id, blob_hash, size, mime_type, category_order)
SELECT name,
COALESCE($2::uuid, folder_id),
user_id,
blob_hash,
size,
mime_type,
category_order
FROM src
RETURNING id::text, name, folder_id::text, size, mime_type,
EXTRACT(EPOCH FROM created_at)::bigint,
EXTRACT(EPOCH FROM updated_at)::bigint,
blob_hash
)
SELECT * FROM new_file
"#,
)
SELECT * FROM new_file
"#,
)
.bind(file_id)
.bind(&target_fid)
.fetch_optional(self.pool.as_ref())
.bind(file_id)
.bind(&target_fid)
.fetch_optional(self.pool.as_ref())
})
.await
.map_err(|e| {
if let sqlx::Error::Database(ref db_err) = e
@@ -557,23 +571,25 @@ impl FileWritePort for FileBlobWriteRepository {
// The write-behind cache will call update_file_content later.
let placeholder_hash = "0000000000000000000000000000000000000000000000000000000000000000";
let row = sqlx::query_as::<_, (String, i64, i64)>(
r#"
INSERT INTO storage.files (name, folder_id, user_id, blob_hash, size, mime_type, category_order)
VALUES ($1, $2::uuid, $3, $4, $5, $6, $7)
RETURNING id::text,
EXTRACT(EPOCH FROM created_at)::bigint,
EXTRACT(EPOCH FROM updated_at)::bigint
"#,
)
.bind(&name)
.bind(&folder_id)
.bind(user_id)
.bind(placeholder_hash)
.bind(size as i64)
.bind(&content_type)
.bind(category_order_for(&name, &content_type))
.fetch_one(self.pool.as_ref())
let row = retry_on_deadlock("files.insert_deferred", || {
sqlx::query_as::<_, (String, i64, i64)>(
r#"
INSERT INTO storage.files (name, folder_id, user_id, blob_hash, size, mime_type, category_order)
VALUES ($1, $2::uuid, $3, $4, $5, $6, $7)
RETURNING id::text,
EXTRACT(EPOCH FROM created_at)::bigint,
EXTRACT(EPOCH FROM updated_at)::bigint
"#,
)
.bind(&name)
.bind(&folder_id)
.bind(user_id)
.bind(placeholder_hash)
.bind(size as i64)
.bind(&content_type)
.bind(category_order_for(&name, &content_type))
.fetch_one(self.pool.as_ref())
})
.await
.map_err(|e| DomainError::internal_error("FileBlobWrite", format!("deferred: {e}")))?;
@@ -12,6 +12,7 @@ use sqlx::PgPool;
use std::sync::Arc;
use uuid::Uuid;
use super::transaction_utils::retry_on_deadlock;
use crate::application::dtos::folder_dto::{FolderResourceCursor, FolderResourceRow};
use crate::common::errors::DomainError;
use crate::domain::entities::folder::Folder;
@@ -450,20 +451,25 @@ impl FolderRepository for FolderDbRepository {
// The BEFORE UPDATE trigger recomputes path/lpath for this row;
// the AFTER UPDATE cascade trigger then batch-updates all
// descendants in a single UPDATE using the GiST lpath index.
let row = sqlx::query_as::<_, FolderRow>(
r#"
UPDATE storage.folders
SET name = $1, updated_at = NOW()
WHERE id = $2::uuid AND NOT is_trashed
RETURNING id::text, name, path, parent_id::text, user_id,
EXTRACT(EPOCH FROM created_at)::bigint,
EXTRACT(EPOCH FROM updated_at)::bigint,
EXTRACT(EPOCH FROM tree_modified_at)::bigint
"#,
)
.bind(&new_name)
.bind(id)
.fetch_optional(self.pool())
// That multi-row rewrite can deadlock against the tree-ETag
// flusher's id-ordered ancestor bump — retry instead of failing
// the user's operation (40P01 only; 23505 still maps below).
let row = retry_on_deadlock("folders.rename", || {
sqlx::query_as::<_, FolderRow>(
r#"
UPDATE storage.folders
SET name = $1, updated_at = NOW()
WHERE id = $2::uuid AND NOT is_trashed
RETURNING id::text, name, path, parent_id::text, user_id,
EXTRACT(EPOCH FROM created_at)::bigint,
EXTRACT(EPOCH FROM updated_at)::bigint,
EXTRACT(EPOCH FROM tree_modified_at)::bigint
"#,
)
.bind(&new_name)
.bind(id)
.fetch_optional(self.pool())
})
.await
.map_err(|e| {
if let sqlx::Error::Database(ref db_err) = e
@@ -486,20 +492,23 @@ impl FolderRepository for FolderDbRepository {
// The BEFORE UPDATE trigger recomputes path/lpath for this row;
// the AFTER UPDATE cascade trigger then batch-updates all
// descendants in a single UPDATE using the GiST lpath index.
let row = sqlx::query_as::<_, FolderRow>(
r#"
UPDATE storage.folders
SET parent_id = $1::uuid, updated_at = NOW()
WHERE id = $2::uuid AND NOT is_trashed
RETURNING id::text, name, path, parent_id::text, user_id,
EXTRACT(EPOCH FROM created_at)::bigint,
EXTRACT(EPOCH FROM updated_at)::bigint,
EXTRACT(EPOCH FROM tree_modified_at)::bigint
"#,
)
.bind(new_parent_id)
.bind(id)
.fetch_optional(self.pool())
// Retried on deadlock vs the tree-ETag flusher (see rename_folder).
let row = retry_on_deadlock("folders.move", || {
sqlx::query_as::<_, FolderRow>(
r#"
UPDATE storage.folders
SET parent_id = $1::uuid, updated_at = NOW()
WHERE id = $2::uuid AND NOT is_trashed
RETURNING id::text, name, path, parent_id::text, user_id,
EXTRACT(EPOCH FROM created_at)::bigint,
EXTRACT(EPOCH FROM updated_at)::bigint,
EXTRACT(EPOCH FROM tree_modified_at)::bigint
"#,
)
.bind(new_parent_id)
.bind(id)
.fetch_optional(self.pool())
})
.await
.map_err(|e| DomainError::internal_error("FolderDb", format!("move: {e}")))?
.ok_or_else(|| DomainError::not_found("Folder", id))?;
@@ -511,24 +520,30 @@ impl FolderRepository for FolderDbRepository {
// Delete all files whose folder is anywhere in the subtree.
// Uses the GiST-indexed ltree `<@` operator — O(log N) vs the
// O(depth × N) recursive CTE it replaces.
sqlx::query(
"DELETE FROM storage.files \
WHERE folder_id IN ( \
SELECT id FROM storage.folders \
WHERE lpath <@ (SELECT lpath FROM storage.folders WHERE id = $1::uuid) \
)",
)
.bind(id)
.execute(self.pool())
// Both statements retried on deadlock vs the tree-ETag flusher's
// id-ordered ancestor bump (multi-row exclusive locks).
retry_on_deadlock("folders.delete_files", || {
sqlx::query(
"DELETE FROM storage.files \
WHERE folder_id IN ( \
SELECT id FROM storage.folders \
WHERE lpath <@ (SELECT lpath FROM storage.folders WHERE id = $1::uuid) \
)",
)
.bind(id)
.execute(self.pool())
})
.await
.map_err(|e| DomainError::internal_error("FolderDb", format!("delete files: {e}")))?;
// Then delete the folder (CASCADE will remove descendant folders)
let result = sqlx::query("DELETE FROM storage.folders WHERE id = $1::uuid")
.bind(id)
.execute(self.pool())
.await
.map_err(|e| DomainError::internal_error("FolderDb", format!("delete: {e}")))?;
let result = retry_on_deadlock("folders.delete", || {
sqlx::query("DELETE FROM storage.folders WHERE id = $1::uuid")
.bind(id)
.execute(self.pool())
})
.await
.map_err(|e| DomainError::internal_error("FolderDb", format!("delete: {e}")))?;
if result.rows_affected() == 0 {
return Err(DomainError::not_found("Folder", id));
@@ -570,22 +585,24 @@ impl FolderRepository for FolderDbRepository {
// Child files and sub-folders are implicitly hidden because their
// ancestor is trashed — list queries already filter NOT is_trashed,
// and folder navigation won't reach a trashed folder's children.
let result = sqlx::query_scalar::<_, i64>(
r#"
WITH trash_folder AS (
UPDATE storage.folders
SET is_trashed = TRUE,
trashed_at = NOW(),
original_parent_id = parent_id,
updated_at = NOW()
WHERE id = $1::uuid AND NOT is_trashed
RETURNING 1
let result = retry_on_deadlock("folders.trash", || {
sqlx::query_scalar::<_, i64>(
r#"
WITH trash_folder AS (
UPDATE storage.folders
SET is_trashed = TRUE,
trashed_at = NOW(),
original_parent_id = parent_id,
updated_at = NOW()
WHERE id = $1::uuid AND NOT is_trashed
RETURNING 1
)
SELECT COUNT(*) FROM trash_folder
"#,
)
SELECT COUNT(*) FROM trash_folder
"#,
)
.bind(folder_id)
.fetch_one(self.pool())
.bind(folder_id)
.fetch_one(self.pool())
})
.await
.map_err(|e| DomainError::internal_error("FolderDb", format!("trash: {e}")))?;
@@ -607,23 +624,25 @@ impl FolderRepository for FolderDbRepository {
// The BEFORE UPDATE trigger recomputes path/lpath when
// original_parent_id is restored; the cascade trigger
// batch-updates all descendants via the GiST lpath index.
let result = sqlx::query_scalar::<_, i64>(
r#"
WITH restore_folder AS (
UPDATE storage.folders
SET is_trashed = FALSE,
trashed_at = NULL,
parent_id = COALESCE(original_parent_id, parent_id),
original_parent_id = NULL,
updated_at = NOW()
WHERE id = $1::uuid AND is_trashed
RETURNING 1
let result = retry_on_deadlock("folders.restore", || {
sqlx::query_scalar::<_, i64>(
r#"
WITH restore_folder AS (
UPDATE storage.folders
SET is_trashed = FALSE,
trashed_at = NULL,
parent_id = COALESCE(original_parent_id, parent_id),
original_parent_id = NULL,
updated_at = NOW()
WHERE id = $1::uuid AND is_trashed
RETURNING 1
)
SELECT COUNT(*) FROM restore_folder
"#,
)
SELECT COUNT(*) FROM restore_folder
"#,
)
.bind(folder_id)
.fetch_one(self.pool())
.bind(folder_id)
.fetch_one(self.pool())
})
.await
.map_err(|e| DomainError::internal_error("FolderDb", format!("restore: {e}")))?;
@@ -636,25 +655,30 @@ impl FolderRepository for FolderDbRepository {
async fn delete_folder_permanently(&self, folder_id: &str) -> Result<(), DomainError> {
// Delete all files whose folder is anywhere in the subtree
// (GiST ltree index, same pattern as delete_folder).
sqlx::query(
"DELETE FROM storage.files \
WHERE folder_id IN ( \
SELECT id FROM storage.folders \
WHERE lpath <@ (SELECT lpath FROM storage.folders WHERE id = $1::uuid) \
)",
)
.bind(folder_id)
.execute(self.pool())
// (GiST ltree index, same pattern as delete_folder — both
// statements retried on deadlock vs the tree-ETag flusher).
retry_on_deadlock("folders.perm_delete_files", || {
sqlx::query(
"DELETE FROM storage.files \
WHERE folder_id IN ( \
SELECT id FROM storage.folders \
WHERE lpath <@ (SELECT lpath FROM storage.folders WHERE id = $1::uuid) \
)",
)
.bind(folder_id)
.execute(self.pool())
})
.await
.map_err(|e| DomainError::internal_error("FolderDb", format!("perm delete files: {e}")))?;
// Then permanently delete folder — CASCADE handles descendant folders
let result = sqlx::query("DELETE FROM storage.folders WHERE id = $1::uuid")
.bind(folder_id)
.execute(self.pool())
.await
.map_err(|e| DomainError::internal_error("FolderDb", format!("perm delete: {e}")))?;
let result = retry_on_deadlock("folders.perm_delete", || {
sqlx::query("DELETE FROM storage.folders WHERE id = $1::uuid")
.bind(folder_id)
.execute(self.pool())
})
.await
.map_err(|e| DomainError::internal_error("FolderDb", format!("perm delete: {e}")))?;
if result.rows_affected() == 0 {
return Err(DomainError::not_found("Folder", folder_id));
+1 -1
View File
@@ -16,7 +16,7 @@ mod session_pg_repository;
mod settings_pg_repository;
mod share_pg_repository;
mod subject_group_pg_repository;
mod transaction_utils;
pub(crate) mod transaction_utils;
mod user_pg_repository;
// ── Blob-storage repositories ──
@@ -1,6 +1,7 @@
use sqlx::{Error as SqlxError, PgPool, Postgres, Transaction};
use std::future::Future;
use std::sync::Arc;
use tracing::{debug, error, info};
use tracing::{debug, error, info, warn};
/// Helper function to execute database operations in a transaction
/// Takes a database pool and a closure that will be executed within a transaction
@@ -56,3 +57,128 @@ where
}
}
}
/// True when the error is a PostgreSQL deadlock abort (SQLSTATE `40P01`).
///
/// Deadlock victims are safe to re-run when the statement is a single
/// autocommit round-trip: the aborted implicit transaction left nothing
/// behind, and PostgreSQL chose this session as the victim precisely so
/// the competing transaction could finish — a retry usually succeeds
/// immediately.
pub fn is_deadlock(err: &SqlxError) -> bool {
matches!(err, SqlxError::Database(db) if db.code().as_deref() == Some("40P01"))
}
/// Re-run `op` while it fails with an error matching `should_retry`, up to
/// 3 retries with a short growing backoff. Errors that don't match the
/// predicate — and the final attempt's error — are returned untouched, so
/// callers' existing error mapping (e.g. `23505` → already-exists) still
/// sees exactly what it expects.
pub async fn retry_when<T, E, F, Fut, P>(
operation_name: &str,
should_retry: P,
op: F,
) -> Result<T, E>
where
F: Fn() -> Fut,
Fut: Future<Output = Result<T, E>>,
P: Fn(&E) -> bool,
{
const BACKOFF_MS: [u64; 3] = [10, 50, 150];
let mut attempt = 0;
loop {
match op().await {
Ok(value) => return Ok(value),
Err(e) if attempt < BACKOFF_MS.len() && should_retry(&e) => {
warn!(
"Retryable failure on {} (attempt {}/{}), backing off {}ms",
operation_name,
attempt + 1,
BACKOFF_MS.len() + 1,
BACKOFF_MS[attempt]
);
tokio::time::sleep(std::time::Duration::from_millis(BACKOFF_MS[attempt])).await;
attempt += 1;
}
Err(e) => return Err(e),
}
}
}
/// [`retry_when`] specialised to PostgreSQL deadlocks (`40P01`) — the only
/// transient SQLSTATE our single-statement write paths can hit.
pub async fn retry_on_deadlock<T, F, Fut>(operation_name: &str, op: F) -> Result<T, SqlxError>
where
F: Fn() -> Fut,
Fut: Future<Output = Result<T, SqlxError>>,
{
retry_when(operation_name, is_deadlock, op).await
}
#[cfg(test)]
mod tests {
use super::retry_when;
use std::sync::atomic::{AtomicU32, Ordering};
#[derive(Debug, PartialEq)]
enum FakeError {
Transient,
Fatal,
}
#[tokio::test]
async fn retries_transient_errors_until_success() {
let calls = AtomicU32::new(0);
let result = retry_when(
"test",
|e| *e == FakeError::Transient,
|| {
let n = calls.fetch_add(1, Ordering::SeqCst);
async move {
if n < 2 {
Err(FakeError::Transient)
} else {
Ok(n)
}
}
},
)
.await;
assert_eq!(result, Ok(2));
assert_eq!(calls.load(Ordering::SeqCst), 3);
}
#[tokio::test]
async fn gives_up_after_max_attempts_returning_last_error() {
let calls = AtomicU32::new(0);
let result: Result<(), _> = retry_when(
"test",
|e| *e == FakeError::Transient,
|| {
calls.fetch_add(1, Ordering::SeqCst);
async { Err(FakeError::Transient) }
},
)
.await;
assert_eq!(result, Err(FakeError::Transient));
// 1 initial attempt + 3 retries
assert_eq!(calls.load(Ordering::SeqCst), 4);
}
#[tokio::test]
async fn non_matching_errors_are_not_retried() {
let calls = AtomicU32::new(0);
let result: Result<(), _> = retry_when(
"test",
|e| *e == FakeError::Transient,
|| {
calls.fetch_add(1, Ordering::SeqCst);
async { Err(FakeError::Fatal) }
},
)
.await;
assert_eq!(result, Err(FakeError::Fatal));
assert_eq!(calls.load(Ordering::SeqCst), 1);
}
}