1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
// use std::{collections::HashMap, thread};
// use crossbeam::channel::{Receiver, Sender};
// use crate::nodes;
// pub struct Scheduler <Ctx>{
// nodes: HashMap<String, Box<dyn nodes::TNode<Value = Ctx> + Sync>>,
// nodes_inputs: HashMap<String, Vec<Receiver<Ctx>>>,
// nodes_outputs: HashMap<String, Vec<Receiver<Ctx>>>,
// nodes_senders: HashMap<String, Vec<Sender<Ctx>>>
// }
// impl<Ctx: Send> Scheduler <Ctx>{
// pub fn new() -> Self {
// Scheduler {
// nodes: HashMap::new(),
// nodes_inputs: HashMap::new(),
// nodes_outputs: HashMap::new(),
// nodes_senders: HashMap::new()
// }
// }
// /// 一个node可以依赖前序多个节点的输入
// pub fn add_node(&mut self, last: bool, pre_name_idxs: &Vec<(&str, usize)>, node: Box<dyn nodes::TNode<Value = Ctx> + Sync>) -> Option<Vec<Receiver<Ctx>>>{
// let mut final_receivers = None;
// for (pre_name, pre_idx) in pre_name_idxs {
// let pre_dep = &self.nodes_outputs.get(*pre_name).unwrap()[*pre_idx];
// self.nodes_inputs.entry(node.name().to_string())
// .and_modify(|v| v.push(pre_dep.clone()))
// .or_insert(vec![pre_dep.clone()]);
// }
// self.nodes_inputs.entry(node.name().to_string()).or_insert(vec![]);
// let mut senders = vec![];
// let mut receivers = vec![];
// for cap in node.channel_caps() {
// let (s, r) = crossbeam::channel::bounded(*cap);
// senders.push(s);
// receivers.push(r);
// }
// self.nodes_senders.insert(node.name().to_string(), senders);
// self.nodes_outputs.insert(node.name().to_string(), receivers);
// if last {
// final_receivers = Some(self.nodes_outputs.get(node.name()).unwrap()
// .iter()
// .map(|v| v.clone())
// .collect::<Vec<Receiver<Ctx>>>());
// }
// self.nodes.insert(node.name().to_string(), node);
// final_receivers
// }
// pub fn start(&mut self) {
// thread::scope(|s| {
// for (_, node) in & self.nodes {
// let senders = self.nodes_senders.get(node.name()).unwrap();
// let inputs = self.nodes_inputs.get(node.name()).unwrap();
// for _ in 0..node.num_workers() {
// let senders = senders.iter()
// .map(|v| v.clone())
// .collect::<Vec<Sender<Ctx>>>();
// let inputs = inputs.iter()
// .map(|v| v.clone())
// .collect::<Vec<Receiver<Ctx>>>();
// s.spawn(move || {
// node.work(inputs, senders);
// });
// }
// self.nodes_senders.get_mut(node.name()).unwrap().clear();
// self.nodes_outputs.get_mut(node.name()).unwrap().clear();
// self.nodes_inputs.get_mut(node.name()).unwrap().clear();
// }
// })
// }
// }