use core_query::Dir;
use core_storage::{Direction, IdMap, Interner, Topology};
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: &Topology, 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(),
}
}
#[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,
}
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,
}
}
}
#[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: &Topology,
idmap: &IdMap,
syms: &Interner,
labels: &[u32],
config: &PageRankConfig,
) -> PageRankReport {
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>> = vec![Vec::new(); n];
for &et in &etypes {
for (i, &id) in node_ids.iter().enumerate() {
match config.direction {
AlgoDir::Out => {
for &nbr in topo.neighbors(et, Direction::Out, id).as_ref() {
if let Some(&j) = id_to_idx.get(&nbr) {
if !send_to[i].contains(&j) {
send_to[i].push(j);
}
}
}
}
AlgoDir::In => {
for &nbr in topo.neighbors(et, Direction::In, id).as_ref() {
if let Some(&j) = id_to_idx.get(&nbr) {
if !send_to[i].contains(&j) {
send_to[i].push(j);
}
}
}
}
AlgoDir::Both => {
for dir in [Direction::Out, Direction::In] {
for &nbr in topo.neighbors(et, dir, id).as_ref() {
if let Some(&j) = id_to_idx.get(&nbr) {
if !send_to[i].contains(&j) {
send_to[i].push(j);
}
}
}
}
}
}
}
}
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_deg = send.len();
if out_deg == 0 {
dangling.push(i);
} else {
let w = 1.0 / out_deg as f64;
for &j in send {
receive_from[j].push((i, w));
}
}
}
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,
}
impl Default for WccConfig {
fn default() -> Self {
Self {
edge_type: None,
budget_ms: 5_000,
}
}
}
#[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: &Topology,
idmap: &IdMap,
syms: &Interner,
labels: &[u32],
config: &WccConfig,
) -> WccReport {
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) {
uf.union(i, j);
}
}
for &nbr in topo.neighbors(et, Direction::In, id).as_ref() {
if let Some(&j) = id_to_idx.get(&nbr) {
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,
}
impl Default for DegreeConfig {
fn default() -> Self {
Self {
edge_type: None,
direction: AlgoDir::Both,
budget_ms: 5_000,
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct DegreeReport {
pub scores: Vec<(String, u64)>,
pub truncated: bool,
}
pub(crate) fn degree_centrality(
topo: &Topology,
idmap: &IdMap,
syms: &Interner,
labels: &[u32],
config: &DegreeConfig,
) -> DegreeReport {
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 {
match config.direction {
AlgoDir::Out => {
degrees[i] += topo.neighbors(et, Direction::Out, id).len() as u64;
}
AlgoDir::In => {
degrees[i] += topo.neighbors(et, Direction::In, id).len() as u64;
}
AlgoDir::Both => {
degrees[i] += topo.neighbors(et, Direction::Out, id).len() as u64;
degrees[i] += topo.neighbors(et, Direction::In, id).len() as u64;
}
}
}
}
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 }
}