use futures::future::{pending, FutureExt}; use tokio::sync::oneshot; use tokio::task::LocalSet; use tokio::time::Duration; use tokio_util::task::JoinMap;
// Spawn `N` tasks that return their index (`i`). fn spawn_index_tasks(map: &mut JoinMap<usize, usize>, n: usize, on: Option<&LocalSet>) { for i in0..n { let rc = std::rc::Rc::new(i); match on {
None => map.spawn_local(i, asyncmove { *rc }),
Some(local) => map.spawn_local_on(i, asyncmove { *rc }, local),
};
}
}
// Spawn `N` “pending” tasks that own a `oneshot::Sender`. // When the task is aborted the sender is dropped, which is observed // via the returned `Receiver`s. fn spawn_pending_tasks(
map: &mut JoinMap<usize, ()>,
receivers: &mut Vec<oneshot::Receiver<()>>,
n: usize,
on: Option<&LocalSet>,
) { for i in0..n { let (tx, rx) = oneshot::channel::<()>();
receivers.push(rx);
let fut = asyncmove {
pending::<()>().await;
drop(tx);
}; match on {
None => map.spawn_local(i, fut),
Some(local) => map.spawn_local_on(i, fut, local),
};
}
}
/// Await every task in JoinMap and assert every task returns its own key. asyncfn drain_joinmap_and_assert(mut map: JoinMap<usize, usize>, n: usize) { letmut seen = vec![false; n]; whilelet Some((k, res)) = map.join_next().await { let v = res.expect("task panicked");
assert_eq!(k, v);
seen[v] = true;
}
assert!(seen.into_iter().all(|b| b));
assert!(map.is_empty());
}
// Await every receiver and assert they all return `Err` because the // corresponding sender (inside an aborted task) was dropped. asyncfn await_receivers_and_assert(receivers: Vec<oneshot::Receiver<()>>) { for rx in receivers {
assert!(
rx.await.is_err(), "task should have been aborted and sender dropped"
);
}
}
let (key, res) = rt().block_on(map.join_next()).unwrap();
assert_eq!(key, "key");
assert!(res.unwrap_err().is_cancelled());
}
// This ensures that `join_next` works correctly when the coop budget is // exhausted. #[tokio::test(flavor = "current_thread")] asyncfn join_map_coop() { // Large enough to trigger coop. const TASK_NUM: u32 = 1000;
for i in0..TASK_NUM {
map.spawn(i, asyncmove {
SEM.add_permits(1);
i
});
}
// Wait for all tasks to complete. // // Since this is a `current_thread` runtime, there's no race condition // between the last permit being added and the task completing. let _ = SEM.acquire_many(TASK_NUM).await.unwrap();
for i in0..5 {
map.spawn(i, futures::future::pending());
} for i in5..10 {
map.spawn(i, async {
tokio::time::sleep(Duration::from_secs(1)).await;
});
}
// The join map will now have 5 pending tasks and 5 ready tasks.
tokio::time::sleep(Duration::from_secs(2)).await;
tokio::select! {
biased;
res = map.join_next() => match res {
Some((_key, res)) => panic!("Task {res:?} exited."),
None => panic!("Phantom task completion."),
},
() = tokio::task::yield_now() => {},
}
send.send(()).unwrap();
let (key, res) = map.join_next().await.unwrap();
assert_eq!(key, 1);
assert_eq!(res.unwrap(), 2);
assert!(map.join_next().await.is_none());
}
#[cfg_attr(not(panic = "unwind"), ignore)] #[tokio::test] asyncfn duplicate_keys_drop() { #[derive(Hash, Debug, PartialEq, Eq)] struct Key; impl Drop for Key { fn drop(&mutself) {
panic!("drop called for key");
}
}
let (send, recv) = oneshot::channel::<()>();
letmut map = JoinMap::new();
map.spawn(Key, async { recv.await.unwrap() });
// replace the task, force it to drop the key and abort the task // we should expect it to panic when dropping the key. let _ = std::panic::catch_unwind(AssertUnwindSafe(|| map.spawn(Key, async {}))).unwrap_err();
// don't panic when this key drops. let (key, _) = map.join_next().await.unwrap();
std::mem::forget(key);
// original task should have been aborted, so the sender should be dangling.
assert!(send.is_closed());
assert!(map.join_next().await.is_none());
}
mod spawn_local { usesuper::*;
#[test] #[should_panic(
expected = "`spawn_local` called from outside of a `task::LocalSet` or `runtime::LocalRuntime`"
)] fn panic_outside_any_runtime() { letmut map = JoinMap::new();
map.spawn_local((), async {});
}
#[tokio::test(flavor = "multi_thread")] #[should_panic(
expected = "`spawn_local` called from outside of a `task::LocalSet` or `runtime::LocalRuntime`"
)] asyncfn panic_in_multi_thread_runtime() { letmut map = JoinMap::new();
map.spawn_local((), async {});
}
#[cfg(tokio_unstable)] mod local_runtime { usesuper::*;
/// Spawn several tasks, and then join all tasks. #[tokio::test(flavor = "local")] asyncfn spawn_then_join_next() { const N: usize = 8;
letmut map = JoinMap::new();
spawn_index_tasks(&mut map, N, None);
/// Spawn several pending-forever tasks, and then drop the [`JoinMap`]. #[tokio::test(flavor = "local")] asyncfn spawn_then_drop() { const N: usize = 8;
/// Spawn several tasks, and then join all tasks. #[tokio::test(flavor = "current_thread")] asyncfn spawn_then_join_next() { const N: usize = 8; let local = LocalSet::new();
local
.run_until(asyncmove { letmut map = JoinMap::new();
spawn_index_tasks(&mut map, N, None);
drain_joinmap_and_assert(map, N).await;
})
.await;
}
/// Spawn several pending-forever tasks, and then shutdown the [`JoinMap`]. #[tokio::test(flavor = "current_thread")] asyncfn spawn_then_shutdown() { const N: usize = 8; let local = LocalSet::new();
/// Spawn several pending-forever tasks, and then drop the [`JoinMap`]. #[tokio::test(flavor = "current_thread")] asyncfn spawn_then_drop() { const N: usize = 8; let local = LocalSet::new();
#[cfg(tokio_unstable)] mod local_runtime { usesuper::*;
/// Spawn several tasks, and then join all tasks. #[tokio::test(flavor = "local")] asyncfn spawn_then_join_next() { const N: usize = 8;
let local = LocalSet::new(); letmut map = JoinMap::new();
spawn_index_tasks(&mut map, N, Some(&local));
assert!(map.join_next().now_or_never().is_none());
local
.run_until(asyncmove {
drain_joinmap_and_assert(map, N).await;
})
.await;
}
}
mod local_set { usesuper::*;
/// Spawn several tasks, and then join all tasks. #[tokio::test(flavor = "current_thread")] asyncfn spawn_then_join_next() { const N: usize = 8; let local = LocalSet::new(); letmut pending_map = JoinMap::new();
spawn_index_tasks(&mut pending_map, N, Some(&local));
assert!(pending_map.join_next().now_or_never().is_none());
local
.run_until(asyncmove {
drain_joinmap_and_assert(pending_map, N).await;
})
.await;
}
/// Spawn several pending-forever tasks, and then shutdown the [`JoinMap`]. #[tokio::test(flavor = "current_thread")] asyncfn spawn_then_shutdown() { const N: usize = 8; let local = LocalSet::new(); letmut map = JoinMap::new(); letmut receivers = Vec::new();
spawn_pending_tasks(&mut map, &mut receivers, N, Some(&local));
assert!(map.join_next().now_or_never().is_none());
local
.run_until(asyncmove {
map.shutdown().await;
assert!(map.is_empty());
await_receivers_and_assert(receivers).await;
})
.await;
}
/// Spawn several pending-forever tasks and then drop the [`JoinMap`] /// before the `LocalSet` is driven and while the `LocalSet` is already driven. #[tokio::test(flavor = "current_thread")] asyncfn spawn_then_drop() { const N: usize = 8;
{ let local = LocalSet::new(); letmut map = JoinMap::new(); letmut receivers = Vec::new();
spawn_pending_tasks(&mut map, &mut receivers, N, Some(&local));
assert!(map.join_next().now_or_never().is_none());
drop(map);
local
.run_until(asyncmove { await_receivers_and_assert(receivers).await })
.await;
}
{ let local = LocalSet::new(); letmut map = JoinMap::new(); letmut receivers = Vec::new();
spawn_pending_tasks(&mut map, &mut receivers, N, Some(&local));
assert!(map.join_next().now_or_never().is_none());
Die Informationen auf dieser Webseite wurden
nach bestem Wissen sorgfältig zusammengestellt. Es wird jedoch weder Vollständigkeit, noch Richtigkeit,
noch Qualität der bereit gestellten Informationen zugesichert.
Bemerkung:
Die farbliche Syntaxdarstellung und die Messung sind noch experimentell.