use std::collections::VecDeque;
use std::rc::Rc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant};
use crate::app::ExternalData;
use crate::reactive::{Dirty, Runtime};
use crate::window::WindowId;
pub trait Waker: Send + Sync + 'static {
fn wake(&self) -> bool;
fn post(&self, window: WindowId, data: ExternalData) -> bool;
}
#[derive(Default)]
pub(crate) struct LocalQueue {
items: Mutex<VecDeque<(WindowId, ExternalData)>>,
}
impl LocalQueue {
fn push(&self, window: WindowId, data: ExternalData) {
self.items
.lock()
.unwrap_or_else(|e| e.into_inner())
.push_back((window, data));
}
fn take(&self) -> Vec<(WindowId, ExternalData)> {
self.items
.lock()
.unwrap_or_else(|e| e.into_inner())
.drain(..)
.collect()
}
}
#[derive(Clone)]
pub(crate) enum WakerSlot {
Local(Arc<LocalQueue>),
Platform(Arc<dyn Waker>),
}
impl Default for WakerSlot {
fn default() -> Self {
Self::Local(Arc::new(LocalQueue::default()))
}
}
impl WakerSlot {
pub(crate) fn wake(&self) -> bool {
match self {
WakerSlot::Local(_) => true,
WakerSlot::Platform(w) => w.wake(),
}
}
pub(crate) fn post(&self, window: WindowId, data: ExternalData) -> bool {
match self {
WakerSlot::Local(q) => {
q.push(window, data);
true
}
WakerSlot::Platform(w) => w.post(window, data),
}
}
}
#[derive(Clone, Default)]
pub struct CancelToken(Arc<AtomicBool>);
impl CancelToken {
pub fn new() -> Self {
Self::default()
}
pub fn cancel(&self) {
self.0.store(true, Ordering::SeqCst);
}
pub fn is_cancelled(&self) -> bool {
self.0.load(Ordering::SeqCst)
}
}
impl std::fmt::Debug for CancelToken {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "CancelToken({})", self.is_cancelled())
}
}
#[derive(Clone)]
pub struct TaskCtx {
id: u64,
window: WindowId,
waker: WakerSlot,
cancel: CancelToken,
}
impl TaskCtx {
pub fn id(&self) -> u64 {
self.id
}
pub fn window(&self) -> WindowId {
self.window
}
pub fn post<T: Send + 'static>(&self, msg: T) -> bool {
self.waker.post(self.window, ExternalData::new(msg))
}
pub fn progress(&self, done: usize, total: usize) -> bool {
self.waker.post(
self.window,
ExternalData::new(TaskProgress {
id: self.id,
done,
total,
detail: None,
}),
)
}
pub fn progress_with(&self, done: usize, total: usize, detail: impl Into<String>) -> bool {
self.waker.post(
self.window,
ExternalData::new(TaskProgress {
id: self.id,
done,
total,
detail: Some(detail.into()),
}),
)
}
pub fn is_cancelled(&self) -> bool {
self.cancel.is_cancelled()
}
pub fn cancel_token(&self) -> CancelToken {
self.cancel.clone()
}
pub fn wake(&self) -> bool {
self.waker.wake()
}
}
impl std::fmt::Debug for TaskCtx {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("TaskCtx")
.field("id", &self.id)
.field("window", &self.window)
.finish_non_exhaustive()
}
}
pub struct TaskHandle {
rt: Runtime,
id: u64,
cancel: CancelToken,
done: Arc<AtomicBool>,
}
impl TaskHandle {
pub fn id(&self) -> u64 {
self.id
}
pub fn cancel(&self) {
self.cancel.cancel();
}
pub fn is_done(&self) -> bool {
self.done.load(Ordering::SeqCst)
}
pub fn is_cancelled(&self) -> bool {
self.cancel.is_cancelled()
}
pub fn is_running(&self) -> bool {
self.rt
.inner
.tasks
.borrow()
.iter()
.any(|t| t.id == self.id)
}
}
impl std::fmt::Debug for TaskHandle {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("TaskHandle")
.field("id", &self.id)
.field("done", &self.is_done())
.finish_non_exhaustive()
}
}
pub struct TaskEvent {
pub id: u64,
pub payload: ExternalData,
}
impl std::fmt::Debug for TaskEvent {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "TaskEvent {{ id: {} }}", self.id)
}
}
pub struct TaskFailed {
pub id: u64,
}
impl std::fmt::Debug for TaskFailed {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "TaskFailed {{ id: {} }}", self.id)
}
}
pub(crate) struct TaskProgress {
pub(crate) id: u64,
pub(crate) done: usize,
pub(crate) total: usize,
pub(crate) detail: Option<String>,
}
pub(crate) struct TaskRecord {
pub(crate) id: u64,
pub(crate) window: WindowId,
pub(crate) cancel: CancelToken,
pub(crate) busy: Option<u64>,
}
pub struct BusyItem {
pub id: u64,
pub window: WindowId,
pub label: String,
pub detail: Option<String>,
pub progress: Option<(usize, usize)>,
pub cancel: Option<Rc<dyn Fn()>>,
pub(crate) since: Instant,
pub(crate) hide_at: Option<Instant>,
}
impl Clone for BusyItem {
fn clone(&self) -> Self {
Self {
id: self.id,
window: self.window,
label: self.label.clone(),
detail: self.detail.clone(),
progress: self.progress,
cancel: self.cancel.clone(),
since: self.since,
hide_at: self.hide_at,
}
}
}
impl BusyItem {
pub fn ratio(&self) -> Option<f32> {
self.progress
.map(|(d, t)| if t == 0 { 0.0 } else { (d as f32 / t as f32).clamp(0.0, 1.0) })
}
pub fn is_cancellable(&self) -> bool {
self.cancel.is_some()
}
}
impl std::fmt::Debug for BusyItem {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("BusyItem")
.field("id", &self.id)
.field("window", &self.window)
.field("label", &self.label)
.field("detail", &self.detail)
.field("progress", &self.progress)
.field("cancellable", &self.is_cancellable())
.finish()
}
}
impl Runtime {
pub fn set_waker(&self, waker: Arc<dyn Waker>) {
*self.inner.waker.borrow_mut() = WakerSlot::Platform(waker);
}
pub fn waker(&self) -> Option<Arc<dyn Waker>> {
match &*self.inner.waker.borrow() {
WakerSlot::Local(_) => None,
WakerSlot::Platform(w) => Some(Arc::clone(w)),
}
}
pub fn is_online(&self) -> bool {
self.waker().is_some()
}
pub fn wake(&self) -> bool {
self.inner.waker.borrow().wake()
}
pub fn take_pending_external(&self) -> Vec<(WindowId, ExternalData)> {
match &*self.inner.waker.borrow() {
WakerSlot::Local(q) => q.take(),
WakerSlot::Platform(_) => Vec::new(),
}
}
pub fn spawn_task<T, F>(&self, window: WindowId, work: F) -> TaskHandle
where
T: Send + 'static,
F: FnOnce(TaskCtx) -> T + Send + 'static,
{
self.spawn_task_inner(window, None, work)
}
pub fn spawn_task_busy<T, F>(
&self,
window: WindowId,
label: impl Into<String>,
work: F,
) -> TaskHandle
where
T: Send + 'static,
F: FnOnce(TaskCtx) -> T + Send + 'static,
{
self.spawn_task_inner(window, Some(label.into()), work)
}
fn spawn_task_inner<T, F>(
&self,
window: WindowId,
busy_label: Option<String>,
work: F,
) -> TaskHandle
where
T: Send + 'static,
F: FnOnce(TaskCtx) -> T + Send + 'static,
{
let id = {
let n = self.inner.next_task_id.get() + 1;
self.inner.next_task_id.set(n);
n
};
let cancel = CancelToken::new();
let done = Arc::new(AtomicBool::new(false));
let waker = self.inner.waker.borrow().clone();
let busy = busy_label.map(|label| {
let bid = {
let n = self.inner.next_task_id.get() + 1;
self.inner.next_task_id.set(n);
n
};
let cancel_from_ui = cancel.clone();
self.inner.busy.borrow_mut().push(BusyItem {
id: bid,
window,
label,
detail: None,
progress: None,
cancel: Some(Rc::new(move || cancel_from_ui.cancel())),
since: Instant::now(),
hide_at: None,
});
bid
});
self.inner.tasks.borrow_mut().push(TaskRecord {
id,
window,
cancel: cancel.clone(),
busy,
});
self.mark(window, Dirty::VIEW | Dirty::PRESENT);
let ctx = TaskCtx {
id,
window,
waker: waker.clone(),
cancel: cancel.clone(),
};
let done_flag = Arc::clone(&done);
std::thread::spawn(move || {
let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| work(ctx)));
done_flag.store(true, Ordering::SeqCst);
let msg = match result {
Ok(v) => ExternalData::new(TaskEvent {
id,
payload: ExternalData::new(v),
}),
Err(_) => ExternalData::new(TaskFailed { id }),
};
let _ = waker.post(window, msg);
});
TaskHandle {
rt: self.clone(),
id,
cancel,
done,
}
}
pub(crate) fn cancel_tasks_of(&self, window: WindowId) {
let mut tasks = self.inner.tasks.borrow_mut();
for t in tasks.iter() {
if t.window == window {
t.cancel.cancel();
}
}
tasks.retain(|t| t.window != window);
self.inner.busy.borrow_mut().retain(|b| b.window != window);
}
pub fn has_tasks(&self, window: WindowId) -> bool {
self.inner.tasks.borrow().iter().any(|t| t.window == window)
}
pub fn set_busy_min_visible(&self, d: std::time::Duration) {
self.inner.busy_min_visible.set(d);
}
pub fn busy_min_visible(&self) -> std::time::Duration {
self.inner.busy_min_visible.get()
}
pub fn begin_busy(&self, window: WindowId, label: impl Into<String>) -> BusyToken {
let id = {
let n = self.inner.next_task_id.get() + 1;
self.inner.next_task_id.set(n);
n
};
self.inner.busy.borrow_mut().push(BusyItem {
id,
window,
label: label.into(),
detail: None,
progress: None,
cancel: None,
since: Instant::now(),
hide_at: None,
});
self.mark(window, Dirty::VIEW | Dirty::PRESENT);
BusyToken {
rt: self.clone(),
window,
id,
finished: false,
}
}
pub fn busy_items(&self, window: WindowId) -> Vec<BusyItem> {
self.inner
.busy
.borrow()
.iter()
.filter(|b| b.window == window)
.cloned()
.collect()
}
pub fn is_busy(&self, window: WindowId) -> bool {
self.inner.busy.borrow().iter().any(|b| b.window == window)
}
pub(crate) fn end_busy(&self, id: u64) {
let hidden = {
let mut items = self.inner.busy.borrow_mut();
let Some(pos) = items.iter().position(|b| b.id == id) else {
return;
};
let min = self.inner.busy_min_visible.get();
let hide_at = items[pos].since + min;
if min > Duration::ZERO && Instant::now() < hide_at {
items[pos].hide_at = Some(hide_at);
return;
}
items.remove(pos)
};
self.mark(hidden.window, Dirty::VIEW | Dirty::PRESENT);
}
pub fn reap_busy(&self, now: Instant) -> bool {
let mut removed: Vec<WindowId> = Vec::new();
{
let mut items = self.inner.busy.borrow_mut();
let mut i = 0;
while i < items.len() {
if items[i].hide_at.is_some_and(|t| t <= now) {
removed.push(items.remove(i).window);
} else {
i += 1;
}
}
}
for w in &removed {
self.mark(*w, Dirty::VIEW | Dirty::PRESENT);
}
!removed.is_empty()
}
pub(crate) fn set_busy_label(&self, id: u64, label: String) {
let mut items = self.inner.busy.borrow_mut();
let Some(b) = items.iter_mut().find(|b| b.id == id) else {
return;
};
if b.label == label {
return;
}
b.label = label;
let window = b.window;
drop(items);
self.mark(window, Dirty::VIEW);
}
pub(crate) fn set_busy_progress(&self, id: u64, done: usize, total: usize) {
let mut items = self.inner.busy.borrow_mut();
let Some(b) = items.iter_mut().find(|b| b.id == id) else {
return;
};
if b.progress == Some((done, total)) {
return;
}
b.progress = Some((done, total));
let window = b.window;
drop(items);
self.mark(window, Dirty::VIEW);
}
pub(crate) fn set_busy_detail(&self, id: u64, detail: Option<String>) {
let mut items = self.inner.busy.borrow_mut();
let Some(b) = items.iter_mut().find(|b| b.id == id) else {
return;
};
if b.detail == detail {
return;
}
b.detail = detail;
let window = b.window;
drop(items);
self.mark(window, Dirty::VIEW);
}
pub(crate) fn set_busy_cancel(&self, id: u64, f: Rc<dyn Fn()>) {
let mut items = self.inner.busy.borrow_mut();
if let Some(b) = items.iter_mut().find(|b| b.id == id) {
b.cancel = Some(f);
}
}
pub(crate) fn busy_id_of_task(&self, task: u64) -> Option<u64> {
self.inner
.tasks
.borrow()
.iter()
.find(|t| t.id == task)
.and_then(|t| t.busy)
}
}
pub(crate) fn on_task_message(rt: &Runtime, window: WindowId, data: &ExternalData) -> bool {
if let Some(p) = data.downcast_ref::<TaskProgress>() {
if let Some(bid) = rt.busy_id_of_task(p.id) {
rt.set_busy_progress(bid, p.done, p.total);
if let Some(d) = &p.detail {
rt.set_busy_detail(bid, Some(d.clone()));
}
}
return true;
}
let finished = data
.downcast_ref::<TaskEvent>()
.map(|ev| ev.id)
.or_else(|| data.downcast_ref::<TaskFailed>().map(|ev| ev.id));
if let Some(id) = finished {
let rec = {
let mut tasks = rt.inner.tasks.borrow_mut();
let idx = tasks.iter().position(|t| t.id == id);
idx.map(|i| tasks.remove(i))
};
if let Some(rec) = rec
&& let Some(bid) = rec.busy
{
rt.end_busy(bid);
}
rt.mark(window, Dirty::VIEW | Dirty::PRESENT);
return false; }
false
}
pub struct BusyToken {
rt: Runtime,
window: WindowId,
id: u64,
finished: bool,
}
impl BusyToken {
pub fn id(&self) -> u64 {
self.id
}
pub fn window(&self) -> WindowId {
self.window
}
pub fn set_label(&self, label: impl Into<String>) {
self.rt.set_busy_label(self.id, label.into());
}
pub fn set_progress(&self, done: usize, total: usize) {
self.rt.set_busy_progress(self.id, done, total);
}
pub fn set_detail(&self, detail: impl Into<String>) {
self.rt.set_busy_detail(self.id, Some(detail.into()));
}
pub fn cancellable(&self, f: impl Fn() + 'static) {
self.rt.set_busy_cancel(self.id, Rc::new(f));
}
pub fn finish(mut self) {
self.rt.end_busy(self.id);
self.finished = true;
}
}
impl Drop for BusyToken {
fn drop(&mut self) {
if !self.finished {
self.rt.end_busy(self.id);
}
}
}
impl std::fmt::Debug for BusyToken {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("BusyToken")
.field("id", &self.id)
.field("window", &self.window)
.finish_non_exhaustive()
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::window::WindowId;
use std::sync::atomic::AtomicU32;
use std::time::{Duration, Instant};
#[derive(Default)]
struct TestWaker {
wakes: AtomicU32,
posts: Mutex<Vec<(WindowId, ExternalData)>>,
}
impl TestWaker {
fn take(&self) -> Vec<(WindowId, ExternalData)> {
self.posts.lock().unwrap().drain(..).collect()
}
}
impl Waker for TestWaker {
fn wake(&self) -> bool {
self.wakes.fetch_add(1, Ordering::SeqCst);
true
}
fn post(&self, window: WindowId, data: ExternalData) -> bool {
self.posts.lock().unwrap().push((window, data));
true
}
}
fn win() -> WindowId {
WindowId::new(1)
}
#[test]
fn local_queue_round_trips_without_a_platform() {
let rt = Runtime::new();
assert!(!rt.is_online(), "默认是本地队列模式");
assert!(rt.wake(), "无头下唤醒是 no-op 成功");
let w = win();
let handle = rt.waker();
assert!(handle.is_none(), "无平台 ⇒ 拿不到平台唤醒器");
let slot = rt.inner.waker.borrow().clone();
assert!(slot.post(w, ExternalData::new(7u8)));
let got = rt.take_pending_external();
assert_eq!(got.len(), 1);
assert_eq!(got[0].0, w);
assert_eq!(got[0].1.downcast_ref::<u8>(), Some(&7));
}
#[test]
fn platform_waker_receives_posts_and_marks_online() {
let rt = Runtime::new();
let tw = Arc::new(TestWaker::default());
rt.set_waker(tw.clone());
assert!(rt.is_online());
assert!(rt.wake());
assert_eq!(tw.wakes.load(Ordering::SeqCst), 1);
let slot = rt.inner.waker.borrow().clone();
assert!(slot.post(win(), ExternalData::new("hi")));
assert_eq!(tw.take().len(), 1);
assert!(rt.take_pending_external().is_empty(), "平台模式不落本地队列");
}
#[test]
fn spawn_task_delivers_a_task_event_with_the_payload() {
let rt = Runtime::new();
let tw = Arc::new(TestWaker::default());
rt.set_waker(tw.clone());
rt.register_window(win());
let handle = rt.spawn_task(win(), |ctx| {
ctx.post("progress-note");
Ok::<_, ()>(21u32 * 2)
});
assert_eq!(handle.id(), 1, "任务 id 自增");
assert!(rt.has_tasks(win()));
let deadline = Instant::now() + Duration::from_secs(5);
while !handle.is_done() && Instant::now() < deadline {
std::thread::sleep(Duration::from_millis(2));
}
assert!(handle.is_done(), "任务应在超时前完成");
let deadline = Instant::now() + Duration::from_secs(5);
let mut msgs = Vec::new();
while msgs.len() < 2 && Instant::now() < deadline {
msgs.extend(tw.take());
std::thread::sleep(Duration::from_millis(2));
}
assert_eq!(msgs.len(), 2, "一条中间消息 + 一条完成事件");
let (w, done) = msgs.pop().unwrap();
assert_eq!(w, win());
assert!(!on_task_message(&rt, w, &done), "TaskEvent 要交给用户");
assert!(!rt.has_tasks(win()), "任务表已清");
let ev = done.downcast::<TaskEvent>().expect("是 TaskEvent");
assert_eq!(ev.id, handle.id());
assert_eq!(ev.payload.downcast::<Result<u32, ()>>(), Some(Ok(42)));
let (_, note) = msgs.pop().unwrap();
assert!(!on_task_message(&rt, win(), ¬e));
assert_eq!(note.downcast::<&str>(), Some("progress-note"));
}
#[test]
fn spawn_task_busy_shows_a_busy_item_and_clears_it_on_completion() {
let rt = Runtime::new();
let tw = Arc::new(TestWaker::default());
rt.set_waker(tw.clone());
rt.register_window(win());
let handle = rt.spawn_task_busy(win(), "正在导出…", |ctx| {
ctx.progress(1, 4);
"done"
});
assert!(rt.is_busy(win()), "任务开始 ⇒ 遮罩出现");
let items = rt.busy_items(win());
assert_eq!(items.len(), 1);
assert_eq!(items[0].label, "正在导出…");
assert_eq!(items[0].ratio(), None, "还没进度 ⇒ 不确定进度");
assert!(items[0].is_cancellable(), "带遮罩的任务默认可取消");
let deadline = Instant::now() + Duration::from_secs(5);
let mut msgs = Vec::new();
while (msgs.len() < 2 || !handle.is_done()) && Instant::now() < deadline {
msgs.extend(tw.take());
std::thread::sleep(Duration::from_millis(2));
}
let progress = msgs
.iter()
.position(|(_, d)| d.downcast_ref::<TaskProgress>().is_some())
.expect("有进度消息");
let (_, p) = msgs.remove(progress);
assert!(on_task_message(&rt, win(), &p), "进度消息被框架完全消费");
assert_eq!(rt.busy_items(win())[0].ratio(), Some(0.25));
let (_, done) = msgs.pop().expect("完成事件");
assert!(!on_task_message(&rt, win(), &done));
assert!(!rt.is_busy(win()), "任务完成 ⇒ 遮罩自动收起");
}
#[test]
fn progress_with_carries_a_human_readable_detail_line() {
let rt = Runtime::new();
let tw = Arc::new(TestWaker::default());
rt.set_waker(tw.clone());
rt.register_window(win());
let _handle = rt.spawn_task_busy(win(), "正在打开 3 个文件…", |ctx| {
ctx.progress_with(1, 9, "第 1 / 3 个文件 · 正在读取 a.pdf");
ctx.progress_with(5, 9, "第 2 / 3 个文件 · 正在合并 b.pdf");
"done"
});
let deadline = Instant::now() + Duration::from_secs(5);
let mut msgs = Vec::new();
while msgs.len() < 3 && Instant::now() < deadline {
msgs.extend(tw.take());
std::thread::sleep(Duration::from_millis(2));
}
let mut seen = Vec::new();
for (_, d) in msgs {
if d.downcast_ref::<TaskProgress>().is_some() {
on_task_message(&rt, win(), &d);
let item = rt.busy_items(win()).into_iter().next().expect("遮罩项在");
seen.push((
item.ratio().map(|r| (r * 100.0).round() as i32),
item.detail.clone(),
));
}
}
assert_eq!(
seen,
vec![
(Some(11), Some("第 1 / 3 个文件 · 正在读取 a.pdf".to_string())),
(Some(56), Some("第 2 / 3 个文件 · 正在合并 b.pdf".to_string())),
],
"进度条与明细都跟着上报走"
);
}
#[test]
fn busy_overlay_is_held_for_the_minimum_visible_time() {
let rt = Runtime::new();
let tw = Arc::new(TestWaker::default());
rt.set_waker(tw.clone());
rt.register_window(win());
assert_eq!(rt.busy_min_visible(), Duration::ZERO, "默认不等待");
rt.set_busy_min_visible(Duration::from_millis(300));
let handle = rt.spawn_task_busy(win(), "正在打开…", |_| 0u32);
let deadline = Instant::now() + Duration::from_secs(5);
let mut msgs = Vec::new();
while msgs.is_empty() && Instant::now() < deadline {
msgs.extend(tw.take());
std::thread::sleep(Duration::from_millis(2));
}
assert!(handle.is_done());
let (_, done) = msgs.pop().expect("完成事件");
on_task_message(&rt, win(), &done);
assert!(!rt.has_tasks(win()), "任务已经结束");
assert!(rt.is_busy(win()), "但遮罩还挂着(最短可见时间)");
assert!(!rt.reap_busy(Instant::now()), "还没到点 ⇒ 不收");
assert!(
rt.reap_busy(Instant::now() + Duration::from_millis(400)),
"过了最短可见时间 ⇒ 收掉"
);
assert!(!rt.is_busy(win()), "遮罩收起");
assert!(
!rt.reap_busy(Instant::now() + Duration::from_secs(1)),
"已经收干净了(幂等)"
);
}
#[test]
fn busy_token_can_set_a_detail_line() {
let rt = Runtime::new();
let w = win();
rt.register_window(w);
let token = rt.begin_busy(w, "正在整理…");
token.set_progress(2, 5);
token.set_detail("正在写第 2 个分片");
let item = rt.busy_items(w).into_iter().next().expect("忙碌项");
assert_eq!(item.ratio(), Some(0.4));
assert_eq!(item.detail.as_deref(), Some("正在写第 2 个分片"));
token.finish();
}
#[test]
fn clicking_the_overlay_cancel_button_cancels_the_task() {
let rt = Runtime::new();
let tw = Arc::new(TestWaker::default());
rt.set_waker(tw.clone());
rt.register_window(win());
let handle = rt.spawn_task_busy(win(), "正在导出…", |ctx| {
let deadline = Instant::now() + Duration::from_secs(5);
while !ctx.is_cancelled() && Instant::now() < deadline {
std::thread::sleep(Duration::from_millis(2));
}
ctx.is_cancelled()
});
let click = rt.busy_items(win())[0].cancel.clone().expect("有取消按钮");
click();
assert!(handle.is_cancelled(), "点击 ⇒ 任务被取消");
let deadline = Instant::now() + Duration::from_secs(5);
let mut msgs = Vec::new();
while msgs.is_empty() && Instant::now() < deadline {
msgs.extend(tw.take());
std::thread::sleep(Duration::from_millis(2));
}
let (_, done) = msgs.pop().expect("完成事件");
on_task_message(&rt, win(), &done);
assert!(!rt.is_busy(win()));
}
#[test]
fn a_panicking_task_still_finishes_and_clears_its_overlay() {
let rt = Runtime::new();
let tw = Arc::new(TestWaker::default());
rt.set_waker(tw.clone());
rt.register_window(win());
let handle = rt.spawn_task_busy(win(), "正在解析…", |_ctx| -> u32 {
panic!("任务体崩了")
});
assert!(rt.is_busy(win()));
let deadline = Instant::now() + Duration::from_secs(5);
let mut msgs = Vec::new();
while msgs.is_empty() && Instant::now() < deadline {
msgs.extend(tw.take());
std::thread::sleep(Duration::from_millis(2));
}
assert!(handle.is_done(), "panic 也要走收尾路径(否则遮罩永远挂着)");
let (_, failed) = msgs.pop().expect("失败事件");
assert!(!on_task_message(&rt, win(), &failed), "失败事件也交给用户");
assert_eq!(
failed.downcast::<TaskFailed>().map(|f| f.id),
Some(handle.id()),
"投递的是 TaskFailed"
);
assert!(!rt.is_busy(win()), "失败 ⇒ 遮罩同样收起");
assert!(!rt.has_tasks(win()), "任务表同样清空");
}
#[test]
fn cancelling_a_task_is_visible_inside_the_worker() {
let rt = Runtime::new();
let tw = Arc::new(TestWaker::default());
rt.set_waker(tw.clone());
rt.register_window(win());
let handle = rt.spawn_task(win(), |ctx| {
let deadline = Instant::now() + Duration::from_secs(5);
while !ctx.is_cancelled() && Instant::now() < deadline {
std::thread::sleep(Duration::from_millis(2));
}
ctx.is_cancelled()
});
handle.cancel();
assert!(handle.is_cancelled());
let deadline = Instant::now() + Duration::from_secs(5);
let mut msgs = Vec::new();
while msgs.is_empty() && Instant::now() < deadline {
msgs.extend(tw.take());
std::thread::sleep(Duration::from_millis(2));
}
let (_, done) = msgs.pop().expect("完成事件");
on_task_message(&rt, win(), &done);
let ev = done.downcast::<TaskEvent>().unwrap();
assert_eq!(ev.payload.downcast::<bool>(), Some(true), "线程里看到了取消");
}
#[test]
fn closing_a_window_cancels_its_tasks_and_drops_its_busy_items() {
let rt = Runtime::new();
let tw = Arc::new(TestWaker::default());
rt.set_waker(tw.clone());
rt.register_window(win());
let handle = rt.spawn_task_busy(win(), "正在加载…", |ctx| {
let deadline = Instant::now() + Duration::from_secs(5);
while !ctx.is_cancelled() && Instant::now() < deadline {
std::thread::sleep(Duration::from_millis(2));
}
});
assert!(rt.is_busy(win()));
rt.cancel_tasks_of(win());
assert!(handle.is_cancelled());
assert!(!rt.has_tasks(win()), "任务表清空");
assert!(!rt.is_busy(win()), "遮罩项一并清掉");
}
#[test]
fn busy_token_is_raii_and_updates_progress() {
let rt = Runtime::new();
let w = win();
rt.register_window(w);
{
let busy = rt.begin_busy(w, "正在合并…");
assert!(rt.is_busy(w));
busy.set_label("正在合并 3 个 PDF…");
busy.set_progress(2, 3);
let items = rt.busy_items(w);
assert_eq!(items[0].label, "正在合并 3 个 PDF…");
assert_eq!(items[0].ratio(), Some(2.0 / 3.0));
assert!(rt.take_dirty(w).contains(Dirty::VIEW));
}
assert!(!rt.is_busy(w), "Drop ⇒ 遮罩收起");
assert!(rt.take_dirty(w).contains(Dirty::VIEW), "收起也要刷新一次");
}
#[test]
fn busy_token_can_be_cancelled_from_the_overlay() {
let rt = Runtime::new();
let w = win();
rt.register_window(w);
let hits = Arc::new(AtomicU32::new(0));
let busy = rt.begin_busy(w, "正在导出…");
let h = hits.clone();
busy.cancellable(move || {
h.fetch_add(1, Ordering::SeqCst);
});
let items = rt.busy_items(w);
assert!(items[0].is_cancellable());
let cb = items[0].cancel.clone().unwrap();
cb();
assert_eq!(hits.load(Ordering::SeqCst), 1);
busy.finish();
assert!(!rt.is_busy(w));
}
#[test]
fn multiple_tasks_stack_busy_items_per_window() {
let rt = Runtime::new();
let tw = Arc::new(TestWaker::default());
rt.set_waker(tw.clone());
let a = WindowId::new(1);
let b = WindowId::new(2);
rt.register_window(a);
rt.register_window(b);
let h1 = rt.spawn_task_busy(a, "A 的任务", |_| ());
let h2 = rt.spawn_task_busy(a, "另一个任务", |_| ());
let _h3 = rt.spawn_task_busy(b, "B 的任务", |_| ());
assert_eq!(rt.busy_items(a).len(), 2, "同窗口多个任务各占一项");
assert_eq!(rt.busy_items(b).len(), 1, "遮罩按窗口隔离");
assert!(h1.is_running() && h2.is_running());
let deadline = Instant::now() + Duration::from_secs(5);
while (h1.is_running() || h2.is_running()) && Instant::now() < deadline {
for (w, d) in tw.take() {
if w == a {
on_task_message(&rt, w, &d);
}
}
std::thread::sleep(Duration::from_millis(2));
}
assert!(rt.busy_items(a).is_empty(), "两个任务都完成后遮罩收起");
assert_eq!(rt.busy_items(b).len(), 1, "B 的任务未被消费 ⇒ 遮罩仍在");
}
}