mirror of
https://github.com/thegeneralist01/archivr
synced 2026-10-09 21:03:17 +02:00
feat: share videos to TV (Chromecast + AirPlay) (#33)
* server: add scoped media-token endpoint for Cast/AirPlay auth bypass
Chromecast and Apple TV fetch media URLs as independent HTTP clients
with no session cookie. The existing serve_artifact handler requires
auth_user.require_auth(), so those devices always received 401.
Changes:
- MediaToken struct stored in AppState (Arc<Mutex<HashMap>>), scoped to
a single (archive_id, entry_uid, artifact_index) tuple with a 2-hour TTL
- POST /api/archives/:id/entries/:uid/artifacts/:idx/media-token
requires an authenticated session, verifies the artifact exists,
prunes expired tokens, mints a 43-char URL-safe token, and returns
{ url, expires_in_secs }
- serve_artifact now accepts an optional ?token= query param; a valid
scoped token bypasses require_auth() while a missing/invalid/expired
token falls through to the normal 401 path
- CSP script-src extended to include https://www.gstatic.com so the
Cast sender SDK script (injected lazily by VideoPreview) is not blocked
- 4 new tests: bare-URL still 401, tokenized fetch succeeds without
session cookie, bogus token 401, wrong-artifact-index 401
* frontend: Cast/AirPlay overlay in VideoPreview
When the user opens a video archive entry, VideoPreview now:
1. Issues a signed media token (POST .../artifacts/:idx/media-token) and
uses the returned signed URL as <video src>. This ensures the video
element's src is one that Cast devices and Apple TV can fetch without a
session cookie.
2. Lazily injects the Google Cast SDK script (cast_sender.js from
gstatic.com, now allowed by the updated CSP). Once the SDK reports
available, a <google-cast-launcher> web component appears as an overlay
button in the top-right corner of the video. Selecting a Cast device
triggers loadMedia() with the signed URL and the artifact's MIME type.
3. Detects AirPlay support (webkitShowPlaybackTargetPicker on
HTMLVideoElement) and shows an AirPlay icon button alongside Cast.
The <video> element carries x-webkit-airplay='allow', so Safari's native
controls also surface the AirPlay option. The explicit overlay button
calls webkitShowPlaybackTargetPicker() for consistent placement.
Both buttons are hidden when the respective APIs are unavailable (HTTP
pages, non-Safari for AirPlay, no Cast extension/devices), so there is no
UI regression for users who don't cast.
PreviewPanel now passes contentType (derived from artifact extension) to
VideoPreview so Cast receives a correct MIME type.
New CSS: .video-tv-controls (absolute overlay), .video-tv-btn (frosted
glass icon button), .video-tv-loading (placeholder during token fetch).
* server: fix serve_artifact auth OR logic — bogus token falls back to session
Previously a request carrying ?token=<expired> was immediately rejected
with 401, even if the user held a valid session cookie. This broke
logged-in browser playback after the 2-hour signed-URL window expired,
because VideoPreview uses the signed URL as <video src>.
Fix: compute token_valid first; if the token is absent or invalid, fall
through to auth_user.require_auth() instead of returning early.
Effect: valid token skips session check, invalid/missing token checks
session, both invalid → 401 as before.
Updated the bogus-token-no-session test docstring to clarify it tests
the no-auth path specifically. Added new test:
media_token_bogus_token_with_session_returns_200 — verifies a logged-in
user can still fetch the artifact via a URL carrying a stale token.
* frontend: guard token-fetch effect against stale async resolution
A slow issueMediaToken() response for video A could resolve after the
user selected video B and call setSignedSrc(urlA), making the
preview/Cast play the wrong file.
Add a cancelled flag set in the effect cleanup; both .then and .catch
check it before touching state, so only the most recent src wins.
* frontend: load Cast media immediately if session already exists
Previously the effect only sent video to the TV on SESSION_STARTED /
SESSION_RESUMED events. Two gaps:
1. If a Cast session was already active when signedSrc became ready
(e.g. the SDK resumed a session before the token fetch finished, or
the user switches videos while already casting), nothing was sent.
2. Same gap if castReady fired after an already-established session.
Fix: extract loadMedia(session) and call it against
ctx.getCurrentSession() immediately when castReady + signedSrc are both
truthy, in addition to keeping the event listener for future connects.
* server: staged file-upload endpoint
POST /api/archives/:id/uploads streams a multipart body to a temp file
under the archive's store/temp/ directory and returns a staged_path the
capture pipeline can move into place.
- Routes: /api/archives/:id/uploads (POST, requires auth)
- Body cap: 10 GiB; chunk-streamed to disk, never buffered in memory
- Path-traversal sanitised on the filename field
- Temp files are cleaned up on error paths (disk-leak fix)
- main.rs wires the new route into the server startup
- Cargo: adds the multipart dependency
* frontend: file upload in Capture dialog
Drag-and-drop or 'Upload file' button stages files for archiving:
- File items sit alongside URL rows in the same list; each shows the
original filename, a live progress bar during upload, and a check badge
when ready. The locator input is replaced entirely — no editable field.
- Archive button is disabled until all uploads finish; each file item
contributes to the Archive N count once its upload is done.
- File items are excluded from sessionStorage persistence (they are
transient — the staged server path would be invalid after a reload).
Staged-file cleanup is handled at every exit path so temp/uploads/ does
not accumulate:
• removeRow on an in-progress item aborts the XHR; removeRow on a done
item calls DELETE /archives/:id/uploads.
• Dialog cancel (Escape / Cancel button) aborts all in-flight XHRs and
DELETEs all completed staged files via the close-event handler.
• handleArchive sets isSubmittingRef=true before dialog.close() so the
close handler skips cleanup — the background capture job handles
staged-file removal on success instead.
• uploadFile() returns { promise, abort } so the component can cancel
the XHR without any visible fetch.
api.js additions: uploadFile (XHR with progress + abort), deleteUpload.
(Static assets rebuilt from combined source to include screensharing
changes from this branch.)
* fix: collection enrollment with default_visibility_bits
Two related fixes from feat-file-uploading:
core: fix collection enrollment using default_visibility_bits instead of
entry.visibility — entries were being enrolled with the entry-level
visibility rather than the collection's configured default.
server: allow changing default_visibility_bits on the default collection
— the PATCH handler was incorrectly blocking updates to the default
collection's visibility configuration.
* server: fix unbounded staged-upload disk growth
Two review findings:
P2 — delete staged file on capture failure (routes.rs)
When perform_capture returns Err, the job was marked failed but
staged_upload_path was never removed. With a 10 GiB body cap a few
failed imports could exhaust archive storage before the next restart.
Mirror the success-path cleanup into the Err arm so the file is removed
immediately regardless of outcome.
P1 — periodic staged-upload pruning (main.rs)
The startup prune of temp/uploads/ only ran once, so uploads abandoned
mid-session (browser crash, navigation away) accumulated forever on a
long-running server. Folded the pruning logic into the existing 24 h
maintenance task alongside session cleanup, so stale dirs are swept
continuously without requiring a restart.
* server+frontend: fix staged-upload disk-growth and prune safety
Server (main.rs + routes.rs):
- Extract prune_stale_upload_dirs() helper called by both startup and
the periodic 24h task, eliminating the duplicated loop.
- Sentinel (.uploading) created in the UUID dir before streaming begins;
removed on successful completion; error path uses remove_dir_all so
the partial file and sentinel are cleaned up together.
The periodic prune skips any dir containing .uploading (active XHR).
- Startup prune passes cleanup_stale_sentinels=true: the server has not
started accepting connections yet so any sentinel is a crash remnant —
it is removed and the dir proceeds to the age check, preventing leaked
dirs from a previous crash accumulating forever.
- Staleness measured from the newest non-sentinel child file mtime so a
just-finished slow upload (dir mtime stale, file mtime fresh) is not
pruned before the user can submit it for capture. Empty dirs fall back
to dir mtime.
- Failed captures (Err branch in spawn_blocking) now also delete the
staged file and UUID dir immediately, matching the success path.
Frontend (api.js + CaptureDialog.jsx):
- submitCapture attaches err.status = res.status on non-2xx responses
so callers can distinguish a definite HTTP rejection from a network
error where the response may have been lost.
- submitBgJob catch deletes the staged file only when e.status is set
(server definitively rejected the POST /captures request). A network
error leaves the file in place because the server may have accepted
the job and the response was lost — deleting would race the capture.
This commit is contained in:
parent
6377daadae
commit
1af920eb63
17 changed files with 1498 additions and 142 deletions
|
|
@ -7,6 +7,64 @@ use std::{net::SocketAddr, path::PathBuf};
|
|||
|
||||
const DEFAULT_BIND: &str = "127.0.0.1:8080";
|
||||
|
||||
/// Prune abandoned staged-upload UUID dirs under `uploads_dir` whose content
|
||||
/// is older than `cutoff`.
|
||||
///
|
||||
/// `cleanup_stale_sentinels`: pass `true` at startup (before the server begins
|
||||
/// accepting connections) so crash-leftover `.uploading` markers are removed and
|
||||
/// those dirs are subject to the normal age check. Pass `false` from the
|
||||
/// in-process periodic task so dirs with a live sentinel (active XHR) are
|
||||
/// skipped entirely.
|
||||
///
|
||||
/// Staleness is measured from the newest non-sentinel child file's mtime so a
|
||||
/// just-finished slow upload is not pruned before the user submits it for
|
||||
/// capture. Empty dirs fall back to the directory mtime.
|
||||
fn prune_stale_upload_dirs(
|
||||
uploads_dir: &std::path::Path,
|
||||
cutoff: std::time::SystemTime,
|
||||
cleanup_stale_sentinels: bool,
|
||||
) {
|
||||
let Ok(entries) = std::fs::read_dir(uploads_dir) else {
|
||||
return;
|
||||
};
|
||||
for entry in entries.flatten() {
|
||||
let path = entry.path();
|
||||
if !path.is_dir() {
|
||||
continue;
|
||||
}
|
||||
let sentinel = path.join(".uploading");
|
||||
if sentinel.exists() {
|
||||
if cleanup_stale_sentinels {
|
||||
// Server just started — no uploads are in flight, so any
|
||||
// sentinel is a crash remnant. Remove it and fall through
|
||||
// to the age check below.
|
||||
let _ = std::fs::remove_file(&sentinel);
|
||||
} else {
|
||||
// An active XHR is writing to this dir — leave it alone.
|
||||
continue;
|
||||
}
|
||||
}
|
||||
// Measure staleness from the newest non-sentinel child file so a
|
||||
// completed slow upload isn't pruned while the user is still on the
|
||||
// capture form. Fall back to dir mtime only when the dir is empty.
|
||||
let newest_child = std::fs::read_dir(&path)
|
||||
.ok()
|
||||
.into_iter()
|
||||
.flatten()
|
||||
.filter_map(|e| e.ok())
|
||||
.filter(|e| e.file_name() != std::ffi::OsStr::new(".uploading"))
|
||||
.filter_map(|e| e.metadata().ok())
|
||||
.filter_map(|m| m.modified().ok())
|
||||
.max();
|
||||
let reference = newest_child
|
||||
.or_else(|| entry.metadata().and_then(|m| m.modified()).ok())
|
||||
.unwrap_or(std::time::SystemTime::UNIX_EPOCH);
|
||||
if reference < cutoff {
|
||||
let _ = std::fs::remove_dir_all(&path);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::main]
|
||||
async fn main() -> Result<()> {
|
||||
let config_path = std::env::args()
|
||||
|
|
@ -30,15 +88,34 @@ async fn main() -> Result<()> {
|
|||
for archive in ®istry.archives {
|
||||
if let Ok(conn) = archivr_core::database::open_or_initialize(&archive.archive_path) {
|
||||
match archivr_core::database::fail_stalled_capture_jobs(&conn) {
|
||||
Ok(n) if n > 0 => eprintln!("info: marked {n} stalled capture job(s) as failed in '{}'", archive.id),
|
||||
Err(e) => eprintln!("warn: stalled job cleanup failed for '{}': {e:#}", archive.id),
|
||||
Ok(n) if n > 0 => eprintln!(
|
||||
"info: marked {n} stalled capture job(s) as failed in '{}'",
|
||||
archive.id
|
||||
),
|
||||
Err(e) => eprintln!(
|
||||
"warn: stalled job cleanup failed for '{}': {e:#}",
|
||||
archive.id
|
||||
),
|
||||
_ => {}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Spawn session cleanup: runs at startup and every 24h.
|
||||
// Prune staged upload dirs older than 24 h. cleanup_stale_sentinels=true
|
||||
// because no uploads are in flight before the server starts listening.
|
||||
let prune_cutoff = std::time::SystemTime::now()
|
||||
.checked_sub(std::time::Duration::from_secs(24 * 60 * 60))
|
||||
.unwrap_or(std::time::SystemTime::UNIX_EPOCH);
|
||||
for archive in ®istry.archives {
|
||||
if let Ok(paths) = archivr_core::archive::read_archive_paths(&archive.archive_path) {
|
||||
let uploads_dir = paths.store_path.join("temp").join("uploads");
|
||||
prune_stale_upload_dirs(&uploads_dir, prune_cutoff, true);
|
||||
}
|
||||
}
|
||||
|
||||
// Spawn maintenance task: session cleanup + staged-upload pruning, every 24 h.
|
||||
let cleanup_auth_path = auth_db_path.clone();
|
||||
let cleanup_registry = registry.clone();
|
||||
tokio::spawn(async move {
|
||||
loop {
|
||||
if let Ok(conn) = archivr_core::database::open_auth_db(&cleanup_auth_path) {
|
||||
|
|
@ -48,6 +125,17 @@ async fn main() -> Result<()> {
|
|||
_ => {}
|
||||
}
|
||||
}
|
||||
// cleanup_stale_sentinels=false: server is live, respect active uploads.
|
||||
let prune_cutoff = std::time::SystemTime::now()
|
||||
.checked_sub(std::time::Duration::from_secs(24 * 60 * 60))
|
||||
.unwrap_or(std::time::SystemTime::UNIX_EPOCH);
|
||||
for archive in &cleanup_registry.archives {
|
||||
if let Ok(paths) = archivr_core::archive::read_archive_paths(&archive.archive_path)
|
||||
{
|
||||
let uploads_dir = paths.store_path.join("temp").join("uploads");
|
||||
prune_stale_upload_dirs(&uploads_dir, prune_cutoff, false);
|
||||
}
|
||||
}
|
||||
tokio::time::sleep(tokio::time::Duration::from_secs(24 * 60 * 60)).await;
|
||||
}
|
||||
});
|
||||
|
|
@ -63,6 +151,10 @@ async fn main() -> Result<()> {
|
|||
|
||||
let listener = tokio::net::TcpListener::bind(addr).await?;
|
||||
println!("archivr-server listening on http://{addr}");
|
||||
axum::serve(listener, app.into_make_service_with_connect_info::<SocketAddr>()).await?;
|
||||
axum::serve(
|
||||
listener,
|
||||
app.into_make_service_with_connect_info::<SocketAddr>(),
|
||||
)
|
||||
.await?;
|
||||
Ok(())
|
||||
}
|
||||
|
|
|
|||
|
|
@ -32,7 +32,7 @@ use std::{
|
|||
use archivr_core::{archive, capture, database, downloader};
|
||||
use axum::{
|
||||
Json, Router,
|
||||
extract::{ConnectInfo, Path, Query, Request, State},
|
||||
extract::{ConnectInfo, DefaultBodyLimit, Multipart, Path, Query, Request, State},
|
||||
http::StatusCode,
|
||||
middleware::Next,
|
||||
response::{IntoResponse, Response},
|
||||
|
|
@ -50,11 +50,21 @@ use rusqlite::OptionalExtension;
|
|||
const LOGIN_WINDOW: Duration = Duration::from_secs(15 * 60);
|
||||
const LOGIN_MAX_ATTEMPTS: usize = 5;
|
||||
|
||||
// Short-lived token granting unauthenticated access to one specific artifact.
|
||||
// Used so Cast / AirPlay devices (which carry no session cookie) can fetch media.
|
||||
pub(crate) struct MediaToken {
|
||||
archive_id: String,
|
||||
entry_uid: String,
|
||||
artifact_index: usize,
|
||||
expires_at: std::time::Instant,
|
||||
}
|
||||
|
||||
#[derive(Clone)]
|
||||
pub struct AppState {
|
||||
registry: Arc<ServerRegistry>,
|
||||
pub auth_db_path: Arc<std::path::PathBuf>,
|
||||
pub login_attempts: Arc<Mutex<HashMap<IpAddr, VecDeque<Instant>>>>,
|
||||
pub media_tokens: Arc<Mutex<HashMap<String, MediaToken>>>,
|
||||
}
|
||||
|
||||
#[derive(Debug, serde::Deserialize, Default)]
|
||||
|
|
@ -133,7 +143,7 @@ async fn security_headers(req: Request, next: Next) -> Response {
|
|||
axum::http::header::HeaderName::from_static("content-security-policy"),
|
||||
axum::http::HeaderValue::from_static(
|
||||
"default-src 'self'; \
|
||||
script-src 'self'; \
|
||||
script-src 'self' https://www.gstatic.com; \
|
||||
style-src 'self' 'unsafe-inline' https://fonts.googleapis.com; \
|
||||
img-src 'self' data: blob: https:; \
|
||||
font-src 'self' https://fonts.gstatic.com; \
|
||||
|
|
@ -249,6 +259,10 @@ pub fn app_with_state(state: AppState) -> Router {
|
|||
"/api/archives/:archive_id/entries/:entry_uid/artifacts/:artifact_index",
|
||||
get(serve_artifact),
|
||||
)
|
||||
.route(
|
||||
"/api/archives/:archive_id/entries/:entry_uid/artifacts/:artifact_index/media-token",
|
||||
post(issue_media_token),
|
||||
)
|
||||
.route(
|
||||
"/api/archives/:archive_id/entries/:entry_uid/rearchive",
|
||||
post(rearchive_handler),
|
||||
|
|
@ -260,6 +274,12 @@ pub fn app_with_state(state: AppState) -> Router {
|
|||
.route("/api/archives/:archive_id/blobs/:sha256", get(serve_blob))
|
||||
.route("/api/archives/:archive_id/runs", get(list_runs))
|
||||
.route("/api/archives/:archive_id/captures", post(capture_handler))
|
||||
.route(
|
||||
"/api/archives/:archive_id/uploads",
|
||||
post(upload_handler)
|
||||
.delete(delete_upload_handler)
|
||||
.layer(DefaultBodyLimit::max(10 * 1024 * 1024 * 1024)),
|
||||
)
|
||||
.route(
|
||||
"/api/archives/:archive_id/captures/probe",
|
||||
get(probe_handler),
|
||||
|
|
@ -386,6 +406,7 @@ pub fn app(registry: ServerRegistry, auth_db_path: std::path::PathBuf) -> Router
|
|||
registry: Arc::new(registry),
|
||||
auth_db_path: Arc::new(auth_db_path),
|
||||
login_attempts: Arc::new(Mutex::new(HashMap::new())),
|
||||
media_tokens: Arc::new(Mutex::new(HashMap::new())),
|
||||
};
|
||||
app_with_state(state)
|
||||
}
|
||||
|
|
@ -421,7 +442,11 @@ async fn list_entry_children(
|
|||
let mounted = mounted_archive(&state, &archive_id)?;
|
||||
let conn = database::open_or_initialize(&mounted.archive_path)?;
|
||||
let caller_bits = auth_to_caller_bits(&auth);
|
||||
Ok(Json(archive::list_child_entries(&conn, &entry_uid, caller_bits)?))
|
||||
Ok(Json(archive::list_child_entries(
|
||||
&conn,
|
||||
&entry_uid,
|
||||
caller_bits,
|
||||
)?))
|
||||
}
|
||||
|
||||
async fn search_entries_handler(
|
||||
|
|
@ -466,14 +491,41 @@ async fn list_runs(
|
|||
let conn = database::open_or_initialize(&mounted.archive_path)?;
|
||||
Ok(Json(archive::list_runs(&conn)?))
|
||||
}
|
||||
const MEDIA_TOKEN_TTL: Duration = Duration::from_secs(2 * 60 * 60); // 2 h
|
||||
|
||||
#[derive(Debug, serde::Deserialize, Default)]
|
||||
struct ArtifactQuery {
|
||||
token: Option<String>,
|
||||
}
|
||||
|
||||
#[derive(serde::Serialize)]
|
||||
struct MediaTokenResponse {
|
||||
url: String,
|
||||
expires_in_secs: u64,
|
||||
}
|
||||
|
||||
async fn serve_artifact(
|
||||
State(state): State<AppState>,
|
||||
auth_user: AuthUser,
|
||||
Path((archive_id, entry_uid, artifact_index)): Path<(String, String, usize)>,
|
||||
Query(params): Query<ArtifactQuery>,
|
||||
req: Request,
|
||||
) -> Result<Response, ApiError> {
|
||||
auth_user.require_auth()?;
|
||||
// Auth: valid scoped token OR authenticated session (OR both).
|
||||
// A token present but invalid/expired falls back to session auth so that
|
||||
// a logged-in browser player keeps working after a token expires.
|
||||
let token_valid = params.token.as_deref().map_or(false, |tok| {
|
||||
let tokens = state.media_tokens.lock();
|
||||
tokens.get(tok).map_or(false, |t| {
|
||||
t.archive_id == archive_id
|
||||
&& t.entry_uid == entry_uid
|
||||
&& t.artifact_index == artifact_index
|
||||
&& t.expires_at > std::time::Instant::now()
|
||||
})
|
||||
});
|
||||
if !token_valid {
|
||||
auth_user.require_auth()?;
|
||||
}
|
||||
let mounted = mounted_archive(&state, &archive_id)?;
|
||||
let paths = archive::read_archive_paths(&mounted.archive_path)?;
|
||||
let conn = database::open_or_initialize(&mounted.archive_path)?;
|
||||
|
|
@ -491,6 +543,51 @@ async fn serve_artifact(
|
|||
.into_response())
|
||||
}
|
||||
|
||||
/// POST /api/archives/:archive_id/entries/:entry_uid/artifacts/:artifact_index/media-token
|
||||
///
|
||||
/// Requires an authenticated session. Returns a short-lived signed URL that
|
||||
/// allows unauthenticated GET of the specified artifact — intended for Cast /
|
||||
/// AirPlay devices that cannot carry the browser's session cookie.
|
||||
async fn issue_media_token(
|
||||
State(state): State<AppState>,
|
||||
auth_user: AuthUser,
|
||||
Path((archive_id, entry_uid, artifact_index)): Path<(String, String, usize)>,
|
||||
) -> Result<Json<MediaTokenResponse>, ApiError> {
|
||||
auth_user.require_auth()?;
|
||||
// Verify the artifact actually exists before issuing a token.
|
||||
let mounted = mounted_archive(&state, &archive_id)?;
|
||||
let conn = database::open_or_initialize(&mounted.archive_path)?;
|
||||
let detail = archive::get_entry_detail(&conn, &entry_uid)?
|
||||
.ok_or(ApiError::not_found("entry not found"))?;
|
||||
if artifact_index >= detail.artifacts.len() {
|
||||
return Err(ApiError::not_found("artifact index out of range"));
|
||||
}
|
||||
let token = auth::generate_token();
|
||||
let now = std::time::Instant::now();
|
||||
{
|
||||
let mut tokens = state.media_tokens.lock();
|
||||
// GC expired tokens on each issuance to keep the map bounded.
|
||||
tokens.retain(|_, t| t.expires_at > now);
|
||||
tokens.insert(
|
||||
token.clone(),
|
||||
MediaToken {
|
||||
archive_id: archive_id.clone(),
|
||||
entry_uid: entry_uid.clone(),
|
||||
artifact_index,
|
||||
expires_at: now + MEDIA_TOKEN_TTL,
|
||||
},
|
||||
);
|
||||
}
|
||||
let url = format!(
|
||||
"/api/archives/{}/entries/{}/artifacts/{}?token={}",
|
||||
archive_id, entry_uid, artifact_index, token
|
||||
);
|
||||
Ok(Json(MediaTokenResponse {
|
||||
url,
|
||||
expires_in_secs: MEDIA_TOKEN_TTL.as_secs(),
|
||||
}))
|
||||
}
|
||||
|
||||
async fn serve_entry_favicon(
|
||||
State(state): State<AppState>,
|
||||
auth_user: AuthUser,
|
||||
|
|
@ -836,6 +933,11 @@ struct UpdateCookieRuleBody {
|
|||
ordinal: Option<i64>,
|
||||
}
|
||||
|
||||
#[derive(Debug, serde::Deserialize)]
|
||||
struct DeleteUploadBody {
|
||||
locator: String,
|
||||
}
|
||||
|
||||
async fn capture_handler(
|
||||
State(state): State<AppState>,
|
||||
auth_user: AuthUser,
|
||||
|
|
@ -866,7 +968,11 @@ async fn capture_handler(
|
|||
.and_then(|n| n.parse::<u32>().ok())
|
||||
.is_some()
|
||||
};
|
||||
if let Some(bad) = body.per_item_quality.values().find(|q| !is_valid_quality(q)) {
|
||||
if let Some(bad) = body
|
||||
.per_item_quality
|
||||
.values()
|
||||
.find(|q| !is_valid_quality(q))
|
||||
{
|
||||
return Err(ApiError::bad_request(&format!(
|
||||
"invalid per_item_quality value {bad:?}: must be \"best\", \"audio\", or a height string like \"1080p\""
|
||||
)));
|
||||
|
|
@ -923,6 +1029,24 @@ async fn capture_handler(
|
|||
|
||||
// Spawn background capture.
|
||||
let locator = body.locator.trim().to_string();
|
||||
// If the locator is a file:// path staged under temp/uploads/, track it for cleanup.
|
||||
// Canonicalize both sides to prevent path-traversal via `..` components in the locator.
|
||||
// Same pattern as artifact serving (line ~631). The staged file must already exist on disk
|
||||
// (it was written by upload_handler), so canonicalize() will resolve symlinks correctly.
|
||||
let staged_upload_path: Option<std::path::PathBuf> = if locator.starts_with("file://") {
|
||||
let file_path = std::path::PathBuf::from(locator.trim_start_matches("file://"));
|
||||
let staging_dir = archive_paths.store_path.join("temp").join("uploads");
|
||||
match (file_path.canonicalize(), staging_dir.canonicalize()) {
|
||||
(Ok(canonical_file), Ok(canonical_staging))
|
||||
if canonical_file.starts_with(&canonical_staging) =>
|
||||
{
|
||||
Some(canonical_file)
|
||||
}
|
||||
_ => None,
|
||||
}
|
||||
} else {
|
||||
None
|
||||
};
|
||||
let quality = body.quality.clone();
|
||||
let archive_path = mounted.archive_path.clone();
|
||||
let job_uid_bg = job_uid.clone();
|
||||
|
|
@ -976,6 +1100,16 @@ async fn capture_handler(
|
|||
notes,
|
||||
)
|
||||
.ok();
|
||||
// Clean up staged upload file — content is now in the raw store.
|
||||
// `staged` is already the canonicalized path (safe to remove_file directly).
|
||||
// Also attempt to remove the now-empty UUID parent dir; fails silently if
|
||||
// non-empty or already gone.
|
||||
if let Some(staged) = staged_upload_path {
|
||||
let _ = std::fs::remove_file(&staged);
|
||||
if let Some(parent) = staged.parent() {
|
||||
let _ = std::fs::remove_dir(parent);
|
||||
}
|
||||
}
|
||||
}
|
||||
Err(e) => {
|
||||
database::update_capture_job_status(
|
||||
|
|
@ -987,6 +1121,15 @@ async fn capture_handler(
|
|||
None,
|
||||
)
|
||||
.ok();
|
||||
// Failed captures never move the file into the raw store,
|
||||
// so clean up the staged upload here rather than waiting
|
||||
// for the next startup or periodic prune.
|
||||
if let Some(staged) = staged_upload_path {
|
||||
let _ = std::fs::remove_file(&staged);
|
||||
if let Some(parent) = staged.parent() {
|
||||
let _ = std::fs::remove_dir(parent);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
});
|
||||
|
|
@ -997,6 +1140,155 @@ async fn capture_handler(
|
|||
))
|
||||
}
|
||||
|
||||
async fn upload_handler(
|
||||
State(state): State<AppState>,
|
||||
auth_user: AuthUser,
|
||||
Path(archive_id): Path<String>,
|
||||
mut multipart: Multipart,
|
||||
) -> Result<(StatusCode, Json<serde_json::Value>), ApiError> {
|
||||
auth_user.require_role(ROLE_USER)?;
|
||||
let mounted = mounted_archive(&state, &archive_id)?;
|
||||
let archive_paths =
|
||||
archive::read_archive_paths(&mounted.archive_path).map_err(ApiError::from)?;
|
||||
|
||||
// Stage under temp/uploads/<uuid>/<safe_name> so Source::Local derives the
|
||||
// entry title from Path::file_name() of the locator rather than the uuid.
|
||||
let staging_base = archive_paths.store_path.join("temp").join("uploads");
|
||||
|
||||
while let Some(mut field) = multipart
|
||||
.next_field()
|
||||
.await
|
||||
.map_err(|e| ApiError::bad_request(&e.to_string()))?
|
||||
{
|
||||
if field.name() != Some("file") {
|
||||
continue;
|
||||
}
|
||||
let filename = field.file_name().unwrap_or("upload").to_string();
|
||||
|
||||
// Sanitize: strip path separators and control chars, truncate to 200 chars.
|
||||
let safe_name: String = filename
|
||||
.chars()
|
||||
.filter(|c| !matches!(*c, '/' | '\\' | '\0'))
|
||||
.collect::<String>()
|
||||
.trim()
|
||||
.to_string();
|
||||
let safe_name = if safe_name.is_empty() {
|
||||
"upload".to_string()
|
||||
} else {
|
||||
safe_name.chars().take(200).collect()
|
||||
};
|
||||
|
||||
let uuid_dir = staging_base.join(uuid::Uuid::new_v4().simple().to_string());
|
||||
tokio::fs::create_dir_all(&uuid_dir).await.map_err(|e| {
|
||||
ApiError::from(anyhow::anyhow!("failed to create upload staging dir: {e}"))
|
||||
})?;
|
||||
|
||||
// Sentinel: exists while the XHR is streaming. The prune task skips any
|
||||
// uuid_dir that contains this file, so a slow upload is never deleted
|
||||
// mid-transfer regardless of wall-clock age.
|
||||
let sentinel = uuid_dir.join(".uploading");
|
||||
tokio::fs::File::create(&sentinel).await.map_err(|e| {
|
||||
ApiError::from(anyhow::anyhow!("failed to create upload sentinel: {e}"))
|
||||
})?;
|
||||
|
||||
let staged_path = uuid_dir.join(&safe_name);
|
||||
|
||||
// Stream chunks directly to disk — never buffers the full file in RAM.
|
||||
// The body limit (10 GiB) is enforced by the DefaultBodyLimit layer.
|
||||
let stream_result: Result<(), ApiError> = async {
|
||||
use tokio::io::AsyncWriteExt as _;
|
||||
let file = tokio::fs::File::create(&staged_path).await.map_err(|e| {
|
||||
ApiError::from(anyhow::anyhow!("failed to create staged file: {e}"))
|
||||
})?;
|
||||
let mut writer = tokio::io::BufWriter::new(file);
|
||||
while let Some(chunk) = field
|
||||
.chunk()
|
||||
.await
|
||||
.map_err(|e| ApiError::bad_request(&e.to_string()))?
|
||||
{
|
||||
writer
|
||||
.write_all(&chunk)
|
||||
.await
|
||||
.map_err(|e| ApiError::from(anyhow::anyhow!("failed to write chunk: {e}")))?;
|
||||
}
|
||||
writer
|
||||
.flush()
|
||||
.await
|
||||
.map_err(|e| ApiError::from(anyhow::anyhow!("failed to flush upload: {e}")))?;
|
||||
Ok(())
|
||||
}
|
||||
.await;
|
||||
if let Err(e) = stream_result {
|
||||
// remove_dir_all cleans up the partial file and the sentinel together.
|
||||
let _ = tokio::fs::remove_dir_all(&uuid_dir).await;
|
||||
return Err(e);
|
||||
}
|
||||
|
||||
// Stream complete — drop the sentinel so the prune task can reclaim the
|
||||
// dir if it is later abandoned without being submitted for capture.
|
||||
let _ = tokio::fs::remove_file(&sentinel).await;
|
||||
|
||||
let size = tokio::fs::metadata(&staged_path)
|
||||
.await
|
||||
.map(|m| m.len() as i64)
|
||||
.unwrap_or(0);
|
||||
|
||||
let locator = format!("file://{}", staged_path.display());
|
||||
return Ok((
|
||||
StatusCode::OK,
|
||||
Json(serde_json::json!({
|
||||
"locator": locator,
|
||||
"filename": filename,
|
||||
"size": size,
|
||||
})),
|
||||
));
|
||||
}
|
||||
|
||||
Err(ApiError::bad_request(
|
||||
"no file field found in multipart upload",
|
||||
))
|
||||
}
|
||||
|
||||
/// `DELETE /api/archives/:archive_id/uploads`
|
||||
///
|
||||
/// Discards a staged upload file that was never submitted for capture —
|
||||
/// called by the frontend when the user removes a file row or cancels the
|
||||
/// dialog. The same canonicalize-then-prefix-check used in `capture_handler`
|
||||
/// prevents path traversal via crafted `file://` locators.
|
||||
async fn delete_upload_handler(
|
||||
State(state): State<AppState>,
|
||||
auth_user: AuthUser,
|
||||
Path(archive_id): Path<String>,
|
||||
Json(body): Json<DeleteUploadBody>,
|
||||
) -> Result<StatusCode, ApiError> {
|
||||
auth_user.require_role(ROLE_USER)?;
|
||||
let mounted = mounted_archive(&state, &archive_id)?;
|
||||
let archive_paths =
|
||||
archive::read_archive_paths(&mounted.archive_path).map_err(ApiError::from)?;
|
||||
|
||||
if !body.locator.starts_with("file://") {
|
||||
return Err(ApiError::bad_request("locator must be a file:// URI"));
|
||||
}
|
||||
let file_path = std::path::PathBuf::from(body.locator.trim_start_matches("file://"));
|
||||
let staging_dir = archive_paths.store_path.join("temp").join("uploads");
|
||||
|
||||
// Canonicalize both sides before the prefix check (path-traversal guard).
|
||||
let (canonical_file, canonical_staging) =
|
||||
match (file_path.canonicalize(), staging_dir.canonicalize()) {
|
||||
(Ok(f), Ok(s)) => (f, s),
|
||||
_ => return Err(ApiError::not_found("staged upload not found")),
|
||||
};
|
||||
if !canonical_file.starts_with(&canonical_staging) {
|
||||
return Err(ApiError::bad_request("locator is not a staged upload"));
|
||||
}
|
||||
|
||||
tokio::fs::remove_file(&canonical_file).await.ok();
|
||||
if let Some(parent) = canonical_file.parent() {
|
||||
tokio::fs::remove_dir(parent).await.ok(); // no-op if non-empty
|
||||
}
|
||||
Ok(StatusCode::NO_CONTENT)
|
||||
}
|
||||
|
||||
async fn get_capture_job_handler(
|
||||
State(state): State<AppState>,
|
||||
auth_user: AuthUser,
|
||||
|
|
@ -1188,8 +1480,11 @@ async fn probe_playlist_handler(
|
|||
return Err(ApiError::bad_request("locator must not be empty"));
|
||||
}
|
||||
// Validate it's a playlist/channel source and expand shorthands.
|
||||
let canonical_url = capture::locator_to_playlist_url(&locator)
|
||||
.ok_or_else(|| ApiError::bad_request("locator is not a YouTube playlist, channel, YTM playlist, or Spotify album/playlist"))?;
|
||||
let canonical_url = capture::locator_to_playlist_url(&locator).ok_or_else(|| {
|
||||
ApiError::bad_request(
|
||||
"locator is not a YouTube playlist, channel, YTM playlist, or Spotify album/playlist",
|
||||
)
|
||||
})?;
|
||||
// Verify archive exists.
|
||||
let _ = mounted_archive(&state, &archive_id)?;
|
||||
// Resolve cookies.
|
||||
|
|
@ -3758,6 +4053,7 @@ mod tests {
|
|||
}),
|
||||
auth_db_path: Arc::new(auth_path),
|
||||
login_attempts: Arc::new(Mutex::new(HashMap::new())),
|
||||
media_tokens: Arc::new(Mutex::new(HashMap::new())),
|
||||
};
|
||||
let bad_creds = serde_json::json!({ "username": "nobody", "password": "wrong" });
|
||||
for _ in 0..LOGIN_MAX_ATTEMPTS {
|
||||
|
|
@ -3821,6 +4117,7 @@ mod tests {
|
|||
}),
|
||||
auth_db_path: Arc::new(auth_path),
|
||||
login_attempts: Arc::new(Mutex::new(HashMap::new())),
|
||||
media_tokens: Arc::new(Mutex::new(HashMap::new())),
|
||||
};
|
||||
let bad_creds = serde_json::json!({ "username": "x", "password": "y" });
|
||||
for _ in 0..LOGIN_MAX_ATTEMPTS {
|
||||
|
|
@ -4589,6 +4886,7 @@ mod tests {
|
|||
}),
|
||||
auth_db_path: Arc::new(auth_path.clone()),
|
||||
login_attempts: Arc::new(Mutex::new(HashMap::new())),
|
||||
media_tokens: Arc::new(Mutex::new(HashMap::new())),
|
||||
};
|
||||
let session_cookie = make_test_session(&auth_path);
|
||||
|
||||
|
|
@ -4887,4 +5185,237 @@ mod tests {
|
|||
"extra disk-only file must be deleted"
|
||||
);
|
||||
}
|
||||
|
||||
// ── Media token tests ────────────────────────────────────────────────────
|
||||
|
||||
// Helper: build a minimal archive + auth setup and return (state, entry_uid, session_cookie).
|
||||
async fn make_media_token_state(
|
||||
dir: &tempfile::TempDir,
|
||||
) -> (AppState, String, std::path::PathBuf, String) {
|
||||
let store_path = dir.path().join("store");
|
||||
let paths =
|
||||
archivr_core::archive::initialize_archive(dir.path(), &store_path, "test", false)
|
||||
.unwrap();
|
||||
// Write artifact file.
|
||||
let artifact_relpath = "raw/m/e/video.mp4";
|
||||
let artifact_dir = store_path.join("raw").join("m").join("e");
|
||||
std::fs::create_dir_all(&artifact_dir).unwrap();
|
||||
std::fs::write(artifact_dir.join("video.mp4"), b"fakevideo").unwrap();
|
||||
// Populate DB.
|
||||
let conn = database::open_or_initialize(&paths.archive_path).unwrap();
|
||||
let user_id = database::ensure_default_user(&conn).unwrap();
|
||||
let sid = database::upsert_source_identity(
|
||||
&conn,
|
||||
"yt",
|
||||
"video",
|
||||
Some("media-token-test"),
|
||||
Some("https://yt.example/v"),
|
||||
"https://yt.example/v",
|
||||
)
|
||||
.unwrap();
|
||||
let run = database::create_archive_run(&conn, user_id, 1).unwrap();
|
||||
let entry = database::create_archived_entry(
|
||||
&conn,
|
||||
&database::NewEntry {
|
||||
source_identity_id: sid,
|
||||
archive_run_id: run.id,
|
||||
parent_entry_id: None,
|
||||
root_entry_id: None,
|
||||
created_by_user_id: user_id,
|
||||
owned_by_user_id: user_id,
|
||||
source_kind: "yt".to_string(),
|
||||
entity_kind: "video".to_string(),
|
||||
title: Some("Test Video".to_string()),
|
||||
visibility: "private".to_string(),
|
||||
representation_kind: "video".to_string(),
|
||||
source_metadata_json: "{}".to_string(),
|
||||
display_metadata_json: None,
|
||||
},
|
||||
)
|
||||
.unwrap();
|
||||
let blob_id = database::upsert_blob(
|
||||
&conn,
|
||||
&database::BlobRecord {
|
||||
sha256: "bbbb2222cccc3333dddd4444aaaa1111bbbb2222cccc3333dddd4444aaaa1111"
|
||||
.to_string(),
|
||||
byte_size: 9,
|
||||
mime_type: Some("video/mp4".to_string()),
|
||||
extension: Some("mp4".to_string()),
|
||||
raw_relpath: artifact_relpath.to_string(),
|
||||
},
|
||||
)
|
||||
.unwrap();
|
||||
database::add_entry_artifact(
|
||||
&conn,
|
||||
&database::NewArtifact {
|
||||
entry_id: entry.id,
|
||||
artifact_role: "primary_media".to_string(),
|
||||
storage_area: "raw".to_string(),
|
||||
relpath: artifact_relpath.to_string(),
|
||||
blob_id: Some(blob_id),
|
||||
logical_path: None,
|
||||
metadata_json: None,
|
||||
},
|
||||
)
|
||||
.unwrap();
|
||||
drop(conn);
|
||||
let auth_path = dir.path().join("auth.sqlite");
|
||||
{
|
||||
let conn = archivr_core::database::open_auth_db(&auth_path).unwrap();
|
||||
archivr_core::database::create_owner(&conn, "testowner", "dummy").unwrap();
|
||||
}
|
||||
let session_cookie = make_test_session(&auth_path);
|
||||
let registry = ServerRegistry {
|
||||
archives: vec![MountedArchive {
|
||||
id: "test".to_string(),
|
||||
label: "Test".to_string(),
|
||||
archive_path: paths.archive_path.clone(),
|
||||
}],
|
||||
bind: None,
|
||||
auth_db_path: None,
|
||||
};
|
||||
let state = AppState {
|
||||
registry: Arc::new(registry),
|
||||
auth_db_path: Arc::new(auth_path),
|
||||
login_attempts: Arc::new(Mutex::new(HashMap::new())),
|
||||
media_tokens: Arc::new(Mutex::new(HashMap::new())),
|
||||
};
|
||||
(state, entry.entry_uid, paths.archive_path, session_cookie)
|
||||
}
|
||||
|
||||
/// Bare artifact URL (no token, no session) must still return 401.
|
||||
#[tokio::test]
|
||||
async fn media_token_bare_artifact_without_auth_returns_401() {
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
let (state, entry_uid, _, _) = make_media_token_state(&dir).await;
|
||||
let uri = format!("/api/archives/test/entries/{}/artifacts/0", entry_uid);
|
||||
let response = app_with_state(state)
|
||||
.oneshot(Request::builder().uri(&uri).body(Body::empty()).unwrap())
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(response.status(), StatusCode::UNAUTHORIZED);
|
||||
}
|
||||
|
||||
/// Authenticated POST to media-token, then unauthenticated GET with token → 200.
|
||||
#[tokio::test]
|
||||
async fn media_token_tokenized_artifact_succeeds_unauthenticated() {
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
let (state, entry_uid, _, session_cookie) = make_media_token_state(&dir).await;
|
||||
// Issue token (authenticated).
|
||||
let token_uri = format!(
|
||||
"/api/archives/test/entries/{}/artifacts/0/media-token",
|
||||
entry_uid
|
||||
);
|
||||
let token_resp = app_with_state(state.clone())
|
||||
.oneshot(
|
||||
Request::builder()
|
||||
.method("POST")
|
||||
.uri(&token_uri)
|
||||
.header("cookie", &session_cookie)
|
||||
.body(Body::empty())
|
||||
.unwrap(),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(token_resp.status(), StatusCode::OK);
|
||||
let body = body_json(token_resp).await;
|
||||
let signed_url = body["url"].as_str().expect("url field missing");
|
||||
assert!(body["expires_in_secs"].as_u64().unwrap() > 0);
|
||||
// Fetch artifact with signed URL — no session cookie.
|
||||
let artifact_resp = app_with_state(state)
|
||||
.oneshot(
|
||||
Request::builder()
|
||||
.uri(signed_url)
|
||||
.body(Body::empty())
|
||||
.unwrap(),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(artifact_resp.status(), StatusCode::OK);
|
||||
}
|
||||
|
||||
/// A bogus token with NO session must return 401 (no valid auth path).
|
||||
#[tokio::test]
|
||||
async fn media_token_invalid_token_returns_401() {
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
let (state, entry_uid, _, _) = make_media_token_state(&dir).await;
|
||||
let uri = format!(
|
||||
"/api/archives/test/entries/{}/artifacts/0?token=not-a-real-token",
|
||||
entry_uid
|
||||
);
|
||||
let response = app_with_state(state)
|
||||
.oneshot(Request::builder().uri(&uri).body(Body::empty()).unwrap())
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(response.status(), StatusCode::UNAUTHORIZED);
|
||||
}
|
||||
|
||||
/// A bogus token WITH a valid session must return 200 — the session fallback
|
||||
/// keeps a logged-in browser player working after a token expires.
|
||||
#[tokio::test]
|
||||
async fn media_token_bogus_token_with_session_returns_200() {
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
let (state, entry_uid, _, session_cookie) = make_media_token_state(&dir).await;
|
||||
let uri = format!(
|
||||
"/api/archives/test/entries/{}/artifacts/0?token=not-a-real-token",
|
||||
entry_uid
|
||||
);
|
||||
let response = app_with_state(state)
|
||||
.oneshot(
|
||||
Request::builder()
|
||||
.uri(&uri)
|
||||
.header("cookie", &session_cookie)
|
||||
.body(Body::empty())
|
||||
.unwrap(),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(response.status(), StatusCode::OK);
|
||||
}
|
||||
|
||||
/// A token issued for artifact 0 must not unlock artifact 1.
|
||||
#[tokio::test]
|
||||
async fn media_token_wrong_artifact_index_returns_401() {
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
let (state, entry_uid, _, session_cookie) = make_media_token_state(&dir).await;
|
||||
// Issue token for artifact 0.
|
||||
let token_uri = format!(
|
||||
"/api/archives/test/entries/{}/artifacts/0/media-token",
|
||||
entry_uid
|
||||
);
|
||||
let token_resp = app_with_state(state.clone())
|
||||
.oneshot(
|
||||
Request::builder()
|
||||
.method("POST")
|
||||
.uri(&token_uri)
|
||||
.header("cookie", &session_cookie)
|
||||
.body(Body::empty())
|
||||
.unwrap(),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(token_resp.status(), StatusCode::OK);
|
||||
let body = body_json(token_resp).await;
|
||||
let token = body["url"]
|
||||
.as_str()
|
||||
.unwrap()
|
||||
.split("token=")
|
||||
.nth(1)
|
||||
.unwrap();
|
||||
// Try to use it for artifact 1.
|
||||
let wrong_uri = format!(
|
||||
"/api/archives/test/entries/{}/artifacts/1?token={}",
|
||||
entry_uid, token
|
||||
);
|
||||
let response = app_with_state(state)
|
||||
.oneshot(
|
||||
Request::builder()
|
||||
.uri(&wrong_uri)
|
||||
.body(Body::empty())
|
||||
.unwrap(),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(response.status(), StatusCode::UNAUTHORIZED);
|
||||
}
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue