files.rs
⎇
Raw
1//! File operations API.
2//!
3//! All operations live on the same URL shape as the listing, distinguished by
4//! method (and content type for POST):
5//!
6//! - `GET /api/files/{root_id}` and `/api/files/{root_id}/{*path}` — list
7//! - `DELETE /api/files/{root_id}/{*path}` — delete a file or folder
8//! - `POST /api/files/{root_id}/{*path}` — create a folder (no body)
9//! - `POST /api/files/{root_id}/{*path}` (JSON body) — rename / move / copy
10//! - `POST /api/files/{root_id}/{*path}` (multipart) — upload into the dir
11
12use std::io::{self, Read};
13use std::path::{Component, PathBuf};
14use std::sync::Arc;
15
16use axum::Json;
17use axum::extract::{Path as AxumPath, Query as AxumQuery, State};
18use axum::http::{StatusCode, header};
19use axum::response::{IntoResponse, Response};
20use futures_util::StreamExt;
21use multer::Multipart;
22use serde::Deserialize;
23use tokio::io::AsyncWriteExt;
24use tokio::sync::mpsc;
25use tokio_stream::wrappers::ReceiverStream;
26
27use crate::api::common::{AuthUser, blocking, target_rel};
28use crate::archive::{self, ArchiveFormat};
29use crate::db::{RootRow, ShareRow};
30use crate::error::{ApiError, AppState};
31use crate::fs::{self, FsError};
32use api_types::{
33 DeleteResp, FilesResp, Mutation, OkResp, Op, P_ACTION, P_OVERWRITE, SaveResp, UploadResp,
34};
35
36/// Upper bound for the in-memory text endpoint (preview, later editor).
37pub(super) const MAX_TEXT_BYTES: u64 = 2 * 1024 * 1024;
38
39// ---------------------------------------------------------------------------
40// Query params
41// ---------------------------------------------------------------------------
42
43/// Query params for `GET /api/files/{root_id}/{*path}`. Without `action` the
44/// route lists the directory; `?action=download|preview|content` serve the
45/// item itself.
46///
47/// The field names are the shared [`P_ACTION`] / [`P_FORMAT`] constants.
48/// `#[serde(rename)]` only takes a literal, so that link cannot be written
49/// here; `tests::query_fields_are_the_shared_constants` pins it instead.
50#[derive(Deserialize, Default)]
51pub struct FileQuery {
52 #[serde(default)]
53 action: Option<String>,
54 #[serde(default)]
55 format: Option<String>,
56}
57
58// ---------------------------------------------------------------------------
59// Listing
60// ---------------------------------------------------------------------------
61
62/// GET /api/files/{root_id}/{*path} — list a directory, or serve the item
63/// itself via `?action=download|preview|content`.
64pub async fn file_get(
65 State(state): State<Arc<AppState>>,
66 auth: AuthUser,
67 path: AxumPath<(i64, String)>,
68 query: AxumQuery<FileQuery>,
69 headers: axum::http::HeaderMap,
70) -> Result<Response, ApiError> {
71 let (root_id, req_rel) = path.0;
72 match query.action.as_deref() {
73 Some(a) if a == api_types::ACTION_DOWNLOAD => {
74 let range = range_header(&headers);
75 download(
76 state,
77 auth,
78 root_id,
79 req_rel,
80 query.format.as_deref(),
81 range,
82 ims_header(&headers),
83 )
84 .await
85 }
86 Some(a) if a == api_types::ACTION_PREVIEW => {
87 preview(
88 state,
89 auth,
90 root_id,
91 req_rel,
92 range_header(&headers),
93 ims_header(&headers),
94 )
95 .await
96 }
97 Some(a) if a == api_types::ACTION_CONTENT => content(state, auth, root_id, req_rel).await,
98 Some(a) if a == api_types::ACTION_THUMB => thumb(state, auth, root_id, req_rel).await,
99 _ => {
100 let json = list_inner(state, auth, root_id, req_rel).await?;
101 Ok(json.into_response())
102 }
103 }
104}
105
106/// GET /api/files/{root_id} — list the root directory itself, or serve the
107/// root item via `?action=download|preview|content` (the root of a *file*
108/// share is the file itself).
109///
110/// Same thing as [`file_get`] with an empty path, so it delegates there.
111pub async fn list_root(
112 state: State<Arc<AppState>>,
113 auth: AuthUser,
114 path: AxumPath<i64>,
115 query: AxumQuery<FileQuery>,
116 headers: axum::http::HeaderMap,
117) -> Result<Response, ApiError> {
118 file_get(
119 state,
120 auth,
121 AxumPath((path.0, String::new())),
122 query,
123 headers,
124 )
125 .await
126}
127
128fn range_header(headers: &axum::http::HeaderMap) -> Option<String> {
129 headers
130 .get(header::RANGE)
131 .and_then(|v| v.to_str().ok())
132 .map(str::to_string)
133}
134
135fn ims_header(headers: &axum::http::HeaderMap) -> Option<String> {
136 headers
137 .get(header::IF_MODIFIED_SINCE)
138 .and_then(|v| v.to_str().ok())
139 .map(str::to_string)
140}
141
142/// Caching policy for a served file.
143///
144/// `private` because the response depends on who asked. `no-cache`, not
145/// `no-store`, so a client can revalidate a large media file and get a `304`
146/// instead of re-downloading it.
147const FILE_CACHE: &str = "private, no-cache";
148
149/// An mtime (unix seconds) as an HTTP-date, the `Last-Modified` format.
150/// None for an unknown mtime (0) and for a file written in the last two
151/// seconds: whole-second dates cannot tell two writes in one second apart.
152fn http_date(mtime: i64) -> Option<String> {
153 if mtime <= 0 || mtime + 2 > chrono::Utc::now().timestamp() {
154 return None;
155 }
156 chrono::DateTime::from_timestamp(mtime, 0)
157 .map(|d| d.format("%a, %d %b %Y %H:%M:%S GMT").to_string())
158}
159
160/// True when the client's `If-Modified-Since` is at or after `mtime`.
161/// HTTP-dates carry whole seconds, so the comparison is second-precision.
162fn unmodified_since(ims: Option<&str>, mtime: i64) -> bool {
163 chrono::DateTime::parse_from_rfc2822(ims.unwrap_or_default())
164 .is_ok_and(|t| t.timestamp() >= mtime)
165}
166
167async fn list_inner(
168 state: Arc<AppState>,
169 auth: AuthUser,
170 root_id: i64,
171 req_rel: String,
172) -> Result<Json<FilesResp>, ApiError> {
173 // A file share's root is the file itself: there is nothing to list.
174 if auth.share.as_ref().is_some_and(|s| s.is_file) {
175 return Err(FsError::NotADirectory.into());
176 }
177 let root = find_root(&auth.roots, root_id)?;
178 let server_root = state.root.clone();
179 let root_rel = root.path.clone();
180 // Resolve and list in one blocking hop: both are filesystem work.
181 let (entries, truncated) = blocking(move || {
182 let full = fs::resolve_path(&server_root, &root_rel, &req_rel)?;
183 fs::list_dir(&full)
184 })
185 .await?;
186
187 Ok(Json(FilesResp { entries, truncated }))
188}
189
190// ---------------------------------------------------------------------------
191// download / preview / content (milestone 4)
192// ---------------------------------------------------------------------------
193
194/// `Content-Disposition` parameters for `name`: an ASCII `filename=` fallback
195/// (non-ASCII and control bytes become `_`) plus the RFC 8187 `filename*=`
196/// that every current browser reads. Never fails header validation.
197fn disposition(kind: &str, name: &str) -> String {
198 let ascii: String = name
199 .chars()
200 .map(|c| match c {
201 '"' | '\\' => '_',
202 c if c.is_ascii_graphic() || c == ' ' => c,
203 _ => '_',
204 })
205 .collect();
206 let mut enc = String::with_capacity(name.len() * 3);
207 for b in name.bytes() {
208 // attr-char per RFC 8187.
209 if b.is_ascii_alphanumeric() || b"!#$&+-.^_`|~".contains(&b) {
210 enc.push(b as char);
211 } else {
212 use std::fmt::Write as _;
213 let _ = write!(enc, "%{b:02X}");
214 }
215 }
216 format!("{kind}; filename=\"{ascii}\"; filename*=UTF-8''{enc}")
217}
218
219/// Resolve the requested item to an absolute path + metadata (blocking).
220async fn resolve_item(
221 state: &AppState,
222 root: &RootRow,
223 req_rel: &str,
224 share: Option<&ShareRow>,
225) -> Result<(std::path::PathBuf, String, bool, u64, i64), ApiError> {
226 let (server_root, root_rel, rel) = (state.root.clone(), root.path.clone(), req_rel.to_string());
227 // A file share's synthetic root *is* the file, so resolve it directly.
228 let share_target = share.filter(|s| s.is_file).map(|s| s.target.clone());
229 blocking(move || {
230 let full = match share_target {
231 Some(target) => fs::resolve_file(&server_root, &target)?,
232 None => fs::resolve_path(&server_root, &root_rel, &rel)?,
233 };
234 let name = full
235 .file_name()
236 .map(|n| n.to_string_lossy().into_owned())
237 .ok_or_else(|| FsError::Invalid("invalid path".to_string()))?;
238 let meta = std::fs::metadata(&full).map_err(|_| FsError::NotFound)?;
239 let mtime = fs::mtime_secs(&meta).unwrap_or(0);
240 Ok::<_, FsError>((full, name, meta.is_dir(), meta.len(), mtime))
241 })
242 .await
243}
244
245/// `GET ...?action=download` — a single file as-is, a folder as an archive
246/// (format chosen by the client).
247async fn download(
248 state: Arc<AppState>,
249 auth: AuthUser,
250 root_id: i64,
251 req_rel: String,
252 format: Option<&str>,
253 range: Option<String>,
254 ims: Option<String>,
255) -> Result<Response, ApiError> {
256 let root = find_root(&auth.roots, root_id)?;
257 let (full, name, is_dir, size, mtime) =
258 resolve_item(&state, root, &req_rel, auth.share.as_ref()).await?;
259
260 if !is_dir {
261 return file_response(
262 &full,
263 &name,
264 size,
265 false,
266 range.as_deref(),
267 mtime,
268 ims.as_deref(),
269 )
270 .await;
271 }
272
273 let fmt = format.and_then(ArchiveFormat::parse).ok_or_else(|| {
274 ApiError::localized(
275 StatusCode::BAD_REQUEST,
276 "format must be one of: zip, tar, tar.gz, tar.zst",
277 "err_bad_format",
278 )
279 })?;
280 let disp = disposition("attachment", &format!("{name}.{}", fmt.extension()));
281 let body = stream_archive(fmt, full, name);
282 Response::builder()
283 .status(StatusCode::OK)
284 .header(header::CONTENT_TYPE, fmt.mime())
285 .header(header::CONTENT_DISPOSITION, disp)
286 .header(header::CACHE_CONTROL, FILE_CACHE)
287 .body(body)
288 .map_err(|e| {
289 ApiError::new(
290 StatusCode::INTERNAL_SERVER_ERROR,
291 format!("bad response: {e}"),
292 )
293 })
294}
295
296/// `GET ...?action=preview` — a single file, inline (for native media).
297async fn preview(
298 state: Arc<AppState>,
299 auth: AuthUser,
300 root_id: i64,
301 req_rel: String,
302 range: Option<String>,
303 ims: Option<String>,
304) -> Result<Response, ApiError> {
305 let root = find_root(&auth.roots, root_id)?;
306 let (full, name, is_dir, size, mtime) =
307 resolve_item(&state, root, &req_rel, auth.share.as_ref()).await?;
308 if is_dir {
309 return Err(ApiError::localized(
310 StatusCode::BAD_REQUEST,
311 "not a file",
312 "err_not_a_file",
313 ));
314 }
315 file_response(
316 &full,
317 &name,
318 size,
319 true,
320 range.as_deref(),
321 mtime,
322 ims.as_deref(),
323 )
324 .await
325}
326
327/// `GET ...?action=content` — raw file bytes for the text preview/editor.
328/// Capped at `MAX_TEXT_BYTES`.
329async fn content(
330 state: Arc<AppState>,
331 auth: AuthUser,
332 root_id: i64,
333 req_rel: String,
334) -> Result<Response, ApiError> {
335 let root = find_root(&auth.roots, root_id)?;
336 let (full, _name, is_dir, size, _mtime) =
337 resolve_item(&state, root, &req_rel, auth.share.as_ref()).await?;
338 if is_dir {
339 return Err(ApiError::localized(
340 StatusCode::BAD_REQUEST,
341 "not a file",
342 "err_not_a_file",
343 ));
344 }
345 if size > MAX_TEXT_BYTES {
346 return Err(ApiError::localized(
347 StatusCode::PAYLOAD_TOO_LARGE,
348 "file too large to preview",
349 "err_too_large_preview",
350 ));
351 }
352 // mtime and bytes from one handle, mtime first: a write in between would
353 // otherwise hand the editor a stale conflict anchor for fresh content.
354 let (mtime, bytes) = blocking(move || -> io::Result<(i64, Vec<u8>)> {
355 let mut f = std::fs::File::open(&full)?;
356 let mtime = fs::mtime_secs(&f.metadata()?).unwrap_or(0);
357 let mut bytes = Vec::with_capacity(size as usize);
358 f.read_to_end(&mut bytes)?;
359 Ok((mtime, bytes))
360 })
361 .await?;
362 Ok((
363 [
364 (
365 header::CONTENT_TYPE,
366 "text/plain; charset=utf-8".to_string(),
367 ),
368 (header::CACHE_CONTROL, FILE_CACHE.to_string()),
369 (
370 axum::http::HeaderName::from_static("x-file-mtime"),
371 mtime.to_string(),
372 ),
373 ],
374 bytes,
375 )
376 .into_response())
377}
378
379/// `GET ...?action=thumb` — a small WebP preview of an image or video.
380///
381/// `404` covers every "no thumbnail here" case: thumbnails switched off, a
382/// folder, a kind we do not render, an undecodable file, a video without
383/// ffmpeg. The grid falls back to its icon for all of them alike.
384async fn thumb(
385 state: Arc<AppState>,
386 auth: AuthUser,
387 root_id: i64,
388 req_rel: String,
389) -> Result<Response, ApiError> {
390 let Some(thumbs) = state.thumbs.as_ref() else {
391 return Err(no_thumb());
392 };
393 let root = find_root(&auth.roots, root_id)?;
394 let (full, _name, is_dir, size, mtime) =
395 resolve_item(&state, root, &req_rel, auth.share.as_ref()).await?;
396 if is_dir {
397 return Err(no_thumb());
398 }
399 // The kind is sniffed from the bytes, not the name, so an mp4 called .txt
400 // still gets a thumbnail and a .jpg full of text does not.
401 let kind = {
402 let full = full.clone();
403 blocking(move || Ok::<_, FsError>(fs::detect_kind(&full, false))).await?
404 };
405 let Some(bytes) = thumbs.get(&full, kind, size, mtime).await else {
406 return Err(no_thumb());
407 };
408 Ok((
409 [
410 (header::CONTENT_TYPE, "image/webp"),
411 // Safe despite the stable path: the client varies the query on
412 // mtime, so new content always means a new URL.
413 (
414 header::CACHE_CONTROL,
415 "private, max-age=31536000, immutable",
416 ),
417 ],
418 bytes,
419 )
420 .into_response())
421}
422
423/// The message never reaches a user: an `<img>` 404 just leaves the tile's
424/// icon showing. That is why it carries no localized code.
425fn no_thumb() -> ApiError {
426 ApiError::new(StatusCode::NOT_FOUND, "no thumbnail".to_string())
427}
428
429/// `PUT ...?action=content` — save a file's text contents (the editor).
430///
431/// Requires a read-write root. The body is the new contents. If the
432/// `X-Expected-Mtime` header is present, the file's current mtime must match
433/// it, otherwise `409 Conflict` (the file changed on disk since it was read).
434/// Returns the file's new mtime so the client can anchor the next check.
435pub async fn file_put(
436 State(state): State<Arc<AppState>>,
437 auth: AuthUser,
438 path: AxumPath<(i64, String)>,
439 query: AxumQuery<FileQuery>,
440 headers: axum::http::HeaderMap,
441 body: axum::body::Bytes,
442) -> Result<Json<SaveResp>, ApiError> {
443 let (root_id, req_rel) = path.0;
444 put_inner(state, auth, root_id, req_rel, query, headers, body).await
445}
446
447/// `PUT .../{root_id}?action=content` — the root item itself. Only reachable
448/// for a *file* share (its root is the file); for folders it resolves to a
449/// directory and is rejected below.
450pub async fn file_put_root(
451 State(state): State<Arc<AppState>>,
452 auth: AuthUser,
453 path: AxumPath<i64>,
454 query: AxumQuery<FileQuery>,
455 headers: axum::http::HeaderMap,
456 body: axum::body::Bytes,
457) -> Result<Json<SaveResp>, ApiError> {
458 put_inner(state, auth, path.0, String::new(), query, headers, body).await
459}
460
461async fn put_inner(
462 state: Arc<AppState>,
463 auth: AuthUser,
464 root_id: i64,
465 req_rel: String,
466 query: AxumQuery<FileQuery>,
467 headers: axum::http::HeaderMap,
468 body: axum::body::Bytes,
469) -> Result<Json<SaveResp>, ApiError> {
470 if query.action.as_deref() != Some(api_types::ACTION_CONTENT) {
471 return Err(ApiError::localized(
472 StatusCode::BAD_REQUEST,
473 "expected action=content",
474 "err_bad_action",
475 ));
476 }
477 let root = require_rw_root(&auth.roots, root_id)?;
478 if body.len() as u64 > MAX_TEXT_BYTES {
479 return Err(ApiError::localized(
480 StatusCode::PAYLOAD_TOO_LARGE,
481 "file too large to save",
482 "err_too_large_save",
483 ));
484 }
485 let expected: Option<i64> = headers
486 .get("x-expected-mtime")
487 .and_then(|v| v.to_str().ok())
488 .and_then(|s| s.parse().ok());
489 // A file share's synthetic root *is* the file, so resolve it directly.
490 let share_target = auth
491 .share
492 .as_ref()
493 .filter(|s| s.is_file)
494 .map(|s| s.target.clone());
495 let (server_root, root_rel, rel) = (state.root.clone(), root.path.clone(), req_rel);
496 let content = body.to_vec();
497 let mtime = blocking(move || match share_target {
498 Some(target) => fs::save_file_at(&server_root, &target, &content, expected),
499 None => fs::save_file(&server_root, &root_rel, &rel, &content, expected),
500 })
501 .await?;
502 Ok(Json(SaveResp { mtime }))
503}
504
505/// Stream a single file to the client with the right disposition.
506async fn file_response(
507 full: &std::path::Path,
508 name: &str,
509 size: u64,
510 inline: bool,
511 range: Option<&str>,
512 mtime: i64,
513 ims: Option<&str>,
514) -> Result<Response, ApiError> {
515 let last_modified = http_date(mtime);
516 // A revalidating client gets the empty 304 instead of the whole file. The
517 // policy is repeated on it, so the stored copy does not lose the directive.
518 if last_modified.is_some() && unmodified_since(ims, mtime) {
519 return Ok((
520 StatusCode::NOT_MODIFIED,
521 [(header::CACHE_CONTROL, FILE_CACHE)],
522 )
523 .into_response());
524 }
525 let mime = mime_guess::from_path(full)
526 .first_or_octet_stream()
527 .to_string();
528 let disp = disposition(if inline { "inline" } else { "attachment" }, name);
529 // A single `bytes=a-b` range (media seeking). Anything else is served whole.
530 let (start, end) = match parse_range(range, size) {
531 Some(Some(r)) => r,
532 Some(None) => {
533 return Response::builder()
534 .status(StatusCode::RANGE_NOT_SATISFIABLE)
535 .header(header::CONTENT_RANGE, format!("bytes */{size}"))
536 .body(axum::body::Body::empty())
537 .map_err(|e| {
538 ApiError::new(
539 StatusCode::INTERNAL_SERVER_ERROR,
540 format!("bad response: {e}"),
541 )
542 });
543 }
544 None => (0, size),
545 };
546 let partial = (start, end) != (0, size);
547 let body = stream_file(full.to_path_buf(), start, end);
548 let mut res = Response::builder()
549 .status(if partial {
550 StatusCode::PARTIAL_CONTENT
551 } else {
552 StatusCode::OK
553 })
554 .header(header::CONTENT_DISPOSITION, disp)
555 .header(header::ACCEPT_RANGES, "bytes")
556 .header(header::CONTENT_LENGTH, end - start)
557 // Always sent, validator or not: with no directive a cache may apply
558 // heuristic freshness to a response that depends on who asked.
559 .header(header::CACHE_CONTROL, FILE_CACHE);
560 if let Some(lm) = last_modified {
561 // The validator the `no-cache` above revalidates against.
562 res = res.header(header::LAST_MODIFIED, lm);
563 }
564 if partial {
565 res = res.header(
566 header::CONTENT_RANGE,
567 format!("bytes {start}-{}/{size}", end - 1),
568 );
569 }
570 // A file the browser would parse as a document (HTML/SVG/XML) is served
571 // under the sandboxed policy, so it can render as a page without being
572 // able to act as the app. Derived from the same `mime` we declare.
573 // Non-scriptable inline files (PDF, …) are frameable by the app itself,
574 // for the preview modal.
575 if crate::api::is_scriptable_mime(&mime) {
576 res = res.header("content-security-policy", crate::api::FILE_CSP);
577 } else if inline {
578 res = res
579 .header("content-security-policy", crate::api::INLINE_CSP)
580 .header(header::X_FRAME_OPTIONS, "SAMEORIGIN");
581 }
582 res.header(header::CONTENT_TYPE, mime)
583 .body(body)
584 .map_err(|e| {
585 ApiError::new(
586 StatusCode::INTERNAL_SERVER_ERROR,
587 format!("bad response: {e}"),
588 )
589 })
590}
591
592/// Parse a `Range` header against `size`. `None` = serve the whole file,
593/// `Some(None)` = unsatisfiable, `Some(Some((start, end)))` = half-open range.
594fn parse_range(range: Option<&str>, size: u64) -> Option<Option<(u64, u64)>> {
595 let spec = range?.strip_prefix("bytes=")?;
596 // ponytail: one range only; multipart/byteranges is not worth it here.
597 let (a, b) = spec.split_once('-')?;
598 let (start, end) = match (a.trim().parse::<u64>().ok(), b.trim().parse::<u64>().ok()) {
599 (Some(s), Some(e)) => (s, e.saturating_add(1).min(size)),
600 (Some(s), None) if b.trim().is_empty() => (s, size),
601 // Suffix form: the last N bytes.
602 (None, Some(n)) if a.trim().is_empty() => (size.saturating_sub(n), size),
603 _ => return None,
604 };
605 if start >= size || start >= end {
606 return Some(None);
607 }
608 Some(Some((start, end)))
609}
610
611/// Stream `path[start..end)` to the client in chunks (blocking reader → channel).
612fn stream_file(path: std::path::PathBuf, start: u64, end: u64) -> axum::body::Body {
613 use std::io::Seek;
614 let (tx, rx) = mpsc::channel::<Vec<u8>>(16);
615 tokio::task::spawn_blocking(move || {
616 let mut f = match std::fs::File::open(&path) {
617 Ok(f) => f,
618 Err(e) => {
619 tracing::warn!(error = %e, path = %path.display(), "download open failed");
620 return;
621 }
622 };
623 if start > 0 && f.seek(io::SeekFrom::Start(start)).is_err() {
624 return;
625 }
626 let mut left = end - start;
627 let mut buf = vec![0u8; 256 * 1024];
628 while left > 0 {
629 let want = buf.len().min(left as usize);
630 match f.read(&mut buf[..want]) {
631 Ok(0) => break,
632 Ok(n) => {
633 left -= n as u64;
634 // Client gone → stop producing.
635 if tx.blocking_send(buf[..n].to_vec()).is_err() {
636 break;
637 }
638 }
639 Err(e) => {
640 tracing::warn!(error = %e, path = %path.display(), "download read failed");
641 break;
642 }
643 }
644 }
645 });
646 let stream = ReceiverStream::new(rx).map(Ok::<_, io::Error>);
647 axum::body::Body::from_stream(stream)
648}
649
650/// Stream an archive of `dir` (top-level entry `top`) to the client.
651fn stream_archive(fmt: ArchiveFormat, dir: std::path::PathBuf, top: String) -> axum::body::Body {
652 let (tx, rx) = mpsc::channel::<Vec<u8>>(16);
653 tokio::task::spawn_blocking(move || {
654 let mut sink = ChanWriter::new(tx);
655 if let Err(e) = archive::build(fmt, &dir, &top, &mut sink) {
656 tracing::warn!(error = %e, dir = %dir.display(), "archive build failed");
657 }
658 // Dropping the sink flushes its buffer and closes the channel.
659 });
660 let stream = ReceiverStream::new(rx).map(Ok::<_, io::Error>);
661 axum::body::Body::from_stream(stream)
662}
663
664/// A `Write` that buffers chunks and forwards them over an mpsc channel — the
665/// bridge between the blocking archive builder and the async response body.
666struct ChanWriter {
667 tx: mpsc::Sender<Vec<u8>>,
668 buf: Vec<u8>,
669}
670
671impl ChanWriter {
672 fn new(tx: mpsc::Sender<Vec<u8>>) -> Self {
673 Self {
674 tx,
675 buf: Vec::with_capacity(64 * 1024),
676 }
677 }
678}
679
680impl io::Write for ChanWriter {
681 fn write(&mut self, b: &[u8]) -> io::Result<usize> {
682 self.buf.extend_from_slice(b);
683 if self.buf.len() >= 64 * 1024 {
684 io::Write::flush(self)?;
685 }
686 Ok(b.len())
687 }
688 fn flush(&mut self) -> io::Result<()> {
689 if !self.buf.is_empty() {
690 let chunk = std::mem::take(&mut self.buf);
691 self.tx
692 .blocking_send(chunk)
693 .map_err(|_| io::Error::new(io::ErrorKind::BrokenPipe, "client disconnected"))?;
694 }
695 Ok(())
696 }
697}
698
699impl Drop for ChanWriter {
700 fn drop(&mut self) {
701 let _ = io::Write::flush(self);
702 }
703}
704
705// ---------------------------------------------------------------------------
706// POST dispatch: mkdir | rename/move/copy | upload
707// ---------------------------------------------------------------------------
708
709/// POST /api/files/{root_id} — top-level operations (upload into the root
710/// directory). The bare root never targets a specific item, so JSON ops with
711/// a missing folder name are rejected by the individual handlers.
712pub async fn dispatch_root(
713 State(state): State<Arc<AppState>>,
714 auth: AuthUser,
715 path: AxumPath<i64>,
716 headers: axum::http::HeaderMap,
717 req: axum::http::Request<axum::body::Body>,
718) -> Result<Response, ApiError> {
719 dispatch_inner(state, auth, path.0, String::new(), headers, req).await
720}
721
722/// POST /api/files/{root_id}/{*path}
723pub async fn dispatch(
724 State(state): State<Arc<AppState>>,
725 auth: AuthUser,
726 path: AxumPath<(i64, String)>,
727 headers: axum::http::HeaderMap,
728 req: axum::http::Request<axum::body::Body>,
729) -> Result<Response, ApiError> {
730 let (root_id, req_rel) = path.0;
731 dispatch_inner(state, auth, root_id, req_rel, headers, req).await
732}
733
734/// Route one POST to upload, mutation or mkdir.
735///
736/// Upload and mutation are recognized by their content type. mkdir carries no
737/// body, so it names itself with `?action=mkdir`. Anything else is rejected:
738/// an unrecognized content type used to fall through to mkdir, which turned a
739/// typo in a header into a silently created folder.
740async fn dispatch_inner(
741 state: Arc<AppState>,
742 auth: AuthUser,
743 root_id: i64,
744 req_rel: String,
745 headers: axum::http::HeaderMap,
746 req: axum::http::Request<axum::body::Body>,
747) -> Result<Response, ApiError> {
748 let ct = headers
749 .get(header::CONTENT_TYPE)
750 .and_then(|v| v.to_str().ok())
751 .unwrap_or("");
752
753 if ct.starts_with("multipart/form-data") {
754 return upload(state, auth, root_id, req_rel, req).await;
755 }
756 if ct.starts_with("application/json") {
757 let bytes = axum::body::to_bytes(req.into_body(), 1_000_000)
758 .await
759 .map_err(|_| {
760 ApiError::localized(
761 StatusCode::BAD_REQUEST,
762 "invalid request body",
763 "err_bad_body",
764 )
765 })?;
766 let body: Mutation = axum::Json::from_bytes(&bytes)
767 .map_err(|_| {
768 ApiError::localized(
769 StatusCode::BAD_REQUEST,
770 "invalid request body",
771 "err_bad_body",
772 )
773 })?
774 .0;
775 return Ok(mutation(state, auth, root_id, req_rel, body)
776 .await?
777 .into_response());
778 }
779 match action_param(req.uri()).as_deref() {
780 Some(api_types::ACTION_MKDIR) => {
781 return Ok(mkdir(state, auth, root_id, req_rel).await?.into_response());
782 }
783 Some(api_types::ACTION_CREATE_FILE) => {
784 return Ok(create_file(state, auth, root_id, req_rel)
785 .await?
786 .into_response());
787 }
788 _ => {}
789 }
790 Err(ApiError::localized(
791 StatusCode::UNSUPPORTED_MEDIA_TYPE,
792 "POST expects a multipart upload, a JSON mutation, or ?action=mkdir",
793 "err_bad_post",
794 ))
795}
796
797// ---------------------------------------------------------------------------
798// create file
799// ---------------------------------------------------------------------------
800
801/// Create an empty file in a writable root.
802async fn create_file(
803 state: Arc<AppState>,
804 auth: AuthUser,
805 root_id: i64,
806 req_rel: String,
807) -> Result<Json<OkResp>, ApiError> {
808 let root = require_rw_root(&auth.roots, root_id)?;
809 if req_rel.trim().is_empty() {
810 return Err(ApiError::localized(
811 StatusCode::BAD_REQUEST,
812 "a file name is required",
813 "err_file_name_required",
814 ));
815 }
816 let (server_root, root_rel, rel) = (state.root.clone(), root.path.clone(), req_rel);
817 blocking(move || fs::create_file(&server_root, &root_rel, &rel)).await?;
818 Ok(Json(OkResp {}))
819}
820
821// ---------------------------------------------------------------------------
822// mkdir
823// ---------------------------------------------------------------------------
824
825async fn mkdir(
826 state: Arc<AppState>,
827 auth: AuthUser,
828 root_id: i64,
829 req_rel: String,
830) -> Result<Json<OkResp>, ApiError> {
831 let root = require_rw_root(&auth.roots, root_id)?;
832 if req_rel.trim().is_empty() {
833 return Err(ApiError::localized(
834 StatusCode::BAD_REQUEST,
835 "a folder name is required",
836 "err_folder_name_required",
837 ));
838 }
839 let (server_root, root_rel, rel) = (state.root.clone(), root.path.clone(), req_rel);
840 blocking(move || fs::mkdir(&server_root, &root_rel, &rel)).await?;
841 Ok(Json(OkResp {}))
842}
843
844// ---------------------------------------------------------------------------
845// rename / move / copy
846// ---------------------------------------------------------------------------
847
848async fn mutation(
849 state: Arc<AppState>,
850 auth: AuthUser,
851 root_id: i64,
852 req_rel: String,
853 body: Mutation,
854) -> Result<Json<OkResp>, ApiError> {
855 match body.op {
856 Op::Rename => {
857 let new_name = body
858 .new_name
859 .as_deref()
860 .ok_or_else(|| {
861 ApiError::localized(
862 StatusCode::BAD_REQUEST,
863 "new_name is required",
864 "err_new_name_required",
865 )
866 })?
867 .to_string();
868 let root = require_rw_root(&auth.roots, root_id)?;
869 let (server_root, root_rel, rel, overwrite) = (
870 state.root.clone(),
871 root.path.clone(),
872 req_rel,
873 body.overwrite,
874 );
875 let vacated = blocking(move || {
876 fs::rename_item(&server_root, &root_rel, &rel, &new_name, overwrite)
877 })
878 .await?;
879 revoke_shares_at(&state, &vacated).await;
880 Ok(Json(OkResp {}))
881 }
882 Op::Move | Op::Copy => {
883 let dst_root_id = body.dst_root_id.ok_or_else(|| {
884 ApiError::localized(
885 StatusCode::BAD_REQUEST,
886 "dst_root_id is required",
887 "err_dst_required",
888 )
889 })?;
890 let dst = body.dst.clone().unwrap_or_default();
891 // Moving or copying out of a folder requires rw there; copying
892 // *from* a read-only root is fine.
893 let op_is_move = body.op == Op::Move;
894 let src_root = if op_is_move {
895 require_rw_root(&auth.roots, root_id)?
896 } else {
897 find_root(&auth.roots, root_id)?
898 };
899 let dst_root = require_rw_root(&auth.roots, dst_root_id)?;
900 let (server_root, src_rel, dst_rel, dst_path, rel, overwrite) = (
901 state.root.clone(),
902 src_root.path.clone(),
903 dst_root.path.clone(),
904 dst,
905 req_rel,
906 body.overwrite,
907 );
908 let vacated = blocking(move || {
909 if op_is_move {
910 fs::move_item(&server_root, &src_rel, &rel, &dst_rel, &dst_path, overwrite)
911 .map(Some)
912 } else {
913 // A copy frees no path, so it revokes nothing.
914 fs::copy_item(&server_root, &src_rel, &rel, &dst_rel, &dst_path, overwrite)
915 .map(|()| None)
916 }
917 })
918 .await?;
919 if let Some(vacated) = vacated {
920 revoke_shares_at(&state, &vacated).await;
921 }
922 Ok(Json(OkResp {}))
923 }
924 }
925}
926
927// ---------------------------------------------------------------------------
928// DELETE
929// ---------------------------------------------------------------------------
930
931pub async fn delete(
932 State(state): State<Arc<AppState>>,
933 auth: AuthUser,
934 path: AxumPath<(i64, String)>,
935) -> Result<Json<DeleteResp>, ApiError> {
936 let (root_id, req_rel) = path.0;
937 let root = require_rw_root(&auth.roots, root_id)?;
938 if req_rel.trim().is_empty() {
939 return Err(ApiError::localized(
940 StatusCode::BAD_REQUEST,
941 "a path inside the folder is required",
942 "err_path_required",
943 ));
944 }
945 let (server_root, root_rel, rel) = (state.root.clone(), root.path.clone(), req_rel);
946 let (is_dir, gone) = blocking(move || fs::remove_item(&server_root, &root_rel, &rel)).await?;
947 revoke_shares_at(&state, &gone).await;
948 Ok(Json(DeleteResp { is_dir }))
949}
950
951// ---------------------------------------------------------------------------
952// Upload (multipart)
953// ---------------------------------------------------------------------------
954
955async fn upload(
956 state: Arc<AppState>,
957 auth: AuthUser,
958 root_id: i64,
959 req_rel: String,
960 req: axum::http::Request<axum::body::Body>,
961) -> Result<Response, ApiError> {
962 let root = require_rw_root(&auth.roots, root_id)?;
963 // `base` is the upload directory; `root_abs` the user's root, which is the
964 // containment boundary (a symlink may legitimately point elsewhere inside it).
965 let (root_abs, base) = {
966 let (server_root, root_rel, rel) = (state.root.clone(), root.path.clone(), req_rel);
967 blocking(move || {
968 Ok::<_, FsError>((
969 fs::resolve_root(&server_root, &root_rel)?,
970 fs::resolve_dir(&server_root, &root_rel, &rel)?,
971 ))
972 })
973 .await?
974 };
975 let boundary = req
976 .headers()
977 .get(header::CONTENT_TYPE)
978 .and_then(|v| v.to_str().ok())
979 .and_then(parse_boundary)
980 .ok_or_else(|| {
981 ApiError::localized(
982 StatusCode::BAD_REQUEST,
983 "expected multipart/form-data with a boundary",
984 "err_bad_multipart",
985 )
986 })?;
987 let overwrite = parse_overwrite(req.uri());
988
989 let stream = req.into_body().into_data_stream();
990 let mut multipart = Multipart::new(stream, boundary);
991 let mut uploaded: usize = 0;
992 let mut skipped: Vec<String> = Vec::new();
993
994 while let Some(mut field) = multipart.next_field().await.map_err(|_| {
995 ApiError::localized(
996 StatusCode::BAD_REQUEST,
997 "invalid upload data",
998 "err_bad_upload",
999 )
1000 })? {
1001 let part_name = field
1002 .name()
1003 .filter(|n| !n.is_empty())
1004 .or_else(|| field.file_name())
1005 .map(str::to_string)
1006 .ok_or_else(|| {
1007 ApiError::localized(
1008 StatusCode::BAD_REQUEST,
1009 "part without a name",
1010 "err_part_no_name",
1011 )
1012 })?;
1013
1014 validate_rel_path(&part_name)?;
1015
1016 let target = base.join(&part_name);
1017 let parent = target
1018 .parent()
1019 .filter(|p| !p.as_os_str().is_empty())
1020 .ok_or_else(|| {
1021 ApiError::localized(
1022 StatusCode::BAD_REQUEST,
1023 "invalid part name",
1024 "err_bad_part_name",
1025 )
1026 })?;
1027 // Create the parent, then canonicalize it and require it to still be
1028 // inside the upload directory. `validate_rel_path` blocks `..`, but a
1029 // symlinked directory on disk would otherwise carry the write outside.
1030 let (p, b) = (parent.to_path_buf(), root_abs.clone());
1031 let file_name = target.file_name().map(|n| n.to_owned());
1032 let (parent, target_state) = tokio::task::spawn_blocking(move || {
1033 let escape = || io::Error::other("upload parent escapes the root");
1034 // Check the nearest existing ancestor *before* creating anything,
1035 // so no directory is ever created outside the root either.
1036 let mut existing = p.as_path();
1037 while !existing.exists() {
1038 existing = existing.parent().ok_or_else(escape)?;
1039 }
1040 if !fs::is_within_or_eq(&b, &existing.canonicalize()?) {
1041 return Err(escape());
1042 }
1043 if !p.is_dir() {
1044 std::fs::create_dir_all(&p)?;
1045 }
1046 let canon = p.canonicalize()?;
1047 if !fs::is_within_or_eq(&b, &canon) {
1048 return Err(escape());
1049 }
1050 // Whether the target already exists, and as what: two stats that
1051 // belong on this thread, not on an async worker.
1052 let state = file_name
1053 .map(|n| canon.join(n))
1054 .map(|t| (t.exists(), t.is_dir()));
1055 Ok::<_, io::Error>((canon, state))
1056 })
1057 .await
1058 .map_err(|_| io::Error::other("join"))?
1059 .map_err(|e| {
1060 tracing::warn!(error = %e, "upload parent rejected");
1061 ApiError::localized(
1062 StatusCode::FORBIDDEN,
1063 "invalid file path in upload",
1064 "err_bad_upload_path",
1065 )
1066 })?;
1067 let (Some(file_name), Some((exists, is_dir))) = (target.file_name(), target_state) else {
1068 return Err(ApiError::localized(
1069 StatusCode::BAD_REQUEST,
1070 "invalid part name",
1071 "err_bad_part_name",
1072 ));
1073 };
1074 let target = parent.join(file_name);
1075
1076 if exists && !overwrite {
1077 // Drain this part and report it as a conflict at the end.
1078 while let Some(_chunk) = field.chunk().await.map_err(|_| {
1079 ApiError::localized(
1080 StatusCode::BAD_REQUEST,
1081 "invalid upload data",
1082 "err_bad_upload",
1083 )
1084 })? {}
1085 skipped.push(part_name);
1086 continue;
1087 }
1088 if exists {
1089 if is_dir {
1090 return Err(ApiError::localized(
1091 StatusCode::CONFLICT,
1092 "a folder with this name already exists",
1093 "err_folder_exists",
1094 ));
1095 }
1096 tokio::fs::remove_file(&target).await.map_err(|_| {
1097 ApiError::localized(
1098 StatusCode::INTERNAL_SERVER_ERROR,
1099 "internal error",
1100 "err_internal",
1101 )
1102 })?;
1103 }
1104
1105 // Stream to a temp file in the same directory, then rename into place.
1106 let suffix = crate::auth::random_token();
1107 let tmp = Scratch(parent.join(format!(".upload-{suffix}")));
1108 let tmp_file = tokio::fs::File::create(&tmp.0).await.map_err(|_| {
1109 ApiError::localized(
1110 StatusCode::INTERNAL_SERVER_ERROR,
1111 "internal error",
1112 "err_internal",
1113 )
1114 })?;
1115 // Buffered: a multipart chunk is often a few kilobytes, and each
1116 // unbuffered write would be its own syscall.
1117 let mut tmp_file = tokio::io::BufWriter::with_capacity(1 << 20, tmp_file);
1118 let write_failed = loop {
1119 match field.chunk().await.map_err(|_| {
1120 ApiError::localized(
1121 StatusCode::BAD_REQUEST,
1122 "invalid upload data",
1123 "err_bad_upload",
1124 )
1125 }) {
1126 Ok(Some(chunk)) => {
1127 if let Err(e) = tmp_file.write_all(&chunk).await {
1128 tracing::warn!(error = %e, "write failed during upload");
1129 break true;
1130 }
1131 }
1132 Ok(None) => break tmp_file.flush().await.is_err(),
1133 Err(e) => return Err(e),
1134 }
1135 };
1136 if write_failed {
1137 return Err(ApiError::localized(
1138 StatusCode::INTERNAL_SERVER_ERROR,
1139 "could not save the file",
1140 "err_save_failed",
1141 ));
1142 }
1143 let tmp2 = tmp.0.clone();
1144 let target2 = target.clone();
1145 // Flatten both errors: the outer `Err` is a panicking or shut-down task,
1146 // the inner one is `rename` refusing. Either way nothing was published.
1147 let renamed = tokio::task::spawn_blocking(move || std::fs::rename(&tmp2, &target2))
1148 .await
1149 .map_err(io::Error::other)
1150 .and_then(|r| r);
1151 if let Err(e) = renamed {
1152 tracing::warn!(error = %e, path = %target.display(), "rename failed during upload");
1153 return Err(ApiError::localized(
1154 StatusCode::INTERNAL_SERVER_ERROR,
1155 "internal error",
1156 "err_internal",
1157 ));
1158 }
1159 tmp.disarm();
1160 uploaded += 1;
1161 }
1162
1163 if uploaded == 0 && skipped.is_empty() {
1164 return Err(ApiError::localized(
1165 StatusCode::BAD_REQUEST,
1166 "no files were uploaded",
1167 "err_no_files_uploaded",
1168 ));
1169 }
1170 if !skipped.is_empty() {
1171 return Err(ApiError::localized(
1172 StatusCode::CONFLICT,
1173 "some files already exist",
1174 "err_files_exist",
1175 )
1176 .with_extra(serde_json::json!({ "skipped": skipped, "uploaded": uploaded })));
1177 }
1178 Ok(Json(UploadResp { uploaded }).into_response())
1179}
1180
1181/// The `.upload-<token>` scratch file of one in-flight upload part. Dropping it
1182/// removes the file, which covers the paths no `return` can see, above all the
1183/// request future being dropped when the client closes the connection. A leaked
1184/// scratch file is never named again and shows up in listings, which include
1185/// hidden entries on purpose. [`Scratch::disarm`] after `rename` keeps the file.
1186struct Scratch(PathBuf);
1187
1188impl Scratch {
1189 /// The scratch file is now the uploaded file: leave it alone.
1190 fn disarm(self) {
1191 std::mem::forget(self);
1192 }
1193}
1194
1195impl Drop for Scratch {
1196 fn drop(&mut self) {
1197 // Plain blocking unlink: `Drop` can run during runtime shutdown, where
1198 // `tokio::spawn` panics. One unlink cannot block meaningfully.
1199 if let Err(e) = std::fs::remove_file(&self.0)
1200 && e.kind() != io::ErrorKind::NotFound
1201 {
1202 tracing::warn!(error = %e, path = %self.0.display(), "could not remove upload scratch file");
1203 }
1204 }
1205}
1206
1207fn parse_boundary(content_type: &str) -> Option<String> {
1208 content_type
1209 .split(';')
1210 .map(|s| s.trim())
1211 .find_map(|s| s.strip_prefix("boundary="))
1212 .map(|b| b.trim_matches('"').to_string())
1213 .filter(|b| !b.is_empty())
1214}
1215
1216/// Read one query parameter from the request URI.
1217fn query_param(uri: &axum::http::Uri, key: &str) -> Option<String> {
1218 uri.query()?.split('&').find_map(|kv| {
1219 let (k, v) = kv.split_once('=')?;
1220 (k == key).then(|| v.to_string())
1221 })
1222}
1223
1224fn action_param(uri: &axum::http::Uri) -> Option<String> {
1225 query_param(uri, P_ACTION)
1226}
1227
1228fn parse_overwrite(uri: &axum::http::Uri) -> bool {
1229 matches!(query_param(uri, P_OVERWRITE).as_deref(), Some("true" | "1"))
1230}
1231
1232fn validate_rel_path(name: &str) -> Result<(), ApiError> {
1233 for c in std::path::Path::new(name).components() {
1234 match c {
1235 Component::Normal(_) => {}
1236 _ => {
1237 return Err(ApiError::localized(
1238 StatusCode::BAD_REQUEST,
1239 "invalid file path in upload",
1240 "err_bad_upload_path",
1241 ));
1242 }
1243 }
1244 }
1245 Ok(())
1246}
1247
1248// ---------------------------------------------------------------------------
1249// Helpers
1250// ---------------------------------------------------------------------------
1251
1252/// Drop every share that named `abs` or anything under it.
1253///
1254/// Called after a delete, a rename, or a move: each one frees a path, and a
1255/// share stores a path, not a file identity. Without this, a *new* item that
1256/// later lands on the freed path would inherit the old link's audience.
1257///
1258/// Only covers changes made through this API. A file moved out from under the
1259/// server (over SSH, say) leaves its shares in place, still pointing at a
1260/// path. Closing that needs inode pinning, which breaks across a restore from
1261/// backup, so it is deliberately not done.
1262///
1263/// Best-effort: the file operation has already succeeded by the time this
1264/// runs, so a database error must not turn it into a 500. The client would
1265/// read that as "the delete failed" and retry, and the retry would 404. The
1266/// failure is logged at `error` instead, and leaves a share pointing at a
1267/// path that no longer holds what it did.
1268pub(crate) async fn revoke_shares_at(state: &AppState, abs: &std::path::Path) {
1269 let target = target_rel(state, abs);
1270 match state.db.revoke_shares_at(&target).await {
1271 Ok(0) => {}
1272 Ok(n) => tracing::info!(target = %target, revoked = n, "shares revoked: path is gone"),
1273 Err(e) => {
1274 tracing::error!(error = %e, target = %target, "could not revoke shares on a freed path")
1275 }
1276 }
1277}
1278
1279fn find_root(roots: &[RootRow], root_id: i64) -> Result<&RootRow, ApiError> {
1280 roots.iter().find(|r| r.id == root_id).ok_or_else(|| {
1281 ApiError::localized(
1282 StatusCode::FORBIDDEN,
1283 "no such folder",
1284 "err_no_such_folder",
1285 )
1286 })
1287}
1288
1289fn require_rw_root(roots: &[RootRow], root_id: i64) -> Result<&RootRow, ApiError> {
1290 let root = find_root(roots, root_id)?;
1291 if !root.mode.is_writable() {
1292 return Err(ApiError::localized(
1293 StatusCode::FORBIDDEN,
1294 "read-only folder",
1295 "err_read_only_folder",
1296 ));
1297 }
1298 Ok(root)
1299}
1300
1301#[cfg(test)]
1302mod tests {
1303 use super::*;
1304 use api_types::{ACTION_DOWNLOAD, P_FORMAT};
1305 use axum::http::Uri;
1306
1307 /// A rename of `P_ACTION` or `P_FORMAT` without the matching field
1308 /// rename would silently stop the server from reading the parameter the
1309 /// client sends. This builds the query string from the constants and
1310 /// runs the real extractor over it.
1311 #[test]
1312 fn query_fields_are_the_shared_constants() {
1313 let uri: Uri = format!("/f/1/a.txt?{P_ACTION}={ACTION_DOWNLOAD}&{P_FORMAT}=zip")
1314 .parse()
1315 .unwrap();
1316 let q: FileQuery = AxumQuery::try_from_uri(&uri).unwrap().0;
1317 assert_eq!(q.action.as_deref(), Some(ACTION_DOWNLOAD));
1318 assert_eq!(q.format.as_deref(), Some("zip"));
1319
1320 // The same constants drive the hand-rolled readers on the POST path.
1321 assert_eq!(action_param(&uri).as_deref(), Some(ACTION_DOWNLOAD));
1322 let uri: Uri = format!("/f/1/a.txt?{P_OVERWRITE}=true").parse().unwrap();
1323 assert!(parse_overwrite(&uri));
1324 }
1325}
1326