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]>
185 lines
5.7 KiB
Rust
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(())
|
|
}
|
|
}
|