use super::builder_markers::{
NoWait, Once, Open, Rate, ReceiveAll, ReceiveAny, Repeat, Repeatable, Set, Unset, WaitFor,
Waits, Wiring,
};
use crate::{
executor::Executor,
futures::task::Task,
modules::{
forward::{Forwarder, cloned_into},
handle_kind::Waiting,
handle_set::HandleSet,
input::{self, Receives, Standalone},
merge_set::MergeSet,
task_handle::TaskHandle,
task_kind::TaskKind,
task_setup::{Deadline, TaskSetup},
},
};
use std::{
marker::PhantomData,
time::{Duration, Instant},
};
type Forwards = Vec<Box<dyn FnOnce(usize)>>;
#[must_use = "a builder does nothing until `spawn` is called"]
pub struct TaskBuilder<F, K = Once, D = Open, C = Open, W = NoWait>
where
F: Task,
W: Wiring,
{
task: F,
setup: TaskSetup,
link: W::Link,
forwards: Forwards,
_state: PhantomData<(K, D, C, W)>,
}
impl<F> TaskBuilder<F, Once, Open, Open, NoWait>
where
F: Task,
{
pub(crate) fn new(task: F) -> Self {
Self {
task,
setup: TaskSetup::default(),
link: (),
forwards: Vec::new(),
_state: PhantomData,
}
}
}
impl<F, K, D, C, W> TaskBuilder<F, K, D, C, W>
where
F: Task,
W: Wiring,
{
pub fn priority(mut self, priority: u8) -> Self {
self.setup.priority = priority;
self
}
pub fn timeout(mut self, timeout: Duration) -> Self {
self.setup.timeout = Some(timeout);
self
}
pub fn after(mut self, delay: Duration) -> TaskBuilder<F, K, Set, Set, W> {
self.setup.start_delay = delay;
self.moved()
}
pub fn give_to<O>(mut self, target: &TaskHandle<O, Waiting<F::Output>>) -> Self
where
F::Output: Clone,
{
let target = target.retyped::<()>();
self.forwards.push(Box::new(move |id| {
let forward = Forwarder::new(target, cloned_into::<F::Output, F::Output>);
Executor::forward(id, Box::new(forward));
}));
self
}
fn moved<K2, D2, C2>(self) -> TaskBuilder<F, K2, D2, C2, W> {
TaskBuilder {
task: self.task,
setup: self.setup,
link: self.link,
forwards: self.forwards,
_state: PhantomData,
}
}
fn rewire<W2>(self, link: W2::Link) -> TaskBuilder<F, K, D, C, W2>
where
W2: Wiring,
{
TaskBuilder {
task: self.task,
setup: self.setup,
link,
forwards: self.forwards,
_state: PhantomData,
}
}
}
impl<F, D, C, W> TaskBuilder<F, Once, D, C, W>
where
F: Task,
W: Wiring,
{
pub fn repeat(mut self) -> TaskBuilder<F, Repeat, D, W::AfterKind<C>, W> {
self.setup.kind = TaskKind::Repeating;
self.moved()
}
pub fn at_rate(mut self, period: Duration) -> TaskBuilder<F, Rate, D, W::AfterKind<C>, W> {
self.setup.kind = TaskKind::Series;
self.setup.interval = period;
self.moved()
}
}
impl<F, D, C> TaskBuilder<F, Once, D, C, NoWait>
where
F: Task,
{
pub fn wait_for<T>(mut self) -> TaskBuilder<F, Once, D, C, WaitFor<T>> {
self.setup.waits = true;
self.rewire(())
}
pub fn receive<H>(mut self, from: H) -> TaskBuilder<F, Once, D, C, ReceiveAll<H>>
where
H: HandleSet,
{
self.setup.waits = true;
self.rewire(from)
}
pub fn receive_any<H>(mut self, from: H) -> TaskBuilder<F, Once, D, C, ReceiveAny<H>>
where
H: Send + 'static,
{
self.setup.waits = true;
self.rewire(from)
}
}
impl<F, D, C, W> TaskBuilder<F, Once, D, C, W>
where
F: Task,
W: Wiring,
{
pub fn count(mut self, arrivals: u32) -> TaskBuilder<F, Once, D, Set, W>
where
C: Unset,
W: Waits,
{
self.setup.gives = arrivals;
self.moved()
}
}
impl<F, D, C, W> TaskBuilder<F, Repeat, D, C, W>
where
F: Task,
W: Wiring,
{
pub fn every(mut self, gap: Duration) -> Self {
self.setup.kind = TaskKind::RepeatEvery;
self.setup.interval = gap;
self
}
}
impl<F, K, C, W> TaskBuilder<F, K, Open, C, W>
where
F: Task,
K: Repeatable,
W: Wiring,
{
pub fn for_duration(mut self, span: Duration) -> TaskBuilder<F, K, Set, C, W> {
self.setup.deadline = Deadline::Span(span);
self.moved()
}
pub fn until(mut self, when: Instant) -> TaskBuilder<F, K, Set, C, W> {
self.setup.deadline = Deadline::At(when);
self.moved()
}
}
impl<F, K, D, C, W> TaskBuilder<F, K, D, C, W>
where
F: Task,
K: Repeatable,
W: Wiring,
{
pub fn count(mut self, runs: u32) -> TaskBuilder<F, K, D, Set, W>
where
C: Unset,
{
self.setup.runs = runs;
self.moved()
}
}
fn link_forwards(forwards: Forwards, id: usize) {
for forward in forwards {
forward(id);
}
}
impl<F, D, C> TaskBuilder<F, Once, D, C, NoWait>
where
F: Task,
F::Input: Standalone,
{
pub fn spawn(self) -> TaskHandle<F::Output> {
let TaskBuilder {
mut task,
setup,
forwards,
..
} = self;
task.give(input::token(), input::standalone());
let handle = Executor::new_task(task, setup);
link_forwards(forwards, handle.id());
handle
}
}
impl<F, D, C> TaskBuilder<F, Repeat, D, C, NoWait>
where
F: Task,
F::Input: Standalone,
{
pub fn spawn(self) -> TaskHandle<F::Output> {
let TaskBuilder {
mut task,
setup,
forwards,
..
} = self;
task.give(input::token(), input::standalone());
let handle = Executor::new_task(task, setup);
link_forwards(forwards, handle.id());
handle
}
}
impl<F, D, C> TaskBuilder<F, Rate, D, C, NoWait>
where
F: Task + Clone,
F::Input: Standalone,
{
pub fn spawn(self) -> TaskHandle<F::Output> {
let TaskBuilder {
mut task,
setup,
forwards,
..
} = self;
task.give(input::token(), input::standalone());
let handle = Executor::new_series(task, setup);
link_forwards(forwards, handle.id());
handle
}
}
impl<F, D, C, T> TaskBuilder<F, Once, D, C, WaitFor<T>>
where
F: Task,
T: Send + 'static,
{
pub fn spawn<M>(self) -> TaskHandle<F::Output, Waiting<T>>
where
F::Input: Receives<T, M>,
M: 'static,
{
let TaskBuilder {
task,
setup,
forwards,
..
} = self;
let handle = Executor::new_waiting::<F, T, M>(task, setup);
link_forwards(forwards, handle.id());
handle
}
}
impl<F, D, C, T> TaskBuilder<F, Repeat, D, C, WaitFor<T>>
where
F: Task,
T: Send + 'static,
{
pub fn spawn<M>(self) -> TaskHandle<F::Output, Waiting<T>>
where
F::Input: Receives<T, M>,
M: 'static,
{
let TaskBuilder {
task,
setup,
forwards,
..
} = self;
let handle = Executor::new_waiting::<F, T, M>(task, setup);
link_forwards(forwards, handle.id());
handle
}
}
impl<F, D, C, T> TaskBuilder<F, Rate, D, C, WaitFor<T>>
where
F: Task + Clone,
T: Send + 'static,
{
pub fn spawn<M>(self) -> TaskHandle<F::Output, Waiting<T>>
where
F::Input: Receives<T, M>,
M: 'static,
{
let TaskBuilder {
task,
setup,
forwards,
..
} = self;
let handle = Executor::new_waiting_series::<F, T, M>(task, setup);
link_forwards(forwards, handle.id());
handle
}
}
impl<F, D, C, H> TaskBuilder<F, Once, D, C, ReceiveAll<H>>
where
F: Task,
H: HandleSet,
{
pub fn spawn<M>(self) -> TaskHandle<F::Output>
where
F::Input: Receives<H::Output, M>,
M: 'static,
{
let TaskBuilder {
task,
setup,
link,
forwards,
..
} = self;
let waiting = Executor::new_waiting::<F, H::Output, M>(task, setup);
link_forwards(forwards, waiting.id());
Executor::receive_all(waiting, link)
}
}
impl<F, D, C, H> TaskBuilder<F, Repeat, D, C, ReceiveAll<H>>
where
F: Task,
H: HandleSet,
{
pub fn spawn<M>(self) -> TaskHandle<F::Output>
where
F::Input: Receives<H::Output, M>,
M: 'static,
{
let TaskBuilder {
task,
setup,
link,
forwards,
..
} = self;
let waiting = Executor::new_waiting::<F, H::Output, M>(task, setup);
link_forwards(forwards, waiting.id());
Executor::receive_all(waiting, link)
}
}
impl<F, D, C, H> TaskBuilder<F, Rate, D, C, ReceiveAll<H>>
where
F: Task + Clone,
H: HandleSet,
{
pub fn spawn<M>(self) -> TaskHandle<F::Output>
where
F::Input: Receives<H::Output, M>,
M: 'static,
{
let TaskBuilder {
task,
setup,
link,
forwards,
..
} = self;
let waiting = Executor::new_waiting_series::<F, H::Output, M>(task, setup);
link_forwards(forwards, waiting.id());
Executor::receive_all(waiting, link)
}
}
impl<F, D, C, H> TaskBuilder<F, Once, D, C, ReceiveAny<H>>
where
F: Task,
H: Send + 'static,
{
pub fn spawn<M>(self) -> TaskHandle<F::Output>
where
H: MergeSet<F::Input, M>,
F::Input: Receives<H::Given, M>,
M: 'static,
{
let TaskBuilder {
task,
setup,
link,
forwards,
..
} = self;
let waiting = Executor::new_waiting::<F, H::Given, M>(task, setup);
link_forwards(forwards, waiting.id());
Executor::receive_any::<_, F::Input, M, H>(waiting, link)
}
}
impl<F, D, C, H> TaskBuilder<F, Repeat, D, C, ReceiveAny<H>>
where
F: Task,
H: Send + 'static,
{
pub fn spawn<M>(self) -> TaskHandle<F::Output>
where
H: MergeSet<F::Input, M>,
F::Input: Receives<H::Given, M>,
M: 'static,
{
let TaskBuilder {
task,
setup,
link,
forwards,
..
} = self;
let waiting = Executor::new_waiting::<F, H::Given, M>(task, setup);
link_forwards(forwards, waiting.id());
Executor::receive_any::<_, F::Input, M, H>(waiting, link)
}
}
impl<F, D, C, H> TaskBuilder<F, Rate, D, C, ReceiveAny<H>>
where
F: Task + Clone,
H: Send + 'static,
{
pub fn spawn<M>(self) -> TaskHandle<F::Output>
where
H: MergeSet<F::Input, M>,
F::Input: Receives<H::Given, M>,
M: 'static,
{
let TaskBuilder {
task,
setup,
link,
forwards,
..
} = self;
let waiting = Executor::new_waiting_series::<F, H::Given, M>(task, setup);
link_forwards(forwards, waiting.id());
Executor::receive_any::<_, F::Input, M, H>(waiting, link)
}
}
#[cfg(test)]
mod type_checks {
use super::*;
use crate::{
Runtime,
compute::Compute,
sleep::{Sleep, SleepMode, SleepTask},
};
fn task() -> SleepTask {
Sleep::sleep(Duration::from_millis(1)).mode(SleepMode::Relaxed)
}
#[allow(dead_code)]
fn one_shots() {
let _ = TaskBuilder::new(task()).spawn();
let _ = TaskBuilder::new(task()).priority(200).spawn();
let _ = TaskBuilder::new(task()).after(Duration::ZERO).spawn();
let _ = TaskBuilder::new(task())
.after(Duration::ZERO)
.priority(200)
.spawn();
}
#[allow(dead_code)]
fn repeats() {
let _ = TaskBuilder::new(task()).repeat().spawn();
let _ = TaskBuilder::new(task())
.repeat()
.every(Duration::ZERO)
.spawn();
let _ = TaskBuilder::new(task()).repeat().count(10).spawn();
let _ = TaskBuilder::new(task())
.repeat()
.for_duration(Duration::ZERO)
.spawn();
let _ = TaskBuilder::new(task())
.repeat()
.until(Instant::now())
.spawn();
}
#[allow(dead_code)]
fn timeouts() {
let _ = TaskBuilder::new(task()).timeout(Duration::ZERO).spawn();
let _ = TaskBuilder::new(task())
.timeout(Duration::ZERO)
.repeat()
.count(10)
.spawn();
let _ = TaskBuilder::new(task())
.repeat()
.every(Duration::ZERO)
.after(Duration::ZERO)
.timeout(Duration::ZERO)
.spawn();
let _ = TaskBuilder::new(task())
.at_rate(Duration::ZERO)
.timeout(Duration::ZERO)
.spawn();
let _ = TaskBuilder::new(Compute::compute(|value: u8| value))
.wait_for::<u8>()
.timeout(Duration::ZERO)
.count(3)
.spawn();
}
#[allow(dead_code)]
fn both_bounds() {
let _ = TaskBuilder::new(task())
.repeat()
.count(10)
.for_duration(Duration::ZERO)
.spawn();
let _ = TaskBuilder::new(task())
.repeat()
.until(Instant::now())
.count(10)
.spawn();
}
#[allow(dead_code)]
fn delays() {
let _ = TaskBuilder::new(task())
.repeat()
.count(10)
.after(Duration::ZERO)
.spawn();
let _ = TaskBuilder::new(task())
.repeat()
.after(Duration::ZERO)
.every(Duration::ZERO)
.spawn();
let _ = TaskBuilder::new(task())
.after(Duration::ZERO)
.repeat()
.every(Duration::ZERO)
.spawn();
}
#[allow(dead_code)]
fn computes() {
let _ = TaskBuilder::new(Compute::compute(|()| 1)).spawn();
let _ = TaskBuilder::new(Compute::compute(|()| 1).blocking())
.priority(200)
.spawn();
let _ = TaskBuilder::new(Compute::compute(|()| 1))
.repeat()
.every(Duration::ZERO)
.count(3)
.spawn();
let _ = TaskBuilder::new(Compute::compute(|()| 1))
.at_rate(Duration::ZERO)
.count(3)
.spawn();
let _ = TaskBuilder::new(Compute::compute(|()| 1))
.after(Duration::ZERO)
.spawn();
let _: i32 = Runtime::block(Compute::compute(|()| 1));
}
#[allow(dead_code)]
fn waits() {
let _ = TaskBuilder::new(Compute::compute(|value: i32| value))
.wait_for::<i32>()
.spawn();
let _ = TaskBuilder::new(Compute::compute(|value: i32| value))
.wait_for::<i32>()
.count(3)
.spawn();
let _ = TaskBuilder::new(Compute::compute(|value: i32| value))
.wait_for::<i32>()
.count(3)
.repeat()
.count(5)
.every(Duration::ZERO)
.spawn();
let _ = TaskBuilder::new(Compute::compute(|value: i32| value))
.wait_for::<i32>()
.at_rate(Duration::ZERO)
.for_duration(Duration::ZERO)
.spawn();
let _ = TaskBuilder::new(Compute::compute(|()| 1))
.wait_for::<()>()
.priority(200)
.after(Duration::ZERO)
.spawn();
let _ = TaskBuilder::new(task())
.wait_for::<String>()
.repeat()
.spawn();
let handle = TaskBuilder::new(Compute::compute(|value: i32| value))
.wait_for::<i32>()
.spawn();
let _ = handle.give(1);
}
#[allow(dead_code)]
fn receives() {
let a = TaskBuilder::new(Compute::compute(|()| 1u8)).spawn();
let b = TaskBuilder::new(Compute::compute(|()| 2u16)).spawn();
let c = TaskBuilder::new(Compute::compute(|()| String::new())).spawn();
let _ = TaskBuilder::new(Compute::compute(|(a, (b, c)): (u8, (u16, String))| {
(a, b, c)
}))
.receive((a.clone(), (b.clone(), c.clone())))
.count(1)
.spawn();
let _ = TaskBuilder::new(Compute::compute(|all: Vec<u8>| all))
.receive(vec![a.clone(), a.clone()])
.repeat()
.count(2)
.spawn();
let _ = TaskBuilder::new(Compute::compute(|pair: [u8; 2]| pair))
.receive([a.clone(), a.clone()])
.at_rate(Duration::ZERO)
.spawn();
let _ = TaskBuilder::new(task())
.receive((a.clone(), c.clone()))
.spawn();
let _ = TaskBuilder::new(Compute::compute(|value: u64| value))
.receive_any((a.clone(), b.clone()))
.spawn();
let _ = TaskBuilder::new(task())
.receive_any(vec![c.clone()])
.spawn();
let waiting = TaskBuilder::new(Compute::compute(|value: u8| value))
.wait_for::<u8>()
.spawn();
let _ = TaskBuilder::new(Compute::compute(|()| 3u8))
.give_to(&waiting)
.give_to(&waiting)
.repeat()
.spawn();
}
#[allow(dead_code)]
fn schedules() {
let _ = TaskBuilder::new(task()).at_rate(Duration::ZERO).spawn();
let _ = TaskBuilder::new(task())
.at_rate(Duration::ZERO)
.count(10)
.spawn();
let _ = TaskBuilder::new(task())
.at_rate(Duration::ZERO)
.for_duration(Duration::ZERO)
.priority(200)
.spawn();
}
}