1use 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#[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#[derive(Debug, Clone, PartialEq, Eq)]
31pub enum TaskOutcome {
32 Done,
34 Failed(String),
36 Cancelled,
39}
40
41#[derive(Debug, Clone, Copy, PartialEq, Eq)]
43#[non_exhaustive]
44pub enum RecvWait {
45 Timeout,
47 Cancelled,
49 Closed,
51}
52
53#[derive(Debug, Clone, PartialEq)]
55pub enum TaskEvent {
56 Started {
58 id: TaskId,
60 label: String,
62 },
63 Progress {
65 id: TaskId,
67 fraction: Option<f32>,
69 note: Option<String>,
71 },
72 Finished {
74 id: TaskId,
76 outcome: TaskOutcome,
78 },
79}
80
81impl TaskEvent {
82 #[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>;
93type Deliver<Msg> = Arc<dyn Fn(Msg) + Send + Sync>;
95type Report = Arc<dyn Fn(TaskEvent) + Send + Sync>;
97
98pub 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 #[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 #[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 #[must_use]
150 pub fn id(&self) -> TaskId {
151 self.id
152 }
153
154 #[must_use]
156 pub fn label(&self) -> &str {
157 &self.label
158 }
159
160 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
183pub(crate) enum Delivery<Msg> {
185 Message(Msg),
187 Ended,
189}
190
191pub(crate) struct TaskClock {
194 fake: bool,
195 state: Mutex<ClockState>,
196 changed: Condvar,
197}
198
199#[derive(Default)]
200struct ClockState {
201 now: Duration,
202 busy: usize,
204 cancelled: HashSet<TaskId>,
205 sleeping: HashMap<TaskId, Duration>,
208 task_time: HashMap<TaskId, Duration>,
211}
212
213const RECV_SLICE: Duration = Duration::from_millis(50);
217
218const 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 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 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 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 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 if state.cancelled.contains(&id) {
363 return Err(RecvWait::Cancelled);
364 }
365 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 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 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
404pub struct TaskCx<Msg> {
406 id: TaskId,
407 clock: Arc<TaskClock>,
408 deliver: Deliver<Msg>,
409 report: Option<Report>,
412}
413
414impl<Msg: Send + 'static> TaskCx<Msg> {
415 #[must_use]
417 pub fn id(&self) -> TaskId {
418 self.id
419 }
420
421 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 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 pub fn send(&self, message: Msg) {
433 (self.deliver)(message);
434 }
435
436 #[must_use]
438 pub fn is_cancelled(&self) -> bool {
439 self.clock.is_cancelled(self.id)
440 }
441
442 #[must_use]
445 pub fn sleep(&self, duration: Duration) -> bool {
446 self.clock.sleep(self.id, duration)
447 }
448
449 #[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 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
484pub(crate) type Spawner = fn(String, Box<dyn FnOnce() + Send>) -> io::Result<()>;
486
487pub(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
492const NO_THREAD: &str = "could not start a thread";
494
495pub(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 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 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 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 let _ = sender.send(Delivery::Ended);
548 }
549 started
550}
551
552#[derive(Debug, Clone, PartialEq)]
554pub struct TaskEntry {
555 pub id: TaskId,
557 pub label: String,
559 pub fraction: Option<f32>,
561 pub note: Option<String>,
563 pub outcome: Option<TaskOutcome>,
565}
566
567#[derive(Debug, Clone, Default, PartialEq)]
569pub struct Tasks {
570 entries: Vec<TaskEntry>,
571}
572
573impl Tasks {
574 #[must_use]
576 pub fn new() -> Self {
577 Self::default()
578 }
579
580 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 #[must_use]
611 pub fn entries(&self) -> &[TaskEntry] {
612 &self.entries
613 }
614
615 #[must_use]
617 pub fn get(&self, id: TaskId) -> Option<&TaskEntry> {
618 self.entries.iter().find(|entry| entry.id == id)
619 }
620
621 #[must_use]
623 pub fn running(&self) -> usize {
624 self.entries.iter().filter(|entry| entry.outcome.is_none()).count()
625 }
626
627 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 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 #[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 Listen,
808 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 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 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}