//! Upload jobs: one job per pick (files or a folder), files sent one request //! each so progress, pause and abort work per file. //! //! The job list is a page-global signal (wasm is single-threaded, so a //! `thread_local` is the whole "manager"). The overlay panel and the uploads //! view both render it; the browser starts jobs; the runner below drives //! them. //! //! Pause aborts the in-flight request and restarts that file on resume. //! HTTP cannot pause one request. //! ponytail: per-file restart on resume; chunked resumable uploads if //! multi-GB files matter. use std::cell::{Cell, RefCell}; use std::collections::{HashSet, VecDeque}; use std::rc::Rc; use std::time::Duration; use gloo_timers::future::TimeoutFuture; use leptos::prelude::*; use wasm_bindgen_futures::spawn_local; use crate::api; use crate::i18n; /// What to do when a file turns out to exist after all (the pre-check said /// no, someone created it during the transfer). #[derive(Clone, Copy, PartialEq, Eq, Debug)] pub enum OnConflict { /// Pause the job on that file and ask in the panel. Ask, Overwrite, Skip, } #[derive(Clone, PartialEq, Debug)] pub enum JobState { /// Files picked, existence check in flight. No bytes move yet. Preparing, Running, Paused, /// Waiting for the user's answer about this file (relative path). Conflict(String), Done, Failed(String), Aborted, } impl JobState { pub fn is_active(&self) -> bool { matches!( self, JobState::Preparing | JobState::Running | JobState::Paused | JobState::Conflict(_) ) } } #[derive(Clone, PartialEq)] pub struct Job { pub id: u64, pub root_id: i64, /// Target directory relative to the root ("" = the root). pub dir: String, /// "name.txt" for one file, the folder name for a folder pick, else /// "N files". pub label: String, /// Number of files in the job. The file list itself stays with the /// runner: a folder pick can bring 80 000 entries, and every progress /// update clones this struct for the row views. pub file_count: usize, pub total: f64, /// Bytes of the files that finished (written or skipped). pub done_bytes: f64, /// Bytes sent of the in-flight file. pub cur_loaded: f64, /// Index of the in-flight file. pub current: usize, /// Relative path of the in-flight file. pub current_name: String, pub state: JobState, pub on_conflict: OnConflict, pub skipped: Vec, /// Files that failed for good, with the server's message. pub errors: Vec<(String, String)>, /// (ms since page load, bytes so far), newest last. ~500 ms apart. pub samples: VecDeque<(f64, f64)>, pub started_ms: f64, pub finished_ms: Option, /// Speed graph open (per row). pub expanded: bool, /// Left the overlay (finished a while ago). Still in the uploads view. pub hidden: bool, } impl Job { pub fn sent(&self) -> f64 { self.done_bytes + self.cur_loaded } pub fn percent(&self) -> f64 { if self.total <= 0.0 { return if self.state == JobState::Done { 100.0 } else { 0.0 }; } (self.sent() / self.total * 100.0).clamp(0.0, 100.0) } /// Bytes per second over the last two samples; 0 while not running. pub fn speed(&self) -> f64 { if self.state != JobState::Running || self.samples.len() < 2 { return 0.0; } let n = self.samples.len(); rate(self.samples[n - 2], self.samples[n - 1]) } pub fn is_finished(&self) -> bool { !self.state.is_active() } } /// Bytes per second for each interval between consecutive `(ms, bytes)` /// samples. One shorter than `samples`, empty for fewer than two. pub fn speeds(samples: &[(f64, f64)]) -> Vec { samples.windows(2).map(|w| rate(w[0], w[1])).collect() } fn rate((t0, b0): (f64, f64), (t1, b1): (f64, f64)) -> f64 { if t1 > t0 { (b1 - b0) / (t1 - t0) * 1000.0 } else { 0.0 } } /// Runtime handles the runner and the panel share, outside the signal: the /// in-flight request and the user's pending answers. #[derive(Default)] struct Control { paused: bool, abort: bool, req: Option, /// The answer for the file the job is waiting on. decision: Option, } /// The page's "folder changed" hook (see [`set_on_done`]). type OnDone = Rc; thread_local! { /// Owner of the global signals below. A signal created inside a /// component belongs to that component and dies with it on navigation; /// this root owner lives as long as the page. static ROOT: Owner = new_detached_owner(); static JOBS: RwSignal, LocalStorage> = ROOT.with(|o| o.with(|| RwSignal::new_local(Vec::new()))); static CONTROLS: RefCell>)>> = const { RefCell::new(Vec::new()) }; static NEXT_ID: Cell = const { Cell::new(1) }; /// Overlay collapsed to its header line. static COLLAPSED: RwSignal = ROOT.with(|o| o.with(|| RwSignal::new(false))); /// Called with (root_id, dir) when a job finishes, so the page can /// re-list that folder if it is the one shown. static ON_DONE: RefCell> = const { RefCell::new(None) }; } /// `Owner::new_root` makes itself the current owner as a side effect. This /// runs lazily from inside a component, so put the previous owner back. fn new_detached_owner() -> Owner { let prev = Owner::current(); let root = Owner::new_root(None); match prev { Some(p) => p.set(), None => root.clone().unset(), } root } const SAMPLE_MS: f64 = 500.0; const SAMPLES_KEPT: usize = 60; /// How long a finished job stays in the overlay. const OVERLAY_LINGER: Duration = Duration::from_secs(5); pub fn jobs() -> RwSignal, LocalStorage> { JOBS.with(|j| *j) } pub fn collapsed() -> RwSignal { COLLAPSED.with(|c| *c) } /// Register the page's "folder changed" hook. One per page; the last one /// wins. Returns the token to hand to [`clear_on_done`]. pub fn set_on_done(f: impl Fn(i64, &str) + 'static) -> u64 { let token = next_id(); ON_DONE.with(|h| *h.borrow_mut() = Some((token, Rc::new(f)))); token } /// Drop the hook again. The page that registered it owns the signals the /// closure captures; a job finishing after that page is gone would panic. /// A token that is no longer the current one is ignored: the next page may /// register its hook before this one is cleaned up. pub fn clear_on_done(token: u64) { ON_DONE.with(|h| { let mut h = h.borrow_mut(); if h.as_ref().is_some_and(|(t, _)| *t == token) { *h = None; } }); } pub fn any_active() -> bool { jobs().with(|j| j.iter().any(|j| j.state.is_active())) } /// Next value of the page-wide counter, for job ids and hook tokens. fn next_id() -> u64 { NEXT_ID.with(|n| n.replace(n.get() + 1)) } fn now_ms() -> f64 { web_sys::window() .and_then(|w| w.performance()) .map(|p| p.now()) .unwrap_or(0.0) } fn control(id: u64) -> Option>> { CONTROLS.with(|c| { c.borrow() .iter() .find(|(i, _)| *i == id) .map(|(_, c)| c.clone()) }) } fn update(id: u64, f: impl FnOnce(&mut Job)) { jobs().update(|list| { if let Some(j) = list.iter_mut().find(|j| j.id == id) { f(j); } }); } fn read(id: u64, f: impl FnOnce(&Job) -> T) -> Option { jobs().with_untracked(|list| list.iter().find(|j| j.id == id).map(f)) } // --------------------------------------------------------------------------- // Starting and controlling jobs // --------------------------------------------------------------------------- /// Create a job in the `Preparing` state so the panel shows it at once. /// The existence check and the conflict dialog run in between; then /// [`begin`] starts the transfer, or [`discard`] drops the job again. pub fn prepare(root_id: i64, dir: String, label: String) -> u64 { let id = next_id(); let started_ms = now_ms(); let job = Job { id, root_id, dir, label, file_count: 0, total: 0.0, done_bytes: 0.0, cur_loaded: 0.0, current: 0, current_name: String::new(), state: JobState::Preparing, on_conflict: OnConflict::Ask, skipped: Vec::new(), errors: Vec::new(), samples: VecDeque::from([(started_ms, 0.0)]), started_ms, finished_ms: None, expanded: false, hidden: false, }; CONTROLS.with(|c| { c.borrow_mut() .push((id, Rc::new(RefCell::new(Control::default())))) }); jobs().update(|list| list.push(job)); id } /// Start the transfer of a prepared job. `files` are `(relative path, /// file)`; `errors` are files rejected before the start (a folder sits where /// the file would go). `on_conflict` is the pre-check dialog's choice for /// files that appear during the transfer; `decided` are the files the /// dialog already covered. They are overwritten or skipped per /// `decided_overwrite` without a first attempt. A job the user aborted /// while preparing stays discarded. pub fn begin( id: u64, files: Vec<(String, web_sys::File)>, errors: Vec<(String, String)>, on_conflict: OnConflict, decided: Vec, decided_overwrite: bool, ) { let Some(ctl) = control(id) else { return }; let Some((root_id, dir)) = read(id, |j| (j.root_id, j.dir.clone())) else { return; }; let total = files.iter().map(|(_, f)| f.size()).sum(); let started_ms = now_ms(); update(id, |j| { j.file_count = files.len(); j.total = total; j.state = JobState::Running; j.on_conflict = on_conflict; j.errors = errors; j.samples = VecDeque::from([(started_ms, 0.0)]); j.started_ms = started_ms; }); spawn_local(run( id, root_id, dir, files, ctl, on_conflict, (decided, decided_overwrite), )); spawn_local(sample(id)); } /// Replace the provisional label of a preparing job. pub fn set_label(id: u64, label: String) { update(id, |j| j.label = label); } /// Drop a job that never started (check failed, dialog cancelled). pub fn discard(id: u64) { jobs().update(|list| list.retain(|j| j.id != id)); CONTROLS.with(|c| c.borrow_mut().retain(|(i, _)| *i != id)); } /// Push one speed sample every `SAMPLE_MS` while the job is active. A /// fixed clock, not "per file" or "per progress event": thousands of small /// files would otherwise redraw the graph hundreds of times a second. async fn sample(id: u64) { loop { TimeoutFuture::new(SAMPLE_MS as u32).await; let active = read(id, |j| j.state.is_active()); if active != Some(true) { return; } update(id, |j| { let sent = j.sent(); j.samples.push_back((now_ms(), sent)); if j.samples.len() > SAMPLES_KEPT { j.samples.pop_front(); } }); } } pub fn pause(id: u64) { let Some(ctl) = control(id) else { return }; let mut c = ctl.borrow_mut(); c.paused = true; if let Some(r) = &c.req { let _ = r.abort(); } drop(c); update(id, |j| { if j.state == JobState::Running { j.state = JobState::Paused; j.cur_loaded = 0.0; } }); } pub fn resume(id: u64) { let Some(ctl) = control(id) else { return }; ctl.borrow_mut().paused = false; update(id, |j| { if j.state == JobState::Paused { j.state = JobState::Running; } }); } pub fn abort(id: u64) { if read(id, |j| j.state == JobState::Preparing) == Some(true) { discard(id); return; } let Some(ctl) = control(id) else { return }; let mut c = ctl.borrow_mut(); c.abort = true; c.paused = false; if let Some(r) = &c.req { let _ = r.abort(); } // The runner marks the job Aborted when it notices. } /// Answer the conflict prompt of a job. `apply_all` makes `choice` the /// job's rule for later conflicts too. pub fn decide(id: u64, choice: OnConflict, apply_all: bool) { let Some(ctl) = control(id) else { return }; let mut c = ctl.borrow_mut(); c.decision = Some(choice); // Pausing while the prompt was up left the flag set: the gate blocks // the next batch, so the row must say Paused, not Running. let paused = c.paused; drop(c); update(id, |j| { if apply_all { j.on_conflict = choice; } if matches!(j.state, JobState::Conflict(_)) { j.state = if paused { JobState::Paused } else { JobState::Running }; } }); } pub fn toggle_expanded(id: u64) { update(id, |j| j.expanded = !j.expanded); } /// Remove a finished job from the log. pub fn clear(id: u64) { jobs().update(|list| list.retain(|j| j.id != id || j.state.is_active())); CONTROLS.with(|c| c.borrow_mut().retain(|(i, _)| *i != id)); } pub fn clear_finished() { let ids: Vec = jobs().with_untracked(|list| { list.iter() .filter(|j| j.is_finished()) .map(|j| j.id) .collect() }); for id in ids { clear(id); } } // --------------------------------------------------------------------------- // Runner // --------------------------------------------------------------------------- /// Files per request: files are grouped until the group holds /// `BATCH_BYTES` or `BATCH_FILES`, whichever comes first. A kernel tree /// (80 000 small files) becomes a few hundred requests instead of 80 000. /// A file above the byte cap travels alone. const BATCH_BYTES: f64 = 1024.0 * 1024.0; const BATCH_FILES: usize = 200; /// One request's worth of files, plus the overwrite flag they all share. type Batch = (Vec<(String, web_sys::File)>, bool); /// Cut `items` into groups at both caps. `size` reads an item's bytes. /// Generic over the item so the caps can be tested without a `web_sys::File` /// (which needs a browser). fn chunk(items: Vec, size: impl Fn(&T) -> f64) -> Vec> { let mut out: Vec> = Vec::new(); let mut bytes = 0.0; for it in items { let s = size(&it); let full = out .last() .is_some_and(|b| bytes + s > BATCH_BYTES || b.len() >= BATCH_FILES); if out.is_empty() || full { out.push(Vec::new()); bytes = 0.0; } bytes += s; out.last_mut().expect("pushed above").push(it); } out } /// Group `files` into request batches. `decided` holds the paths the /// pre-check dialog listed. With `decided_overwrite` they go into batches /// that carry the overwrite flag. Otherwise they come back as `(path, size)` /// and are never sent. fn plan_batches( files: Vec<(String, web_sys::File)>, decided: &HashSet, decided_overwrite: bool, ) -> (VecDeque, Vec<(String, f64)>) { let mut skipped = Vec::new(); let (mut overwrite, mut plain) = (Vec::new(), Vec::new()); for (rel, file) in files { match (decided.contains(&rel), decided_overwrite) { (true, false) => skipped.push((rel, file.size())), (true, true) => overwrite.push((rel, file)), (false, _) => plain.push((rel, file)), } } let size = |(_, f): &(String, web_sys::File)| f.size(); let queue = chunk(overwrite, size) .into_iter() .map(|b| (b, true)) .chain(chunk(plain, size).into_iter().map(|b| (b, false))) .collect(); (queue, skipped) } async fn run( id: u64, root_id: i64, dir: String, files: Vec<(String, web_sys::File)>, ctl: Rc>, on_conflict: OnConflict, decided: (Vec, bool), ) { let mut rule = on_conflict; let (decided, decided_overwrite) = decided; let decided: HashSet = decided.into_iter().collect(); let (mut queue, pre_skipped) = plan_batches(files, &decided, decided_overwrite); let mut done_count = pre_skipped.len(); if !pre_skipped.is_empty() { update(id, |j| { for (rel, size) in pre_skipped { j.skipped.push(rel); j.done_bytes += size; } }); } while let Some(batch) = queue.pop_front() { // Pause / abort gate. // ponytail: 200 ms poll, wake channel if latency matters. while ctl.borrow().paused && !ctl.borrow().abort { TimeoutFuture::new(200).await; } if ctl.borrow().abort { finish(id, JobState::Aborted, &dir); return; } let (files, batch_ow) = batch; let total: f64 = files.iter().map(|(_, f)| f.size()).sum(); update(id, |j| { j.current = done_count; j.current_name = files[0].0.clone(); j.cur_loaded = 0.0; }); let on_progress = move |loaded: f64| { update(id, |j| j.cur_loaded = loaded); }; let started = api::start_upload(root_id, &dir, files.clone(), batch_ow, on_progress); let (req, fut) = match started { Ok(x) => x, Err(e) => { finish(id, JobState::Failed(e.to_string()), &dir); return; } }; ctl.borrow_mut().req = Some(req); let res = fut.await; ctl.borrow_mut().req = None; match res { Ok(()) => { done_count += files.len(); update(id, |j| { j.done_bytes += total; j.cur_loaded = 0.0; }); } // Our own abort (pause or stop) cannot be told from a network // failure by the response alone, so the control flags decide. // The gate at the top of the loop handles it; the whole batch // goes again, and it is at most ~1 MiB. Err(_) if ctl.borrow().paused || ctl.borrow().abort => { update(id, |j| j.cur_loaded = 0.0); queue.push_front((files, batch_ow)); } Err(e) if e.status() == Some(409) && e.skipped().is_some() => { // Some files appeared during the transfer. The server wrote // the others; only the listed ones need an answer. let hit: HashSet<&str> = e .skipped() .unwrap_or_default() .iter() .map(String::as_str) .collect(); let (conflicts, written): (Vec<_>, Vec<_>) = files .into_iter() .partition(|(rel, _)| hit.contains(rel.as_str())); done_count += written.len(); let written_bytes: f64 = written.iter().map(|(_, f)| f.size()).sum(); update(id, |j| { j.done_bytes += written_bytes; j.cur_loaded = 0.0; }); let mut overwrite = Vec::new(); for (rel, file) in conflicts { let size = file.size(); if batch_ow { // An overwrite request only ever reports a skip when // a folder sits where the file would go. Sending it // again would loop, so it fails here. done_count += 1; update(id, |j| { j.errors .push((rel, i18n::t(i18n::k::ERR_FOLDER_EXISTS).to_string())); j.done_bytes += size; }); continue; } let choice = match rule { OnConflict::Ask => { update(id, |j| j.state = JobState::Conflict(rel.clone())); // A stalled job must be seen: open the overlay. collapsed().set(false); let choice = wait_decision(&ctl).await; if ctl.borrow().abort { finish(id, JobState::Aborted, &dir); return; } // `decide` may have made this the job's rule. rule = read(id, |j| j.on_conflict).unwrap_or(rule); choice } r => r, }; match choice { OnConflict::Skip => { done_count += 1; update(id, |j| { j.skipped.push(rel); j.done_bytes += size; }); } _ => overwrite.push((rel, file)), } } if !overwrite.is_empty() { queue.push_front((overwrite, true)); } } Err(e) if e.status() == Some(409) && files.len() > 1 => { // The server rejected the whole request with a conflict it // could not attribute. Retry one by one so the error lands // on the file that caused it. // ponytail: files before the culprit were already written; // their single retries come back as "appeared during the // transfer" and may prompt. Rare: the pre-check catches // folders unless one was created mid-upload. let _ = e; for f in files.into_iter().rev() { queue.push_front((vec![f], batch_ow)); } } Err(e) if e.status() == Some(409) => { // A folder sits where the file would go. Per-file // failure; the job goes on. done_count += 1; let (rel, file) = &files[0]; let size = file.size(); update(id, |j| { j.errors.push((rel.clone(), e.to_string())); j.done_bytes += size; j.cur_loaded = 0.0; }); } Err(e) => { finish(id, JobState::Failed(e.to_string()), &dir); return; } } } finish(id, JobState::Done, &dir); } async fn wait_decision(ctl: &Rc>) -> OnConflict { // ponytail: 200 ms poll, wake channel if latency matters. loop { if ctl.borrow().abort { return OnConflict::Skip; } if let Some(d) = ctl.borrow_mut().decision.take() { return d; } TimeoutFuture::new(200).await; } } fn finish(id: u64, state: JobState, dir: &str) { let root_id = read(id, |j| j.root_id); update(id, |j| { j.state = state; j.cur_loaded = 0.0; j.finished_ms = Some(now_ms()); }); if let Some(root_id) = root_id { ON_DONE.with(|h| { if let Some((_, f)) = h.borrow().as_ref() { f(root_id, dir); } }); } spawn_local(async move { gloo_timers::future::sleep(OVERLAY_LINGER).await; update(id, |j| j.hidden = true); }); } /// "12.3 MB/s". pub fn format_speed(bytes_per_s: f64) -> String { i18n::t_fmt( i18n::k::UPLOAD_SPEED, &crate::util::format_size(bytes_per_s.max(0.0) as u64), ) } #[cfg(test)] mod tests { use super::{BATCH_BYTES, BATCH_FILES, chunk}; /// Fake sizes: a real `web_sys::File` needs a browser. #[test] fn chunk_respects_both_caps() { let sizes = |v: Vec| chunk(v, |s: &f64| *s); let lens = |b: Vec>| b.iter().map(Vec::len).collect::>(); // The byte cap cuts before the group would exceed it. assert_eq!( lens(sizes(vec![BATCH_BYTES * 0.6, BATCH_BYTES * 0.6, 1.0])), vec![1, 2] ); // A file above the cap travels alone. assert_eq!( lens(sizes(vec![1.0, BATCH_BYTES * 3.0, 1.0])), vec![1, 1, 1] ); // The file cap cuts even when the bytes are nothing. assert_eq!( lens(sizes(vec![0.0; BATCH_FILES * 2 + 1])), vec![BATCH_FILES, BATCH_FILES, 1] ); assert!(sizes(Vec::new()).is_empty()); } }