use crate::config_parser::ConfigurationValue;
use std::boxed::Box;
use std::cell::{RefCell};
use crate::{Message,Plugs};
use ::rand::{Rng,StdRng};
use crate::pattern::{Pattern,new_pattern,PatternBuilderArgument};
use std::rc::Rc;
use std::collections::{BTreeSet,BTreeMap,VecDeque};
use crate::topology::Topology;
use quantifiable_derive::Quantifiable;use crate::quantify::Quantifiable;
use std::fmt::Debug;
#[derive(Debug)]
pub enum TrafficError
{
OriginOutsideTraffic,
SelfMessage,
}
pub trait Traffic : Quantifiable + Debug
{
fn generate_message(&mut self, origin:usize, cycle:usize, topology:&Box<dyn Topology>, rng: &RefCell<StdRng>) -> Result<Rc<Message>,TrafficError>;
fn probability_per_cycle(&self, server:usize) -> f32;
fn try_consume(&mut self, message: Rc<Message>, cycle:usize, topology:&Box<dyn Topology>, rng: &RefCell<StdRng>) -> bool;
fn is_finished(&self) -> bool;
fn should_generate(&self, server:usize, rng: &RefCell<StdRng>) -> bool
{
let p=self.probability_per_cycle(server);
let r=rng.borrow_mut().gen_range(0f32,1f32);
r<p
}
}
#[non_exhaustive]
#[derive(Debug)]
pub struct TrafficBuilderArgument<'a>
{
pub cv: &'a ConfigurationValue,
pub plugs: &'a Plugs,
pub topology: &'a Box<dyn Topology>,
pub rng: &'a RefCell<StdRng>,
}
pub fn new_traffic(arg:TrafficBuilderArgument) -> Box<dyn Traffic>
{
if let &ConfigurationValue::Object(ref cv_name, ref _cv_pairs)=arg.cv
{
match arg.plugs.traffics.get(cv_name)
{
Some(builder) => return builder(arg),
_ => (),
};
match cv_name.as_ref()
{
"HomogeneousTraffic" => Box::new(Homogeneous::new(arg)),
"TrafficSum" => Box::new(Sum::new(arg)),
"ShiftedTraffic" => Box::new(Shifted::new(arg)),
"ProductTraffic" => Box::new(ProductTraffic::new(arg)),
"SubRangeTraffic" => Box::new(SubRangeTraffic::new(arg)),
"Burst" => Box::new(Burst::new(arg)),
"Reactive" => Box::new(Reactive::new(arg)),
_ => panic!("Unknown traffic {}",cv_name),
}
}
else
{
panic!("Trying to create a traffic from a non-Object");
}
}
#[derive(Quantifiable)]
#[derive(Debug)]
pub struct Homogeneous
{
servers: usize,
pattern: Box<dyn Pattern>,
message_size: usize,
load: f32,
generated_messages: BTreeSet<*const Message>,
}
impl Traffic for Homogeneous
{
fn generate_message(&mut self, origin:usize, cycle:usize, topology:&Box<dyn Topology>, rng: &RefCell<StdRng>) -> Result<Rc<Message>,TrafficError>
{
if origin>=self.servers
{
return Err(TrafficError::OriginOutsideTraffic);
}
let destination=self.pattern.get_destination(origin,topology,rng);
if origin==destination
{
return Err(TrafficError::SelfMessage);
}
let message=Rc::new(Message{
origin: origin,
destination,
size:self.message_size,
creation_cycle: cycle,
});
self.generated_messages.insert(message.as_ref() as *const Message);
Ok(message)
}
fn probability_per_cycle(&self, _server:usize) -> f32
{
let r=self.load/self.message_size as f32;
if r>1.0
{
1.0
}
else
{
r
}
}
fn try_consume(&mut self, message: Rc<Message>, _cycle:usize, _topology:&Box<dyn Topology>, _rng: &RefCell<StdRng>) -> bool
{
let message_ptr=message.as_ref() as *const Message;
self.generated_messages.remove(&message_ptr)
}
fn is_finished(&self) -> bool
{
false
}
}
impl Homogeneous
{
pub fn new(arg:TrafficBuilderArgument) -> Homogeneous
{
let mut servers=None;
let mut load=None;
let mut pattern=None;
let mut message_size=None;
if let &ConfigurationValue::Object(ref cv_name, ref cv_pairs)=arg.cv
{
if cv_name!="HomogeneousTraffic"
{
panic!("A Homogeneous must be created from a `HomogeneousTraffic` object not `{}`",cv_name);
}
for &(ref name,ref value) in cv_pairs
{
match name.as_ref()
{
"pattern" => pattern=Some(new_pattern(PatternBuilderArgument{cv:value,plugs:arg.plugs})),
"servers" => match value
{
&ConfigurationValue::Number(f) => servers=Some(f as usize),
_ => panic!("bad value for servers"),
}
"load" => match value
{
&ConfigurationValue::Number(f) => load=Some(f as f32),
_ => panic!("bad value for load ({:?})",value),
}
"message_size" => match value
{
&ConfigurationValue::Number(f) => message_size=Some(f as usize),
_ => panic!("bad value for message_size"),
}
_ => panic!("Nothing to do with field {} in HomogeneousTraffic",name),
}
}
}
else
{
panic!("Trying to create a Homogeneous from a non-Object");
}
let servers=servers.expect("There were no servers");
let message_size=message_size.expect("There were no message_size");
let load=load.expect("There were no load");
let mut pattern=pattern.expect("There were no pattern");
pattern.initialize(servers, servers, arg.topology, arg.rng);
Homogeneous{
servers,
pattern,
message_size,
load,
generated_messages: BTreeSet::new(),
}
}
}
#[derive(Quantifiable)]
#[derive(Debug)]
pub struct Sum
{
list: Vec<Box<dyn Traffic>>,
}
impl Traffic for Sum
{
fn generate_message(&mut self, origin:usize, cycle:usize, topology:&Box<dyn Topology>, rng: &RefCell<StdRng>) -> Result<Rc<Message>,TrafficError>
{
let probs:Vec<f32> =self.list.iter().map(|t|t.probability_per_cycle(origin)).collect();
let mut r=rng.borrow_mut().gen_range(0f32,probs.iter().sum());
for i in 0..self.list.len()
{
if r<probs[i]
{
return self.list[i].generate_message(origin,cycle,topology,rng);
}
else
{
r-=probs[i];
}
}
panic!("failed probability");
}
fn probability_per_cycle(&self,server:usize) -> f32
{
self.list.iter().map(|t|t.probability_per_cycle(server)).sum()
}
fn try_consume(&mut self, message: Rc<Message>, cycle:usize, topology:&Box<dyn Topology>, rng: &RefCell<StdRng>) -> bool
{
for traffic in self.list.iter_mut()
{
if traffic.try_consume(message.clone(),cycle,topology,rng)
{
return true;
}
}
return false;
}
fn is_finished(&self) -> bool
{
for traffic in self.list.iter()
{
if !traffic.is_finished()
{
return false;
}
}
return true;
}
}
impl Sum
{
pub fn new(arg:TrafficBuilderArgument) -> Sum
{
let mut list=None;
if let &ConfigurationValue::Object(ref cv_name, ref cv_pairs)=arg.cv
{
if cv_name!="TrafficSum"
{
panic!("A Sum must be created from a `TrafficSum` object not `{}`",cv_name);
}
for &(ref name,ref value) in cv_pairs
{
match name.as_ref()
{
"list" => match value
{
&ConfigurationValue::Array(ref a) => list=Some(a.iter().map(|v|new_traffic(TrafficBuilderArgument{cv:v,..arg})).collect()),
_ => panic!("bad value for list"),
}
_ => panic!("Nothing to do with field {} in TrafficSum",name),
}
}
}
else
{
panic!("Trying to create a Sum from a non-Object");
}
let list=list.expect("There were no list");
Sum{
list
}
}
}
#[derive(Quantifiable)]
#[derive(Debug)]
pub struct Shifted
{
shift: usize,
traffic: Box<dyn Traffic>,
generated_messages: BTreeMap<*const Message,Rc<Message>>,
}
impl Traffic for Shifted
{
fn generate_message(&mut self, origin:usize, cycle:usize, topology:&Box<dyn Topology>, rng: &RefCell<StdRng>) -> Result<Rc<Message>,TrafficError>
{
if origin<self.shift
{
return Err(TrafficError::OriginOutsideTraffic);
}
let inner_message=self.traffic.generate_message(origin-self.shift,cycle,topology,rng)?;
let outer_message=Rc::new(Message{
origin,
destination:inner_message.destination+self.shift,
size:inner_message.size,
creation_cycle: cycle,
});
self.generated_messages.insert(outer_message.as_ref() as *const Message,inner_message);
Ok(outer_message)
}
fn probability_per_cycle(&self,server:usize) -> f32
{
self.traffic.probability_per_cycle(server)
}
fn try_consume(&mut self, message: Rc<Message>, cycle:usize, topology:&Box<dyn Topology>, rng: &RefCell<StdRng>) -> bool
{
let message_ptr=message.as_ref() as *const Message;
let outer_message=match self.generated_messages.remove(&message_ptr)
{
None => return false,
Some(m) => m,
};
if !self.traffic.try_consume(outer_message,cycle,topology,rng)
{
panic!("Shifted traffic consumed a message but its child did not.");
}
true
}
fn is_finished(&self) -> bool
{
self.traffic.is_finished()
}
}
impl Shifted
{
pub fn new(arg:TrafficBuilderArgument) -> Shifted
{
let mut shift=None;
let mut traffic=None;
if let &ConfigurationValue::Object(ref cv_name, ref cv_pairs)=arg.cv
{
if cv_name!="ShiftedTraffic"
{
panic!("A Shifted must be created from a `ShiftedTraffic` object not `{}`",cv_name);
}
for &(ref name,ref value) in cv_pairs
{
match name.as_ref()
{
"traffic" => traffic=Some(new_traffic(TrafficBuilderArgument{cv:value,..arg})),
"shift" => match value
{
&ConfigurationValue::Number(f) => shift=Some(f as usize),
_ => panic!("bad value for shift"),
}
_ => panic!("Nothing to do with field {} in ShiftedTraffic",name),
}
}
}
else
{
panic!("Trying to create a Shifted from a non-Object");
}
let shift=shift.expect("There were no shift");
let traffic=traffic.expect("There were no traffic");
Shifted{
shift,
traffic,
generated_messages: BTreeMap::new(),
}
}
}
#[derive(Quantifiable)]
#[derive(Debug)]
pub struct ProductTraffic
{
block_size: usize,
block_traffic: Box<dyn Traffic>,
global_pattern: Box<dyn Pattern>,
generated_messages: BTreeMap<*const Message,Rc<Message>>,
}
impl Traffic for ProductTraffic
{
fn generate_message(&mut self, origin:usize, cycle:usize, topology:&Box<dyn Topology>, rng: &RefCell<StdRng>) -> Result<Rc<Message>,TrafficError>
{
let local=origin % self.block_size;
let global=origin / self.block_size;
let global_dest=self.global_pattern.get_destination(global,topology,rng);
let inner_message=self.block_traffic.generate_message(local,cycle,topology,rng)?;
let outer_message=Rc::new(Message{
origin,
destination:global_dest*self.block_size+inner_message.destination,
size:inner_message.size,
creation_cycle: cycle,
});
self.generated_messages.insert(outer_message.as_ref() as *const Message,inner_message);
Ok(outer_message)
}
fn probability_per_cycle(&self,server:usize) -> f32
{
self.block_traffic.probability_per_cycle(server)
}
fn try_consume(&mut self, message: Rc<Message>, cycle:usize, topology:&Box<dyn Topology>, rng: &RefCell<StdRng>) -> bool
{
let message_ptr=message.as_ref() as *const Message;
let outer_message=match self.generated_messages.remove(&message_ptr)
{
None => return false,
Some(m) => m,
};
if !self.block_traffic.try_consume(outer_message,cycle,topology,rng)
{
panic!("ProductTraffic traffic consumed a message but its child did not.");
}
true
}
fn is_finished(&self) -> bool
{
self.block_traffic.is_finished()
}
}
impl ProductTraffic
{
pub fn new(arg:TrafficBuilderArgument) -> ProductTraffic
{
let mut block_size=None;
let mut block_traffic=None;
let mut global_pattern=None;
if let &ConfigurationValue::Object(ref cv_name, ref cv_pairs)=arg.cv
{
if cv_name!="ProductTraffic"
{
panic!("A ProductTraffic must be created from a `ProductTraffic` object not `{}`",cv_name);
}
for &(ref name,ref value) in cv_pairs
{
match name.as_ref()
{
"block_traffic" => block_traffic=Some(new_traffic(TrafficBuilderArgument{cv:value,..arg})),
"global_pattern" => global_pattern=Some(new_pattern(PatternBuilderArgument{cv:value,plugs:arg.plugs})),
"block_size" => match value
{
&ConfigurationValue::Number(f) => block_size=Some(f as usize),
_ => panic!("bad value for block_size"),
}
_ => panic!("Nothing to do with field {} in ShiftedTraffic",name),
}
}
}
else
{
panic!("Trying to create a ProductTraffic from a non-Object");
}
let block_size=block_size.expect("There were no block_size");
let block_traffic=block_traffic.expect("There were no block_traffic");
let mut global_pattern=global_pattern.expect("There were no global_pattern");
let global_size=arg.topology.num_servers()/block_size;
global_pattern.initialize(global_size,global_size,arg.topology,arg.rng);
ProductTraffic{
block_size,
block_traffic,
global_pattern,
generated_messages: BTreeMap::new(),
}
}
}
#[derive(Quantifiable)]
#[derive(Debug)]
pub struct SubRangeTraffic
{
start: usize,
end: usize,
traffic: Box<dyn Traffic>,
}
impl Traffic for SubRangeTraffic
{
fn generate_message(&mut self, origin:usize, cycle:usize, topology:&Box<dyn Topology>, rng: &RefCell<StdRng>) -> Result<Rc<Message>,TrafficError>
{
if origin<self.start || origin>=self.end
{
return Err(TrafficError::OriginOutsideTraffic);
}
self.traffic.generate_message(origin,cycle,topology,rng)
}
fn probability_per_cycle(&self,server:usize) -> f32
{
self.traffic.probability_per_cycle(server)
}
fn try_consume(&mut self, message: Rc<Message>, cycle:usize, topology:&Box<dyn Topology>, rng: &RefCell<StdRng>) -> bool
{
self.traffic.try_consume(message,cycle,topology,rng)
}
fn is_finished(&self) -> bool
{
self.traffic.is_finished()
}
}
impl SubRangeTraffic
{
pub fn new(arg:TrafficBuilderArgument) -> SubRangeTraffic
{
let mut start=None;
let mut end=None;
let mut traffic=None;
if let &ConfigurationValue::Object(ref cv_name, ref cv_pairs)=arg.cv
{
if cv_name!="SubRangeTraffic"
{
panic!("A SubRangeTraffic must be created from a `SubRangeTraffic` object not `{}`",cv_name);
}
for &(ref name,ref value) in cv_pairs
{
match name.as_ref()
{
"traffic" => traffic=Some(new_traffic(TrafficBuilderArgument{cv:value,..arg})),
"start" => match value
{
&ConfigurationValue::Number(f) => start=Some(f as usize),
_ => panic!("bad value for start"),
}
"end" => match value
{
&ConfigurationValue::Number(f) => end=Some(f as usize),
_ => panic!("bad value for end"),
}
_ => panic!("Nothing to do with field {} in SubRangeTraffic",name),
}
}
}
else
{
panic!("Trying to create a SubRangeTraffic from a non-Object");
}
let start=start.expect("There were no start");
let end=end.expect("There were no end");
let traffic=traffic.expect("There were no traffic");
SubRangeTraffic{
start,
end,
traffic,
}
}
}
#[derive(Quantifiable)]
#[derive(Debug)]
pub struct Burst
{
servers: usize,
pattern: Box<dyn Pattern>,
message_size: usize,
pending_messages: Vec<usize>,
generated_messages: BTreeSet<*const Message>,
}
impl Traffic for Burst
{
fn generate_message(&mut self, origin:usize, cycle:usize, topology:&Box<dyn Topology>, rng: &RefCell<StdRng>) -> Result<Rc<Message>,TrafficError>
{
if origin>=self.servers
{
return Err(TrafficError::OriginOutsideTraffic);
}
self.pending_messages[origin]-=1;
let destination=self.pattern.get_destination(origin,topology,rng);
if origin==destination
{
return Err(TrafficError::SelfMessage);
}
let message=Rc::new(Message{
origin: origin,
destination,
size:self.message_size,
creation_cycle: cycle,
});
self.generated_messages.insert(message.as_ref() as *const Message);
Ok(message)
}
fn probability_per_cycle(&self, server:usize) -> f32
{
if self.pending_messages[server]>0
{
1.0
}
else
{
0.0
}
}
fn try_consume(&mut self, message: Rc<Message>, _cycle:usize, _topology:&Box<dyn Topology>, _rng: &RefCell<StdRng>) -> bool
{
let message_ptr=message.as_ref() as *const Message;
self.generated_messages.remove(&message_ptr)
}
fn is_finished(&self) -> bool
{
if self.generated_messages.len()>0
{
return false;
}
for &pm in self.pending_messages.iter()
{
if pm>0
{
return false;
}
}
true
}
}
impl Burst
{
pub fn new(arg:TrafficBuilderArgument) -> Burst
{
let mut servers=None;
let mut messages_per_server=None;
let mut pattern=None;
let mut message_size=None;
if let &ConfigurationValue::Object(ref cv_name, ref cv_pairs)=arg.cv
{
if cv_name!="Burst"
{
panic!("A Burst must be created from a `Burst` object not `{}`",cv_name);
}
for &(ref name,ref value) in cv_pairs
{
match name.as_ref()
{
"pattern" => pattern=Some(new_pattern(PatternBuilderArgument{cv:value,plugs:arg.plugs})),
"servers" => match value
{
&ConfigurationValue::Number(f) => servers=Some(f as usize),
_ => panic!("bad value for servers"),
}
"messages_per_server" => match value
{
&ConfigurationValue::Number(f) => messages_per_server=Some(f as usize),
_ => panic!("bad value for messages_per_server ({:?})",value),
}
"message_size" => match value
{
&ConfigurationValue::Number(f) => message_size=Some(f as usize),
_ => panic!("bad value for message_size"),
}
_ => panic!("Nothing to do with field {} in Burst",name),
}
}
}
else
{
panic!("Trying to create a Burst from a non-Object");
}
let servers=servers.expect("There were no servers");
let message_size=message_size.expect("There were no message_size");
let messages_per_server=messages_per_server.expect("There were no messages_per_server");
let mut pattern=pattern.expect("There were no pattern");
pattern.initialize(servers, servers, arg.topology, arg.rng);
Burst{
servers,
pattern,
message_size,
pending_messages:vec![messages_per_server;servers],
generated_messages: BTreeSet::new(),
}
}
}
#[derive(Quantifiable)]
#[derive(Debug)]
pub struct Reactive
{
action_traffic: Box<dyn Traffic>,
reaction_traffic: Box<dyn Traffic>,
pending_messages: Vec<VecDeque<Rc<Message>>>,
}
impl Traffic for Reactive
{
fn generate_message(&mut self, origin:usize, cycle:usize, topology:&Box<dyn Topology>, rng: &RefCell<StdRng>) -> Result<Rc<Message>,TrafficError>
{
if origin<self.pending_messages.len()
{
if let Some(message)=self.pending_messages[origin].pop_front()
{
return Ok(message);
}
}
return self.action_traffic.generate_message(origin,cycle,topology,rng);
}
fn probability_per_cycle(&self, server:usize) -> f32
{
if server<self.pending_messages.len()
{
if self.pending_messages[server].len()>0
{
return 1.0;
}
}
return self.action_traffic.probability_per_cycle(server);
}
fn try_consume(&mut self, message: Rc<Message>, cycle:usize, topology:&Box<dyn Topology>, rng: &RefCell<StdRng>) -> bool
{
if self.action_traffic.try_consume(message.clone(),cycle,topology,rng)
{
if self.reaction_traffic.should_generate(message.origin,rng)
{
match self.reaction_traffic.generate_message(message.origin,cycle,topology,rng)
{
Ok(response_message) =>
{
if self.pending_messages.len()<message.origin+1
{
self.pending_messages.resize(message.origin+1,VecDeque::new());
}
self.pending_messages[message.origin].push_back(response_message);
},
Err(error) => panic!("An error happened when generating response traffic: {:?}",error),
};
}
return true;
}
self.reaction_traffic.try_consume(message,cycle,topology,rng)
}
fn is_finished(&self) -> bool
{
if !self.action_traffic.is_finished() || !self.reaction_traffic.is_finished()
{
return false;
}
for pm in self.pending_messages.iter()
{
if pm.len()>0
{
return false;
}
}
return true;
}
}
impl Reactive
{
pub fn new(arg:TrafficBuilderArgument) -> Reactive
{
let mut action_traffic=None;
let mut reaction_traffic=None;
if let &ConfigurationValue::Object(ref cv_name, ref cv_pairs)=arg.cv
{
if cv_name!="Reactive"
{
panic!("A Reactive must be created from a `Reactive` object not `{}`",cv_name);
}
for &(ref name,ref value) in cv_pairs
{
match name.as_ref()
{
"action_traffic" => action_traffic=Some(new_traffic(TrafficBuilderArgument{cv:value,..arg})),
"reaction_traffic" => reaction_traffic=Some(new_traffic(TrafficBuilderArgument{cv:value,..arg})),
_ => panic!("Nothing to do with field {} in Reactive",name),
}
}
}
else
{
panic!("Trying to create a Reactive from a non-Object");
}
let action_traffic=action_traffic.expect("There were no action_traffic");
let reaction_traffic=reaction_traffic.expect("There were no reaction_traffic");
Reactive{
action_traffic,
reaction_traffic,
pending_messages:vec![],
}
}
}