extern crate timely;
extern crate graph_map;
extern crate differential_dataflow;
use std::hash::Hash;
use timely::dataflow::*;
use differential_dataflow::Collection;
use differential_dataflow::lattice::Lattice;
use differential_dataflow::operators::*;
use differential_dataflow::operators::arrange::ArrangeBySelf;
use differential_dataflow::operators::arrange::ArrangeByKey;
use graph_map::GraphMMap;
type Node = u32;
type Edge = (Node, Node);
fn main() {
let filename = std::env::args().nth(1).unwrap();
let batch: usize = std::env::args().nth(2).unwrap().parse().unwrap();
timely::execute_from_args(std::env::args().skip(2), move |worker| {
let peers = worker.peers();
let index = worker.index();
let mut edges = worker.dataflow::<_,_,_>(|scope| {
use differential_dataflow::input::Input;
let (edges_handle, edges) = scope.new_collection();
triangles(&edges).map(|_| ()).count().inspect(|x| println!("{:?}", x));
edges_handle
});
let graph = GraphMMap::new(&filename);
let nodes = graph.nodes();
for node in 0 .. nodes {
if node % peers == index {
edges.advance_to(node / batch);
for &dest in graph.edges(node) {
edges.insert((node as u32, dest));
}
if node % batch == 0 {
edges.flush();
worker.step();
println!("inserted through: {:?}", node);
}
}
}
}).unwrap();
}
fn triangles<G: Scope>(edges: &Collection<G, Edge>) -> Collection<G, (Node, Node, Node)>
where G::Timestamp: Lattice+Hash+Ord {
let edges = edges.filter(|&(src, dst)| src < dst);
let as_self = edges.arrange_by_self();
let forward = edges.arrange_by_key();
let reverse = edges.map_in_place(|x| ::std::mem::swap(&mut x.0, &mut x.1))
.arrange_by_key();
let counts = edges.map(|(src, _dst)| src)
.arrange_by_self();
let cand_count1 = forward.join_core(&counts, |&src, &dst, &()| Some(((src, dst), 1)));
let cand_count2 = reverse.join_core(&counts, |&dst, &src, &()| Some(((src, dst), 2)));
let winners = cand_count1.concat(&cand_count2)
.group(|_srcdst, counts, output| {
let mut min_cnt = isize::max_value();
let mut min_idx = usize::max_value();
for &(&idx, cnt) in counts.iter() {
if min_cnt > cnt {
min_idx = idx;
min_cnt = cnt;
}
}
output.push((min_idx, 1));
});
let winners1 = winners.flat_map(|((src, dst), index)| if index == 1 { Some((src, dst)) } else { None })
.join_core(&forward, |&src, &dst, &ext| Some(((dst, ext), src)))
.join_core(&as_self, |&(dst, ext), &src, &()| Some(((dst, ext), src)))
.map(|((dst, ext), src)| (src, dst, ext));
let winners2 = winners.flat_map(|((src, dst), index)| if index == 2 { Some((dst, src)) } else { None })
.join_core(&forward, |&dst, &src, &ext| Some(((src, ext), dst)))
.join_core(&as_self, |&(src, ext), &dst, &()| Some(((src, ext), dst)))
.map(|((src, ext), dst)| (src, dst, ext));
winners1.concat(&winners2)
}