Skip to main content

simple/
simple.rs

1use taskorch::{Pool, Queue, TaskBuildNew};
2// [A]      => [B1, B2] ## 1->N
3// [B1, B2] => [Exit]   ## N->1
4fn main() {
5    println!("----- test task orch -----");
6
7    // Step#1. create a Pool
8    let mut pool = Pool::new();
9
10    // Step#2. create a queue
11    let qid = pool.insert_queue(&Queue::new()).unwrap();
12    let submitter = pool.task_submitter(qid).unwrap();
13
14    // Step#3. create tasks
15
16    // an indepent task
17    let task = (|| println!("task='free':  Hello, 1 2 3 ..")).into_task();
18    let _ = submitter.submit(task);
19
20    // an exit task with cond(#0 i32, #1 str)
21    let exit = submitter
22        .submit(
23            (|a: i32, msg: &str| println!("task='exit': received ({a},{msg:?}) and EXIT"))
24                .into_exit_task(),
25        )
26        .take();
27
28    // N->1 : pass i32 to exit-task.p0
29    let b1 = (|a: i32| {
30        println!("task='B1':  pass ('{a}') to task='exit'");
31        a
32    })
33    .into_task()
34    .bind_to(exit.input_ca::<0>());
35    let b1 = submitter.submit(b1).take();
36
37    // N->1 : pass str to exit task.p1
38    let b2 = (|msg: &'static str| {
39        println!("task='B2':  recv ('{msg}') and then pass ('{msg}') to task='exit'");
40        msg
41    })
42    .into_task()
43    .bind_to(exit.input_ca::<1>());
44    let b2 = submitter.submit(b2).take();
45
46    // 1->N : map result to task-b1 and task-b2
47    let b3 = (||())
48        .into_task()
49        .map_tuple_with(move |_: ()| {
50            println!("task='A': map `()=>(i32,&str)` and then pass (10,'exit') to task=['B1','B2']");
51            (10, "exit")
52        })
53        .bind_all_to((b1.input_ca::<0>(), b2.input_ca::<0>()));
54    let _ = submitter.submit(b3);
55
56    // Step#4. start a thread and run
57    pool.spawn_thread_for(qid);
58
59    // Step#5. wait until all finished
60    pool.join();
61}