//! 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 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 { inner: Arc>>>>, } struct UploadSession { filename: String, expected_size: Option, offset: u64, handle: Option, } #[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 }, } impl UploadSessions { pub async fn begin( &self, store: &IsoStore, filename: &str, expected_size: Option, ) -> Result { 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 .lock() .await .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 { let Some(session_lock) = self.inner.lock().await.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.lock().await.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.lock().await.remove(upload_id); return Err(e); } }; self.inner.lock().await.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.lock().await.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(()) } }