Skip to main content

qframe/runtime/
task.rs

1//! Background tasks with progress, cancellation and an outcome, and the model that tracks them
2//! for display.
3//!
4//! A [`Task`] runs on its own thread and talks to the application only through messages, so
5//! drawing never waits for it. The application keeps a [`Tasks`] model up to date from
6//! [`TaskEvent`]s and shows it, for example with [`TaskList`](crate::widgets::TaskList).
7
8use std::collections::{HashMap, HashSet};
9use std::io;
10use std::panic::{AssertUnwindSafe, catch_unwind};
11use std::sync::atomic::{AtomicU64, Ordering};
12use std::sync::mpsc::{Receiver, RecvTimeoutError, Sender, TryRecvError};
13use std::sync::{Arc, Condvar, Mutex, MutexGuard, PoisonError};
14use std::time::{Duration, Instant};
15
16use super::command::MapFn;
17
18/// Identifies one task for its whole life. Ids are unique within the process.
19#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord)]
20pub struct TaskId(u64);
21
22impl TaskId {
23    fn next() -> Self {
24        static NEXT: AtomicU64 = AtomicU64::new(1);
25        Self(NEXT.fetch_add(1, Ordering::Relaxed))
26    }
27}
28
29/// How a task ended.
30#[derive(Debug, Clone, PartialEq, Eq)]
31pub enum TaskOutcome {
32    /// The work returned `Ok`; its message was delivered just before this outcome.
33    Done,
34    /// The work returned `Err` with this reason, or panicked.
35    Failed(String),
36    /// [`Command::cancel_task`](crate::runtime::Command::cancel_task) asked it to stop; its
37    /// result was dropped.
38    Cancelled,
39}
40
41/// Why [`TaskCx::recv_timeout`] returned without an item.
42#[derive(Debug, Clone, Copy, PartialEq, Eq)]
43#[non_exhaustive]
44pub enum RecvWait {
45    /// The time ran out before an item came.
46    Timeout,
47    /// The task was asked to stop.
48    Cancelled,
49    /// Every sender is gone and nothing is left to receive.
50    Closed,
51}
52
53/// What happened to a task, delivered through [`Task::on_event`].
54#[derive(Debug, Clone, PartialEq)]
55pub enum TaskEvent {
56    /// The task was started with this label.
57    Started {
58        /// The task.
59        id: TaskId,
60        /// Its label.
61        label: String,
62    },
63    /// The task reported progress. `None` fields keep their previous value.
64    Progress {
65        /// The task.
66        id: TaskId,
67        /// Completed share from 0 to 1.
68        fraction: Option<f32>,
69        /// A short note about the current step.
70        note: Option<String>,
71    },
72    /// The task ended.
73    Finished {
74        /// The task.
75        id: TaskId,
76        /// How.
77        outcome: TaskOutcome,
78    },
79}
80
81impl TaskEvent {
82    /// The task the event is about.
83    #[must_use]
84    pub fn id(&self) -> TaskId {
85        match self {
86            Self::Started { id, .. } | Self::Progress { id, .. } | Self::Finished { id, .. } => *id,
87        }
88    }
89}
90
91type Work<Msg> = Box<dyn FnOnce(&TaskCx<Msg>) -> Result<Msg, String> + Send>;
92type EventMessage<Msg> = Arc<dyn Fn(TaskEvent) -> Msg + Send + Sync>;
93/// Hands a message of running work to the event loop.
94type Deliver<Msg> = Arc<dyn Fn(Msg) + Send + Sync>;
95/// Hands a task event, already turned into a message, to the event loop.
96type Report = Arc<dyn Fn(TaskEvent) + Send + Sync>;
97
98/// Work to run in the background with progress, cancellation and an outcome.
99///
100/// ```
101/// use std::time::Duration;
102/// use qframe::runtime::{Command, Task, TaskEvent};
103///
104/// enum Msg {
105///     Task(TaskEvent),
106///     Built(String),
107/// }
108///
109/// let task = Task::new("Build image", |cx| {
110///     for step in 0..4 {
111///         if !cx.sleep(Duration::from_millis(300)) {
112///             return Err("stopped".into());
113///         }
114///         cx.progress((step + 1) as f32 / 4.0);
115///     }
116///     Ok(Msg::Built("sha256:4f2a".into()))
117/// })
118/// .on_event(Msg::Task);
119/// let id = task.id(); // keep it to cancel the task later
120/// let command: Command<Msg> = Command::task(task);
121/// # let _ = (id, command);
122/// ```
123pub struct Task<Msg> {
124    id: TaskId,
125    label: String,
126    work: Work<Msg>,
127    on_event: Option<EventMessage<Msg>>,
128}
129
130impl<Msg: Send + 'static> Task<Msg> {
131    /// A task labelled `label` running `work`. `Ok` delivers its message; `Err` fails the task
132    /// with a reason.
133    #[must_use]
134    pub fn new(
135        label: impl Into<String>,
136        work: impl FnOnce(&TaskCx<Msg>) -> Result<Msg, String> + Send + 'static,
137    ) -> Self {
138        Self { id: TaskId::next(), label: label.into(), work: Box::new(work), on_event: None }
139    }
140
141    /// Turns start, progress and the outcome into messages, e.g. to update a [`Tasks`] model.
142    #[must_use]
143    pub fn on_event(mut self, message: impl Fn(TaskEvent) -> Msg + Send + Sync + 'static) -> Self {
144        self.on_event = Some(Arc::new(message));
145        self
146    }
147
148    /// The task's id, known before it starts.
149    #[must_use]
150    pub fn id(&self) -> TaskId {
151        self.id
152    }
153
154    /// The task's label.
155    #[must_use]
156    pub fn label(&self) -> &str {
157        &self.label
158    }
159
160    /// The same task delivering `map(message)` for every message it would deliver: its result,
161    /// what the work sends while it runs and its events.
162    pub(crate) fn map<B: Send + 'static>(self, map: MapFn<Msg, B>) -> Task<B> {
163        let Self { id, label, work, on_event } = self;
164        let on_event = on_event.map(|message| {
165            let map = Arc::clone(&map);
166            Arc::new(move |event| map(message(event))) as EventMessage<B>
167        });
168        let work: Work<B> = Box::new(move |cx: &TaskCx<B>| {
169            let deliver = Arc::clone(&cx.deliver);
170            let inner_map = Arc::clone(&map);
171            let inner = TaskCx {
172                id: cx.id,
173                clock: Arc::clone(&cx.clock),
174                deliver: Arc::new(move |message| deliver(inner_map(message))),
175                report: cx.report.clone(),
176            };
177            work(&inner).map(|message| map(message))
178        });
179        Task { id, label, work, on_event }
180    }
181}
182
183/// A message from a background thread to the event loop.
184pub(crate) enum Delivery<Msg> {
185    /// Apply this message.
186    Message(Msg),
187    /// A piece of background work ended.
188    Ended,
189}
190
191/// Time as tasks see it: the real clock, or the test harness's fake clock that only moves when
192/// the test advances it.
193pub(crate) struct TaskClock {
194    fake: bool,
195    state: Mutex<ClockState>,
196    changed: Condvar,
197}
198
199#[derive(Default)]
200struct ClockState {
201    now: Duration,
202    /// Tasks that are working rather than sleeping.
203    busy: usize,
204    cancelled: HashSet<TaskId>,
205    /// Fake clock only: tasks asleep and when they wake. Whoever wakes a sleeper (the clock
206    /// moving or a cancel) counts it as busy right away, so settling never misses it.
207    sleeping: HashMap<TaskId, Duration>,
208    /// Fake clock only: where each task is in time. A task woken by a big clock jump carries on
209    /// from the moment its sleep ended, so a loop of sleeps keeps its rhythm.
210    task_time: HashMap<TaskId, Duration>,
211}
212
213/// How long a task waiting for a channel item listens before it looks again whether it was
214/// cancelled. A channel of the standard library cannot be woken by anything but its own senders,
215/// so a cancel reaches a waiting task within this.
216const RECV_SLICE: Duration = Duration::from_millis(50);
217
218/// How long the harness waits for tasks to reach a sleep before it gives up.
219const SETTLE_LIMIT: Duration = Duration::from_secs(10);
220
221impl TaskClock {
222    pub(crate) fn new(fake: bool) -> Arc<Self> {
223        Arc::new(Self { fake, state: Mutex::new(ClockState::default()), changed: Condvar::new() })
224    }
225
226    fn lock(&self) -> MutexGuard<'_, ClockState> {
227        self.state.lock().unwrap_or_else(PoisonError::into_inner)
228    }
229
230    /// Asks task `id` to stop; its sleeps return at once.
231    pub(crate) fn cancel(&self, id: TaskId) {
232        let mut state = self.lock();
233        state.cancelled.insert(id);
234        if state.sleeping.remove(&id).is_some() {
235            state.busy += 1;
236            let now = state.now;
237            state.task_time.insert(id, now);
238        }
239        drop(state);
240        self.changed.notify_all();
241    }
242
243    /// Moves the fake clock to `now` and waits until every task is asleep or finished.
244    ///
245    /// # Panics
246    ///
247    /// Panics when a task keeps working for longer than ten seconds, which in a test means it
248    /// never sleeps or blocks forever.
249    pub(crate) fn settle(&self, now: Duration) {
250        let mut state = self.lock();
251        state.now = state.now.max(now);
252        let current = state.now;
253        let due: Vec<(TaskId, Duration)> =
254            state.sleeping.iter().filter(|(_, until)| **until <= current).map(|(id, until)| (*id, *until)).collect();
255        for (id, until) in due {
256            state.sleeping.remove(&id);
257            state.task_time.insert(id, until);
258            state.busy += 1;
259        }
260        self.changed.notify_all();
261        let started = Instant::now();
262        while state.busy > 0 {
263            let waited = started.elapsed();
264            assert!(waited < SETTLE_LIMIT, "a background task kept working for {SETTLE_LIMIT:?} without sleeping");
265            state = self.changed.wait_timeout(state, SETTLE_LIMIT - waited).unwrap_or_else(PoisonError::into_inner).0;
266        }
267    }
268
269    fn is_cancelled(&self, id: TaskId) -> bool {
270        self.lock().cancelled.contains(&id)
271    }
272
273    fn begin(&self, id: TaskId) {
274        let mut state = self.lock();
275        state.busy += 1;
276        let now = state.now;
277        state.task_time.insert(id, now);
278    }
279
280    fn end(&self, id: TaskId) {
281        let mut state = self.lock();
282        state.busy = state.busy.saturating_sub(1);
283        state.cancelled.remove(&id);
284        state.task_time.remove(&id);
285        drop(state);
286        self.changed.notify_all();
287    }
288
289    /// Sleeps task `id` for `duration`; returns `false` when it was cancelled.
290    fn sleep(&self, id: TaskId, duration: Duration) -> bool {
291        let mut state = self.lock();
292        if self.fake {
293            let until = state.task_time.get(&id).copied().unwrap_or(state.now) + duration;
294            if until <= state.now {
295                state.task_time.insert(id, until);
296            } else if !state.cancelled.contains(&id) {
297                state.sleeping.insert(id, until);
298                state.busy = state.busy.saturating_sub(1);
299                self.changed.notify_all();
300                while state.sleeping.contains_key(&id) {
301                    state = self.changed.wait(state).unwrap_or_else(PoisonError::into_inner);
302                }
303            }
304        } else {
305            let deadline = Instant::now() + duration;
306            while !state.cancelled.contains(&id) {
307                let left = deadline.saturating_duration_since(Instant::now());
308                if left.is_zero() {
309                    break;
310                }
311                state = self.changed.wait_timeout(state, left).unwrap_or_else(PoisonError::into_inner).0;
312            }
313        }
314        !state.cancelled.contains(&id)
315    }
316}
317
318impl TaskClock {
319    /// Waits for task `id`'s next item from `rx`, for at most `timeout` when one is given.
320    ///
321    /// On the fake clock a waiting task rests like a sleeper, so the harness does not wait for
322    /// it: `timeout` runs on the fake clock, and a cancel or the clock passing the deadline wakes
323    /// it and counts it as busy, as they do a sleeper. An item wakes it too, counting it as busy
324    /// itself unless the clock or a cancel already did.
325    fn receive<T>(&self, id: TaskId, rx: &Receiver<T>, timeout: Option<Duration>) -> Result<T, RecvWait> {
326        if self.is_cancelled(id) {
327            return Err(RecvWait::Cancelled);
328        }
329        match rx.try_recv() {
330            Ok(item) => return Ok(item),
331            Err(TryRecvError::Disconnected) => return Err(RecvWait::Closed),
332            Err(TryRecvError::Empty) => {}
333        }
334        if timeout.is_some_and(|timeout| timeout.is_zero()) {
335            return Err(RecvWait::Timeout);
336        }
337        if self.fake { self.receive_fake(id, rx, timeout) } else { self.receive_real(id, rx, timeout) }
338    }
339
340    fn receive_real<T>(&self, id: TaskId, rx: &Receiver<T>, timeout: Option<Duration>) -> Result<T, RecvWait> {
341        let deadline = timeout.and_then(|timeout| Instant::now().checked_add(timeout));
342        loop {
343            if self.is_cancelled(id) {
344                return Err(RecvWait::Cancelled);
345            }
346            let left = deadline.map_or(RECV_SLICE, |deadline| deadline.saturating_duration_since(Instant::now()));
347            if left.is_zero() {
348                return Err(RecvWait::Timeout);
349            }
350            match rx.recv_timeout(left.min(RECV_SLICE)) {
351                Ok(item) => return Ok(item),
352                Err(RecvTimeoutError::Timeout) => {}
353                Err(RecvTimeoutError::Disconnected) => return Err(RecvWait::Closed),
354            }
355        }
356    }
357
358    fn receive_fake<T>(&self, id: TaskId, rx: &Receiver<T>, timeout: Option<Duration>) -> Result<T, RecvWait> {
359        let mut state = self.lock();
360        // A cancel that came since the first look found no sleeper to wake, so nothing would wake
361        // this one.
362        if state.cancelled.contains(&id) {
363            return Err(RecvWait::Cancelled);
364        }
365        // A wait without a deadline rests until an item, a cancel or a closed channel.
366        let until = match timeout {
367            Some(timeout) => state.task_time.get(&id).copied().unwrap_or(state.now).saturating_add(timeout),
368            None => Duration::MAX,
369        };
370        // A task behind the clock has already waited the time out, as a sleep would have.
371        if until <= state.now {
372            state.task_time.insert(id, until);
373            return Err(RecvWait::Timeout);
374        }
375        state.sleeping.insert(id, until);
376        state.busy = state.busy.saturating_sub(1);
377        drop(state);
378        self.changed.notify_all();
379        loop {
380            let heard = rx.recv_timeout(RECV_SLICE);
381            let mut state = self.lock();
382            let woken = !state.sleeping.contains_key(&id);
383            match heard {
384                Err(RecvTimeoutError::Timeout) if !woken => continue,
385                // Nobody woke the task, so it counts itself busy again, from the present moment.
386                Ok(_) | Err(RecvTimeoutError::Disconnected) if !woken => {
387                    state.sleeping.remove(&id);
388                    state.busy += 1;
389                    let now = state.now;
390                    state.task_time.insert(id, now);
391                }
392                _ => {}
393            }
394            return match heard {
395                _ if state.cancelled.contains(&id) => Err(RecvWait::Cancelled),
396                Ok(item) => Ok(item),
397                Err(RecvTimeoutError::Disconnected) => Err(RecvWait::Closed),
398                Err(RecvTimeoutError::Timeout) => Err(RecvWait::Timeout),
399            };
400        }
401    }
402}
403
404/// What running work can do: report progress, send messages, notice cancellation and sleep.
405pub struct TaskCx<Msg> {
406    id: TaskId,
407    clock: Arc<TaskClock>,
408    deliver: Deliver<Msg>,
409    /// Events go out as messages of the application's type, which a task built for another
410    /// message type ([`Command::map`](crate::runtime::Command::map)) does not know.
411    report: Option<Report>,
412}
413
414impl<Msg: Send + 'static> TaskCx<Msg> {
415    /// The running task.
416    #[must_use]
417    pub fn id(&self) -> TaskId {
418        self.id
419    }
420
421    /// Reports the completed share, from 0 to 1.
422    pub fn progress(&self, fraction: f32) {
423        self.event(TaskEvent::Progress { id: self.id, fraction: Some(fraction.clamp(0.0, 1.0)), note: None });
424    }
425
426    /// Describes the current step, e.g. "pushing layers".
427    pub fn note(&self, note: impl Into<String>) {
428        self.event(TaskEvent::Progress { id: self.id, fraction: None, note: Some(note.into()) });
429    }
430
431    /// Delivers `message` to the application while the work goes on, e.g. a log line.
432    pub fn send(&self, message: Msg) {
433        (self.deliver)(message);
434    }
435
436    /// Whether the application asked this task to stop. Long work should check it and return.
437    #[must_use]
438    pub fn is_cancelled(&self) -> bool {
439        self.clock.is_cancelled(self.id)
440    }
441
442    /// Waits `duration`, waking early when the task is cancelled. Returns `false` when it was
443    /// cancelled. In tests the harness's fake clock decides when the sleep ends.
444    #[must_use]
445    pub fn sleep(&self, duration: Duration) -> bool {
446        self.clock.sleep(self.id, duration)
447    }
448
449    /// Waits for the next item from `rx`, the receiving end of a channel another thread feeds,
450    /// such as a connection handing over what it read. Returns `None` when the task is cancelled
451    /// (which wins over an item already waiting) or when every sender is gone and nothing is
452    /// left. A cancel reaches a waiting task within 50 ms: a standard channel wakes only for its
453    /// senders, so the wait looks at the cancel in slices.
454    ///
455    /// In a [`Harness`](crate::runtime::Harness) a task waiting here rests like one asleep, so no
456    /// step of the harness waits for it. An item sent from the test wakes it on its own thread;
457    /// the harness applies what the task then sends at a later step, such as
458    /// [`advance`](crate::runtime::Harness::advance).
459    #[must_use]
460    pub fn recv<T>(&self, rx: &Receiver<T>) -> Option<T> {
461        self.clock.receive(self.id, rx, None).ok()
462    }
463
464    /// Like [`TaskCx::recv`], giving up after `timeout`, and telling why no item came. In a
465    /// [`Harness`](crate::runtime::Harness) the timeout runs on the fake clock, like
466    /// [`TaskCx::sleep`].
467    ///
468    /// # Errors
469    ///
470    /// [`RecvWait::Timeout`] when `timeout` passed without an item, [`RecvWait::Cancelled`]
471    /// when the task was cancelled and [`RecvWait::Closed`] when every sender is gone and
472    /// nothing is left to receive.
473    pub fn recv_timeout<T>(&self, rx: &Receiver<T>, timeout: Duration) -> Result<T, RecvWait> {
474        self.clock.receive(self.id, rx, Some(timeout))
475    }
476
477    fn event(&self, event: TaskEvent) {
478        if let Some(report) = &self.report {
479            report(event);
480        }
481    }
482}
483
484/// Starts a thread named `name` running `run`. A failed start drops `run`.
485pub(crate) type Spawner = fn(String, Box<dyn FnOnce() + Send>) -> io::Result<()>;
486
487/// The [`Spawner`] of the runtime: a real thread.
488pub(crate) fn spawn_thread(name: String, run: Box<dyn FnOnce() + Send>) -> io::Result<()> {
489    std::thread::Builder::new().name(name).spawn(run).map(drop)
490}
491
492/// The reason of a task whose thread could not start.
493const NO_THREAD: &str = "could not start a thread";
494
495/// Starts `task` on a thread. The `Started` message is returned for the caller to apply at once.
496/// When no thread can start, the task fails at once and still ends.
497pub(crate) fn spawn<Msg: Send + 'static>(
498    task: Task<Msg>,
499    clock: &Arc<TaskClock>,
500    sender: &Sender<Delivery<Msg>>,
501    spawner: Spawner,
502) -> Option<Msg> {
503    let Task { id, label, work, on_event } = task;
504    let started = on_event.as_ref().map(|message| message(TaskEvent::Started { id, label: label.clone() }));
505    let failed = on_event.clone();
506    let outlet = sender.clone();
507    let deliver: Deliver<Msg> = Arc::new(move |message| {
508        // A message from a task that outlived the loop has nowhere to be handled. The task is
509        // still ended below, so nothing waits for it either way.
510        let _ = outlet.send(Delivery::Message(message));
511    });
512    let report = on_event.map(|message| {
513        let deliver = Arc::clone(&deliver);
514        Arc::new(move |event| deliver(message(event))) as Report
515    });
516    let cx = TaskCx { id, clock: Arc::clone(clock), deliver, report };
517    let ended = sender.clone();
518    clock.begin(id);
519    let run = Box::new(move || {
520        let result = catch_unwind(AssertUnwindSafe(|| work(&cx)))
521            .unwrap_or_else(|_| Err(format!("the task `{label}` panicked")));
522        let outcome = match result {
523            _ if cx.is_cancelled() => TaskOutcome::Cancelled,
524            Ok(message) => {
525                cx.send(message);
526                TaskOutcome::Done
527            }
528            Err(reason) => TaskOutcome::Failed(reason),
529        };
530        // The event's message is the application's code (and a conversion of `Command::map`);
531        // if it panics there is nothing left to tell, but the task must still end, or the
532        // runtime would wait for it forever.
533        let _ = catch_unwind(AssertUnwindSafe(|| cx.event(TaskEvent::Finished { id, outcome })));
534        let _ = ended.send(Delivery::Ended);
535        cx.clock.end(id);
536    });
537    if spawner(format!("quvyta-task-{}", id.0), run).is_err() {
538        // The work went with the closure; report the failure so the task never stays running.
539        clock.end(id);
540        if let Some(message) = failed {
541            let outcome = TaskOutcome::Failed(NO_THREAD.to_owned());
542            let _ = sender.send(Delivery::Message(message(TaskEvent::Finished { id, outcome })));
543        }
544        // Both sends fail only once the loop has gone, and a loop that has gone is not counting
545        // tasks any more. What must not happen is the task ending without this, which would
546        // leave the loop waiting for a task that never ran.
547        let _ = sender.send(Delivery::Ended);
548    }
549    started
550}
551
552/// One task as the application shows it.
553#[derive(Debug, Clone, PartialEq)]
554pub struct TaskEntry {
555    /// The task.
556    pub id: TaskId,
557    /// Its label.
558    pub label: String,
559    /// Completed share from 0 to 1, when the task reports one.
560    pub fraction: Option<f32>,
561    /// The latest note.
562    pub note: Option<String>,
563    /// How it ended; `None` while running.
564    pub outcome: Option<TaskOutcome>,
565}
566
567/// Tasks an application shows, kept up to date with [`Tasks::apply`]. Newest last.
568#[derive(Debug, Clone, Default, PartialEq)]
569pub struct Tasks {
570    entries: Vec<TaskEntry>,
571}
572
573impl Tasks {
574    /// No tasks.
575    #[must_use]
576    pub fn new() -> Self {
577        Self::default()
578    }
579
580    /// Records `event`.
581    pub fn apply(&mut self, event: &TaskEvent) {
582        match event {
583            TaskEvent::Started { id, label } => {
584                self.entries.retain(|entry| entry.id != *id);
585                self.entries.push(TaskEntry {
586                    id: *id,
587                    label: label.clone(),
588                    fraction: None,
589                    note: None,
590                    outcome: None,
591                });
592            }
593            TaskEvent::Progress { id, fraction, note } => {
594                if let Some(entry) = self.entries.iter_mut().find(|entry| entry.id == *id) {
595                    entry.fraction = fraction.or(entry.fraction);
596                    if note.is_some() {
597                        entry.note.clone_from(note);
598                    }
599                }
600            }
601            TaskEvent::Finished { id, outcome } => {
602                if let Some(entry) = self.entries.iter_mut().find(|entry| entry.id == *id) {
603                    entry.outcome = Some(outcome.clone());
604                }
605            }
606        }
607    }
608
609    /// Every task, oldest first.
610    #[must_use]
611    pub fn entries(&self) -> &[TaskEntry] {
612        &self.entries
613    }
614
615    /// The task `id`.
616    #[must_use]
617    pub fn get(&self, id: TaskId) -> Option<&TaskEntry> {
618        self.entries.iter().find(|entry| entry.id == id)
619    }
620
621    /// How many tasks are still running.
622    #[must_use]
623    pub fn running(&self) -> usize {
624        self.entries.iter().filter(|entry| entry.outcome.is_none()).count()
625    }
626
627    /// Forgets finished tasks.
628    pub fn clear_finished(&mut self) {
629        self.entries.retain(|entry| entry.outcome.is_none());
630    }
631}
632
633#[cfg(test)]
634mod tests {
635    use super::*;
636    use crate::runtime::engine::{Engine, TaskMode};
637    use crate::runtime::{App, Command, Harness};
638    use crate::widget::View;
639    use crate::widgets::Text;
640
641    #[derive(Default)]
642    struct Pipeline {
643        tasks: Tasks,
644        built: Option<String>,
645        lines: Vec<String>,
646        build: Option<TaskId>,
647    }
648
649    enum Msg {
650        Build,
651        Cancel,
652        Fail,
653        Panic,
654        Task(TaskEvent),
655        Built(String),
656        Line(String),
657    }
658
659    impl App for Pipeline {
660        type Msg = Msg;
661        fn update(&mut self, msg: Msg) -> Command<Msg> {
662            match msg {
663                Msg::Build => {
664                    let task = Task::new("Build image", |cx| {
665                        cx.note("resolving layers");
666                        for step in 0..4 {
667                            if !cx.sleep(Duration::from_millis(100)) {
668                                return Err("stopped".into());
669                            }
670                            cx.progress((step + 1) as f32 / 4.0);
671                            cx.send(Msg::Line(format!("layer {step}")));
672                        }
673                        Ok(Msg::Built("sha256:4f2a".into()))
674                    })
675                    .on_event(Msg::Task);
676                    self.build = Some(task.id());
677                    return Command::task(task);
678                }
679                Msg::Cancel => return self.build.map_or_else(Command::none, Command::cancel_task),
680                Msg::Fail => {
681                    return Command::task(
682                        Task::new("Sync registry", |cx| {
683                            let _ = cx.sleep(Duration::from_millis(50));
684                            Err("registry timed out".into())
685                        })
686                        .on_event(Msg::Task),
687                    );
688                }
689                Msg::Panic => {
690                    return Command::task(
691                        Task::new("Broken", |_| -> Result<Msg, String> { panic!("boom") }).on_event(Msg::Task),
692                    );
693                }
694                Msg::Task(event) => self.tasks.apply(&event),
695                Msg::Built(digest) => self.built = Some(digest),
696                Msg::Line(line) => self.lines.push(line),
697            }
698            Command::none()
699        }
700        fn view(&self, ui: &mut View<'_, Msg>) {
701            ui.add(Text::new(format!("running {}", self.tasks.running())));
702        }
703    }
704
705    #[test]
706    fn progress_follows_the_fake_clock_and_completes() {
707        let mut h = Harness::new(Pipeline::default(), 20, 1);
708        h.send(Msg::Build);
709        assert_eq!(h.screen(), "running 1\n");
710        let entry = h.app().tasks.entries()[0].clone();
711        assert_eq!(entry.label, "Build image");
712        assert_eq!(entry.note.as_deref(), Some("resolving layers"));
713        assert_eq!(entry.fraction, None);
714        h.advance(Duration::from_millis(100));
715        assert_eq!(h.app().tasks.entries()[0].fraction, Some(0.25));
716        assert_eq!(h.app().lines, ["layer 0"]);
717        h.advance(Duration::from_millis(250));
718        assert_eq!(h.app().tasks.entries()[0].fraction, Some(0.75));
719        h.advance(Duration::from_millis(100));
720        assert_eq!(h.app().built.as_deref(), Some("sha256:4f2a"));
721        assert_eq!(h.app().tasks.entries()[0].outcome, Some(TaskOutcome::Done));
722        assert_eq!(h.screen(), "running 0\n");
723    }
724
725    #[test]
726    fn cancelling_wakes_the_sleep_and_drops_the_result() {
727        let mut h = Harness::new(Pipeline::default(), 20, 1);
728        h.send(Msg::Build).advance(Duration::from_millis(150)).send(Msg::Cancel);
729        let entry = &h.app().tasks.entries()[0];
730        assert_eq!(entry.outcome, Some(TaskOutcome::Cancelled));
731        assert_eq!(entry.fraction, Some(0.25));
732        assert!(h.app().built.is_none());
733    }
734
735    fn no_thread(_: String, _: Box<dyn FnOnce() + Send>) -> io::Result<()> {
736        Err(io::Error::other("no threads left"))
737    }
738
739    #[test]
740    fn a_task_whose_thread_cannot_start_fails_and_ends() {
741        let mut engine = Engine::new(Pipeline::default(), crate::env::Env::builtin(), TaskMode::Threads);
742        engine.spawner = no_thread;
743        engine.update(Msg::Build);
744        assert_eq!(engine.poll_tasks(), 2, "Finished, then Ended");
745        let entry = &engine.app.tasks.entries()[0];
746        assert_eq!(entry.outcome, Some(TaskOutcome::Failed("could not start a thread".into())));
747        assert_eq!((engine.app.tasks.running(), engine.pending_tasks), (0, 0));
748        assert!(engine.app.built.is_none());
749    }
750
751    /// Starts, on `Some(())`, a task whose event message panics once the task finishes.
752    struct Fragile;
753
754    impl App for Fragile {
755        type Msg = Option<()>;
756        fn update(&mut self, start: Option<()>) -> Command<Option<()>> {
757            if start.is_none() {
758                return Command::none();
759            }
760            Command::task(Task::new("Fragile", |_| Ok(None)).on_event(|event| match event {
761                TaskEvent::Finished { .. } => panic!("the message of the outcome failed"),
762                _ => None,
763            }))
764        }
765        fn view(&self, ui: &mut View<'_, Option<()>>) {
766            ui.add(Text::new("fragile"));
767        }
768    }
769
770    #[test]
771    fn a_task_whose_last_event_message_panics_still_ends() {
772        let mut engine = Engine::new(Fragile, crate::env::Env::builtin(), TaskMode::Threads);
773        engine.update(Some(()));
774        let started = Instant::now();
775        while engine.pending_tasks > 0 {
776            assert!(started.elapsed() < Duration::from_secs(10), "the runtime waits for the task forever");
777            engine.poll_tasks();
778            std::thread::sleep(Duration::from_millis(5));
779        }
780    }
781
782    #[test]
783    fn failures_and_panics_become_outcomes() {
784        let mut h = Harness::new(Pipeline::default(), 20, 1);
785        h.send(Msg::Fail).send(Msg::Panic);
786        assert_eq!(h.app().tasks.running(), 1, "the failing task still sleeps");
787        assert_eq!(h.app().tasks.entries()[1].outcome, Some(TaskOutcome::Failed("the task `Broken` panicked".into())));
788        h.advance(Duration::from_millis(50));
789        assert_eq!(h.app().tasks.entries()[0].outcome, Some(TaskOutcome::Failed("registry timed out".into())));
790        let mut tasks = h.app().tasks.clone();
791        tasks.clear_finished();
792        assert!(tasks.entries().is_empty());
793    }
794
795    /// Listens to a channel in a task and hands every line it hears to the application.
796    #[derive(Default)]
797    struct Listener {
798        rx: Option<std::sync::mpsc::Receiver<String>>,
799        tasks: Tasks,
800        heard: Vec<String>,
801        ended: Option<String>,
802        listening: Option<TaskId>,
803    }
804
805    enum Heard {
806        /// Listens until the channel closes or the task is cancelled.
807        Listen,
808        /// Waits for one line at most this long and tells why none came.
809        Wait(Duration),
810        Cancel,
811        Task(TaskEvent),
812        Line(String),
813        Ended(String),
814    }
815
816    impl App for Listener {
817        type Msg = Heard;
818        fn update(&mut self, msg: Heard) -> Command<Heard> {
819            match msg {
820                Heard::Listen => {
821                    let rx = self.rx.take().expect("one channel per test");
822                    let task = Task::new("Listen", move |cx| {
823                        while let Some(line) = cx.recv(&rx) {
824                            cx.send(Heard::Line(line));
825                        }
826                        Ok(Heard::Ended("closed".into()))
827                    })
828                    .on_event(Heard::Task);
829                    self.listening = Some(task.id());
830                    return Command::task(task);
831                }
832                Heard::Wait(timeout) => {
833                    let rx = self.rx.take().expect("one channel per test");
834                    let task = Task::new("Wait", move |cx| {
835                        let why = match cx.recv_timeout(&rx, timeout) {
836                            Ok(line) => line,
837                            Err(why) => format!("{why:?}"),
838                        };
839                        Ok(Heard::Ended(why))
840                    })
841                    .on_event(Heard::Task);
842                    self.listening = Some(task.id());
843                    return Command::task(task);
844                }
845                Heard::Cancel => return self.listening.map_or_else(Command::none, Command::cancel_task),
846                Heard::Task(event) => self.tasks.apply(&event),
847                Heard::Line(line) => self.heard.push(line),
848                Heard::Ended(why) => self.ended = Some(why),
849            }
850            Command::none()
851        }
852        fn view(&self, ui: &mut View<'_, Heard>) {
853            ui.add(Text::new(format!("heard {}", self.heard.len())));
854        }
855    }
856
857    /// A listener on real threads with the sending end of its channel.
858    fn listening_engine() -> (Engine<Listener>, std::sync::mpsc::Sender<String>) {
859        let (tx, rx) = std::sync::mpsc::channel();
860        let app = Listener { rx: Some(rx), ..Listener::default() };
861        (Engine::new(app, crate::env::Env::builtin(), TaskMode::Threads), tx)
862    }
863
864    /// Polls `engine` until `done` holds, failing after a generous bound.
865    fn poll_until(engine: &mut Engine<Listener>, what: &str, done: impl Fn(&Listener) -> bool) {
866        let started = Instant::now();
867        while !done(&engine.app) {
868            assert!(started.elapsed() < Duration::from_secs(20), "{what} never happened");
869            engine.poll_tasks();
870            std::thread::sleep(Duration::from_millis(5));
871        }
872    }
873
874    #[test]
875    fn a_task_receives_what_another_thread_sends_and_ends_when_the_channel_closes() {
876        let (mut engine, tx) = listening_engine();
877        engine.update(Heard::Listen);
878        let sender = std::thread::spawn(move || {
879            std::thread::sleep(Duration::from_millis(30));
880            tx.send("GET /index.html 200".into()).expect("the task listens");
881            tx.send("GET /style.css 304".into()).expect("the task listens");
882        });
883        poll_until(&mut engine, "the lines", |app| app.heard.len() == 2);
884        assert_eq!(engine.app.heard, ["GET /index.html 200", "GET /style.css 304"]);
885        sender.join().expect("the sender ends");
886        poll_until(&mut engine, "the end", |app| app.ended.is_some());
887        assert_eq!(engine.app.ended.as_deref(), Some("closed"), "the dropped sender ended the wait");
888        assert_eq!(engine.app.tasks.entries()[0].outcome, Some(TaskOutcome::Done));
889    }
890
891    #[test]
892    fn cancelling_ends_a_wait_for_the_channel_promptly() {
893        let (mut engine, tx) = listening_engine();
894        engine.update(Heard::Listen);
895        std::thread::sleep(Duration::from_millis(100));
896        let asked = Instant::now();
897        engine.update(Heard::Cancel);
898        poll_until(&mut engine, "the cancel", |app| app.tasks.running() == 0);
899        assert!(asked.elapsed() < Duration::from_secs(5), "the cancel took {:?}", asked.elapsed());
900        assert_eq!(engine.app.tasks.entries()[0].outcome, Some(TaskOutcome::Cancelled));
901        assert!(tx.send("late".into()).is_err(), "the task and its receiver are gone");
902    }
903
904    #[test]
905    fn a_wait_with_a_timeout_says_why_no_item_came() {
906        let (mut engine, tx) = listening_engine();
907        engine.update(Heard::Wait(Duration::from_millis(40)));
908        poll_until(&mut engine, "the timeout", |app| app.ended.is_some());
909        assert_eq!(engine.app.ended.as_deref(), Some("Timeout"));
910        drop(tx);
911
912        let (mut engine, tx) = listening_engine();
913        drop(tx);
914        engine.update(Heard::Wait(Duration::from_secs(30)));
915        poll_until(&mut engine, "the closed channel", |app| app.ended.is_some());
916        assert_eq!(engine.app.ended.as_deref(), Some("Closed"));
917    }
918
919    #[test]
920    fn in_the_harness_a_waiting_task_rests_and_still_hears_and_cancels() {
921        let (tx, rx) = std::sync::mpsc::channel();
922        let started = Instant::now();
923        let mut h = Harness::new(Listener { rx: Some(rx), ..Listener::default() }, 20, 1);
924        h.send(Heard::Listen);
925        assert_eq!(h.app().tasks.running(), 1, "the harness did not wait for the listening task");
926        tx.send("GET /health 200".into()).expect("the task listens");
927        while h.app().heard.is_empty() {
928            assert!(started.elapsed() < Duration::from_secs(20), "the line never arrived");
929            std::thread::sleep(Duration::from_millis(5));
930            h.advance(Duration::ZERO);
931        }
932        assert_eq!(h.screen(), "heard 1\n");
933        h.send(Heard::Cancel);
934        assert_eq!(h.app().tasks.entries()[0].outcome, Some(TaskOutcome::Cancelled));
935        assert!(started.elapsed() < Duration::from_secs(20), "the harness hung on the wait");
936    }
937
938    #[test]
939    fn in_the_harness_a_timeout_follows_the_fake_clock() {
940        let (tx, rx) = std::sync::mpsc::channel::<String>();
941        let mut h = Harness::new(Listener { rx: Some(rx), ..Listener::default() }, 20, 1);
942        h.send(Heard::Wait(Duration::from_secs(60)));
943        h.advance(Duration::from_secs(59));
944        assert_eq!(h.app().ended, None, "a minute has not passed on the fake clock");
945        h.advance(Duration::from_secs(1));
946        assert_eq!(h.app().ended.as_deref(), Some("Timeout"));
947        drop(tx);
948    }
949}