Files
OpenPXE/crates/http-api/src/uploads.rs
T
Miles WardandClaude Opus 4.8 674a69f93b v0.5.4: code-cleanup pass (AppError, figment config, encoding dedup, typed status, deps)
Final cleanup before hardware testing. No behaviour changes; 248 tests green,
clippy clean.

#1  AppError newtype (http-api/src/error.rs) with one IntoResponse mapping
    (NotFound→404, Invalid→400, _→500) + From<core::Error>/From<io::Error>.
    Converted the clearly-safe handlers (sso_put, unattended_upload,
    branding_clear) to `?`; intentionally left handlers with bespoke
    status semantics (Invalid→404 on category, 409 on duplicate share /
    open upload) explicit so no asserted status changes.
#2  figment-based Config::load (defaults → TOML → env). Keeps the historical
    flat OPENPXE_* names (Unraid/entrypoint compatible) AND adds the nested
    OPENPXE_SECTION__FIELD form; now covers every field (apply_env had
    silently skipped unattended_dir + bind addrs). 6 Jail tests prove
    backward-compat. Removed the hand-rolled apply_env.
#3  thiserror 1→2; dropped unused mime/mime_guess/once_cell deps.
#4  Re-evaluated: Duration::from_hours/from_mins are stable on the pinned
    1.95 toolchain and clippy prefers them — kept the readable form
    (the "unstable" premise didn't hold; MSRV is intentionally 1.95).
#5  insta snapshot of the rendered iPXE menu (version-filtered) + wiremock
    coverage of the SAML metadata-URL fetch (200 + non-2xx).
#6  api_status → typed StatusResponse struct (was a 25-key json! blob) with
    a full_flow guard test asserting every UI key + the started_at string
    shape. Deferred the /api/docs typed conversion (lowest value, highest
    churn, zero functional benefit).
#7  pct_encode/xml_escape de-duplicated into openpxe_core::encoding (were
    copied across app.rs + the SAML modules). No new crates.
#8  UploadSessions registry → parking_lot::RwLock (sync, never held across
    .await); per-session lock stays tokio::Mutex.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-06-03 03:33:05 -04:00

185 lines
5.7 KiB
Rust

//! Chunked upload sessions for browser-driven ISO uploads.
//!
//! The legacy multipart endpoint still exists for simple API clients, but
//! browsers get a better failure mode with raw chunks: progress advances after
//! each acknowledged write, partial files appear in the ISO directory
//! immediately, and reverse proxies are less likely to buffer an entire DVD
//! image before OpenPXE sees byte one.
use bytes::Bytes;
use openpxe_core::{Error, Result};
use openpxe_iso_store::{IsoMeta, IsoStore, UploadHandle};
use parking_lot::RwLock;
use serde::Serialize;
use std::collections::HashMap;
use std::sync::Arc;
use tokio::sync::Mutex;
use uuid::Uuid;
const DEFAULT_CHUNK_SIZE: u64 = 8 * 1024 * 1024;
#[derive(Clone, Default)]
pub struct UploadSessions {
// v0.5.4: the registry is a sync `parking_lot::RwLock` — it's only ever
// briefly read/inserted/removed to look up a session, never held across
// an `.await`. The per-session lock below stays a `tokio::sync::Mutex`
// because `write_chunk` / `finish` are awaited while it's held.
inner: Arc<RwLock<HashMap<String, Arc<Mutex<UploadSession>>>>>,
}
struct UploadSession {
filename: String,
expected_size: Option<u64>,
offset: u64,
handle: Option<UploadHandle>,
}
#[derive(Debug, Clone, Serialize)]
pub struct UploadStarted {
pub upload_id: String,
pub iso_id: String,
pub filename: String,
pub offset: u64,
pub chunk_size: u64,
}
#[derive(Debug, Clone, Serialize)]
#[serde(rename_all = "snake_case")]
pub enum UploadAppend {
Progress { offset: u64 },
Complete { offset: u64, iso: Box<IsoMeta> },
}
impl UploadSessions {
pub async fn begin(
&self,
store: &IsoStore,
filename: &str,
expected_size: Option<u64>,
) -> Result<UploadStarted> {
if !filename.to_ascii_lowercase().ends_with(".iso") {
return Err(Error::Invalid("only .iso uploads accepted".to_string()));
}
let handle = store.begin_upload(filename).await?;
let iso_id = handle.id.clone();
let upload_id = Uuid::new_v4().to_string();
let session = UploadSession {
filename: filename.to_string(),
expected_size,
offset: 0,
handle: Some(handle),
};
self.inner
.write()
.insert(upload_id.clone(), Arc::new(Mutex::new(session)));
Ok(UploadStarted {
upload_id,
iso_id,
filename: filename.to_string(),
offset: 0,
chunk_size: DEFAULT_CHUNK_SIZE,
})
}
pub async fn append(
&self,
store: &IsoStore,
upload_id: &str,
offset: u64,
chunk: Bytes,
complete: bool,
) -> Result<UploadAppend> {
let Some(session_lock) = self.inner.read().get(upload_id).cloned() else {
return Err(Error::Invalid(format!("no such upload '{upload_id}'")));
};
let mut session = session_lock.lock().await;
if session.offset != offset {
return Err(Error::Invalid(format!(
"expected offset {}, got {offset}",
session.offset
)));
}
let new_offset = session
.offset
.checked_add(chunk.len() as u64)
.ok_or_else(|| Error::Invalid("upload offset overflow".to_string()))?;
if let Some(expected) = session.expected_size {
if new_offset > expected {
return Err(Error::Invalid(format!(
"chunk exceeds declared upload size {expected}"
)));
}
}
let Some(handle) = session.handle.as_mut() else {
return Err(Error::Invalid("upload already completed".to_string()));
};
if let Err(e) = handle.write_chunk(&chunk).await {
let handle = session.handle.take();
drop(session);
self.inner.write().remove(upload_id);
if let Some(handle) = handle {
let _ = handle.abort().await;
}
return Err(e);
}
session.offset = new_offset;
if !complete {
return Ok(UploadAppend::Progress { offset: new_offset });
}
if let Some(expected) = session.expected_size {
if new_offset != expected {
return Err(Error::Invalid(format!(
"final chunk ended at {new_offset}, expected {expected}"
)));
}
}
let Some(handle) = session.handle.take() else {
return Err(Error::Invalid("upload already completed".to_string()));
};
let filename = session.filename.clone();
drop(session);
tracing::info!(
target: "openpxe::http::upload",
upload_id,
filename = %filename,
received_bytes = new_offset,
"chunked upload body complete; introspecting"
);
let meta = match handle.finish(store).await {
Ok(meta) => meta,
Err(e) => {
self.inner.write().remove(upload_id);
return Err(e);
}
};
self.inner.write().remove(upload_id);
Ok(UploadAppend::Complete {
offset: new_offset,
iso: Box::new(meta),
})
}
pub async fn abort(&self, upload_id: &str) -> Result<()> {
let Some(session_lock) = self.inner.write().remove(upload_id) else {
return Err(Error::Invalid(format!("no such upload '{upload_id}'")));
};
let mut session = session_lock.lock().await;
if let Some(handle) = session.handle.take() {
handle.abort().await?;
}
Ok(())
}
}