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,
    // last group runs on the main thread
    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);
    }
}