feat(drive): add quota calculus per drive
- quota per drive
- add 2 internal API endpoints to test purpose
disabled by default, enable it via `OXICLOUD_ENABLE_ADMIN_INTERNAL_ENDPOINTS=true`
this enable:
/api/admin/internal/trigger-sweep
to sweep the trash and recalculated quota
/api/admin/internal/trigger-gc
to garbage orphan blobs
use full for end to end tests and validate lifecycles
This commit is contained in:
@@ -453,6 +453,42 @@ pub trait StorageUsagePort: Send + Sync + 'static {
|
||||
|
||||
/// Returns (used_bytes, quota_bytes) for a user.
|
||||
async fn get_user_storage_info(&self, user_id: Uuid) -> Result<(i64, i64), DomainError>;
|
||||
|
||||
/// Incrementally adjust one drive's cached `storage.drives.used_bytes`
|
||||
/// by `delta` bytes — O(1), the per-upload counterpart to the
|
||||
/// O(N) full recompute below. Mirrors `add_user_storage_usage_delta`
|
||||
/// in shape: single statement, `GREATEST(0, …)` clamp so a late or
|
||||
/// duplicate adjustment can never drive the counter negative.
|
||||
/// Deletes/trash do not decrement here (mirroring user-quota
|
||||
/// design); the periodic reconciliation sweep is the correctness
|
||||
/// backstop.
|
||||
async fn add_drive_storage_usage_delta(
|
||||
&self,
|
||||
drive_id: Uuid,
|
||||
delta: i64,
|
||||
) -> Result<(), DomainError>;
|
||||
|
||||
/// Reconcile every drive's cached `used_bytes` against the actual
|
||||
/// sum of its non-trashed files in one set-based UPDATE. Same
|
||||
/// shape as `update_all_users_storage_usage`: `LEFT JOIN` over a
|
||||
/// `GROUP BY drive_id` aggregate, with an `IS DISTINCT FROM`
|
||||
/// guard so idle drives don't churn dead tuples. Runs from the
|
||||
/// same reconciliation ticker.
|
||||
async fn update_all_drives_storage_usage(&self) -> Result<(), DomainError>;
|
||||
|
||||
/// Pre-upload quota check on a single drive.
|
||||
///
|
||||
/// Returns `Ok(())` when `used_bytes + additional_bytes` fits under
|
||||
/// `quota_bytes`, or `Err(QuotaExceeded)` otherwise.
|
||||
/// `quota_bytes IS NULL` short-circuits to `Ok(())` — unlimited
|
||||
/// drive. Single read-only `SELECT` on `storage.drives`; the
|
||||
/// check/write window is a soft cap by design (same semantics as
|
||||
/// the user-quota path), bounded by the sweep interval.
|
||||
async fn check_drive_quota(
|
||||
&self,
|
||||
drive_id: Uuid,
|
||||
additional_bytes: u64,
|
||||
) -> Result<(), DomainError>;
|
||||
}
|
||||
|
||||
/// Generic storage service interface for calendar and contact services
|
||||
|
||||
@@ -331,6 +331,21 @@ impl DeltaUploadService {
|
||||
.check_storage_quota(caller_id, total_size)
|
||||
.await?;
|
||||
|
||||
// ── Per-drive quota (D4) ─────────────────────────────────
|
||||
// Mirrors the per-user check above on the same `total_size`.
|
||||
// Only on CREATE — Update replaces an existing row's content;
|
||||
// tight size-delta accounting on update is a follow-up (today
|
||||
// the periodic sweep reconciles drift either way). The
|
||||
// single-statement `check_drive_quota_by_folder` lookup is a
|
||||
// PK probe; cost matches the existing per-user check.
|
||||
if let CommitMode::Create { folder_id, .. } = &mode {
|
||||
let folder_uuid = Uuid::parse_str(folder_id)
|
||||
.map_err(|_| DomainError::not_found("Folder", folder_id.clone()))?;
|
||||
self.quota
|
||||
.check_drive_quota_by_folder(folder_uuid, total_size)
|
||||
.await?;
|
||||
}
|
||||
|
||||
// ── Whole-file fast path: caller already owns this exact content ──
|
||||
// Mirrors the instant-upload endpoint: a reference bump, no chunk
|
||||
// work at all. Ownership is required — an existing-but-foreign
|
||||
|
||||
@@ -299,23 +299,46 @@ impl FileUploadService {
|
||||
let Some(storage_service) = &self.storage_usage_service else {
|
||||
return;
|
||||
};
|
||||
let Some(owner) = file
|
||||
let delta = file.size as i64;
|
||||
|
||||
// Per-user delta — unchanged from `b5b80549` / `fbbae541`.
|
||||
if let Some(owner) = file
|
||||
.owner_id
|
||||
.as_deref()
|
||||
.and_then(|s| Uuid::parse_str(s).ok())
|
||||
else {
|
||||
return;
|
||||
};
|
||||
let delta = file.size as i64;
|
||||
let service_clone = Arc::clone(storage_service);
|
||||
tokio::spawn(async move {
|
||||
if let Err(e) = service_clone
|
||||
.add_user_storage_usage_delta(owner, delta)
|
||||
.await
|
||||
{
|
||||
warn!("Failed to bump storage usage for {owner}: {e}");
|
||||
}
|
||||
});
|
||||
{
|
||||
let service_clone = Arc::clone(storage_service);
|
||||
tokio::spawn(async move {
|
||||
if let Err(e) = service_clone
|
||||
.add_user_storage_usage_delta(owner, delta)
|
||||
.await
|
||||
{
|
||||
warn!("Failed to bump storage usage for {owner}: {e}");
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
// Per-drive delta (D4) — same fire-and-forget shape, resolves
|
||||
// the drive id from the file's parent folder in one SQL
|
||||
// statement. `storage.drives.used_bytes` is what the per-drive
|
||||
// quota check and the picker quota bar read; drift from
|
||||
// deletes / trash is reconciled by the same sweep that handles
|
||||
// user-side drift.
|
||||
if let Some(folder) = file
|
||||
.folder_id
|
||||
.as_deref()
|
||||
.and_then(|s| Uuid::parse_str(s).ok())
|
||||
{
|
||||
let service_clone = Arc::clone(storage_service);
|
||||
tokio::spawn(async move {
|
||||
if let Err(e) = service_clone
|
||||
.add_drive_storage_usage_delta_by_folder(folder, delta)
|
||||
.await
|
||||
{
|
||||
warn!("Failed to bump drive usage for folder {folder}: {e}");
|
||||
}
|
||||
});
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -128,6 +128,157 @@ impl StorageUsageService {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Incrementally adjust one drive's cached `storage.drives.used_bytes`
|
||||
/// by `delta` bytes — same shape as
|
||||
/// [`Self::add_user_storage_usage_delta`]: single statement, no
|
||||
/// read-then-write window, `GREATEST(0, …)` clamp so a late or
|
||||
/// duplicate adjustment can never drive the counter negative.
|
||||
/// Deletes / trash do not decrement here; the periodic reconciliation
|
||||
/// sweep ([`Self::update_all_drives_storage_usage`]) remains the
|
||||
/// correctness backstop.
|
||||
pub async fn add_drive_storage_usage_delta(
|
||||
&self,
|
||||
drive_id: Uuid,
|
||||
delta: i64,
|
||||
) -> Result<(), DomainError> {
|
||||
sqlx::query(
|
||||
"UPDATE storage.drives
|
||||
SET used_bytes = GREATEST(0, used_bytes + $2)
|
||||
WHERE id = $1",
|
||||
)
|
||||
.bind(drive_id)
|
||||
.bind(delta)
|
||||
.execute(self.pool.as_ref())
|
||||
.await
|
||||
.map_err(|e| DomainError::internal_error("StorageUsage", format!("drive delta: {e}")))?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Same as [`Self::add_drive_storage_usage_delta`] but resolves
|
||||
/// the drive id from a parent folder id in a single statement.
|
||||
/// Avoids a separate `SELECT drive_id FROM storage.folders` round
|
||||
/// trip at the upload hook site (where the folder id is what's
|
||||
/// naturally on the FileDto). The nested SELECT is point-lookup
|
||||
/// on the folder PK; clamp + idempotency properties are
|
||||
/// unchanged.
|
||||
pub async fn add_drive_storage_usage_delta_by_folder(
|
||||
&self,
|
||||
folder_id: Uuid,
|
||||
delta: i64,
|
||||
) -> Result<(), DomainError> {
|
||||
sqlx::query(
|
||||
"UPDATE storage.drives
|
||||
SET used_bytes = GREATEST(0, used_bytes + $2)
|
||||
WHERE id = (SELECT drive_id FROM storage.folders WHERE id = $1)",
|
||||
)
|
||||
.bind(folder_id)
|
||||
.bind(delta)
|
||||
.execute(self.pool.as_ref())
|
||||
.await
|
||||
.map_err(|e| {
|
||||
DomainError::internal_error("StorageUsage", format!("drive delta by folder: {e}"))
|
||||
})?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Pre-upload quota check on a single drive.
|
||||
///
|
||||
/// Read-only `SELECT (used_bytes, quota_bytes) FROM storage.drives`;
|
||||
/// returns `QuotaExceeded` when the projected `used_bytes +
|
||||
/// additional_bytes` would breach `quota_bytes`. A `NULL`
|
||||
/// `quota_bytes` short-circuits to `Ok(())` (unlimited drive —
|
||||
/// admin override / future system drives).
|
||||
///
|
||||
/// Soft cap by design: the check/write window matches the
|
||||
/// user-quota path, bounded by the sweep interval. The clamp on
|
||||
/// `add_drive_storage_usage_delta` and the set-based reconciliation
|
||||
/// keep the counter honest; small over-quota slippage during the
|
||||
/// window is acceptable.
|
||||
pub async fn check_drive_quota(
|
||||
&self,
|
||||
drive_id: Uuid,
|
||||
additional_bytes: u64,
|
||||
) -> Result<(), DomainError> {
|
||||
let row: Option<(i64, Option<i64>)> = sqlx::query_as(
|
||||
"SELECT used_bytes, quota_bytes FROM storage.drives WHERE id = $1",
|
||||
)
|
||||
.bind(drive_id)
|
||||
.fetch_optional(self.pool.as_ref())
|
||||
.await
|
||||
.map_err(|e| DomainError::internal_error("StorageUsage", format!("drive quota lookup: {e}")))?;
|
||||
|
||||
let Some((used, quota)) = row else {
|
||||
// Anti-enum at the upload edge would normally map to 404,
|
||||
// but at this layer we surface the typed not-found and let
|
||||
// the caller decide how to react. In practice the upload
|
||||
// path resolves the drive id from a folder/file lookup
|
||||
// first, so this branch fires only on a deleted-drive race.
|
||||
return Err(DomainError::not_found("Drive", drive_id.to_string()));
|
||||
};
|
||||
let Some(quota) = quota else {
|
||||
return Ok(()); // unlimited
|
||||
};
|
||||
// Saturate on the i64 + u64 sum so a hostile / corrupt counter
|
||||
// can't silently overflow into a negative comparison.
|
||||
let projected = (used as i128) + (additional_bytes as i128);
|
||||
if projected > quota as i128 {
|
||||
return Err(DomainError::new(
|
||||
crate::common::errors::ErrorKind::QuotaExceeded,
|
||||
"Drive",
|
||||
format!(
|
||||
"Drive quota exceeded: {} + {} > {} bytes",
|
||||
used, additional_bytes, quota
|
||||
),
|
||||
));
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Same as [`Self::check_drive_quota`] but resolves the drive id
|
||||
/// from a parent folder id. Mirrors
|
||||
/// [`Self::add_drive_storage_usage_delta_by_folder`] so the upload
|
||||
/// handler (which holds `folder_id` from the multipart form) can
|
||||
/// gate the write in one round trip. Returns
|
||||
/// `DomainError::not_found("Folder", …)` if the folder id doesn't
|
||||
/// resolve — the upload pipeline would 404 on that anyway.
|
||||
pub async fn check_drive_quota_by_folder(
|
||||
&self,
|
||||
folder_id: Uuid,
|
||||
additional_bytes: u64,
|
||||
) -> Result<(), DomainError> {
|
||||
let row: Option<(i64, Option<i64>)> = sqlx::query_as(
|
||||
"SELECT d.used_bytes, d.quota_bytes
|
||||
FROM storage.drives d
|
||||
JOIN storage.folders f ON f.drive_id = d.id
|
||||
WHERE f.id = $1",
|
||||
)
|
||||
.bind(folder_id)
|
||||
.fetch_optional(self.pool.as_ref())
|
||||
.await
|
||||
.map_err(|e| {
|
||||
DomainError::internal_error("StorageUsage", format!("drive quota by folder: {e}"))
|
||||
})?;
|
||||
|
||||
let Some((used, quota)) = row else {
|
||||
return Err(DomainError::not_found("Folder", folder_id.to_string()));
|
||||
};
|
||||
let Some(quota) = quota else {
|
||||
return Ok(()); // unlimited
|
||||
};
|
||||
let projected = (used as i128) + (additional_bytes as i128);
|
||||
if projected > quota as i128 {
|
||||
return Err(DomainError::new(
|
||||
crate::common::errors::ErrorKind::QuotaExceeded,
|
||||
"Drive",
|
||||
format!(
|
||||
"Drive quota exceeded: {} + {} > {} bytes",
|
||||
used, additional_bytes, quota
|
||||
),
|
||||
));
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Spawn a background task that periodically reconciles every user's cached
|
||||
/// `storage_used_bytes` against the actual sum of their files.
|
||||
///
|
||||
@@ -153,7 +304,13 @@ impl StorageUsageService {
|
||||
ticker.tick().await;
|
||||
debug!("Running scheduled storage-usage reconciliation");
|
||||
if let Err(e) = service.update_all_users_storage_usage().await {
|
||||
error!("Scheduled storage-usage reconciliation failed: {}", e);
|
||||
error!("Scheduled user storage-usage reconciliation failed: {}", e);
|
||||
}
|
||||
// Drive sweep runs alongside the user sweep — same
|
||||
// cadence, same maintenance pool. Failure is logged
|
||||
// but doesn't skip the next tick.
|
||||
if let Err(e) = service.update_all_drives_storage_usage().await {
|
||||
error!("Scheduled drive storage-usage reconciliation failed: {}", e);
|
||||
}
|
||||
}
|
||||
});
|
||||
@@ -266,6 +423,63 @@ impl StorageUsagePort for StorageUsageService {
|
||||
let user = self.user_repository.get_user_by_id(user_id).await?;
|
||||
Ok((user.storage_used_bytes(), user.storage_quota_bytes()))
|
||||
}
|
||||
|
||||
async fn add_drive_storage_usage_delta(
|
||||
&self,
|
||||
drive_id: Uuid,
|
||||
delta: i64,
|
||||
) -> Result<(), DomainError> {
|
||||
StorageUsageService::add_drive_storage_usage_delta(self, drive_id, delta).await
|
||||
}
|
||||
|
||||
/// Reconcile every drive's cached `used_bytes` in ONE set-based UPDATE.
|
||||
///
|
||||
/// Same shape as the per-user sweep above: `LEFT JOIN` over the
|
||||
/// `storage.files` aggregate keyed on `drive_id`, `IS DISTINCT
|
||||
/// FROM` guard to skip no-op rewrites so idle drives don't churn
|
||||
/// dead tuples. Runs from the same reconciliation ticker as the
|
||||
/// user sweep; failure is logged but doesn't stop the next tick.
|
||||
async fn update_all_drives_storage_usage(&self) -> Result<(), DomainError> {
|
||||
debug!("Starting drive storage-usage reconciliation sweep");
|
||||
let result = sqlx::query(
|
||||
r#"
|
||||
UPDATE storage.drives d
|
||||
SET used_bytes = COALESCE(t.total, 0)
|
||||
FROM storage.drives d2
|
||||
LEFT JOIN (
|
||||
SELECT drive_id, SUM(size)::bigint AS total
|
||||
FROM storage.files
|
||||
WHERE NOT is_trashed
|
||||
GROUP BY drive_id
|
||||
) t ON t.drive_id = d2.id
|
||||
WHERE d.id = d2.id
|
||||
AND d.used_bytes IS DISTINCT FROM COALESCE(t.total, 0)
|
||||
"#,
|
||||
)
|
||||
.execute(self.pool.as_ref())
|
||||
.await
|
||||
.map_err(|e| {
|
||||
error!("Drive storage-usage reconciliation sweep failed: {}", e);
|
||||
DomainError::internal_error(
|
||||
"StorageUsage",
|
||||
format!("drive reconciliation sweep: {e}"),
|
||||
)
|
||||
})?;
|
||||
|
||||
info!(
|
||||
"Drive storage-usage reconciliation corrected {} drive(s)",
|
||||
result.rows_affected()
|
||||
);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn check_drive_quota(
|
||||
&self,
|
||||
drive_id: Uuid,
|
||||
additional_bytes: u64,
|
||||
) -> Result<(), DomainError> {
|
||||
StorageUsageService::check_drive_quota(self, drive_id, additional_bytes).await
|
||||
}
|
||||
}
|
||||
|
||||
// Make StorageUsageService cloneable to support spawning concurrent tasks
|
||||
|
||||
Reference in New Issue
Block a user