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