use des_sim::context::{EventContext, SourceContext, UserContext};
use des_sim::execution::Engine;
use des_sim::execution::runner::Runner;
use des_sim::execution::runner::instance::{ParallelModel, ParallelRunner};
use des_sim::modeling::event::{Event, EventPriority};
use des_sim::modeling::hook::instance::{ModelSummary, TraceHook};
use des_sim::modeling::model::Model;
use des_sim::modeling::source::Source;
use des_sim::primitive::time::Duration;
use std::collections::VecDeque;
use std::fmt;
use std::sync::mpsc::Sender;
#[cfg(test)]
mod tests {
#[test]
fn example_runs() {
super::main();
}
}
#[derive(Debug, Clone)]
pub enum MyEvent {
JobArrived { job_id: u32 },
JobProcessed { job_id: u32 },
JobProcessNext,
}
#[derive(Debug)]
pub struct ServerModel {
pub name: &'static str,
pub queue: VecDeque<u32>,
pub is_busy: bool,
}
#[derive(Debug)]
pub enum ServerCommand {
EnqueueJob { job_id: u32 },
ProcessNextOrIdle,
}
impl Model<MyEvent> for ServerModel {
fn handle_event(
&mut self,
_context: &mut EventContext<MyEvent, Self>,
_event: &Event<MyEvent>,
) {
}
}
impl ParallelModel<MyEvent, ServerCommand> for ServerModel {
fn handle_event_parallel(&self, event: Event<MyEvent>, sender: Sender<ServerCommand>) {
match event.payload {
MyEvent::JobArrived { job_id } => {
if self.is_busy {
std::thread::sleep(std::time::Duration::from_millis(700));
} else {
std::thread::sleep(std::time::Duration::from_millis(200));
}
sender.send(ServerCommand::EnqueueJob { job_id }).unwrap();
}
MyEvent::JobProcessed { job_id: _ } => {
sender.send(ServerCommand::ProcessNextOrIdle).unwrap();
}
MyEvent::JobProcessNext => {
sender.send(ServerCommand::ProcessNextOrIdle).unwrap();
}
}
}
fn apply_command(&mut self, context: &mut EventContext<MyEvent, Self>, command: ServerCommand) {
match command {
ServerCommand::EnqueueJob { job_id } => {
const WORKER_COUNT: usize = 3;
let can_process_next = self.queue.len() < WORKER_COUNT;
self.queue.push_back(job_id);
if !self.is_busy {
self.is_busy = true;
}
if can_process_next {
context.schedule_event(
Duration::ticks(0),
EventPriority::minimum(),
MyEvent::JobProcessNext,
);
}
}
ServerCommand::ProcessNextOrIdle => {
if let Some(next_id) = self.queue.pop_front() {
self.is_busy = true;
context.schedule_event(
Duration::ticks(2),
EventPriority::minimum(),
MyEvent::JobProcessed { job_id: next_id },
);
} else {
self.is_busy = false;
}
}
}
}
}
impl ModelSummary for ServerModel {
fn summary(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("ServerModel")
.field("name", &self.name)
.field("queue_len", &self.queue.len())
.field("busy", &self.is_busy)
.finish()
}
}
#[derive(Debug)]
pub struct JobGenerator {
next_job_id: u32,
interval: Duration,
}
impl Source<MyEvent, ServerModel> for JobGenerator {
fn on_registered(
&mut self,
context: &mut dyn UserContext<MyEvent, ServerModel>,
_model: &ServerModel,
) -> Option<Duration> {
let job_id = self.next_job_id;
self.next_job_id += 1;
context.schedule_event(
Duration::ticks(0),
EventPriority::minimum(),
MyEvent::JobArrived { job_id },
);
Some(self.interval)
}
fn fire(
&mut self,
context: &mut SourceContext<MyEvent, ServerModel>,
_model: &ServerModel,
) -> Option<Duration> {
for _ in 0..5 {
let job_id = self.next_job_id;
self.next_job_id += 1;
context.schedule_event(
Duration::ticks(0),
EventPriority::minimum(),
MyEvent::JobArrived { job_id },
);
}
Some(self.interval)
}
}
fn main() {
env_logger::Builder::from_env(env_logger::Env::default().default_filter_or("trace"))
.format(|buf, record| {
use std::io::Write;
writeln!(
buf,
"[{}] {:<5} {}",
chrono::Local::now().format("%H:%M:%S"),
record.level(),
record.args()
)
})
.init();
let mut engine = Engine::new();
engine
.add_hook(TraceHook)
.add_source(
"Job Generator x4",
JobGenerator {
next_job_id: 0,
interval: Duration::ticks(4),
},
)
.add_source(
"Job Generator x6",
JobGenerator {
next_job_id: 0,
interval: Duration::ticks(6),
},
);
let model = ServerModel {
name: "Sample Server",
queue: Default::default(),
is_busy: false,
};
let mut runner = ParallelRunner::new(true, EventPriority::maximum());
let result = runner.run_do_ticks(engine, model, 60, false);
println!("\nSimulation Result: {:?}", result);
}