use core_query::Dir;
use core_storage::v8::seam::TopologyView;
use core_storage::{Direction, EdgePropsView, IdMap, Interner, Value};
use serde::{Deserialize, Serialize};
use std::collections::BTreeMap;
use std::time::{Duration, Instant};
fn live_nodes(idmap: &IdMap, labels: &[u32]) -> (Vec<u32>, Vec<String>) {
let n = idmap.len() as u32;
let mut ids = Vec::new();
let mut keys = Vec::new();
for id in 0..n {
let Some(key) = idmap.key_of(id) else {
continue;
};
let Some(&sym) = labels.get(id as usize) else {
continue;
};
if sym == u32::MAX {
continue; }
ids.push(id);
keys.push(key.to_string());
}
(ids, keys)
}
fn resolve_etype(syms: &Interner, edge_type: Option<&str>) -> Option<Option<u32>> {
match edge_type {
None => Some(None), Some(name) => {
let sym = syms.get(name)?; Some(Some(sym))
}
}
}
fn etypes_filtered(topo: &TopologyView, filter: Option<u32>) -> Vec<u32> {
match filter {
Some(sym) => {
let all: Vec<u32> = topo.etypes().collect();
if all.contains(&sym) {
vec![sym]
} else {
vec![]
}
}
None => topo.etypes().collect(),
}
}
fn resolve_etypes_multi(syms: &Interner, topo: &TopologyView, names: &[String]) -> Vec<u32> {
if names.is_empty() {
return topo.etypes().collect();
}
let mut out = Vec::new();
for name in names {
if let Some(sym) = syms.get(name) {
if !out.contains(&sym) {
out.push(sym);
}
}
}
out
}
fn live_nodes_for_label(
idmap: &IdMap,
syms: &Interner,
labels: &[u32],
label: Option<&str>,
) -> (Vec<u32>, Vec<String>) {
let want = match label {
None => None,
Some(name) => match syms.get(name) {
Some(sym) => Some(sym),
None => return (Vec::new(), Vec::new()),
},
};
let n = idmap.len() as u32;
let mut ids = Vec::new();
let mut keys = Vec::new();
for id in 0..n {
let Some(key) = idmap.key_of(id) else {
continue;
};
let Some(&sym) = labels.get(id as usize) else {
continue;
};
if sym == u32::MAX {
continue; }
if let Some(want_sym) = want {
if sym != want_sym {
continue;
}
}
ids.push(id);
keys.push(key.to_string());
}
(ids, keys)
}
fn edge_weight(
edge_props: &EdgePropsView,
etype: u32,
src: u32,
dst: u32,
weight_prop: Option<&str>,
min_weight: Option<f64>,
) -> Option<f64> {
let w = match weight_prop {
None => 1.0,
Some(prop) => match edge_props.get(etype, src, dst, prop) {
Some(Value::Float(f)) => f,
Some(Value::Int(i)) => i as f64,
_ => 1.0,
},
};
match min_weight {
Some(min) if w < min => None,
_ => Some(w),
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(default)]
pub struct PageRankConfig {
pub damping: f64,
pub max_iters: u32,
pub tol: f64,
pub edge_type: Option<String>,
pub direction: AlgoDir,
pub budget_ms: u64,
pub weight_prop: Option<String>,
pub min_weight: Option<f64>,
}
impl Default for PageRankConfig {
fn default() -> Self {
Self {
damping: 0.85,
max_iters: 50,
tol: 1e-6,
edge_type: None,
direction: AlgoDir::Out,
budget_ms: 5_000,
weight_prop: None,
min_weight: None,
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct PageRankReport {
pub scores: Vec<(String, f64)>,
pub converged: bool,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "lowercase")]
pub enum AlgoDir {
Out,
In,
Both,
}
impl From<Dir> for AlgoDir {
fn from(d: Dir) -> Self {
match d {
Dir::Out => AlgoDir::Out,
Dir::In => AlgoDir::In,
Dir::Both => AlgoDir::Both,
}
}
}
pub(crate) fn pagerank(
topo: &TopologyView,
idmap: &IdMap,
syms: &Interner,
labels: &[u32],
edge_props: &EdgePropsView,
config: &PageRankConfig,
) -> PageRankReport {
let weighted = config.weight_prop.is_some() || config.min_weight.is_some();
let deadline = if config.budget_ms > 0 {
Some(Instant::now() + Duration::from_millis(config.budget_ms))
} else {
None
};
let (node_ids, node_keys) = live_nodes(idmap, labels);
let n = node_ids.len();
if n == 0 {
return PageRankReport {
scores: Vec::new(),
converged: true,
};
}
let max_id = topo.etypes().count(); let _ = max_id;
let mut id_to_idx: BTreeMap<u32, usize> = BTreeMap::new();
for (i, &id) in node_ids.iter().enumerate() {
id_to_idx.insert(id, i);
}
let etype_filter = match resolve_etype(syms, config.edge_type.as_deref()) {
None => {
let score = 1.0 / n as f64;
let mut scores: Vec<(String, f64)> =
node_keys.iter().map(|k| (k.clone(), score)).collect();
scores.sort_by(|(ka, sa), (kb, sb)| {
sb.partial_cmp(sa)
.unwrap_or(std::cmp::Ordering::Equal)
.then(ka.cmp(kb))
});
return PageRankReport {
scores,
converged: true,
};
}
Some(f) => f,
};
let etypes = etypes_filtered(topo, etype_filter);
let mut send_to: Vec<Vec<(usize, f64)>> = vec![Vec::new(); n];
for &et in &etypes {
for (i, &id) in node_ids.iter().enumerate() {
let dirs: &[Direction] = match config.direction {
AlgoDir::Out => &[Direction::Out],
AlgoDir::In => &[Direction::In],
AlgoDir::Both => &[Direction::Out, Direction::In],
};
for &dir in dirs {
for &nbr in topo.neighbors(et, dir, id).as_ref() {
let Some(&j) = id_to_idx.get(&nbr) else {
continue;
};
if weighted {
let Some(w) = edge_weight(
edge_props,
et,
id,
nbr,
config.weight_prop.as_deref(),
config.min_weight,
) else {
continue; };
if let Some(entry) = send_to[i].iter_mut().find(|(k, _)| *k == j) {
entry.1 += w;
} else {
send_to[i].push((j, w));
}
} else if !send_to[i].iter().any(|(k, _)| *k == j) {
send_to[i].push((j, 1.0));
}
}
}
}
}
let mut receive_from: Vec<Vec<(usize, f64)>> = vec![Vec::new(); n];
let mut dangling: Vec<usize> = Vec::new();
for (i, send) in send_to.iter().enumerate() {
let out_weight: f64 = send.iter().map(|(_, w)| w).sum();
if send.is_empty() || out_weight <= 0.0 {
dangling.push(i);
} else {
for &(j, w) in send {
receive_from[j].push((i, w / out_weight));
}
}
}
let nf = n as f64;
let d = config.damping;
let teleport = (1.0 - d) / nf;
let mut pr: Vec<f64> = vec![1.0 / nf; n];
let mut converged = false;
for _iter in 0..config.max_iters {
if let Some(dl) = deadline {
if Instant::now() >= dl {
break;
}
}
let dangling_sum: f64 = dangling.iter().map(|&i| pr[i]).sum::<f64>() * d / nf;
let mut new_pr = vec![teleport + dangling_sum; n];
for j in 0..n {
let received: f64 = receive_from[j].iter().map(|&(i, w)| pr[i] * w).sum();
new_pr[j] += d * received;
}
let delta: f64 = pr
.iter()
.zip(new_pr.iter())
.map(|(a, b)| (a - b).abs())
.sum();
pr = new_pr;
if delta < config.tol {
converged = true;
break;
}
}
let mut scores: Vec<(String, f64)> = node_keys.into_iter().zip(pr).collect();
scores.sort_by(|(ka, sa), (kb, sb)| {
sb.partial_cmp(sa)
.unwrap_or(std::cmp::Ordering::Equal)
.then(ka.cmp(kb))
});
PageRankReport { scores, converged }
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(default)]
pub struct WccConfig {
pub edge_type: Option<String>,
pub budget_ms: u64,
pub weight_prop: Option<String>,
pub min_weight: Option<f64>,
}
impl Default for WccConfig {
fn default() -> Self {
Self {
edge_type: None,
budget_ms: 5_000,
weight_prop: None,
min_weight: None,
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct WccReport {
pub components: Vec<(String, String)>,
pub truncated: bool,
}
struct UnionFind {
parent: Vec<usize>,
rank: Vec<u8>,
}
impl UnionFind {
fn new(n: usize) -> Self {
Self {
parent: (0..n).collect(),
rank: vec![0; n],
}
}
fn find(&mut self, mut x: usize) -> usize {
while self.parent[x] != x {
self.parent[x] = self.parent[self.parent[x]]; x = self.parent[x];
}
x
}
fn union(&mut self, a: usize, b: usize) {
let ra = self.find(a);
let rb = self.find(b);
if ra == rb {
return;
}
match self.rank[ra].cmp(&self.rank[rb]) {
std::cmp::Ordering::Less => self.parent[ra] = rb,
std::cmp::Ordering::Greater => self.parent[rb] = ra,
std::cmp::Ordering::Equal => {
self.parent[rb] = ra;
self.rank[ra] += 1;
}
}
}
}
pub(crate) fn wcc(
topo: &TopologyView,
idmap: &IdMap,
syms: &Interner,
labels: &[u32],
edge_props: &EdgePropsView,
config: &WccConfig,
) -> WccReport {
let weighted = config.weight_prop.is_some() || config.min_weight.is_some();
let deadline = if config.budget_ms > 0 {
Some(Instant::now() + Duration::from_millis(config.budget_ms))
} else {
None
};
let (node_ids, node_keys) = live_nodes(idmap, labels);
let n = node_ids.len();
if n == 0 {
return WccReport {
components: Vec::new(),
truncated: false,
};
}
let mut id_to_idx: BTreeMap<u32, usize> = BTreeMap::new();
for (i, &id) in node_ids.iter().enumerate() {
id_to_idx.insert(id, i);
}
let etype_filter = match resolve_etype(syms, config.edge_type.as_deref()) {
None => {
let mut components: Vec<(String, String)> =
node_keys.iter().map(|k| (k.clone(), k.clone())).collect();
components.sort();
return WccReport {
components,
truncated: false,
};
}
Some(f) => f,
};
let etypes = etypes_filtered(topo, etype_filter);
let mut uf = UnionFind::new(n);
let mut truncated = false;
'outer: for &et in &etypes {
for (i, &id) in node_ids.iter().enumerate() {
if let Some(dl) = deadline {
if Instant::now() >= dl {
truncated = true;
break 'outer;
}
}
for &nbr in topo.neighbors(et, Direction::Out, id).as_ref() {
if let Some(&j) = id_to_idx.get(&nbr) {
if weighted
&& edge_weight(
edge_props,
et,
id,
nbr,
config.weight_prop.as_deref(),
config.min_weight,
)
.is_none()
{
continue; }
uf.union(i, j);
}
}
for &nbr in topo.neighbors(et, Direction::In, id).as_ref() {
if let Some(&j) = id_to_idx.get(&nbr) {
if weighted
&& edge_weight(
edge_props,
et,
nbr,
id,
config.weight_prop.as_deref(),
config.min_weight,
)
.is_none()
{
continue; }
uf.union(i, j);
}
}
}
}
let mut root_min_key: BTreeMap<usize, &str> = BTreeMap::new();
for (i, key_str) in node_keys.iter().enumerate() {
let root = uf.find(i);
let key = key_str.as_str();
let entry = root_min_key.entry(root).or_insert(key);
if key < *entry {
*entry = key;
}
}
let mut components: Vec<(String, String)> = node_keys
.iter()
.enumerate()
.map(|(i, key_str)| {
let root = uf.find(i);
let comp_id = root_min_key[&root].to_string();
(key_str.clone(), comp_id)
})
.collect();
components.sort_by(|(ka, ca), (kb, cb)| ca.cmp(cb).then(ka.cmp(kb)));
WccReport {
components,
truncated,
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(default)]
pub struct DegreeConfig {
pub edge_type: Option<String>,
pub direction: AlgoDir,
pub budget_ms: u64,
pub weight_prop: Option<String>,
pub min_weight: Option<f64>,
}
impl Default for DegreeConfig {
fn default() -> Self {
Self {
edge_type: None,
direction: AlgoDir::Both,
budget_ms: 5_000,
weight_prop: None,
min_weight: None,
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct DegreeReport {
pub scores: Vec<(String, u64)>,
pub truncated: bool,
}
pub(crate) fn degree_centrality(
topo: &TopologyView,
idmap: &IdMap,
syms: &Interner,
labels: &[u32],
edge_props: &EdgePropsView,
config: &DegreeConfig,
) -> DegreeReport {
let weighted = config.weight_prop.is_some() || config.min_weight.is_some();
let deadline = if config.budget_ms > 0 {
Some(Instant::now() + Duration::from_millis(config.budget_ms))
} else {
None
};
let (node_ids, node_keys) = live_nodes(idmap, labels);
let n = node_ids.len();
if n == 0 {
return DegreeReport {
scores: Vec::new(),
truncated: false,
};
}
let etype_filter = match resolve_etype(syms, config.edge_type.as_deref()) {
None => {
let scores = node_keys.iter().map(|k| (k.clone(), 0u64)).collect();
return DegreeReport {
scores,
truncated: false,
};
}
Some(f) => f,
};
let etypes = etypes_filtered(topo, etype_filter);
let mut degrees: Vec<u64> = vec![0u64; n];
let mut truncated = false;
for (i, &id) in node_ids.iter().enumerate() {
if let Some(dl) = deadline {
if Instant::now() >= dl {
truncated = true;
break;
}
}
for &et in &etypes {
let dirs: &[Direction] = match config.direction {
AlgoDir::Out => &[Direction::Out],
AlgoDir::In => &[Direction::In],
AlgoDir::Both => &[Direction::Out, Direction::In],
};
for &dir in dirs {
if !weighted {
degrees[i] += topo.neighbors(et, dir, id).len() as u64;
continue;
}
for &nbr in topo.neighbors(et, dir, id).as_ref() {
let (src, dst) = match dir {
Direction::Out => (id, nbr),
Direction::In => (nbr, id),
};
if edge_weight(
edge_props,
et,
src,
dst,
config.weight_prop.as_deref(),
config.min_weight,
)
.is_some()
{
degrees[i] += 1;
}
}
}
}
}
let mut scores: Vec<(String, u64)> = node_keys.into_iter().zip(degrees).collect();
scores.sort_by(|(ka, da), (kb, db)| db.cmp(da).then(ka.cmp(kb)));
DegreeReport { scores, truncated }
}
type WeightedAdj = Vec<Vec<(usize, f64)>>;
struct AggregatedLevel {
renumbered: Vec<usize>,
n: usize,
adj: WeightedAdj,
self_weight: Vec<f64>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(default)]
pub struct LouvainConfig {
pub edge_types: Vec<String>,
pub weight_prop: Option<String>,
pub min_weight: Option<f64>,
pub resolution: f64,
pub max_passes: u32,
pub max_sweeps: u32,
pub budget_ms: u64,
pub node_label: Option<String>,
}
impl Default for LouvainConfig {
fn default() -> Self {
Self {
edge_types: Vec::new(),
weight_prop: None,
min_weight: None,
resolution: 1.0,
max_passes: 10,
max_sweeps: 20,
budget_ms: 5_000,
node_label: None,
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub struct Community {
pub id: u32,
pub members: Vec<String>,
pub internal_weight: f64,
pub cohesion: f64,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub struct CommunityReport {
pub communities: Vec<Community>,
pub modularity: f64,
pub truncated: bool,
}
fn local_moving(
n: usize,
adj: &WeightedAdj,
self_weight: &[f64],
resolution: f64,
max_sweeps: u32,
deadline: Option<Instant>,
) -> (Vec<usize>, bool) {
let k: Vec<f64> = (0..n)
.map(|i| adj[i].iter().map(|&(_, w)| w).sum::<f64>() + 2.0 * self_weight[i])
.collect();
let m: f64 = k.iter().sum::<f64>() / 2.0;
let mut community_of: Vec<usize> = (0..n).collect();
if m <= 0.0 {
return (community_of, false);
}
let mut tot: Vec<f64> = k.clone();
for _sweep in 0..max_sweeps {
if let Some(dl) = deadline {
if Instant::now() >= dl {
return (community_of, true);
}
}
let mut improved = false;
for i in 0..n {
let ci = community_of[i];
tot[ci] -= k[i];
let mut neighbor_weights: BTreeMap<usize, f64> = BTreeMap::new();
for &(j, w) in &adj[i] {
if j == i {
continue; }
*neighbor_weights.entry(community_of[j]).or_insert(0.0) += w;
}
let gain = |c: usize, w_in: f64| -> f64 {
w_in / m - resolution * tot[c] * k[i] / (2.0 * m * m)
};
let mut best_c = ci;
let mut best_gain = gain(ci, neighbor_weights.get(&ci).copied().unwrap_or(0.0));
for (&c, &w_in) in &neighbor_weights {
if c == ci {
continue;
}
let g = gain(c, w_in);
if g > best_gain + 1e-12 {
best_gain = g;
best_c = c;
}
}
tot[best_c] += k[i];
if best_c != ci {
community_of[i] = best_c;
improved = true;
}
}
if !improved {
break;
}
}
(community_of, false)
}
fn aggregate(
n: usize,
adj: &WeightedAdj,
self_weight: &[f64],
community_of: &[usize],
) -> Option<AggregatedLevel> {
let mut remap: BTreeMap<usize, usize> = BTreeMap::new();
let mut next_id = 0usize;
let mut renumbered: Vec<usize> = vec![0; n];
for (i, item) in renumbered.iter_mut().enumerate() {
let c = community_of[i];
let idx = *remap.entry(c).or_insert_with(|| {
let id = next_id;
next_id += 1;
id
});
*item = idx;
}
let new_n = next_id;
if new_n == n {
return None; }
let mut new_self_weight = vec![0.0; new_n];
let mut new_adj_map: Vec<BTreeMap<usize, f64>> = vec![BTreeMap::new(); new_n];
for i in 0..n {
let ci = renumbered[i];
new_self_weight[ci] += self_weight[i];
for &(j, w) in &adj[i] {
if j < i {
continue; }
let cj = renumbered[j];
if ci == cj {
new_self_weight[ci] += w;
} else {
*new_adj_map[ci].entry(cj).or_insert(0.0) += w;
*new_adj_map[cj].entry(ci).or_insert(0.0) += w;
}
}
}
let new_adj: WeightedAdj = new_adj_map
.into_iter()
.map(|map| map.into_iter().collect())
.collect();
Some(AggregatedLevel {
renumbered,
n: new_n,
adj: new_adj,
self_weight: new_self_weight,
})
}
pub(crate) fn louvain(
topo: &TopologyView,
idmap: &IdMap,
syms: &Interner,
labels: &[u32],
edge_props: &EdgePropsView,
config: &LouvainConfig,
) -> CommunityReport {
let deadline = if config.budget_ms > 0 {
Some(Instant::now() + Duration::from_millis(config.budget_ms))
} else {
None
};
let (raw_ids, raw_keys) =
live_nodes_for_label(idmap, syms, labels, config.node_label.as_deref());
let mut order: Vec<usize> = (0..raw_ids.len()).collect();
order.sort_by(|&a, &b| raw_keys[a].cmp(&raw_keys[b]));
let node_ids: Vec<u32> = order.iter().map(|&i| raw_ids[i]).collect();
let node_keys: Vec<String> = order.iter().map(|&i| raw_keys[i].clone()).collect();
let n0 = node_ids.len();
if n0 == 0 {
return CommunityReport {
communities: Vec::new(),
modularity: 0.0,
truncated: false,
};
}
let mut id_to_idx: BTreeMap<u32, usize> = BTreeMap::new();
for (i, &id) in node_ids.iter().enumerate() {
id_to_idx.insert(id, i);
}
let etypes = resolve_etypes_multi(syms, topo, &config.edge_types);
let mut edge_weight_map: BTreeMap<(usize, usize), f64> = BTreeMap::new();
for &et in &etypes {
for (i, &id) in node_ids.iter().enumerate() {
for &nbr in topo.neighbors(et, Direction::Out, id).as_ref() {
if nbr == id {
continue; }
let Some(&j) = id_to_idx.get(&nbr) else {
continue; };
let Some(w) = edge_weight(
edge_props,
et,
id,
nbr,
config.weight_prop.as_deref(),
config.min_weight,
) else {
continue; };
let key = if i < j { (i, j) } else { (j, i) };
*edge_weight_map.entry(key).or_insert(0.0) += w;
}
}
}
let m: f64 = edge_weight_map.values().sum();
let mut adj: WeightedAdj = vec![Vec::new(); n0];
for (&(a, b), &w) in &edge_weight_map {
adj[a].push((b, w));
adj[b].push((a, w));
}
let mut self_weight: Vec<f64> = vec![0.0; n0];
let mut owner: Vec<usize> = (0..n0).collect();
let mut truncated = false;
let mut n = n0;
if m > 0.0 {
'passes: for _pass in 0..config.max_passes {
let (community_of, hit_budget) = local_moving(
n,
&adj,
&self_weight,
config.resolution,
config.max_sweeps,
deadline,
);
if hit_budget {
owner = owner.iter().map(|&o| community_of[o]).collect();
truncated = true;
break 'passes;
}
let Some(level) = aggregate(n, &adj, &self_weight, &community_of) else {
owner = owner.iter().map(|&o| community_of[o]).collect();
break 'passes;
};
owner = owner.iter().map(|&o| level.renumbered[o]).collect();
n = level.n;
adj = level.adj;
self_weight = level.self_weight;
}
}
let mut internal: BTreeMap<usize, f64> = BTreeMap::new();
let mut leaving: BTreeMap<usize, f64> = BTreeMap::new();
for (&(a, b), &w) in &edge_weight_map {
let ca = owner[a];
let cb = owner[b];
if ca == cb {
*internal.entry(ca).or_insert(0.0) += w;
} else {
*leaving.entry(ca).or_insert(0.0) += w;
*leaving.entry(cb).or_insert(0.0) += w;
}
}
let mut members_by_community: BTreeMap<usize, Vec<String>> = BTreeMap::new();
for (i, key) in node_keys.iter().enumerate() {
members_by_community
.entry(owner[i])
.or_default()
.push(key.clone());
}
let modularity = if m > 0.0 {
members_by_community
.keys()
.map(|c| {
let internal_w = internal.get(c).copied().unwrap_or(0.0);
let leaving_w = leaving.get(c).copied().unwrap_or(0.0);
let sigma_tot = 2.0 * internal_w + leaving_w;
internal_w / m - config.resolution * (sigma_tot * sigma_tot) / (4.0 * m * m)
})
.sum()
} else {
0.0
};
let mut communities: Vec<Community> = members_by_community
.into_iter()
.map(|(c, mut members)| {
members.sort();
let internal_w = internal.get(&c).copied().unwrap_or(0.0);
let leaving_w = leaving.get(&c).copied().unwrap_or(0.0);
let cohesion = if internal_w + leaving_w > 0.0 {
internal_w / (internal_w + leaving_w)
} else {
1.0
};
Community {
id: 0, members,
internal_weight: internal_w,
cohesion,
}
})
.collect();
communities.sort_by(|a, b| {
b.members
.len()
.cmp(&a.members.len())
.then_with(|| a.members[0].cmp(&b.members[0]))
});
for (i, c) in communities.iter_mut().enumerate() {
c.id = i as u32;
}
CommunityReport {
communities,
modularity,
truncated,
}
}