uploads.rs
⎇
Raw
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
14use std::cell::{Cell, RefCell};
15use std::collections::{HashSet, VecDeque};
16use std::rc::Rc;
17use std::time::Duration;
18
19use gloo_timers::future::TimeoutFuture;
20use leptos::prelude::*;
21use wasm_bindgen_futures::spawn_local;
22
23use crate::api;
24use 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)]
29pub 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)]
37pub 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
49impl 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)]
59pub 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
95impl 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.
127pub fn speeds(samples: &[(f64, f64)]) -> Vec<f64> {
128 samples.windows(2).map(|w| rate(w[0], w[1])).collect()
129}
130
131fn 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)]
142struct 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`]).
151type OnDone = Rc<dyn Fn(i64, &str)>;
152
153thread_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.
171fn 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
181const SAMPLE_MS: f64 = 500.0;
182const SAMPLES_KEPT: usize = 60;
183/// How long a finished job stays in the overlay.
184const OVERLAY_LINGER: Duration = Duration::from_secs(5);
185
186pub fn jobs() -> RwSignal<Vec<Job>, LocalStorage> {
187 JOBS.with(|j| *j)
188}
189
190pub 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`].
196pub 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.
206pub 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
215pub 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.
220fn next_id() -> u64 {
221 NEXT_ID.with(|n| n.replace(n.get() + 1))
222}
223
224fn 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
231fn 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
240fn 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
248fn 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.
259pub 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.
298pub 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.
334pub 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).
339pub 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.
347async 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
364pub 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
380pub 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
390pub 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.
407pub 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
429pub fn toggle_expanded(id: u64) {
430 update(id, |j| j.expanded = !j.expanded);
431}
432
433/// Remove a finished job from the log.
434pub 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
439pub 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.
459const BATCH_BYTES: f64 = 1024.0 * 1024.0;
460const BATCH_FILES: usize = 200;
461
462/// One request's worth of files, plus the overwrite flag they all share.
463type 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).
468fn 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.
490fn 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
513async 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
680async 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
693fn 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".
714pub 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)]
722mod 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