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
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
pub mod prelude;
mod handler;
mod message;
mod schedule;
mod util;
mod wait;
use crate::handler::{
InitMessageHandlerBuilder, MessageGroup, MessageGroupBuilder, OpenMessageHandlerBuilder,
};
use crate::message::MessageRegistry;
use crate::schedule::{MessageScheduler, ShutdownSwitch};
use crate::wait::WaitingStrategy;
use std::sync::Arc;
use std::time::Duration;
pub struct RuntimeBuilder {
cpu_affinity: bool,
min_bus_capacity: usize,
shutdown_timeout: Duration,
waiting_strategy: WaitingStrategy,
registry: MessageRegistry,
groups: Vec<MessageGroup>,
}
impl RuntimeBuilder {
pub fn with_cpu_affinity(mut self, cpu_affinity: bool) -> Self {
self.cpu_affinity = cpu_affinity;
self
}
pub fn with_waiting_strategy(mut self, waiting_strategy: WaitingStrategy) -> Self {
self.waiting_strategy = waiting_strategy;
self
}
pub fn add<T: 'static, FH>(self, handler: FH) -> Self
where
FH: FnOnce(OpenMessageHandlerBuilder<T>) -> InitMessageHandlerBuilder<T>,
{
self.add_group(|mut group| {
let hb = group.register(handler);
group.init(move |recv, mut context| {
let mut h = hb.finish(&context).unwrap();
recv.stream(move |message| {
h.handle(&mut context, message);
});
})
})
}
pub fn add_group<F>(mut self, group: F) -> Self
where
F: FnOnce(MessageGroupBuilder) -> MessageGroup,
{
let group = group(MessageGroupBuilder::new(&mut self.registry));
self.groups.push(group);
self
}
pub fn finish_main<T: 'static, FH>(self, handler: FH) -> Runtime
where
FH: FnOnce(OpenMessageHandlerBuilder<T>) -> InitMessageHandlerBuilder<T>,
{
self.add(handler).finish()
}
pub fn finish_main_group<F>(self, group: F) -> Runtime
where
F: FnOnce(MessageGroupBuilder) -> MessageGroup,
{
self.add_group(group).finish()
}
pub fn finish(self) -> Runtime {
let registry = Arc::new(self.registry);
let scheduler = MessageScheduler::new(
registry,
self.groups,
self.cpu_affinity,
self.min_bus_capacity,
self.shutdown_timeout,
self.waiting_strategy,
);
Runtime { scheduler }
}
}
pub struct Runtime {
scheduler: MessageScheduler,
}
impl Runtime {
pub fn builder(min_bus_capacity: usize) -> RuntimeBuilder {
RuntimeBuilder {
cpu_affinity: false,
min_bus_capacity,
shutdown_timeout: Duration::from_secs(10),
waiting_strategy: Default::default(),
registry: MessageRegistry::default(),
groups: Default::default(),
}
}
pub fn get_shutdown_switch(&self) -> ShutdownSwitch {
self.scheduler.get_shutdown_switch()
}
pub fn start<E: 'static + Send + Sync>(self, init: E) {
log::info!("start runtime");
self.scheduler.schedule(init);
}
}