use super::{
Attr, FnOnceFuture, JoinHandle, Runtime, SetAffinity, Sleep, TaskRef, WaitEvent, Worker, Yield,
};
use crate::{Error, Result};
use core::future::Future;
use core::time::Duration;
pub fn block_on_with<T>(future: T, attr: &Attr) -> Result<T::Output>
where
T: Future,
{
if let Ok(task) = sched(future, attr) {
JoinHandle::<T::Output>::new(task).join()
} else {
Err(Error::default())
}
}
pub fn block_on<T>(future: T) -> Result<T::Output>
where
T: Future,
{
block_on_with(future, &Attr::default())
}
pub fn spawn_with<T>(future: T, attr: &Attr) -> JoinHandle<T::Output>
where
T: Future + Send + 'static,
T::Output: Send + 'static,
{
if let Ok(task) = sched(future, attr) {
JoinHandle::<T::Output>::new(task)
} else {
JoinHandle::<T::Output>::null()
}
}
pub fn spawn<T>(future: T) -> JoinHandle<T::Output>
where
T: Future + Send + 'static,
T::Output: Send + 'static,
{
spawn_with(future, &Attr::default())
}
pub fn spawn_fn<F, R>(f: F) -> JoinHandle<R>
where
F: FnOnce() -> R + Send + 'static,
R: Send + 'static,
{
spawn(FnOnceFuture::new(f))
}
pub fn spawn_fn_with<F, R>(f: F, attr: &Attr) -> JoinHandle<R>
where
F: FnOnce() -> R + Send + 'static,
R: Send + 'static,
{
spawn_with(FnOnceFuture::new(f), attr)
}
pub fn spawn_local_with<T>(future: T, attr: &Attr) -> JoinHandle<T::Output>
where
T: Future + 'static,
T::Output: 'static,
{
if let Some(worker) = Worker::current() {
if let Ok(task) = worker.spawn_local(future, attr) {
return JoinHandle::<T::Output>::new(task);
}
}
JoinHandle::<T::Output>::null()
}
pub fn spawn_local<T>(future: T) -> JoinHandle<T::Output>
where
T: Future + 'static,
T::Output: 'static,
{
spawn_local_with(future, &Attr::default())
}
pub fn spawn_fn_local_with<F, R>(f: F, attr: &Attr) -> JoinHandle<R>
where
F: FnOnce() -> R + 'static,
R: 'static,
{
spawn_local_with(FnOnceFuture::new(f), attr)
}
pub fn spawn_fn_local<F, R>(f: F) -> JoinHandle<R>
where
F: FnOnce() -> R + 'static,
R: 'static,
{
spawn_local(FnOnceFuture::new(f))
}
pub async fn sleep(timeout: Duration) {
Sleep::new(timeout).await
}
pub async fn yield_now() {
Yield::new().await
}
pub async fn wait_event(fd: i32, events: u32) -> Result<u32> {
WaitEvent::new(fd, events).await
}
pub async fn set_affinity() {
SetAffinity.await
}
fn sched<T: Future>(future: T, attr: &Attr) -> Result<TaskRef> {
if let Some(worker) = Worker::current() {
if worker.group_id() == attr.group_id {
return worker.spawn(future, attr);
}
}
Runtime::get(attr.group_id).spawn(future, attr)
}
#[cfg(test)]
mod test {
use crate::runtime::*;
use core::time::Duration;
#[test]
fn test_future() {
let _ = Builder::new().nth(1).build();
async fn test_foo() -> i32 {
100
}
let val = spawn(test_foo()).join().unwrap();
assert_eq!(val, 100);
}
#[test]
fn test_sleep() {
let _ = Builder::new().nth(1).build();
async fn test_sleep(val: i32) -> i32 {
sleep(core::time::Duration::new(1, 0)).await;
yield_now().await;
val + 100
}
let val = spawn(test_sleep(100)).join().unwrap();
assert_eq!(val, 200);
}
#[test]
fn test_ready() {
let _ = Builder::new().nth(1).build();
let val = spawn(async {
let mut retn: i32 = 0;
async { 2_i32 }.ready(|result: i32| retn = result).await;
retn
})
.join()
.unwrap();
assert_eq!(val, 2);
}
#[test]
fn test_select_any() {
let _ = Builder::new().nth(1).build();
let (r1, r2) = spawn(async {
let mut r1: i32 = 0;
let mut r2: i32 = 0;
let t1 = async { 1_i32 }.ready(|result: i32| r1 = result);
let t2 = async { 2_i32 }.ready(|result: i32| r2 = result);
t1.or(t2).await;
(r1, r2)
})
.join()
.unwrap();
assert_eq!(r1, 1);
assert_eq!(r2, 0);
}
#[test]
fn test_select_all() {
let _ = Builder::new().nth(1).build();
let (r1, r2) = spawn(async {
let mut r1: i32 = 0;
let mut r2: i32 = 0;
let t1 = async { 1_i32 }.ready(|result: i32| r1 = result);
let t2 = async { 2_i32 }.ready(|result: i32| r2 = result);
t1.and(t2).await;
(r1, r2)
})
.join()
.unwrap();
assert_eq!(r1, 1);
assert_eq!(r2, 2);
}
#[test]
fn test_or_ready() {
let _ = Builder::new().nth(1).build();
let (r1, r2, r3) = spawn(async {
async fn foo() -> i32 {
1
}
async fn bar() -> &'static str {
"hello"
}
async fn baz() -> i32 {
2
}
let mut foo_retn = 0;
let mut bar_retn = "";
let mut baz_retn = 0;
foo()
.ready(|val| foo_retn = val)
.or(bar().ready(|val| bar_retn = val))
.or(baz().ready(|val| baz_retn = val))
.await;
(foo_retn, bar_retn, baz_retn)
})
.join()
.unwrap();
assert_eq!(r1, 1);
assert_eq!(r2, "");
assert_eq!(r3, 0);
}
#[test]
fn test_and_ready() {
let _ = Builder::new().nth(1).build();
let (r1, r2, r3) = spawn(async {
async fn foo() -> i32 {
1
}
async fn bar() -> &'static str {
"hello"
}
async fn baz() -> i32 {
2
}
let mut foo_retn = 0;
let mut bar_retn = "";
let mut baz_retn = 0;
foo()
.ready(|val| foo_retn = val)
.and(bar().ready(|val| bar_retn = val))
.or(baz().ready(|val| baz_retn = val))
.await;
(foo_retn, bar_retn, baz_retn)
})
.join()
.unwrap();
assert_eq!(r1, 1);
assert_eq!(r2, "hello");
assert_eq!(r3, 0);
}
#[test]
fn test_deadline() {
let _ = Builder::new().nth(1).build();
let beg = crate::time::now();
async fn foo() -> Option<i32> {
async fn fun() -> i32 {
sleep(Duration::new(3, 0)).await;
return 1;
}
fun().deadline(Duration::new(1, 0)).await
}
let x = spawn(foo()).join().unwrap();
assert_eq!(None, x);
let end = crate::time::now();
let elapse = end - beg;
assert!(elapse < Duration::new(2, 1000));
}
}