//! File operations API. //! //! All operations live on the same URL shape as the listing, distinguished by //! method (and content type for POST): //! //! - `GET /api/files/{root_id}` and `/api/files/{root_id}/{*path}` — list //! - `DELETE /api/files/{root_id}/{*path}` — delete a file or folder //! - `POST /api/files/{root_id}/{*path}` — create a folder (no body) //! - `POST /api/files/{root_id}/{*path}` (JSON body) — rename / move / copy //! - `POST /api/files/{root_id}/{*path}` (multipart) — upload into the dir use std::io::{self, Read}; use std::path::Component; use std::sync::Arc; use std::time::UNIX_EPOCH; use axum::extract::{Path as AxumPath, Query as AxumQuery, State}; use axum::http::{header, StatusCode}; use axum::response::{IntoResponse, Response}; use axum::Json; use futures_util::StreamExt; use multer::Multipart; use serde::Deserialize; use tokio::io::AsyncWriteExt; use tokio::sync::mpsc; use tokio_stream::wrappers::ReceiverStream; use crate::api::common::AuthUser; use crate::archive::{self, ArchiveFormat}; use crate::db::{RootRow, ShareRow}; use crate::error::{ApiError, AppState}; use crate::fs::{self, FsError}; /// Upper bound for the in-memory text endpoint (preview, later editor). const MAX_TEXT_BYTES: u64 = 2 * 1024 * 1024; // --------------------------------------------------------------------------- // Bodies / query params // --------------------------------------------------------------------------- #[derive(Deserialize)] pub struct MutationBody { pub op: String, // "rename" | "move" | "copy" #[serde(default)] pub new_name: Option, #[serde(default)] pub dst_root_id: Option, #[serde(default)] pub dst: Option, #[serde(default)] pub overwrite: bool, } /// Query params for `GET /api/files/{root_id}/{*path}`. Without `action` the /// route lists the directory; `?action=download|preview|content` serve the /// item itself. #[derive(Deserialize, Default)] pub struct FileQuery { #[serde(default)] action: Option, #[serde(default)] format: Option, } // --------------------------------------------------------------------------- // Listing // --------------------------------------------------------------------------- /// GET /api/files/{root_id}/{*path} — list a directory, or serve the item /// itself via `?action=download|preview|content`. pub async fn file_get( State(state): State>, auth: AuthUser, path: AxumPath<(i64, String)>, query: AxumQuery, ) -> Result { let (root_id, req_rel) = path.0; match query.action.as_deref() { Some("download") => download(state, auth, root_id, req_rel, query.format.as_deref()).await, Some("preview") => preview(state, auth, root_id, req_rel).await, Some("content") => content(state, auth, root_id, req_rel).await, _ => { let json = list_inner(state, auth, root_id, req_rel).await?; Ok(json.into_response()) } } } /// GET /api/files/{root_id} — list the root directory itself, or serve the /// root item via `?action=download|preview|content` (the root of a *file* /// share is the file itself). pub async fn list_root( State(state): State>, auth: AuthUser, path: AxumPath, query: AxumQuery, ) -> Result { let root_id = path.0; match query.action.as_deref() { Some("download") => { download(state, auth, root_id, String::new(), query.format.as_deref()).await } Some("preview") => preview(state, auth, root_id, String::new()).await, Some("content") => content(state, auth, root_id, String::new()).await, _ => { let json = list_inner(state, auth, root_id, String::new()).await?; Ok(json.into_response()) } } } async fn list_inner( state: Arc, auth: AuthUser, root_id: i64, req_rel: String, ) -> Result, ApiError> { let root = find_root(&auth.roots, root_id)?; let server_root = state.root.clone(); let root_rel = root.path.clone(); let full = tokio::task::spawn_blocking(move || fs::resolve_path(&server_root, &root_rel, &req_rel)) .await .map_err(|_| ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, "internal error"))??; let entries = tokio::task::spawn_blocking(move || fs::list_dir(&full)) .await .map_err(|_| ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, "internal error"))??; Ok(Json(serde_json::json!({ "entries": entries }))) } // --------------------------------------------------------------------------- // download / preview / content (milestone 4) // --------------------------------------------------------------------------- /// Escape a file name for a `Content-Disposition` header value. fn disp_name(name: &str) -> String { name.replace('\\', "\\\\") .replace('"', "\\\"") .replace(['\n', '\r'], "_") } /// Resolve the requested item to an absolute path + metadata (blocking). async fn resolve_item( state: &AppState, root: &RootRow, req_rel: &str, share: Option<&ShareRow>, ) -> Result<(std::path::PathBuf, String, bool, u64), ApiError> { let (server_root, root_rel, rel) = (state.root.clone(), root.path.clone(), req_rel.to_string()); // A file share's synthetic root *is* the file, so resolve it directly. let share_target = share.filter(|s| s.is_file).map(|s| s.target.clone()); let inner = tokio::task::spawn_blocking(move || { let full = match share_target { Some(target) => fs::resolve_file(&server_root, &target)?, None => fs::resolve_path(&server_root, &root_rel, &rel)?, }; let name = full .file_name() .map(|n| n.to_string_lossy().into_owned()) .ok_or_else(|| FsError::Invalid("invalid path".to_string()))?; let meta = std::fs::metadata(&full).map_err(|_| FsError::NotFound)?; Ok::<_, FsError>((full, name, meta.is_dir(), meta.len())) }) .await .map_err(|_| ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, "internal error"))?; Ok(inner?) } /// `GET ...?action=download` — a single file as-is, a folder as an archive /// (format chosen by the client). async fn download( state: Arc, auth: AuthUser, root_id: i64, req_rel: String, format: Option<&str>, ) -> Result { let root = find_root(&auth.roots, root_id)?; let (full, name, is_dir, size) = resolve_item(&state, root, &req_rel, auth.share.as_ref()).await?; if !is_dir { return file_response(&full, &name, size, false).await; } let fmt = format.and_then(ArchiveFormat::parse).ok_or_else(|| { ApiError::new( StatusCode::BAD_REQUEST, "format must be one of: zip, tar, tar.gz, tar.zst", ) })?; let disp = format!( "attachment; filename=\"{}.{}\"", disp_name(&name), fmt.extension() ); let body = stream_archive(fmt, full, name); Response::builder() .status(StatusCode::OK) .header(header::CONTENT_TYPE, fmt.mime()) .header(header::CONTENT_DISPOSITION, disp) .body(body) .map_err(|e| { ApiError::new( StatusCode::INTERNAL_SERVER_ERROR, format!("bad response: {e}"), ) }) } /// `GET ...?action=preview` — a single file, inline (for native media). async fn preview( state: Arc, auth: AuthUser, root_id: i64, req_rel: String, ) -> Result { let root = find_root(&auth.roots, root_id)?; let (full, name, is_dir, size) = resolve_item(&state, root, &req_rel, auth.share.as_ref()).await?; if is_dir { return Err(ApiError::new(StatusCode::BAD_REQUEST, "not a file")); } file_response(&full, &name, size, true).await } /// `GET ...?action=content` — raw file bytes for the text preview/editor. /// Capped at `MAX_TEXT_BYTES`. async fn content( state: Arc, auth: AuthUser, root_id: i64, req_rel: String, ) -> Result { let root = find_root(&auth.roots, root_id)?; let (full, _name, is_dir, size) = resolve_item(&state, root, &req_rel, auth.share.as_ref()).await?; if is_dir { return Err(ApiError::new(StatusCode::BAD_REQUEST, "not a file")); } if size > MAX_TEXT_BYTES { return Err(ApiError::new( StatusCode::PAYLOAD_TOO_LARGE, "file too large to preview", )); } let bytes = tokio::fs::read(&full) .await .map_err(|_| ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, "internal error"))?; let meta = tokio::fs::metadata(&full) .await .map_err(|_| ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, "internal error"))?; let mtime = meta .modified() .ok() .and_then(|t| t.duration_since(UNIX_EPOCH).ok()) .map(|d| d.as_secs()) .unwrap_or(0); Ok(( [ ( header::CONTENT_TYPE, "text/plain; charset=utf-8".to_string(), ), ( axum::http::HeaderName::from_static("x-file-mtime"), mtime.to_string(), ), ], bytes, ) .into_response()) } /// `PUT ...?action=content` — save a file's text contents (the editor). /// /// Requires a read-write root. The body is the new contents. If the /// `X-Expected-Mtime` header is present, the file's current mtime must match /// it, otherwise `409 Conflict` (the file changed on disk since it was read). /// Returns the file's new mtime so the client can anchor the next check. pub async fn file_put( State(state): State>, auth: AuthUser, path: AxumPath<(i64, String)>, query: AxumQuery, headers: axum::http::HeaderMap, body: axum::body::Bytes, ) -> Result, ApiError> { let (root_id, req_rel) = path.0; if query.action.as_deref() != Some("content") { return Err(ApiError::new( StatusCode::BAD_REQUEST, "expected action=content", )); } let root = require_rw_root(&auth.roots, root_id)?; if body.len() as u64 > MAX_TEXT_BYTES { return Err(ApiError::new( StatusCode::PAYLOAD_TOO_LARGE, "file too large to save", )); } let expected: Option = headers .get("x-expected-mtime") .and_then(|v| v.to_str().ok()) .and_then(|s| s.parse().ok()); let (server_root, root_rel, rel) = (state.root.clone(), root.path.clone(), req_rel); let content = body.to_vec(); let mtime = tokio::task::spawn_blocking(move || { fs::save_file(&server_root, &root_rel, &rel, &content, expected) }) .await .map_err(|_| ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, "internal error"))??; Ok(Json(serde_json::json!({ "ok": true, "mtime": mtime }))) } /// Stream a single file to the client with the right disposition. async fn file_response( full: &std::path::Path, name: &str, size: u64, inline: bool, ) -> Result { let mime = mime_guess::from_path(full) .first_or_octet_stream() .to_string(); let disp = if inline { format!("inline; filename=\"{}\"", disp_name(name)) } else { format!("attachment; filename=\"{}\"", disp_name(name)) }; let body = stream_file(full.to_path_buf()); Response::builder() .status(StatusCode::OK) .header(header::CONTENT_TYPE, mime) .header(header::CONTENT_DISPOSITION, disp) .header(header::CONTENT_LENGTH, size) .body(body) .map_err(|e| { ApiError::new( StatusCode::INTERNAL_SERVER_ERROR, format!("bad response: {e}"), ) }) } /// Stream `path` to the client in chunks (blocking reader → channel). fn stream_file(path: std::path::PathBuf) -> axum::body::Body { let (tx, rx) = mpsc::channel::>(16); tokio::task::spawn_blocking(move || { let mut f = match std::fs::File::open(&path) { Ok(f) => f, Err(e) => { tracing::warn!(error = %e, path = %path.display(), "download open failed"); return; } }; let mut buf = vec![0u8; 256 * 1024]; loop { match f.read(&mut buf) { Ok(0) => break, Ok(n) => { // Client gone → stop producing. if tx.blocking_send(buf[..n].to_vec()).is_err() { break; } } Err(e) => { tracing::warn!(error = %e, path = %path.display(), "download read failed"); break; } } } }); let stream = ReceiverStream::new(rx).map(Ok::<_, io::Error>); axum::body::Body::from_stream(stream) } /// Stream an archive of `dir` (top-level entry `top`) to the client. fn stream_archive(fmt: ArchiveFormat, dir: std::path::PathBuf, top: String) -> axum::body::Body { let (tx, rx) = mpsc::channel::>(16); tokio::task::spawn_blocking(move || { let mut sink = ChanWriter::new(tx); if let Err(e) = archive::build(fmt, &dir, &top, &mut sink) { tracing::warn!(error = %e, dir = %dir.display(), "archive build failed"); } // Dropping the sink flushes its buffer and closes the channel. }); let stream = ReceiverStream::new(rx).map(Ok::<_, io::Error>); axum::body::Body::from_stream(stream) } /// A `Write` that buffers chunks and forwards them over an mpsc channel — the /// bridge between the blocking archive builder and the async response body. struct ChanWriter { tx: mpsc::Sender>, buf: Vec, } impl ChanWriter { fn new(tx: mpsc::Sender>) -> Self { Self { tx, buf: Vec::with_capacity(64 * 1024), } } } impl io::Write for ChanWriter { fn write(&mut self, b: &[u8]) -> io::Result { self.buf.extend_from_slice(b); if self.buf.len() >= 64 * 1024 { io::Write::flush(self)?; } Ok(b.len()) } fn flush(&mut self) -> io::Result<()> { if !self.buf.is_empty() { let chunk = std::mem::take(&mut self.buf); self.tx .blocking_send(chunk) .map_err(|_| io::Error::new(io::ErrorKind::BrokenPipe, "client disconnected"))?; } Ok(()) } } impl Drop for ChanWriter { fn drop(&mut self) { let _ = io::Write::flush(self); } } // --------------------------------------------------------------------------- // POST dispatch: mkdir | rename/move/copy | upload // --------------------------------------------------------------------------- /// POST /api/files/{root_id} — top-level operations (upload into the root /// directory). The bare root never targets a specific item, so JSON ops with /// a missing folder name are rejected by the individual handlers. pub async fn dispatch_root( State(state): State>, auth: AuthUser, path: AxumPath, headers: axum::http::HeaderMap, req: axum::http::Request, ) -> Result, ApiError> { dispatch_inner(state, auth, path.0, String::new(), headers, req).await } /// POST /api/files/{root_id}/{*path} pub async fn dispatch( State(state): State>, auth: AuthUser, path: AxumPath<(i64, String)>, headers: axum::http::HeaderMap, req: axum::http::Request, ) -> Result, ApiError> { let (root_id, req_rel) = path.0; dispatch_inner(state, auth, root_id, req_rel, headers, req).await } async fn dispatch_inner( state: Arc, auth: AuthUser, root_id: i64, req_rel: String, headers: axum::http::HeaderMap, req: axum::http::Request, ) -> Result, ApiError> { let ct = headers .get(header::CONTENT_TYPE) .and_then(|v| v.to_str().ok()) .unwrap_or(""); if ct.starts_with("multipart/form-data") { return upload(state, auth, root_id, req_rel, req).await; } if ct.starts_with("application/json") { let bytes = axum::body::to_bytes(req.into_body(), 1_000_000) .await .map_err(|_| ApiError::new(StatusCode::BAD_REQUEST, "invalid request body"))?; let body: MutationBody = axum::Json::from_bytes(&bytes) .map_err(|_| ApiError::new(StatusCode::BAD_REQUEST, "invalid request body"))? .0; return mutation(state, auth, root_id, req_rel, body).await; } mkdir(state, auth, root_id, req_rel).await } // --------------------------------------------------------------------------- // mkdir // --------------------------------------------------------------------------- async fn mkdir( state: Arc, auth: AuthUser, root_id: i64, req_rel: String, ) -> Result, ApiError> { let root = require_rw_root(&auth.roots, root_id)?; if req_rel.trim().is_empty() { return Err(ApiError::new( StatusCode::BAD_REQUEST, "a folder name is required", )); } let (server_root, root_rel, rel) = (state.root.clone(), root.path.clone(), req_rel); tokio::task::spawn_blocking(move || fs::mkdir(&server_root, &root_rel, &rel)) .await .map_err(|_| ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, "internal error"))??; Ok(Json(serde_json::json!({ "ok": true }))) } // --------------------------------------------------------------------------- // rename / move / copy // --------------------------------------------------------------------------- async fn mutation( state: Arc, auth: AuthUser, root_id: i64, req_rel: String, body: MutationBody, ) -> Result, ApiError> { match body.op.as_str() { "rename" => { let new_name = body .new_name .as_deref() .ok_or_else(|| ApiError::new(StatusCode::BAD_REQUEST, "new_name is required"))? .to_string(); let root = require_rw_root(&auth.roots, root_id)?; let (server_root, root_rel, rel, overwrite) = ( state.root.clone(), root.path.clone(), req_rel, body.overwrite, ); tokio::task::spawn_blocking(move || { fs::rename_item(&server_root, &root_rel, &rel, &new_name, overwrite) }) .await .map_err(|_| ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, "internal error"))??; Ok(Json(serde_json::json!({ "ok": true }))) } "move" | "copy" => { let dst_root_id = body .dst_root_id .ok_or_else(|| ApiError::new(StatusCode::BAD_REQUEST, "dst_root_id is required"))?; let dst = body.dst.clone().unwrap_or_default(); // Moving or copying out of a folder requires rw there; copying // *from* a read-only root is fine. let src_root = if body.op == "move" { require_rw_root(&auth.roots, root_id)? } else { find_root(&auth.roots, root_id)? }; let dst_root = require_rw_root(&auth.roots, dst_root_id)?; let (server_root, src_rel, dst_rel, dst_path, rel, overwrite) = ( state.root.clone(), src_root.path.clone(), dst_root.path.clone(), dst, req_rel, body.overwrite, ); let op_is_move = body.op == "move"; tokio::task::spawn_blocking(move || { if op_is_move { fs::move_item(&server_root, &src_rel, &rel, &dst_rel, &dst_path, overwrite) } else { fs::copy_item(&server_root, &src_rel, &rel, &dst_rel, &dst_path, overwrite) } }) .await .map_err(|_| ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, "internal error"))??; Ok(Json(serde_json::json!({ "ok": true }))) } _ => Err(ApiError::new( StatusCode::BAD_REQUEST, "unknown op (expected rename, move or copy)", )), } } // --------------------------------------------------------------------------- // DELETE // --------------------------------------------------------------------------- pub async fn delete( State(state): State>, auth: AuthUser, path: AxumPath<(i64, String)>, ) -> Result, ApiError> { let (root_id, req_rel) = path.0; let root = require_rw_root(&auth.roots, root_id)?; if req_rel.trim().is_empty() { return Err(ApiError::new( StatusCode::BAD_REQUEST, "a path inside the folder is required", )); } let (server_root, root_rel, rel) = (state.root.clone(), root.path.clone(), req_rel); let is_dir = tokio::task::spawn_blocking(move || fs::remove_item(&server_root, &root_rel, &rel)) .await .map_err(|_| ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, "internal error"))??; Ok(Json(serde_json::json!({ "ok": true, "is_dir": is_dir }))) } // --------------------------------------------------------------------------- // Upload (multipart) // --------------------------------------------------------------------------- async fn upload( state: Arc, auth: AuthUser, root_id: i64, req_rel: String, req: axum::http::Request, ) -> Result, ApiError> { let root = require_rw_root(&auth.roots, root_id)?; let base = { let (server_root, root_rel, rel) = (state.root.clone(), root.path.clone(), req_rel); tokio::task::spawn_blocking(move || fs::resolve_dir(&server_root, &root_rel, &rel)) .await .map_err(|_| ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, "internal error"))?? }; let boundary = req .headers() .get(header::CONTENT_TYPE) .and_then(|v| v.to_str().ok()) .and_then(parse_boundary) .ok_or_else(|| { ApiError::new( StatusCode::BAD_REQUEST, "expected multipart/form-data with a boundary", ) })?; let overwrite = parse_overwrite(req.uri()); let stream = req.into_body().into_data_stream(); let mut multipart = Multipart::new(stream, boundary); let mut uploaded: usize = 0; let mut skipped: Vec = Vec::new(); while let Some(mut field) = multipart .next_field() .await .map_err(|_| ApiError::new(StatusCode::BAD_REQUEST, "invalid upload data"))? { let part_name = field .name() .filter(|n| !n.is_empty()) .or_else(|| field.file_name()) .map(str::to_string) .ok_or_else(|| ApiError::new(StatusCode::BAD_REQUEST, "part without a name"))?; validate_rel_path(&part_name)?; let target = base.join(&part_name); let parent = target .parent() .filter(|p| !p.as_os_str().is_empty()) .ok_or_else(|| ApiError::new(StatusCode::BAD_REQUEST, "invalid part name"))?; if !parent.is_dir() { let p = parent.to_path_buf(); tokio::task::spawn_blocking(move || std::fs::create_dir_all(&p)) .await .map_err(|_| { ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, "internal error") })??; } if target.exists() && !overwrite { // Drain this part and report it as a conflict at the end. while let Some(_chunk) = field .chunk() .await .map_err(|_| ApiError::new(StatusCode::BAD_REQUEST, "invalid upload data"))? { } skipped.push(part_name); continue; } if target.exists() { if target.is_dir() { return Err(ApiError::new( StatusCode::CONFLICT, "a folder with this name already exists", )); } tokio::fs::remove_file(&target) .await .map_err(|_| ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, "internal error"))?; } // Stream to a temp file in the same directory, then rename into place. let suffix = crate::auth::random_token(); let tmp = parent.join(format!(".upload-{suffix}")); let mut tmp_file = tokio::fs::File::create(&tmp) .await .map_err(|_| ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, "internal error"))?; let write_failed = loop { match field .chunk() .await .map_err(|_| ApiError::new(StatusCode::BAD_REQUEST, "invalid upload data")) { Ok(Some(chunk)) => { if let Err(e) = tmp_file.write_all(&chunk).await { tracing::warn!(error = %e, "write failed during upload"); break true; } } Ok(None) => break false, Err(e) => return Err(e), } }; if write_failed { let _ = tokio::fs::remove_file(&tmp).await; return Err(ApiError::new( StatusCode::INTERNAL_SERVER_ERROR, "could not save the file", )); } let tmp2 = tmp.clone(); let target2 = target.clone(); let renamed = tokio::task::spawn_blocking(move || std::fs::rename(&tmp2, &target2)) .await .map_err(|_| ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, "internal error"))?; renamed.map_err(|e| { let _ = std::fs::remove_file(&tmp); tracing::warn!(error = %e, "rename failed during upload"); ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, "internal error") })?; uploaded += 1; } if uploaded == 0 && skipped.is_empty() { return Err(ApiError::new( StatusCode::BAD_REQUEST, "no files were uploaded", )); } if !skipped.is_empty() { return Err( ApiError::new(StatusCode::CONFLICT, "some files already exist") .with_extra(serde_json::json!({ "skipped": skipped, "uploaded": uploaded })), ); } Ok(Json( serde_json::json!({ "ok": true, "uploaded": uploaded }), )) } fn parse_boundary(content_type: &str) -> Option { content_type .split(';') .map(|s| s.trim()) .find_map(|s| s.strip_prefix("boundary=")) .map(|b| b.trim_matches('"').to_string()) .filter(|b| !b.is_empty()) } fn parse_overwrite(uri: &axum::http::Uri) -> bool { uri.query() .map(|q| { q.split('&') .any(|kv| kv == "overwrite=true" || kv == "overwrite=1") }) .unwrap_or(false) } fn validate_rel_path(name: &str) -> Result<(), ApiError> { for c in std::path::Path::new(name).components() { match c { Component::Normal(_) => {} _ => { return Err(ApiError::new( StatusCode::BAD_REQUEST, "invalid file path in upload", )) } } } Ok(()) } // --------------------------------------------------------------------------- // Helpers // --------------------------------------------------------------------------- fn find_root(roots: &[RootRow], root_id: i64) -> Result<&RootRow, ApiError> { roots .iter() .find(|r| r.id == root_id) .ok_or_else(|| ApiError::new(StatusCode::FORBIDDEN, "no such folder")) } fn require_rw_root(roots: &[RootRow], root_id: i64) -> Result<&RootRow, ApiError> { let root = find_root(roots, root_id)?; if root.mode != "rw" { return Err(ApiError::new(StatusCode::FORBIDDEN, "read-only folder")); } Ok(root) }