use log::info;
use petgraph::Direction;
use std::collections::HashMap;
use std::collections::HashSet;
use std::path::Path;
use std::path::PathBuf;
use crate::graph::Edge;
use crate::graph::Graph;
use crate::graph::Node;
use crate::graph::NodeIndex;
use crate::graph::PetGraph;
use crate::Error;
use crate::Result;
use crate::misc;
#[derive(Clone)]
pub(crate) struct RmqOptions {
pub(crate) get_rmq_uri: String,
pub(crate) work_exchange: String,
pub(crate) retry_exchange: String,
pub(crate) retry_queue_prefix: String,
pub(crate) retry_queue_suffix: String,
}
pub(crate) fn generate<P: AsRef<Path>, S: AsRef<str>>(
graph: Graph,
dir: P,
rmq_options: RmqOptions,
init_input_queue: bool,
init_output_queue: bool,
main_init: S,
) -> Result<()> {
let g = &graph.g;
let dir = dir.as_ref();
let main_init = main_init.as_ref();
let mut outputs = HashMap::new();
for node_i in g.node_indices() {
let node = &g[node_i];
match node {
Node::Start { .. } => (),
Node::Terminate { .. } => (),
Node::Aggregate { name, behaviour_module } => {
let content = generate_aggregate(
g,
node_i,
behaviour_module.into(),
graph.accept_failure.clone(),
rmq_options.clone(),
)?;
update_outputs(&mut outputs, dir, name, content);
}
Node::FanOut { name } => {
let content = generate_fan_out(
g,
node_i,
graph.accept_failure.clone(),
rmq_options.clone()
)?;
update_outputs(&mut outputs, dir, name, content);
}
Node::UserHandler {
name,
behaviour_module,
} => {
let content = generate_user_handler(
g,
node_i,
behaviour_module.into(),
graph.accept_failure.clone(),
rmq_options.clone(),
)?;
update_outputs(&mut outputs, dir, name, content);
}
Node::Poll {
name,
behaviour_module,
} => {
let content = generate_poll(
g,
node_i,
behaviour_module.into(),
graph.accept_failure.clone(),
rmq_options.clone(),
)?;
update_outputs(&mut outputs, dir, name, content);
}
}
}
let content = generate_init_exchanges_and_queues(
g,
rmq_options.clone(),
init_input_queue,
init_output_queue,
)?;
update_outputs(&mut outputs, dir, "init_exchanges_and_queues", content);
let content = generate_main(&outputs, main_init)?;
update_outputs(&mut outputs, dir, "main", content);
let graph_for_display = map_to_string(g);
let dot = petgraph::dot::Dot::with_attr_getters(
&graph_for_display,
&[],
&|_, _| String::from(r#"arrowhead = "onormal""#),
&|_, _| String::from(r#"shape = "box" style = "rounded""#)
);
let dot_file_path = dir.join("graph.dot");
info!("writing dot graph to {}", &dot_file_path.display());
std::fs::write(&dot_file_path, format!("{}", dot))?;
let f = std::fs::OpenOptions::new()
.create(false)
.read(true)
.write(false)
.open(&dot_file_path)?
.sync_all()?;
let svg_file_path = dir.join("graph.svg");
let gen_svg_output = std::process::Command::new("dot")
.arg("-Tsvg")
.arg(
dot_file_path
.to_str()
.ok_or(Error::InvalidFileName(dot_file_path.display().to_string()))?,
)
.arg("-o")
.arg(
svg_file_path
.to_str()
.ok_or(Error::InvalidFileName(svg_file_path.display().to_string()))?,
)
.output()?;
if gen_svg_output.status.success() {
()
} else {
return Err(Error::ErrorGeneratingSvg);
}
let mut mods = Vec::new();
for (file_path, content) in outputs.iter() {
info!("writing to {}", file_path.display());
std::fs::write(file_path, content)?;
match file_path.file_stem() {
None => return Err(Error::InvalidFileName((file_path).display().to_string())),
Some(file_stem) => match (*file_stem).to_str() {
None => return Err(Error::InvalidFileName((file_path).display().to_string())),
Some(s) => mods.push(format!("pub(crate) mod {};", s)),
},
}
}
Ok(())
}
fn generate_aggregate(
g: &PetGraph,
node_i: NodeIndex,
behaviour_module: String,
accept_failure: String,
rmq_options: RmqOptions,
) -> Result<String> {
let Edge {
queue: input_queue,
msg_type: type_input,
retry_interval_in_seconds: _,
} = expect_one_input_edge_for_aggregation_node(g, node_i)?;
let output_queue = expect_optional_outgoing_edge(g, node_i)?.map(|e| e.queue.clone());
super::aggregate::generate(
input_queue,
behaviour_module,
output_queue,
accept_failure,
type_input,
rmq_options,
)
}
fn generate_fan_out(
g: &PetGraph,
node_i: NodeIndex,
accept_failure: String,
rmq_options: RmqOptions,
) -> Result<String> {
let Edge {
queue: input_queue,
msg_type: type_input,
retry_interval_in_seconds: _,
} = expect_one_incoming_edge(g, node_i)?;
let mut output_queues = Vec::new();
let out_edges = g.edges_directed(node_i, Direction::Outgoing);
for out_edge in out_edges {
let output_queue = out_edge.weight().queue.clone();
output_queues.push(output_queue)
}
super::fan_out::generate(
input_queue.clone(),
output_queues,
accept_failure,
type_input.clone(),
rmq_options.clone(),
)
}
fn generate_user_handler(
g: &PetGraph,
node_i: NodeIndex,
module: String,
accept_failure: String,
rmq_options: RmqOptions,
) -> Result<String> {
let Edge {
queue: input_queue,
msg_type: type_input,
retry_interval_in_seconds: _,
} = expect_one_incoming_edge(g, node_i)?;
let output_queue = expect_optional_outgoing_edge(g, node_i)?.map(|e| e.queue.clone());
super::user_handler::generate(
input_queue.clone(),
output_queue,
module,
accept_failure,
type_input.clone(),
rmq_options,
)
}
fn generate_poll(
g: &PetGraph,
node_i: NodeIndex,
module: String,
accept_failure: String,
rmq_options: RmqOptions,
) -> Result<String> {
let Edge {
queue: input_queue,
msg_type: type_input,
retry_interval_in_seconds: _,
} = expect_one_incoming_edge(g, node_i)?;
let output_queue = expect_optional_outgoing_edge(g, node_i)?.map(|e| e.queue.clone());
super::poll::generate(
input_queue.clone(),
output_queue,
module,
accept_failure,
type_input.clone(),
rmq_options,
)
}
fn expect_one_incoming_edge(g: &PetGraph, node_i: NodeIndex) -> Result<&Edge> {
let in_edges: Vec<_> = g.edges_directed(node_i, Direction::Incoming).collect();
if in_edges.len() != 1 {
Err(Error::IllFormedNode {
node: format!("{:?}", g[node_i]),
})
} else {
Ok(in_edges[0].weight())
}
}
fn expect_optional_outgoing_edge(g: &PetGraph, i: NodeIndex) -> Result<Option<&Edge>> {
let edges: Vec<_> = g.edges_directed(i, Direction::Outgoing).collect();
if edges.len() == 0 {
Ok(None)
} else if edges.len() == 1 {
Ok(Some(edges[0].weight()))
} else {
Err(Error::IllFormedNode {
node: format!("{:?}", g[i]),
})
}
}
fn expect_one_input_edge_for_aggregation_node(g: &PetGraph, node_i: NodeIndex) -> Result<Edge> {
let incoming_edges = g.edges_directed(node_i, Direction::Incoming);
let mut edges = HashSet::new();
for edge in incoming_edges {
edges.insert(edge.weight().clone());
}
let edges: Vec<_> = edges.into_iter().collect();
if edges.len() == 1 {
Ok(edges[0].clone())
} else {
Err(Error::IllFormedNode {
node: format!("{:?}", g[node_i]),
})
}
}
fn update_outputs<P: AsRef<Path>, S: AsRef<str>>(
outputs: &mut HashMap<PathBuf, String>,
dir: P,
name: S,
content: String,
) {
let dir = dir.as_ref();
let basename = format!("{}.rs", name.as_ref());
let file_path = dir.join(basename);
outputs.insert(file_path, content);
}
fn map_to_string(old: &PetGraph) -> petgraph::Graph<String, String> {
let all_msg_types: Vec<String> = old.edge_weights().map(|e| e.msg_type.clone()).collect();
let msg_prefix = misc::longest_common_prefix(all_msg_types);
old.map(
|node_index, node| node.name(),
|edge_index, edge| format!("{}\n{}", edge.queue, edge.msg_type.trim_start_matches(&msg_prefix)),
)
}
fn generate_init_exchanges_and_queues(graph: &PetGraph, rmq_options: RmqOptions, init_input_queue: bool, init_output_queue: bool) -> Result<String> {
let mut all_queues: Vec<_> = graph
.edge_weights()
.map(|edge|
(
(&edge.queue).clone(),
format!("{}{}{}", &rmq_options.retry_queue_prefix, &edge.queue, &rmq_options.retry_queue_suffix),
edge.retry_interval_in_seconds,
)
).collect();
let mut input_queues = {
let mut qs = Vec::new();
for node_i in graph.node_indices() {
let node = &graph[node_i];
match node {
Node::Start {..} => {
let edge = expect_optional_outgoing_edge(graph, node_i)?
.ok_or(Error::IllFormedNode {node: format!("{:?}", node)})?;
qs.push(edge.queue.clone());
},
Node::Aggregate {..} | Node::FanOut {..} | Node::UserHandler{..} | Node::Poll {..} | Node::Terminate {..} => ()
}
}
qs
};
let mut output_queues = {
let mut qs = Vec::new();
for node_i in graph.node_indices() {
let node = &graph[node_i];
match node {
Node::Terminate {..} => {
let edge = expect_one_incoming_edge(graph, node_i)?;
qs.push(edge.queue.clone());
},
Node::Aggregate {..} | Node::FanOut {..} | Node::UserHandler{..} | Node::Poll {..} | Node::Start {..} => ()
}
}
qs
};
let mut unwanted_queues = Vec::new();
if !init_input_queue {
unwanted_queues.append(&mut input_queues);
}
if !init_output_queue {
unwanted_queues.append(&mut output_queues);
}
let mut wanted_queues = Vec::new();
for q in all_queues {
if !unwanted_queues.contains(&q.0) {
wanted_queues.push(q);
}
}
wanted_queues.sort();
wanted_queues.dedup();
super::init_exchanges_and_queues::generate(rmq_options, wanted_queues)
}
fn generate_main<S: AsRef<str>>(outputs: &HashMap<PathBuf, String>, main_init: S) -> Result<String> {
let main_init = main_init.as_ref();
let mut modules = Vec::new();
for file_path in outputs.keys() {
let file_stem = file_path
.file_stem()
.ok_or(Error::InvalidFileName(file_path.display().to_string()))?
.to_str()
.ok_or(Error::InvalidFileName(file_path.display().to_string()))?;
if file_stem == "main" {
continue
}
modules.push(String::from(file_stem));
}
modules.sort();
super::main::generate(modules, main_init)
}