uploads.rs
| 1 | //! Upload jobs: one job per pick (files or a folder), files sent one request |
| 2 | //! each so progress, pause and abort work per file. |
| 3 | //! |
| 4 | //! The job list is a page-global signal (wasm is single-threaded, so a |
| 5 | //! `thread_local` is the whole "manager"). The overlay panel and the uploads |
| 6 | //! view both render it; the browser starts jobs; the runner below drives |
| 7 | //! them. |
| 8 | //! |
| 9 | //! Pause aborts the in-flight request and restarts that file on resume. |
| 10 | //! HTTP cannot pause one request. |
| 11 | //! ponytail: per-file restart on resume; chunked resumable uploads if |
| 12 | //! multi-GB files matter. |
| 13 | |
| 14 | use std::cell::{Cell, RefCell}; |
| 15 | use std::collections::{HashSet, VecDeque}; |
| 16 | use std::rc::Rc; |
| 17 | use std::time::Duration; |
| 18 | |
| 19 | use gloo_timers::future::TimeoutFuture; |
| 20 | use leptos::prelude::*; |
| 21 | use wasm_bindgen_futures::spawn_local; |
| 22 | |
| 23 | use crate::api; |
| 24 | use crate::i18n; |
| 25 | |
| 26 | /// What to do when a file turns out to exist after all (the pre-check said |
| 27 | /// no, someone created it during the transfer). |
| 28 | #[derive(Clone, Copy, PartialEq, Eq, Debug)] |
| 29 | pub enum OnConflict { |
| 30 | /// Pause the job on that file and ask in the panel. |
| 31 | Ask, |
| 32 | Overwrite, |
| 33 | Skip, |
| 34 | } |
| 35 | |
| 36 | #[derive(Clone, PartialEq, Debug)] |
| 37 | pub enum JobState { |
| 38 | /// Files picked, existence check in flight. No bytes move yet. |
| 39 | Preparing, |
| 40 | Running, |
| 41 | Paused, |
| 42 | /// Waiting for the user's answer about this file (relative path). |
| 43 | Conflict(String), |
| 44 | Done, |
| 45 | Failed(String), |
| 46 | Aborted, |
| 47 | } |
| 48 | |
| 49 | impl JobState { |
| 50 | pub fn is_active(&self) -> bool { |
| 51 | matches!( |
| 52 | self, |
| 53 | JobState::Preparing | JobState::Running | JobState::Paused | JobState::Conflict(_) |
| 54 | ) |
| 55 | } |
| 56 | } |
| 57 | |
| 58 | #[derive(Clone, PartialEq)] |
| 59 | pub struct Job { |
| 60 | pub id: u64, |
| 61 | pub root_id: i64, |
| 62 | /// Target directory relative to the root ("" = the root). |
| 63 | pub dir: String, |
| 64 | /// "name.txt" for one file, the folder name for a folder pick, else |
| 65 | /// "N files". |
| 66 | pub label: String, |
| 67 | /// Number of files in the job. The file list itself stays with the |
| 68 | /// runner: a folder pick can bring 80 000 entries, and every progress |
| 69 | /// update clones this struct for the row views. |
| 70 | pub file_count: usize, |
| 71 | pub total: f64, |
| 72 | /// Bytes of the files that finished (written or skipped). |
| 73 | pub done_bytes: f64, |
| 74 | /// Bytes sent of the in-flight file. |
| 75 | pub cur_loaded: f64, |
| 76 | /// Index of the in-flight file. |
| 77 | pub current: usize, |
| 78 | /// Relative path of the in-flight file. |
| 79 | pub current_name: String, |
| 80 | pub state: JobState, |
| 81 | pub on_conflict: OnConflict, |
| 82 | pub skipped: Vec<String>, |
| 83 | /// Files that failed for good, with the server's message. |
| 84 | pub errors: Vec<(String, String)>, |
| 85 | /// (ms since page load, bytes so far), newest last. ~500 ms apart. |
| 86 | pub samples: VecDeque<(f64, f64)>, |
| 87 | pub started_ms: f64, |
| 88 | pub finished_ms: Option<f64>, |
| 89 | /// Speed graph open (per row). |
| 90 | pub expanded: bool, |
| 91 | /// Left the overlay (finished a while ago). Still in the uploads view. |
| 92 | pub hidden: bool, |
| 93 | } |
| 94 | |
| 95 | impl Job { |
| 96 | pub fn sent(&self) -> f64 { |
| 97 | self.done_bytes + self.cur_loaded |
| 98 | } |
| 99 | |
| 100 | pub fn percent(&self) -> f64 { |
| 101 | if self.total <= 0.0 { |
| 102 | return if self.state == JobState::Done { |
| 103 | 100.0 |
| 104 | } else { |
| 105 | 0.0 |
| 106 | }; |
| 107 | } |
| 108 | (self.sent() / self.total * 100.0).clamp(0.0, 100.0) |
| 109 | } |
| 110 | |
| 111 | /// Bytes per second over the last two samples; 0 while not running. |
| 112 | pub fn speed(&self) -> f64 { |
| 113 | if self.state != JobState::Running || self.samples.len() < 2 { |
| 114 | return 0.0; |
| 115 | } |
| 116 | let n = self.samples.len(); |
| 117 | rate(self.samples[n - 2], self.samples[n - 1]) |
| 118 | } |
| 119 | |
| 120 | pub fn is_finished(&self) -> bool { |
| 121 | !self.state.is_active() |
| 122 | } |
| 123 | } |
| 124 | |
| 125 | /// Bytes per second for each interval between consecutive `(ms, bytes)` |
| 126 | /// samples. One shorter than `samples`, empty for fewer than two. |
| 127 | pub fn speeds(samples: &[(f64, f64)]) -> Vec<f64> { |
| 128 | samples.windows(2).map(|w| rate(w[0], w[1])).collect() |
| 129 | } |
| 130 | |
| 131 | fn rate((t0, b0): (f64, f64), (t1, b1): (f64, f64)) -> f64 { |
| 132 | if t1 > t0 { |
| 133 | (b1 - b0) / (t1 - t0) * 1000.0 |
| 134 | } else { |
| 135 | 0.0 |
| 136 | } |
| 137 | } |
| 138 | |
| 139 | /// Runtime handles the runner and the panel share, outside the signal: the |
| 140 | /// in-flight request and the user's pending answers. |
| 141 | #[derive(Default)] |
| 142 | struct Control { |
| 143 | paused: bool, |
| 144 | abort: bool, |
| 145 | req: Option<web_sys::XmlHttpRequest>, |
| 146 | /// The answer for the file the job is waiting on. |
| 147 | decision: Option<OnConflict>, |
| 148 | } |
| 149 | |
| 150 | /// The page's "folder changed" hook (see [`set_on_done`]). |
| 151 | type OnDone = Rc<dyn Fn(i64, &str)>; |
| 152 | |
| 153 | thread_local! { |
| 154 | /// Owner of the global signals below. A signal created inside a |
| 155 | /// component belongs to that component and dies with it on navigation; |
| 156 | /// this root owner lives as long as the page. |
| 157 | static ROOT: Owner = new_detached_owner(); |
| 158 | static JOBS: RwSignal<Vec<Job>, LocalStorage> = |
| 159 | ROOT.with(|o| o.with(|| RwSignal::new_local(Vec::new()))); |
| 160 | static CONTROLS: RefCell<Vec<(u64, Rc<RefCell<Control>>)>> = const { RefCell::new(Vec::new()) }; |
| 161 | static NEXT_ID: Cell<u64> = const { Cell::new(1) }; |
| 162 | /// Overlay collapsed to its header line. |
| 163 | static COLLAPSED: RwSignal<bool> = ROOT.with(|o| o.with(|| RwSignal::new(false))); |
| 164 | /// Called with (root_id, dir) when a job finishes, so the page can |
| 165 | /// re-list that folder if it is the one shown. |
| 166 | static ON_DONE: RefCell<Option<(u64, OnDone)>> = const { RefCell::new(None) }; |
| 167 | } |
| 168 | |
| 169 | /// `Owner::new_root` makes itself the current owner as a side effect. This |
| 170 | /// runs lazily from inside a component, so put the previous owner back. |
| 171 | fn new_detached_owner() -> Owner { |
| 172 | let prev = Owner::current(); |
| 173 | let root = Owner::new_root(None); |
| 174 | match prev { |
| 175 | Some(p) => p.set(), |
| 176 | None => root.clone().unset(), |
| 177 | } |
| 178 | root |
| 179 | } |
| 180 | |
| 181 | const SAMPLE_MS: f64 = 500.0; |
| 182 | const SAMPLES_KEPT: usize = 60; |
| 183 | /// How long a finished job stays in the overlay. |
| 184 | const OVERLAY_LINGER: Duration = Duration::from_secs(5); |
| 185 | |
| 186 | pub fn jobs() -> RwSignal<Vec<Job>, LocalStorage> { |
| 187 | JOBS.with(|j| *j) |
| 188 | } |
| 189 | |
| 190 | pub fn collapsed() -> RwSignal<bool> { |
| 191 | COLLAPSED.with(|c| *c) |
| 192 | } |
| 193 | |
| 194 | /// Register the page's "folder changed" hook. One per page; the last one |
| 195 | /// wins. Returns the token to hand to [`clear_on_done`]. |
| 196 | pub fn set_on_done(f: impl Fn(i64, &str) + 'static) -> u64 { |
| 197 | let token = next_id(); |
| 198 | ON_DONE.with(|h| *h.borrow_mut() = Some((token, Rc::new(f)))); |
| 199 | token |
| 200 | } |
| 201 | |
| 202 | /// Drop the hook again. The page that registered it owns the signals the |
| 203 | /// closure captures; a job finishing after that page is gone would panic. |
| 204 | /// A token that is no longer the current one is ignored: the next page may |
| 205 | /// register its hook before this one is cleaned up. |
| 206 | pub fn clear_on_done(token: u64) { |
| 207 | ON_DONE.with(|h| { |
| 208 | let mut h = h.borrow_mut(); |
| 209 | if h.as_ref().is_some_and(|(t, _)| *t == token) { |
| 210 | *h = None; |
| 211 | } |
| 212 | }); |
| 213 | } |
| 214 | |
| 215 | pub fn any_active() -> bool { |
| 216 | jobs().with(|j| j.iter().any(|j| j.state.is_active())) |
| 217 | } |
| 218 | |
| 219 | /// Next value of the page-wide counter, for job ids and hook tokens. |
| 220 | fn next_id() -> u64 { |
| 221 | NEXT_ID.with(|n| n.replace(n.get() + 1)) |
| 222 | } |
| 223 | |
| 224 | fn now_ms() -> f64 { |
| 225 | web_sys::window() |
| 226 | .and_then(|w| w.performance()) |
| 227 | .map(|p| p.now()) |
| 228 | .unwrap_or(0.0) |
| 229 | } |
| 230 | |
| 231 | fn control(id: u64) -> Option<Rc<RefCell<Control>>> { |
| 232 | CONTROLS.with(|c| { |
| 233 | c.borrow() |
| 234 | .iter() |
| 235 | .find(|(i, _)| *i == id) |
| 236 | .map(|(_, c)| c.clone()) |
| 237 | }) |
| 238 | } |
| 239 | |
| 240 | fn update(id: u64, f: impl FnOnce(&mut Job)) { |
| 241 | jobs().update(|list| { |
| 242 | if let Some(j) = list.iter_mut().find(|j| j.id == id) { |
| 243 | f(j); |
| 244 | } |
| 245 | }); |
| 246 | } |
| 247 | |
| 248 | fn read<T>(id: u64, f: impl FnOnce(&Job) -> T) -> Option<T> { |
| 249 | jobs().with_untracked(|list| list.iter().find(|j| j.id == id).map(f)) |
| 250 | } |
| 251 | |
| 252 | // --------------------------------------------------------------------------- |
| 253 | // Starting and controlling jobs |
| 254 | // --------------------------------------------------------------------------- |
| 255 | |
| 256 | /// Create a job in the `Preparing` state so the panel shows it at once. |
| 257 | /// The existence check and the conflict dialog run in between; then |
| 258 | /// [`begin`] starts the transfer, or [`discard`] drops the job again. |
| 259 | pub fn prepare(root_id: i64, dir: String, label: String) -> u64 { |
| 260 | let id = next_id(); |
| 261 | let started_ms = now_ms(); |
| 262 | let job = Job { |
| 263 | id, |
| 264 | root_id, |
| 265 | dir, |
| 266 | label, |
| 267 | file_count: 0, |
| 268 | total: 0.0, |
| 269 | done_bytes: 0.0, |
| 270 | cur_loaded: 0.0, |
| 271 | current: 0, |
| 272 | current_name: String::new(), |
| 273 | state: JobState::Preparing, |
| 274 | on_conflict: OnConflict::Ask, |
| 275 | skipped: Vec::new(), |
| 276 | errors: Vec::new(), |
| 277 | samples: VecDeque::from([(started_ms, 0.0)]), |
| 278 | started_ms, |
| 279 | finished_ms: None, |
| 280 | expanded: false, |
| 281 | hidden: false, |
| 282 | }; |
| 283 | CONTROLS.with(|c| { |
| 284 | c.borrow_mut() |
| 285 | .push((id, Rc::new(RefCell::new(Control::default())))) |
| 286 | }); |
| 287 | jobs().update(|list| list.push(job)); |
| 288 | id |
| 289 | } |
| 290 | |
| 291 | /// Start the transfer of a prepared job. `files` are `(relative path, |
| 292 | /// file)`; `errors` are files rejected before the start (a folder sits where |
| 293 | /// the file would go). `on_conflict` is the pre-check dialog's choice for |
| 294 | /// files that appear during the transfer; `decided` are the files the |
| 295 | /// dialog already covered. They are overwritten or skipped per |
| 296 | /// `decided_overwrite` without a first attempt. A job the user aborted |
| 297 | /// while preparing stays discarded. |
| 298 | pub fn begin( |
| 299 | id: u64, |
| 300 | files: Vec<(String, web_sys::File)>, |
| 301 | errors: Vec<(String, String)>, |
| 302 | on_conflict: OnConflict, |
| 303 | decided: Vec<String>, |
| 304 | decided_overwrite: bool, |
| 305 | ) { |
| 306 | let Some(ctl) = control(id) else { return }; |
| 307 | let Some((root_id, dir)) = read(id, |j| (j.root_id, j.dir.clone())) else { |
| 308 | return; |
| 309 | }; |
| 310 | let total = files.iter().map(|(_, f)| f.size()).sum(); |
| 311 | let started_ms = now_ms(); |
| 312 | update(id, |j| { |
| 313 | j.file_count = files.len(); |
| 314 | j.total = total; |
| 315 | j.state = JobState::Running; |
| 316 | j.on_conflict = on_conflict; |
| 317 | j.errors = errors; |
| 318 | j.samples = VecDeque::from([(started_ms, 0.0)]); |
| 319 | j.started_ms = started_ms; |
| 320 | }); |
| 321 | spawn_local(run( |
| 322 | id, |
| 323 | root_id, |
| 324 | dir, |
| 325 | files, |
| 326 | ctl, |
| 327 | on_conflict, |
| 328 | (decided, decided_overwrite), |
| 329 | )); |
| 330 | spawn_local(sample(id)); |
| 331 | } |
| 332 | |
| 333 | /// Replace the provisional label of a preparing job. |
| 334 | pub fn set_label(id: u64, label: String) { |
| 335 | update(id, |j| j.label = label); |
| 336 | } |
| 337 | |
| 338 | /// Drop a job that never started (check failed, dialog cancelled). |
| 339 | pub fn discard(id: u64) { |
| 340 | jobs().update(|list| list.retain(|j| j.id != id)); |
| 341 | CONTROLS.with(|c| c.borrow_mut().retain(|(i, _)| *i != id)); |
| 342 | } |
| 343 | |
| 344 | /// Push one speed sample every `SAMPLE_MS` while the job is active. A |
| 345 | /// fixed clock, not "per file" or "per progress event": thousands of small |
| 346 | /// files would otherwise redraw the graph hundreds of times a second. |
| 347 | async fn sample(id: u64) { |
| 348 | loop { |
| 349 | TimeoutFuture::new(SAMPLE_MS as u32).await; |
| 350 | let active = read(id, |j| j.state.is_active()); |
| 351 | if active != Some(true) { |
| 352 | return; |
| 353 | } |
| 354 | update(id, |j| { |
| 355 | let sent = j.sent(); |
| 356 | j.samples.push_back((now_ms(), sent)); |
| 357 | if j.samples.len() > SAMPLES_KEPT { |
| 358 | j.samples.pop_front(); |
| 359 | } |
| 360 | }); |
| 361 | } |
| 362 | } |
| 363 | |
| 364 | pub fn pause(id: u64) { |
| 365 | let Some(ctl) = control(id) else { return }; |
| 366 | let mut c = ctl.borrow_mut(); |
| 367 | c.paused = true; |
| 368 | if let Some(r) = &c.req { |
| 369 | let _ = r.abort(); |
| 370 | } |
| 371 | drop(c); |
| 372 | update(id, |j| { |
| 373 | if j.state == JobState::Running { |
| 374 | j.state = JobState::Paused; |
| 375 | j.cur_loaded = 0.0; |
| 376 | } |
| 377 | }); |
| 378 | } |
| 379 | |
| 380 | pub fn resume(id: u64) { |
| 381 | let Some(ctl) = control(id) else { return }; |
| 382 | ctl.borrow_mut().paused = false; |
| 383 | update(id, |j| { |
| 384 | if j.state == JobState::Paused { |
| 385 | j.state = JobState::Running; |
| 386 | } |
| 387 | }); |
| 388 | } |
| 389 | |
| 390 | pub fn abort(id: u64) { |
| 391 | if read(id, |j| j.state == JobState::Preparing) == Some(true) { |
| 392 | discard(id); |
| 393 | return; |
| 394 | } |
| 395 | let Some(ctl) = control(id) else { return }; |
| 396 | let mut c = ctl.borrow_mut(); |
| 397 | c.abort = true; |
| 398 | c.paused = false; |
| 399 | if let Some(r) = &c.req { |
| 400 | let _ = r.abort(); |
| 401 | } |
| 402 | // The runner marks the job Aborted when it notices. |
| 403 | } |
| 404 | |
| 405 | /// Answer the conflict prompt of a job. `apply_all` makes `choice` the |
| 406 | /// job's rule for later conflicts too. |
| 407 | pub fn decide(id: u64, choice: OnConflict, apply_all: bool) { |
| 408 | let Some(ctl) = control(id) else { return }; |
| 409 | let mut c = ctl.borrow_mut(); |
| 410 | c.decision = Some(choice); |
| 411 | // Pausing while the prompt was up left the flag set: the gate blocks |
| 412 | // the next batch, so the row must say Paused, not Running. |
| 413 | let paused = c.paused; |
| 414 | drop(c); |
| 415 | update(id, |j| { |
| 416 | if apply_all { |
| 417 | j.on_conflict = choice; |
| 418 | } |
| 419 | if matches!(j.state, JobState::Conflict(_)) { |
| 420 | j.state = if paused { |
| 421 | JobState::Paused |
| 422 | } else { |
| 423 | JobState::Running |
| 424 | }; |
| 425 | } |
| 426 | }); |
| 427 | } |
| 428 | |
| 429 | pub fn toggle_expanded(id: u64) { |
| 430 | update(id, |j| j.expanded = !j.expanded); |
| 431 | } |
| 432 | |
| 433 | /// Remove a finished job from the log. |
| 434 | pub fn clear(id: u64) { |
| 435 | jobs().update(|list| list.retain(|j| j.id != id || j.state.is_active())); |
| 436 | CONTROLS.with(|c| c.borrow_mut().retain(|(i, _)| *i != id)); |
| 437 | } |
| 438 | |
| 439 | pub fn clear_finished() { |
| 440 | let ids: Vec<u64> = jobs().with_untracked(|list| { |
| 441 | list.iter() |
| 442 | .filter(|j| j.is_finished()) |
| 443 | .map(|j| j.id) |
| 444 | .collect() |
| 445 | }); |
| 446 | for id in ids { |
| 447 | clear(id); |
| 448 | } |
| 449 | } |
| 450 | |
| 451 | // --------------------------------------------------------------------------- |
| 452 | // Runner |
| 453 | // --------------------------------------------------------------------------- |
| 454 | |
| 455 | /// Files per request: files are grouped until the group holds |
| 456 | /// `BATCH_BYTES` or `BATCH_FILES`, whichever comes first. A kernel tree |
| 457 | /// (80 000 small files) becomes a few hundred requests instead of 80 000. |
| 458 | /// A file above the byte cap travels alone. |
| 459 | const BATCH_BYTES: f64 = 1024.0 * 1024.0; |
| 460 | const BATCH_FILES: usize = 200; |
| 461 | |
| 462 | /// One request's worth of files, plus the overwrite flag they all share. |
| 463 | type Batch = (Vec<(String, web_sys::File)>, bool); |
| 464 | |
| 465 | /// Cut `items` into groups at both caps. `size` reads an item's bytes. |
| 466 | /// Generic over the item so the caps can be tested without a `web_sys::File` |
| 467 | /// (which needs a browser). |
| 468 | fn chunk<T>(items: Vec<T>, size: impl Fn(&T) -> f64) -> Vec<Vec<T>> { |
| 469 | let mut out: Vec<Vec<T>> = Vec::new(); |
| 470 | let mut bytes = 0.0; |
| 471 | for it in items { |
| 472 | let s = size(&it); |
| 473 | let full = out |
| 474 | .last() |
| 475 | .is_some_and(|b| bytes + s > BATCH_BYTES || b.len() >= BATCH_FILES); |
| 476 | if out.is_empty() || full { |
| 477 | out.push(Vec::new()); |
| 478 | bytes = 0.0; |
| 479 | } |
| 480 | bytes += s; |
| 481 | out.last_mut().expect("pushed above").push(it); |
| 482 | } |
| 483 | out |
| 484 | } |
| 485 | |
| 486 | /// Group `files` into request batches. `decided` holds the paths the |
| 487 | /// pre-check dialog listed. With `decided_overwrite` they go into batches |
| 488 | /// that carry the overwrite flag. Otherwise they come back as `(path, size)` |
| 489 | /// and are never sent. |
| 490 | fn plan_batches( |
| 491 | files: Vec<(String, web_sys::File)>, |
| 492 | decided: &HashSet<String>, |
| 493 | decided_overwrite: bool, |
| 494 | ) -> (VecDeque<Batch>, Vec<(String, f64)>) { |
| 495 | let mut skipped = Vec::new(); |
| 496 | let (mut overwrite, mut plain) = (Vec::new(), Vec::new()); |
| 497 | for (rel, file) in files { |
| 498 | match (decided.contains(&rel), decided_overwrite) { |
| 499 | (true, false) => skipped.push((rel, file.size())), |
| 500 | (true, true) => overwrite.push((rel, file)), |
| 501 | (false, _) => plain.push((rel, file)), |
| 502 | } |
| 503 | } |
| 504 | let size = |(_, f): &(String, web_sys::File)| f.size(); |
| 505 | let queue = chunk(overwrite, size) |
| 506 | .into_iter() |
| 507 | .map(|b| (b, true)) |
| 508 | .chain(chunk(plain, size).into_iter().map(|b| (b, false))) |
| 509 | .collect(); |
| 510 | (queue, skipped) |
| 511 | } |
| 512 | |
| 513 | async fn run( |
| 514 | id: u64, |
| 515 | root_id: i64, |
| 516 | dir: String, |
| 517 | files: Vec<(String, web_sys::File)>, |
| 518 | ctl: Rc<RefCell<Control>>, |
| 519 | on_conflict: OnConflict, |
| 520 | decided: (Vec<String>, bool), |
| 521 | ) { |
| 522 | let mut rule = on_conflict; |
| 523 | let (decided, decided_overwrite) = decided; |
| 524 | let decided: HashSet<String> = decided.into_iter().collect(); |
| 525 | let (mut queue, pre_skipped) = plan_batches(files, &decided, decided_overwrite); |
| 526 | let mut done_count = pre_skipped.len(); |
| 527 | if !pre_skipped.is_empty() { |
| 528 | update(id, |j| { |
| 529 | for (rel, size) in pre_skipped { |
| 530 | j.skipped.push(rel); |
| 531 | j.done_bytes += size; |
| 532 | } |
| 533 | }); |
| 534 | } |
| 535 | while let Some(batch) = queue.pop_front() { |
| 536 | // Pause / abort gate. |
| 537 | // ponytail: 200 ms poll, wake channel if latency matters. |
| 538 | while ctl.borrow().paused && !ctl.borrow().abort { |
| 539 | TimeoutFuture::new(200).await; |
| 540 | } |
| 541 | if ctl.borrow().abort { |
| 542 | finish(id, JobState::Aborted, &dir); |
| 543 | return; |
| 544 | } |
| 545 | let (files, batch_ow) = batch; |
| 546 | let total: f64 = files.iter().map(|(_, f)| f.size()).sum(); |
| 547 | update(id, |j| { |
| 548 | j.current = done_count; |
| 549 | j.current_name = files[0].0.clone(); |
| 550 | j.cur_loaded = 0.0; |
| 551 | }); |
| 552 | let on_progress = move |loaded: f64| { |
| 553 | update(id, |j| j.cur_loaded = loaded); |
| 554 | }; |
| 555 | let started = api::start_upload(root_id, &dir, files.clone(), batch_ow, on_progress); |
| 556 | let (req, fut) = match started { |
| 557 | Ok(x) => x, |
| 558 | Err(e) => { |
| 559 | finish(id, JobState::Failed(e.to_string()), &dir); |
| 560 | return; |
| 561 | } |
| 562 | }; |
| 563 | ctl.borrow_mut().req = Some(req); |
| 564 | let res = fut.await; |
| 565 | ctl.borrow_mut().req = None; |
| 566 | match res { |
| 567 | Ok(()) => { |
| 568 | done_count += files.len(); |
| 569 | update(id, |j| { |
| 570 | j.done_bytes += total; |
| 571 | j.cur_loaded = 0.0; |
| 572 | }); |
| 573 | } |
| 574 | // Our own abort (pause or stop) cannot be told from a network |
| 575 | // failure by the response alone, so the control flags decide. |
| 576 | // The gate at the top of the loop handles it; the whole batch |
| 577 | // goes again, and it is at most ~1 MiB. |
| 578 | Err(_) if ctl.borrow().paused || ctl.borrow().abort => { |
| 579 | update(id, |j| j.cur_loaded = 0.0); |
| 580 | queue.push_front((files, batch_ow)); |
| 581 | } |
| 582 | Err(e) if e.status() == Some(409) && e.skipped().is_some() => { |
| 583 | // Some files appeared during the transfer. The server wrote |
| 584 | // the others; only the listed ones need an answer. |
| 585 | let hit: HashSet<&str> = e |
| 586 | .skipped() |
| 587 | .unwrap_or_default() |
| 588 | .iter() |
| 589 | .map(String::as_str) |
| 590 | .collect(); |
| 591 | let (conflicts, written): (Vec<_>, Vec<_>) = files |
| 592 | .into_iter() |
| 593 | .partition(|(rel, _)| hit.contains(rel.as_str())); |
| 594 | done_count += written.len(); |
| 595 | let written_bytes: f64 = written.iter().map(|(_, f)| f.size()).sum(); |
| 596 | update(id, |j| { |
| 597 | j.done_bytes += written_bytes; |
| 598 | j.cur_loaded = 0.0; |
| 599 | }); |
| 600 | let mut overwrite = Vec::new(); |
| 601 | for (rel, file) in conflicts { |
| 602 | let size = file.size(); |
| 603 | if batch_ow { |
| 604 | // An overwrite request only ever reports a skip when |
| 605 | // a folder sits where the file would go. Sending it |
| 606 | // again would loop, so it fails here. |
| 607 | done_count += 1; |
| 608 | update(id, |j| { |
| 609 | j.errors |
| 610 | .push((rel, i18n::t(i18n::k::ERR_FOLDER_EXISTS).to_string())); |
| 611 | j.done_bytes += size; |
| 612 | }); |
| 613 | continue; |
| 614 | } |
| 615 | let choice = match rule { |
| 616 | OnConflict::Ask => { |
| 617 | update(id, |j| j.state = JobState::Conflict(rel.clone())); |
| 618 | // A stalled job must be seen: open the overlay. |
| 619 | collapsed().set(false); |
| 620 | let choice = wait_decision(&ctl).await; |
| 621 | if ctl.borrow().abort { |
| 622 | finish(id, JobState::Aborted, &dir); |
| 623 | return; |
| 624 | } |
| 625 | // `decide` may have made this the job's rule. |
| 626 | rule = read(id, |j| j.on_conflict).unwrap_or(rule); |
| 627 | choice |
| 628 | } |
| 629 | r => r, |
| 630 | }; |
| 631 | match choice { |
| 632 | OnConflict::Skip => { |
| 633 | done_count += 1; |
| 634 | update(id, |j| { |
| 635 | j.skipped.push(rel); |
| 636 | j.done_bytes += size; |
| 637 | }); |
| 638 | } |
| 639 | _ => overwrite.push((rel, file)), |
| 640 | } |
| 641 | } |
| 642 | if !overwrite.is_empty() { |
| 643 | queue.push_front((overwrite, true)); |
| 644 | } |
| 645 | } |
| 646 | Err(e) if e.status() == Some(409) && files.len() > 1 => { |
| 647 | // The server rejected the whole request with a conflict it |
| 648 | // could not attribute. Retry one by one so the error lands |
| 649 | // on the file that caused it. |
| 650 | // ponytail: files before the culprit were already written; |
| 651 | // their single retries come back as "appeared during the |
| 652 | // transfer" and may prompt. Rare: the pre-check catches |
| 653 | // folders unless one was created mid-upload. |
| 654 | let _ = e; |
| 655 | for f in files.into_iter().rev() { |
| 656 | queue.push_front((vec![f], batch_ow)); |
| 657 | } |
| 658 | } |
| 659 | Err(e) if e.status() == Some(409) => { |
| 660 | // A folder sits where the file would go. Per-file |
| 661 | // failure; the job goes on. |
| 662 | done_count += 1; |
| 663 | let (rel, file) = &files[0]; |
| 664 | let size = file.size(); |
| 665 | update(id, |j| { |
| 666 | j.errors.push((rel.clone(), e.to_string())); |
| 667 | j.done_bytes += size; |
| 668 | j.cur_loaded = 0.0; |
| 669 | }); |
| 670 | } |
| 671 | Err(e) => { |
| 672 | finish(id, JobState::Failed(e.to_string()), &dir); |
| 673 | return; |
| 674 | } |
| 675 | } |
| 676 | } |
| 677 | finish(id, JobState::Done, &dir); |
| 678 | } |
| 679 | |
| 680 | async fn wait_decision(ctl: &Rc<RefCell<Control>>) -> OnConflict { |
| 681 | // ponytail: 200 ms poll, wake channel if latency matters. |
| 682 | loop { |
| 683 | if ctl.borrow().abort { |
| 684 | return OnConflict::Skip; |
| 685 | } |
| 686 | if let Some(d) = ctl.borrow_mut().decision.take() { |
| 687 | return d; |
| 688 | } |
| 689 | TimeoutFuture::new(200).await; |
| 690 | } |
| 691 | } |
| 692 | |
| 693 | fn finish(id: u64, state: JobState, dir: &str) { |
| 694 | let root_id = read(id, |j| j.root_id); |
| 695 | update(id, |j| { |
| 696 | j.state = state; |
| 697 | j.cur_loaded = 0.0; |
| 698 | j.finished_ms = Some(now_ms()); |
| 699 | }); |
| 700 | if let Some(root_id) = root_id { |
| 701 | ON_DONE.with(|h| { |
| 702 | if let Some((_, f)) = h.borrow().as_ref() { |
| 703 | f(root_id, dir); |
| 704 | } |
| 705 | }); |
| 706 | } |
| 707 | spawn_local(async move { |
| 708 | gloo_timers::future::sleep(OVERLAY_LINGER).await; |
| 709 | update(id, |j| j.hidden = true); |
| 710 | }); |
| 711 | } |
| 712 | |
| 713 | /// "12.3 MB/s". |
| 714 | pub fn format_speed(bytes_per_s: f64) -> String { |
| 715 | i18n::t_fmt( |
| 716 | i18n::k::UPLOAD_SPEED, |
| 717 | &crate::util::format_size(bytes_per_s.max(0.0) as u64), |
| 718 | ) |
| 719 | } |
| 720 | |
| 721 | #[cfg(test)] |
| 722 | mod tests { |
| 723 | use super::{BATCH_BYTES, BATCH_FILES, chunk}; |
| 724 | |
| 725 | /// Fake sizes: a real `web_sys::File` needs a browser. |
| 726 | #[test] |
| 727 | fn chunk_respects_both_caps() { |
| 728 | let sizes = |v: Vec<f64>| chunk(v, |s: &f64| *s); |
| 729 | let lens = |b: Vec<Vec<f64>>| b.iter().map(Vec::len).collect::<Vec<_>>(); |
| 730 | |
| 731 | // The byte cap cuts before the group would exceed it. |
| 732 | assert_eq!( |
| 733 | lens(sizes(vec![BATCH_BYTES * 0.6, BATCH_BYTES * 0.6, 1.0])), |
| 734 | vec![1, 2] |
| 735 | ); |
| 736 | // A file above the cap travels alone. |
| 737 | assert_eq!( |
| 738 | lens(sizes(vec![1.0, BATCH_BYTES * 3.0, 1.0])), |
| 739 | vec![1, 1, 1] |
| 740 | ); |
| 741 | // The file cap cuts even when the bytes are nothing. |
| 742 | assert_eq!( |
| 743 | lens(sizes(vec![0.0; BATCH_FILES * 2 + 1])), |
| 744 | vec![BATCH_FILES, BATCH_FILES, 1] |
| 745 | ); |
| 746 | assert!(sizes(Vec::new()).is_empty()); |
| 747 | } |
| 748 | } |
| 749 |