Skip to main content

spsc/
spsc.rs

1
2use taskorch::{Pool, Queue, TaskBuildNew, TaskSubmitter};
3
4// Thread 1: Task execution (consumer) role
5// Thread 2: Task generation (producer) role
6
7fn main() {
8    println!("----- test task orch -----");
9
10    // Step#1. create a Pool
11    let mut pool = Pool::new();
12
13    // Step#2. create a queue
14    let qid1 = pool.insert_queue(&Queue::new()).unwrap();
15    let submitter1 = pool.task_submitter(qid1).unwrap();
16    // Step#4. start a thread and run
17    pool.spawn_thread_for(qid1);
18
19    // Step#3. create tasks
20    consume_task_prompt(&submitter1);
21
22    std::thread::spawn(||{
23        produce_task(submitter1);
24    });
25
26    // Step#5. wait until all finished
27    pool.join();
28}
29
30fn consume_task_prompt(submitter:&TaskSubmitter) {
31    submitter.submit((||println!("Init: waiting task to do")).into_task());
32}
33
34fn produce_task(submitter:TaskSubmitter) {
35    prompt("hello");
36    submitter.submit((||println!("consmue task='hello': hello everyone!")).into_task());
37
38    prompt("exit");
39    let id_exit = submitter.submit((|a:i32|println!("consume task='exit': recv cond={a} and exit.")).into_exit_task())
40        .take();
41
42    prompt("add");
43    let id_add = submitter.submit(
44        (|a:i32,b:i32|{
45            println!("consume task='add': ({a},{b}) and then pass (r={}) to Task='exit'",a+b);
46            a+b
47        }).into_task().bind_to(id_exit.input_ca::<0>())
48    ).take();
49
50    prompt("params");
51    let _ = submitter.submit(
52        (||{println!("consume task='params': pass (1,2) to task='add'");1},10.into())
53        .into_task()
54        .map_tuple_with(move|_:i32| (1,2))
55        .bind_all_to((id_add.input_ca::<0>(),id_add.input_ca::<1>()))
56    );
57}
58
59fn prompt(taskname:&'static str) {
60    println!("produce task='{taskname}'.");
61}