feat(upload): cover chunk upload + add support of different digest hash
Prefer stream storage rather using buffered (in memory)
note: on many unix like tmpfs are in-memory, sungle PUT are sized limited
Storage map (NC stands for Nextcloud gateway)
┌───────────────────────────────────────────────────────┬────────────────────────────────────────────────────────────────────┬─────────────────────────────────────────────────┐
│ Streaming surface │ Destination │ Configurable via │
├───────────────────────────────────────────────────────┼────────────────────────────────────────────────────────────────────┼─────────────────────────────────────────────────┤
│ REST chunked PUT /api/uploads/{id} chunk │ {storage_path}/.uploads/{upload_id}/chunk_{NNNNNN} │ OXICLOUD_STORAGE_PATH (the .uploads subdir is │
│ │ │ hard-wired) │
├───────────────────────────────────────────────────────┼────────────────────────────────────────────────────────────────────┼─────────────────────────────────────────────────┤
│ REST chunked assemble (during /complete) │ {storage_path}/.uploads/{upload_id}/assembled │ same │
├───────────────────────────────────────────────────────┼────────────────────────────────────────────────────────────────────┼─────────────────────────────────────────────────┤
│ NC chunked PUT /dav/uploads/.../{chunk} │ {storage_path}/.uploads/nextcloud/{user}/{upload_id}/{chunk_name} │ same │
├───────────────────────────────────────────────────────┼────────────────────────────────────────────────────────────────────┼─────────────────────────────────────────────────┤
│ NC chunked assemble (during MOVE) │ {storage_path}/.uploads/nextcloud/{user}/{upload_id}/.assembled │ same │
├───────────────────────────────────────────────────────┼────────────────────────────────────────────────────────────────────┼─────────────────────────────────────────────────┤
│ NC single-file PUT /dav/files/.../{path} (via │ OXICLOUD_UPLOAD_TMPDIR if set, else OS default temp (/tmp on │ OXICLOUD_UPLOAD_TMPDIR │
│ spool_body_to_temp) │ Linux) │ │
├───────────────────────────────────────────────────────┼────────────────────────────────────────────────────────────────────┼─────────────────────────────────────────────────┤
│ REST WebDAV PUT /webdav/{path} (via │ same as above │ OXICLOUD_UPLOAD_TMPDIR │
│ spool_body_to_temp) │ │ │
├───────────────────────────────────────────────────────┼────────────────────────────────────────────────────────────────────┼─────────────────────────────────────────────────┤
│ REST multipart upload /api/files/upload │ {storage_path}/.dedup_temp/upload-{uuid} │ OXICLOUD_STORAGE_PATH (hard-wired subdir) │
├───────────────────────────────────────────────────────┼────────────────────────────────────────────────────────────────────┼─────────────────────────────────────────────────┤
│ WOPI PutFile │ OS default temp via NamedTempFile::new() (no override) │ (none — bug worth tracking) │
├───────────────────────────────────────────────────────┼────────────────────────────────────────────────────────────────────┼─────────────────────────────────────────────────┤
│ Final blob storage (after fsync + rename) │ {storage_path}/.blobs/{ab}/{abc…}.blob │ OXICLOUD_STORAGE_PATH │
└───────────────────────────────────────────────────────┴────────────────────────────────────────────────────────────────────┴─────────────────────────────────────────────────┘
one caveat: a malicious user can create many chunked upload and saturate local storage
This commit is contained in:
@@ -13,11 +13,11 @@ use axum::{
|
||||
http::{HeaderMap, StatusCode, header},
|
||||
response::{IntoResponse, Response},
|
||||
};
|
||||
use bytes::Bytes;
|
||||
use serde::{Deserialize, Serialize};
|
||||
use std::sync::Arc;
|
||||
use utoipa::ToSchema;
|
||||
|
||||
use crate::application::ports::chunked_upload_ports::ChecksumAlg;
|
||||
use crate::application::ports::chunked_upload_ports::ChunkedUploadPort;
|
||||
use crate::application::ports::chunked_upload_ports::DEFAULT_CHUNK_SIZE;
|
||||
use crate::application::ports::file_ports::FileUploadUseCase;
|
||||
@@ -27,6 +27,7 @@ use crate::common::di::AppState;
|
||||
use crate::domain::services::authorization::Permission;
|
||||
use crate::interfaces::errors::AppError;
|
||||
use crate::interfaces::middleware::auth::AuthUser;
|
||||
use crate::interfaces::upload_spool::stream_body_to_path;
|
||||
|
||||
/// Request body for creating an upload session
|
||||
#[derive(Debug, Deserialize, ToSchema)]
|
||||
@@ -38,11 +39,17 @@ pub struct CreateUploadRequest {
|
||||
pub chunk_size: Option<usize>,
|
||||
}
|
||||
|
||||
/// Query params for chunk upload
|
||||
/// Query params for chunk upload.
|
||||
///
|
||||
/// `checksumalg` is parsed via [`ChecksumAlg::parse`] and defaults to
|
||||
/// `Md5` when absent — matching the legacy `Content-MD5` contract that
|
||||
/// older clients rely on. Unknown algorithm names produce a 400 with the
|
||||
/// offending value echoed back.
|
||||
#[derive(Debug, Deserialize)]
|
||||
pub struct ChunkUploadParams {
|
||||
pub chunk_index: usize,
|
||||
pub checksum: Option<String>,
|
||||
pub checksumalg: Option<String>,
|
||||
}
|
||||
|
||||
/// Final response after completing upload
|
||||
@@ -200,58 +207,12 @@ impl ChunkedUploadHandler {
|
||||
}
|
||||
}
|
||||
|
||||
/// PATCH /api/uploads/:upload_id - Upload a chunk
|
||||
///
|
||||
/// Query params:
|
||||
/// - chunk_index: The index of the chunk (0-based)
|
||||
/// - checksum: Optional MD5 checksum for verification
|
||||
///
|
||||
/// Body: Raw bytes of the chunk
|
||||
pub(super) async fn upload_chunk_impl(
|
||||
State(state): State<Arc<AppState>>,
|
||||
auth_user: AuthUser,
|
||||
Path(upload_id): Path<String>,
|
||||
Query(params): Query<ChunkUploadParams>,
|
||||
headers: HeaderMap,
|
||||
body: Bytes,
|
||||
) -> impl IntoResponse {
|
||||
let chunked_service = &state.core.chunked_upload_service;
|
||||
|
||||
// Extract checksum from header or query param
|
||||
let checksum = params.checksum.or_else(|| {
|
||||
headers
|
||||
.get("Content-MD5")
|
||||
.and_then(|v| v.to_str().ok())
|
||||
.map(|s| s.to_string())
|
||||
});
|
||||
|
||||
match chunked_service
|
||||
.upload_chunk(&upload_id, auth_user.id, params.chunk_index, body, checksum)
|
||||
.await
|
||||
{
|
||||
Ok(response) => {
|
||||
let mut resp = Response::builder()
|
||||
.status(StatusCode::OK)
|
||||
.header(header::CONTENT_TYPE, "application/json")
|
||||
.header("Upload-Offset", response.bytes_received.to_string())
|
||||
.header(
|
||||
"Upload-Progress",
|
||||
format!("{:.2}", response.progress * 100.0),
|
||||
);
|
||||
|
||||
if response.is_complete {
|
||||
resp = resp.header("Upload-Complete", "true");
|
||||
}
|
||||
|
||||
resp.body(axum::body::Body::from(
|
||||
serde_json::to_string(&response).unwrap(),
|
||||
))
|
||||
.unwrap()
|
||||
.into_response()
|
||||
}
|
||||
Err(e) => AppError::from(e).into_response(),
|
||||
}
|
||||
}
|
||||
// PATCH /api/uploads/:upload_id — moved entirely to the free
|
||||
// function `upload_chunk` below so the body can be streamed
|
||||
// (axum::body::Body) instead of materialised as `Bytes` here.
|
||||
// The port-level `ChunkedUploadPort::upload_chunk` (Bytes-based)
|
||||
// remains for tests and any future caller that genuinely has the
|
||||
// bytes already in memory.
|
||||
|
||||
/// HEAD /api/uploads/:upload_id - Get upload status
|
||||
///
|
||||
@@ -448,40 +409,114 @@ pub async fn create_upload(
|
||||
pub async fn upload_chunk(
|
||||
State(state): State<Arc<AppState>>,
|
||||
auth_user: AuthUser,
|
||||
path: Path<String>,
|
||||
query: Query<ChunkUploadParams>,
|
||||
Path(upload_id): Path<String>,
|
||||
Query(params): Query<ChunkUploadParams>,
|
||||
headers: HeaderMap,
|
||||
request: Request,
|
||||
) -> impl IntoResponse {
|
||||
// Cap the chunk body at `storage.chunk_max_bytes` (env
|
||||
// `OXICLOUD_CHUNK_MAX_BYTES`, default 100 MB). Previous code used
|
||||
// `usize::MAX` and `unwrap_or_default()` — two compounding bugs:
|
||||
// - No upper bound → an oversized chunk OOMs the server.
|
||||
// - Silent fallback to an empty body on transport error → the
|
||||
// inner size check would either reject (good case) or — if the
|
||||
// declared chunk_size was 0 (illegal but conceivable) — accept
|
||||
// an empty upload as success. Either way the client got no
|
||||
// actionable error.
|
||||
let chunked_service = &state.core.chunked_upload_service;
|
||||
let max_chunk = state.core.config.storage.chunk_max_bytes;
|
||||
let body = match axum::body::to_bytes(request.into_body(), max_chunk).await {
|
||||
Ok(b) => b,
|
||||
|
||||
// ── Resolve the client's checksum + algorithm ────────────────────
|
||||
// Wire shape: `?checksum=<hex>&checksumalg=<name>` (or `Content-MD5`
|
||||
// header for older clients). When `checksumalg` is omitted we
|
||||
// default to MD5, matching the legacy contract — switching the
|
||||
// default would silently break any client still relying on
|
||||
// `Content-MD5` semantics.
|
||||
let expected_checksum = params.checksum.clone().or_else(|| {
|
||||
headers
|
||||
.get("Content-MD5")
|
||||
.and_then(|v| v.to_str().ok())
|
||||
.map(|s| s.to_string())
|
||||
});
|
||||
let alg = match params.checksumalg.as_deref() {
|
||||
Some(name) => match ChecksumAlg::parse(name) {
|
||||
Some(a) => a,
|
||||
None => {
|
||||
return AppError::bad_request(format!(
|
||||
"Unsupported checksumalg: {name} (supported: md5, sha256, blake3)"
|
||||
))
|
||||
.into_response();
|
||||
}
|
||||
},
|
||||
None => ChecksumAlg::Md5,
|
||||
};
|
||||
// Only compute the hash when the client supplied an `expected_checksum`
|
||||
// to verify against — saves ~30 ms per chunk for clients that don't.
|
||||
let alg_to_compute = expected_checksum.as_ref().map(|_| alg);
|
||||
|
||||
// ── Phase 1: prepare ─────────────────────────────────────────────
|
||||
// Validates session ownership + chunk index, returns the on-disk
|
||||
// path and the chunk's declared size. The handler streams the body
|
||||
// to that path; service finalises bookkeeping after the write.
|
||||
let (chunk_path, _expected_size) = match chunked_service
|
||||
.prepare_chunk(&upload_id, auth_user.id, params.chunk_index)
|
||||
.await
|
||||
{
|
||||
Ok(p) => p,
|
||||
Err(e) => return AppError::from(e).into_response(),
|
||||
};
|
||||
|
||||
// ── Phase 2: stream the body straight to disk ────────────────────
|
||||
// Peak heap ~one HTTP frame (~64 KB) regardless of chunk size or
|
||||
// `chunk_max_bytes`. Optional incremental hashing happens here so
|
||||
// verification doesn't require reading the chunk file back.
|
||||
let streamed = match stream_body_to_path(
|
||||
request.into_body(),
|
||||
&chunk_path,
|
||||
max_chunk,
|
||||
alg_to_compute,
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(s) => s,
|
||||
Err(e) => {
|
||||
tracing::warn!(
|
||||
error = %e,
|
||||
upload_id = %path.0,
|
||||
error = ?e,
|
||||
upload_id = %upload_id,
|
||||
chunk_index = params.chunk_index,
|
||||
max_chunk,
|
||||
"Chunked upload PATCH rejected — body read failed (size cap or transport error)"
|
||||
"Chunked upload PATCH rejected — streaming write failed (cap, transport, or IO)"
|
||||
);
|
||||
return AppError::payload_too_large(format!(
|
||||
"Chunk read failed (cap {} bytes): {}",
|
||||
max_chunk, e
|
||||
))
|
||||
.into_response();
|
||||
return e.into_response();
|
||||
}
|
||||
};
|
||||
ChunkedUploadHandler::upload_chunk_impl(State(state), auth_user, path, query, headers, body)
|
||||
|
||||
// ── Phase 3: commit ──────────────────────────────────────────────
|
||||
// Size + checksum verification + session state update. Same RAM-only
|
||||
// DashMap shard ownership pattern as the legacy `upload_chunk_inner`
|
||||
// (held only for ~µs; bitmask persist done after release).
|
||||
let response = match chunked_service
|
||||
.commit_chunk(
|
||||
&upload_id,
|
||||
auth_user.id,
|
||||
params.chunk_index,
|
||||
streamed.bytes_written,
|
||||
streamed.checksum_hex,
|
||||
expected_checksum,
|
||||
)
|
||||
.await
|
||||
.into_response()
|
||||
{
|
||||
Ok(r) => r,
|
||||
Err(e) => return AppError::from(e).into_response(),
|
||||
};
|
||||
|
||||
let mut resp = Response::builder()
|
||||
.status(StatusCode::OK)
|
||||
.header(header::CONTENT_TYPE, "application/json")
|
||||
.header("Upload-Offset", response.bytes_received.to_string())
|
||||
.header(
|
||||
"Upload-Progress",
|
||||
format!("{:.2}", response.progress * 100.0),
|
||||
);
|
||||
if response.is_complete {
|
||||
resp = resp.header("Upload-Complete", "true");
|
||||
}
|
||||
resp.body(axum::body::Body::from(
|
||||
serde_json::to_string(&response).unwrap(),
|
||||
))
|
||||
.unwrap()
|
||||
.into_response()
|
||||
}
|
||||
|
||||
#[utoipa::path(
|
||||
|
||||
Reference in New Issue
Block a user