use petgraph;
use std::path::Path;
pub use petgraph::graph::EdgeIndex;
pub use petgraph::graph::NodeIndex;
use super::generate;
pub(crate) type PetGraph = petgraph::Graph<Node, Edge>;
use crate::Result;
#[derive(Clone, Debug)]
pub(crate) enum Node {
Start {
name: String,
},
Terminate {
name: String,
},
Aggregate {
name: String,
behaviour_module: String,
},
FanOut {
name: String,
},
UserHandler {
name: String,
behaviour_module: String,
},
Poll {
name: String,
behaviour_module: String,
},
}
impl Node {
pub(crate) fn name(&self) -> String {
match self {
Node::Start { name, .. } => name.clone(),
Node::Terminate { name, .. } => name.clone(),
Node::Aggregate { name, .. } => name.clone(),
Node::FanOut { name, .. } => name.clone(),
Node::UserHandler { name, .. } => name.clone(),
Node::Poll { name, .. } => name.clone(),
}
}
}
#[derive(Clone, Debug, Ord, PartialOrd, PartialEq, Eq, Hash)]
pub(crate) struct Edge {
pub(crate) queue: String,
pub(crate) msg_type: String,
pub(crate) retry_interval_in_seconds: u32,
}
pub struct Graph {
pub(crate) g: PetGraph,
pub(crate) accept_failure: String,
pub(crate) type_error: String,
}
impl Graph {
pub fn new<S: Into<String>>(accept_failure: S, type_error: S) -> Graph {
Graph {
g: petgraph::Graph::new(),
accept_failure: accept_failure.into(),
type_error: type_error.into(),
}
}
pub fn start<S: Into<String>>(&mut self, name: S) -> NodeIndex {
self.g.add_node(Node::Start { name: name.into() })
}
pub fn process<S: Into<String>>(
&mut self,
input: NodeIndex,
queue: S,
type_input: S,
name: S,
behaviour_module: S,
retry_interval_in_seconds: u32,
) -> NodeIndex {
let handler_node = Node::UserHandler {
name: name.into(),
behaviour_module: behaviour_module.into(),
};
let handler_node_i = self.g.add_node(handler_node);
let edge = Edge {
queue: queue.into(),
msg_type: type_input.into(),
retry_interval_in_seconds,
};
self.g.add_edge(input, handler_node_i, edge);
handler_node_i
}
pub fn aggregate<S: Into<String>>(
&mut self,
inputs: Vec<NodeIndex>,
queue: S,
type_input: S,
name: S,
behaviour_module: S,
retry_interval_in_seconds: u32,
) -> NodeIndex {
let type_input = type_input.into();
let wait_node_i = self.g.add_node(Node::Aggregate {
name: name.into(),
behaviour_module: behaviour_module.into(),
});
let queue = queue.into();
for input_i in inputs {
self.g.add_edge(
input_i,
wait_node_i,
Edge {
queue: queue.clone(),
msg_type: type_input.clone(),
retry_interval_in_seconds,
},
);
}
wait_node_i
}
pub fn fan_out<S: Into<String>>(
&mut self,
input: NodeIndex,
queue: S,
type_input: S,
name: S,
retry_interval_in_seconds: u32,
) -> NodeIndex {
let fan_out_i = self.g.add_node(Node::FanOut { name: name.into() });
self.g.add_edge(
input,
fan_out_i,
Edge {
queue: queue.into(),
msg_type: type_input.into(),
retry_interval_in_seconds,
},
);
fan_out_i
}
pub fn poll<S: Into<String>>(
&mut self,
input: NodeIndex,
queue: S,
type_input: S,
name: S,
behaviour_module: S,
retry_interval_in_seconds: u32,
) -> NodeIndex {
let poll_node = Node::Poll {
name: name.into(),
behaviour_module: behaviour_module.into(),
};
let poll_node_i = self.g.add_node(poll_node);
let edge = Edge {
queue: queue.into(),
msg_type: type_input.into(),
retry_interval_in_seconds,
};
self.g.add_edge(input, poll_node_i, edge);
poll_node_i
}
pub fn terminate<S: Into<String>>(&mut self, input: NodeIndex, queue: S, type_input: S, name: S, retry_interval_in_seconds: u32) -> () {
let terminate_node = self.g.add_node(Node::Terminate { name: name.into() });
self.g.add_edge(input, terminate_node, Edge {
queue: queue.into(),
msg_type: type_input.into(),
retry_interval_in_seconds
});
}
pub fn generate<P: AsRef<Path>, S: AsRef<str>>(
self,
output_dir: P,
get_rmq_uri: S,
work_exchange: S,
retry_exchange: S,
retry_queue_prefix: S,
retry_queue_suffix: S,
init_input_queue: bool,
init_output_queue: bool,
main_init: S,
) -> Result<()> {
let rmq_options = generate::graph::RmqOptions {
get_rmq_uri: get_rmq_uri.as_ref().into(),
work_exchange: work_exchange.as_ref().into(),
retry_exchange: retry_exchange.as_ref().into(),
retry_queue_prefix: retry_queue_prefix.as_ref().into(),
retry_queue_suffix: retry_queue_suffix.as_ref().into(),
};
generate::graph::generate(self, output_dir, rmq_options, init_input_queue, init_output_queue, main_init)
}
}