#![allow(clippy::needless_range_loop)]
use anyhow::anyhow;
use lax::Lapack;
use log::*;
use std::collections::HashSet;
use std::fs::OpenOptions;
use std::path::Path;
use std::str::FromStr;
use std::io::{BufRead, BufReader, Read};
use csv::ReaderBuilder;
use num_traits::float::*;
use annembed::tools::svdapprox::*;
use indexmap::IndexSet;
use sprs::{CsMat, TriMatI};
use petgraph::graph::{Graph, IndexType};
use petgraph::graphmap::{GraphMap, NodeTrait};
#[allow(unused)]
use petgraph::{Directed, EdgeType};
pub type NodeIndexation<N> = IndexSet<N>;
pub(crate) fn get_header_size(filepath: &Path) -> anyhow::Result<usize> {
log::debug!("get_header_size");
let fileres = OpenOptions::new().read(true).open(filepath);
if fileres.is_err() {
log::error!(
"fn get_header_size : could not open file {:?}",
filepath.as_os_str()
);
println!(
"directed_from_csv could not open file {:?}",
filepath.as_os_str()
);
return Err(anyhow!(
"fn get_header_size : could not open file {}",
filepath.display()
));
}
let mut file = fileres?;
let mut nb_header_lines = 0;
let mut c = [0];
let mut more = true;
while more {
file.read_exact(&mut c)?;
if ['#', '%'].contains(&(c[0] as char)) {
nb_header_lines += 1;
loop {
file.read_exact(&mut c)?;
if c[0] == b'\n' {
break;
}
}
} else {
more = false;
log::debug!("file has {} nb headers lines", nb_header_lines);
}
}
Ok(nb_header_lines)
}
pub fn unweighted_csv_to_graphmap<N, Ty>(
filepath: &Path,
delim: u8,
) -> anyhow::Result<GraphMap<N, (), Ty>>
where
N: NodeTrait + std::hash::Hash + std::cmp::Eq + FromStr + std::fmt::Display,
Ty: EdgeType,
{
let nb_headers_line = get_header_size(filepath)?;
log::info!(
"directed_from_csv , got header nb lines {}",
nb_headers_line
);
let fileres = OpenOptions::new().read(true).open(filepath);
if fileres.is_err() {
log::error!(
"ProcessingState reload_json : reload could not open file {:?}",
filepath.as_os_str()
);
println!(
"unweighted_csv_to_graphmap could not open file {:?}",
filepath.as_os_str()
);
return Err(anyhow!(
"directed_from_csv could not open file {}",
filepath.display()
));
}
let mut file = fileres?;
let mut nb_skipped = 0;
let mut c = [0];
loop {
file.read_exact(&mut c)?;
if c[0] == b'\n' {
nb_skipped += 1;
}
if nb_skipped == nb_headers_line {
break;
}
}
let mut rdr = ReaderBuilder::new()
.delimiter(delim)
.flexible(false)
.has_headers(false)
.from_reader(file);
let nb_nodes_guess = 50_000; let mut graph = GraphMap::<N, (), Ty>::with_capacity(nb_nodes_guess, 500_000);
let mut nb_record = 0;
let mut nb_fields = 0;
let mut node1: N;
let mut node2: N;
for result in rdr.records() {
let record = result?;
if log::log_enabled!(Level::Info) && nb_record <= 5 {
log::info!("{:?}", record);
}
if nb_record == 0 {
nb_fields = record.len();
log::info!("nb fields = {}", nb_fields);
if nb_fields != 2 {
log::error!(
"unweighted_csv_to_graphmap got nb_fileds different from 2, check the delimitor , got {:?} as delimitor ",
delim as char
);
return Err(anyhow!(
"found only one field in record, check the delimitor , got {:?} as delimitor ",
delim as char
));
}
} else if record.len() != nb_fields {
println!(
"non constant number of fields at record {} first record has {}",
nb_record, nb_fields
);
return Err(anyhow!(
"non constant number of fields at record {} first record has {}",
nb_record,
nb_fields
));
}
let field = record.get(0).unwrap();
if let Ok(idx) = field.parse::<N>() {
node1 = idx;
graph.add_node(idx);
} else {
return Err(anyhow!(
"error decoding field 1 of record {}",
nb_record + 1
));
}
let field = record.get(1).unwrap();
if let Ok(idx) = field.parse::<N>() {
node2 = idx;
graph.add_node(idx);
} else {
return Err(anyhow!(
"error decoding field 2 of record {}",
nb_record + 1
));
}
graph.add_edge(node1, node2, ());
nb_record += 1;
if log::log_enabled!(Level::Info) && nb_record <= 5 {
log::info!("{:?}", record);
log::info!(" node1 {}, node2 {}", node1, node2);
}
} log::info!(
"directed_unweighted_csv_to_graph read nb record : {}",
nb_record
);
Ok(graph)
}
pub fn weighted_csv_to_graphmap<N, W, Ty>(
filepath: &Path,
delim: u8,
) -> anyhow::Result<GraphMap<N, W, Ty>>
where
N: NodeTrait + std::hash::Hash + std::cmp::Eq + FromStr + std::fmt::Display,
W: FromStr + Float,
Ty: EdgeType,
{
let nb_headers_line = get_header_size(filepath)?;
log::info!(
"weighted_csv_to_graphmap , got header nb lines {}",
nb_headers_line
);
let fileres = OpenOptions::new().read(true).open(filepath);
if fileres.is_err() {
log::error!(
"ProcessingState reload_json : reload could not open file {:?}",
filepath.as_os_str()
);
println!(
"weighted_csv_to_graphmap could not open file {:?}",
filepath.as_os_str()
);
return Err(anyhow!(
"weighted_csv_to_graphmap could not open file {}",
filepath.display()
));
}
let mut file = fileres?;
let mut nb_skipped = 0;
let mut c = [0];
loop {
file.read_exact(&mut c)?;
if c[0] == b'\n' {
nb_skipped += 1;
}
if nb_skipped == nb_headers_line {
break;
}
}
let mut rdr = ReaderBuilder::new()
.delimiter(delim)
.flexible(false)
.has_headers(false)
.from_reader(file);
let nb_nodes_guess = 50_000; let mut graph = GraphMap::<N, W, Ty>::with_capacity(nb_nodes_guess, 500_000);
let mut nb_record = 0;
let mut nb_fields = 0;
let mut node1: N;
let mut node2: N;
let mut weight: W;
for result in rdr.records() {
let record = result?;
if log::log_enabled!(Level::Info) && nb_record <= 5 {
log::info!("{:?}", record);
}
if nb_record == 0 {
nb_fields = record.len();
log::info!("nb fields = {}", nb_fields);
if nb_fields < 2 {
log::error!(
"csv::weighted_csv_to_graphmap : got nb_fileds different from 2, check the delimitor , got {:?} as delimitor ",
delim as char
);
return Err(anyhow!(
"found only one field in record, check the delimitor , got {:?} as delimitor ",
delim as char
));
}
} else {
if record.len() != nb_fields {
println!(
"non constant number of fields at record {} first record has {}",
nb_record, nb_fields
);
return Err(anyhow!(
"csv::weighted_csv_to_graphmap : non constant number of fields at record {} first record has {}",
nb_record,
nb_fields
));
}
}
let field = record.get(0).unwrap();
if let Ok(idx) = field.parse::<N>() {
node1 = idx;
graph.add_node(idx);
} else {
return Err(anyhow!(
"error decoding field 1 of record {}",
nb_record + 1
));
}
let field = record.get(1).unwrap();
if let Ok(idx) = field.parse::<N>() {
node2 = idx;
graph.add_node(idx);
} else {
return Err(anyhow!(
"error decoding field 2 of record {}",
nb_record + 1
));
}
if nb_fields == 3 {
let field = record.get(2).unwrap();
if let Ok(w) = field.parse::<W>() {
weight = w;
} else {
log::debug!("error decoding field 3 of record {}", nb_record + 1);
return Err(anyhow!(
"error decoding field 3 of record {}",
nb_record + 1
));
}
} else {
weight = W::one();
}
graph.add_edge(node1, node2, weight);
nb_record += 1;
if log::log_enabled!(Level::Info) && nb_record <= 5 {
log::info!("{:?}", record);
log::info!(" node1 {}, node2 {}", node1, node2);
}
} log::info!("weighted_csv_to_graphmap read nb record : {}", nb_record);
Ok(graph)
}
pub fn get_graph_indexation<N, E, Ty, Ix>(graph: &Graph<N, E, Ty, Ix>) -> NodeIndexation<N>
where
N: std::hash::Hash + std::cmp::Eq + Clone + core::fmt::Debug,
Ty: EdgeType,
Ix: IndexType,
{
let mut indexset = IndexSet::<N>::with_capacity(graph.node_count());
let mut i_node = 0;
let mut nodes_weights = graph.node_weights();
while let Some(n) = nodes_weights.next() {
if !indexset.insert(n.clone()) {
log::error!("could not insert node : {:?}, node rank : {:?}", n, i_node);
}
i_node += 1;
}
indexset
}
pub fn csv_to_csrmat<F>(
filepath: &Path,
directed: bool,
delim: u8,
) -> anyhow::Result<(MatRepr<F>, NodeIndexation<usize>)>
where
F: FromStr
+ Float
+ Lapack
+ ndarray::ScalarOperand
+ sprs::MulAcc
+ for<'r> std::ops::MulAssign<&'r F>
+ Default
+ Sync,
{
let res_csv = csv_to_trimat(filepath, directed, delim);
if res_csv.is_ok() {
let csrmat: CsMat<F> = res_csv.as_ref().unwrap().0.to_csr();
log::debug!(
" csrmat dims nb_rows {}, nb_cols {} ",
csrmat.rows(),
csrmat.cols()
);
Ok((MatRepr::from_csrmat(csrmat), res_csv.unwrap().1))
} else {
Err(res_csv.unwrap_err())
}
}
pub fn csv_to_csrmat_delimiters<F>(
filepath: &Path,
directed: bool,
) -> anyhow::Result<(MatRepr<F>, NodeIndexation<usize>)>
where
F: FromStr
+ Float
+ Lapack
+ ndarray::ScalarOperand
+ sprs::MulAcc
+ for<'r> std::ops::MulAssign<&'r F>
+ Sync
+ Default,
{
log::info!("\n\n csv_to_csrmat_delimiters, loading file {:?}", filepath);
let delimiters = ['\t', ',', ' ', ';'];
let mut res: anyhow::Result<(MatRepr<F>, NodeIndexation<usize>)> =
Err(anyhow!("not initializd"));
for delim in delimiters {
log::debug!(
"embedder trying reading {:?} with delimiter {:?}",
&filepath,
delim
);
res = csv_to_csrmat::<F>(filepath, directed, delim as u8);
if res.is_err() {
log::error!(
"embedder failed in csv_to_csrmat_delimiters, reading {:?}, with delimiter {:?} ",
&filepath,
delim
);
} else {
return res;
}
}
if res.is_err() {
log::error!("error : {:?}", res.as_ref().err());
log::error!(
"embedder failed in csv_to_csrmat_delimiters, reading {:?}, tested delimiers {:?}",
&filepath,
delimiters
);
std::process::exit(1);
};
res
}
pub fn csv_to_trimat<F>(
filepath: &Path,
directed: bool,
delim: u8,
) -> anyhow::Result<(TriMatI<F, usize>, NodeIndexation<usize>)>
where
F: FromStr
+ Float
+ Lapack
+ ndarray::ScalarOperand
+ sprs::MulAcc
+ for<'r> std::ops::MulAssign<&'r F>
+ Default,
{
let nb_headers_line = get_header_size(filepath)?;
log::info!(
"directed_from_csv , got header nb lines {}",
nb_headers_line
);
let mut nodeindex = NodeIndexation::<usize>::with_capacity(500000);
let mut hset = HashSet::<(usize, usize)>::new();
let fileres = OpenOptions::new().read(true).open(filepath);
if fileres.is_err() {
log::error!(
"ProcessingState reload_json : reload could not open file {:?}",
filepath.as_os_str()
);
println!(
"directed_from_csv could not open file {:?}",
filepath.as_os_str()
);
return Err(anyhow!(
"directed_from_csv could not open file {}",
filepath.display()
));
}
let file = fileres?;
let mut bufreader = BufReader::new(file);
let mut headerline = String::new();
for _ in 0..nb_headers_line {
bufreader.read_line(&mut headerline)?;
}
let nb_edges_guess = 500_000; let mut rows = Vec::<usize>::with_capacity(nb_edges_guess);
let mut cols = Vec::<usize>::with_capacity(nb_edges_guess);
let mut values = Vec::<F>::with_capacity(nb_edges_guess);
let mut node1: usize; let mut node2: usize;
let mut node_id1: usize; let mut node_id2: usize;
let mut weight: F;
let mut rowmax: usize = 0;
let mut colmax: usize = 0;
let mut nb_record = 0; let mut num_record: usize = 0;
let mut nb_fields = 0;
let mut nb_self_loop = 0;
let mut nb_nodes: usize = 0;
let nb_warnings = 10;
let mut nb_potential_asymetry: usize = 0;
let mut last_edge_inserted = (0usize, 0usize);
let mut rdr = ReaderBuilder::new()
.delimiter(delim)
.flexible(false)
.has_headers(false)
.from_reader(bufreader);
for result in rdr.records() {
num_record += 1;
let record = result?;
if log::log_enabled!(Level::Info) && nb_record <= 2 {
log::info!(" record num {:?}, {:?}", nb_record, record);
}
if nb_record == 0 {
nb_fields = record.len();
log::info!("nb fields = {}", nb_fields);
if nb_fields < 2 {
log::error!(
"found only one field in record, check the delimitor , got {:?} as delimitor ",
delim as char
);
return Err(anyhow!(
"found only one field in record, check the delimitor , got {:?} as delimitor ",
delim as char
));
}
} else {
if record.len() != nb_fields {
println!(
"non constant number of fields at record {} first record has {}",
num_record, nb_fields
);
return Err(anyhow!(
"non constant number of fields at record {} first record has {}",
num_record,
nb_fields
));
}
}
let field = record.get(0).unwrap();
if let Ok(node) = field.parse::<usize>() {
node_id1 = node;
let already = nodeindex.get_index_of(&node_id1);
match already {
Some(idx) => node1 = idx,
None => {
node1 = nodeindex.insert_full(node_id1).0;
log::trace!("inserting node num : {}, rank : {}", node_id1, node1);
nb_nodes += 1;
}
}
rowmax = rowmax.max(node1);
} else {
log::debug!("error decoding field 1 of record {}", num_record);
return Err(anyhow!("error decoding field 1 of record {}", num_record));
}
let field = record.get(1).unwrap();
if let Ok(node) = field.parse::<usize>() {
node_id2 = node;
let already = nodeindex.get_index_of(&node_id2);
match already {
Some(idx) => node2 = idx,
None => {
node2 = nodeindex.insert_full(node_id2).0;
log::trace!("inserting node num : {}, rank : {}", node_id2, node2);
nb_nodes += 1;
}
}
colmax = colmax.max(node2);
} else {
log::debug!("error decoding field 2 of record {}", num_record);
return Err(anyhow!("error decoding field 2 of record {}", num_record));
}
if node_id1 == node_id2 {
nb_self_loop += 1;
log::error!(
"csv_to_trimat got diagonal term ({},{}) at record {:?} record num {}",
node_id1,
node_id2,
record,
num_record
);
continue;
}
if !directed && !hset.insert((node1, node2)) {
if nb_potential_asymetry <= nb_warnings {
println!(
"2-uple ({:?}, {:?}) already present, record {}",
node_id1, node_id2, num_record
);
log::error!(
"2-uple ({:?}, {:?}) already present, record {}",
node_id1,
node_id2,
num_record
);
log::error!("last edge inserted : {:?}", last_edge_inserted);
log::error!("record read : {:?}", record);
}
}
if !directed {
rowmax = rowmax.max(node2);
colmax = colmax.max(node1);
}
if nb_fields == 3 {
let field = record.get(2).unwrap();
if let Ok(w) = field.parse::<F>() {
weight = w;
} else {
log::debug!("error decoding field 3 of record {}", nb_record + 1);
return Err(anyhow!(
"error decoding field 3 of record {}",
nb_record + 1
));
}
} else {
weight = F::one();
}
rows.push(node1);
cols.push(node2);
values.push(weight);
log::trace!("to insert : (node1, node2) : ({}, {})", node1, node2);
nb_record += 1;
if !directed {
if !hset.insert((node2, node1)) {
nb_potential_asymetry += 1;
if nb_potential_asymetry <= nb_warnings {
log::error!(
"undirected mode 2-uple ({:?}, {:?}) symetric edge already present, record {}",
node_id2,
node_id1,
nb_record
);
log::error!("last edge inserted : {:?}", last_edge_inserted);
log::error!("record read : {:?}", record);
}
}
rows.push(node2);
cols.push(node1);
values.push(weight);
}
last_edge_inserted.0 = node_id1;
last_edge_inserted.1 = node_id2;
if log::log_enabled!(Level::Info) && nb_record <= 5 {
log::info!("{:?}", record);
log::info!(" node1 {:?}, node2 {:?}", node1, node2);
}
} assert_eq!(rows.len(), cols.len());
let trimat =
TriMatI::<F, usize>::from_triplets((nodeindex.len(), nodeindex.len()), rows, cols, values);
log::debug!("trimat shape {:?}", trimat.shape());
assert_eq!(trimat.shape().0, nodeindex.len());
assert_eq!(nb_nodes, nodeindex.len());
log::info!(
"\n\n csv file read!, nb nodes {}, nb edges loaded {}, nb_record read {}",
nodeindex.len(),
nb_record,
num_record
);
log::info!(
"rowmax : {}, colmax : {}, nb_edges : {}",
rowmax,
colmax,
trimat.nnz()
);
log::info!(" nb diagonal terms filtered : {}", nb_self_loop);
if nb_potential_asymetry > 0 {
log::error!(
"\n\n csv_to_trimat : CHECK SYMETRY number of couples definded more than one time : {}",
nb_potential_asymetry
);
println!(
"\n\n csv_to_trimat : CHECK SYMETRY number of couples definded more than one time : {}",
nb_potential_asymetry
);
println!("csv_to_trimat took the first couple ")
}
if nb_self_loop > 0 {
println!("nb diagonal terms filtered : {}", nb_self_loop);
}
if log::log_enabled!(log::Level::Trace) {
log::trace!("dump of indexation set");
let iter = nodeindex.iter();
for (rank, id) in iter.enumerate() {
println!(" rank : {}, node id : {} ", rank, id);
}
}
Ok((trimat, nodeindex))
}
pub fn csv_to_trimat_delimiters<F>(
filepath: &Path,
directed: bool,
) -> anyhow::Result<(TriMatI<F, usize>, NodeIndexation<usize>)>
where
F: FromStr
+ Float
+ Lapack
+ ndarray::ScalarOperand
+ sprs::MulAcc
+ for<'r> std::ops::MulAssign<&'r F>
+ Default,
{
log::debug!("in csv_to_trimat_delimiters");
let delimiters = ['\t', ',', ' ', ';'];
let mut res: anyhow::Result<(TriMatI<F, usize>, NodeIndexation<usize>)> =
Err(anyhow!("res not initialized"));
for delim in delimiters {
log::debug!(
"embedder trying reading {:?} with delimiter {:?}",
&filepath,
delim
);
res = csv_to_trimat::<F>(filepath, directed, delim as u8);
if res.is_err() {
log::error!(
"embedder failed in csv_to_trimat_delimiters, reading {:?}, trying delimiter {:?} ",
&filepath,
delim
);
} else {
return res;
}
}
if res.is_err() {
log::error!("error : {:?}", res.as_ref().err());
log::error!(
"embedder failed in csv_to_csrmat_delimiters, reading {:?}, tested delimiers {:?}",
&filepath,
delimiters
);
std::process::exit(1);
};
res
}
#[cfg(test)]
mod tests {
use super::*;
fn log_init_test() {
let _ = env_logger::builder().is_test(true).try_init();
}
#[test]
fn test_directed_unweighted_csv_to_graphmap() {
log_init_test();
let path = Path::new(crate::DATADIR).join("wiki-Vote.txt");
let header_size = get_header_size(&path);
assert_eq!(header_size.unwrap(), 4);
println!(
"\n\n test_directed_unweighted_csv_to_graph data : {:?}",
path
);
let graph = unweighted_csv_to_graphmap::<u32, Directed>(&path, b'\t');
if let Err(err) = &graph {
eprintln!("ERROR: {}", err);
}
assert!(graph.is_ok());
}
#[test]
fn test_weighted_csv_to_trimat() {
println!("\n\n test_weighted_csv_to_trimat");
log_init_test();
let path = Path::new(crate::DATADIR)
.join("moreno_lesmis")
.join("out.moreno_lesmis_lesmis");
log::debug!("\n\n test_weighted_csv_to_trimat, loading file {:?}", path);
let header_size = get_header_size(&path);
assert_eq!(header_size.unwrap(), 2);
println!("\n\n test_weighted_csv_to_trimat, data : {:?}", path);
let trimat_res = csv_to_trimat::<f64>(&path, false, b' ');
if let Err(err) = &trimat_res {
eprintln!("ERROR: {}", err);
assert_eq!(1, 0);
}
let mut trimat_iter = trimat_res.as_ref().unwrap().0.triplet_iter();
let nodeset = &trimat_res.as_ref().unwrap().1;
for _ in 0..5 {
let triplet = trimat_iter.next().unwrap();
let node1 = nodeset.get_index(triplet.1.0).unwrap();
let node2 = nodeset.get_index(triplet.1.1).unwrap();
let value = triplet.0;
log::debug!("node1 {}, node2 {}, value {} ", node1, node2, value);
}
} }