1use core_query::Dir;
33use core_storage::v8::seam::TopologyView;
34use core_storage::{Direction, IdMap, Interner};
35use serde::{Deserialize, Serialize};
36use std::collections::BTreeMap;
37use std::time::{Duration, Instant};
38
39fn live_nodes(idmap: &IdMap, labels: &[u32]) -> (Vec<u32>, Vec<String>) {
48 let n = idmap.len() as u32;
49 let mut ids = Vec::new();
50 let mut keys = Vec::new();
51 for id in 0..n {
52 let Some(key) = idmap.key_of(id) else {
53 continue;
54 };
55 let Some(&sym) = labels.get(id as usize) else {
56 continue;
57 };
58 if sym == u32::MAX {
59 continue; }
61 ids.push(id);
62 keys.push(key.to_string());
63 }
64 (ids, keys)
65}
66
67fn resolve_etype(syms: &Interner, edge_type: Option<&str>) -> Option<Option<u32>> {
73 match edge_type {
74 None => Some(None), Some(name) => {
76 let sym = syms.get(name)?; Some(Some(sym))
78 }
79 }
80}
81
82fn etypes_filtered(topo: &TopologyView, filter: Option<u32>) -> Vec<u32> {
84 match filter {
85 Some(sym) => {
86 let all: Vec<u32> = topo.etypes().collect();
88 if all.contains(&sym) {
89 vec![sym]
90 } else {
91 vec![]
92 }
93 }
94 None => topo.etypes().collect(),
95 }
96}
97
98#[derive(Debug, Clone, Serialize, Deserialize)]
104#[serde(default)]
105pub struct PageRankConfig {
106 pub damping: f64,
109 pub max_iters: u32,
111 pub tol: f64,
113 pub edge_type: Option<String>,
115 pub direction: AlgoDir,
119 pub budget_ms: u64,
122}
123
124impl Default for PageRankConfig {
125 fn default() -> Self {
126 Self {
127 damping: 0.85,
128 max_iters: 50,
129 tol: 1e-6,
130 edge_type: None,
131 direction: AlgoDir::Out,
132 budget_ms: 5_000,
133 }
134 }
135}
136
137#[derive(Debug, Clone, Serialize, Deserialize)]
139pub struct PageRankReport {
140 pub scores: Vec<(String, f64)>,
143 pub converged: bool,
147}
148
149#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
151#[serde(rename_all = "lowercase")]
152pub enum AlgoDir {
153 Out,
155 In,
157 Both,
159}
160
161impl From<Dir> for AlgoDir {
162 fn from(d: Dir) -> Self {
163 match d {
164 Dir::Out => AlgoDir::Out,
165 Dir::In => AlgoDir::In,
166 Dir::Both => AlgoDir::Both,
167 }
168 }
169}
170
171pub(crate) fn pagerank(
175 topo: &TopologyView,
176 idmap: &IdMap,
177 syms: &Interner,
178 labels: &[u32],
179 config: &PageRankConfig,
180) -> PageRankReport {
181 let deadline = if config.budget_ms > 0 {
182 Some(Instant::now() + Duration::from_millis(config.budget_ms))
183 } else {
184 None
185 };
186
187 let (node_ids, node_keys) = live_nodes(idmap, labels);
188 let n = node_ids.len();
189
190 if n == 0 {
191 return PageRankReport {
192 scores: Vec::new(),
193 converged: true,
194 };
195 }
196
197 let max_id = topo.etypes().count(); let _ = max_id;
200 let mut id_to_idx: BTreeMap<u32, usize> = BTreeMap::new();
201 for (i, &id) in node_ids.iter().enumerate() {
202 id_to_idx.insert(id, i);
203 }
204
205 let etype_filter = match resolve_etype(syms, config.edge_type.as_deref()) {
207 None => {
208 let score = 1.0 / n as f64;
210 let mut scores: Vec<(String, f64)> =
211 node_keys.iter().map(|k| (k.clone(), score)).collect();
212 scores.sort_by(|(ka, sa), (kb, sb)| {
213 sb.partial_cmp(sa)
214 .unwrap_or(std::cmp::Ordering::Equal)
215 .then(ka.cmp(kb))
216 });
217 return PageRankReport {
218 scores,
219 converged: true,
220 };
221 }
222 Some(f) => f,
223 };
224
225 let etypes = etypes_filtered(topo, etype_filter);
226
227 let mut send_to: Vec<Vec<usize>> = vec![Vec::new(); n];
231
232 for &et in &etypes {
233 for (i, &id) in node_ids.iter().enumerate() {
234 match config.direction {
235 AlgoDir::Out => {
236 for &nbr in topo.neighbors(et, Direction::Out, id).as_ref() {
238 if let Some(&j) = id_to_idx.get(&nbr) {
239 if !send_to[i].contains(&j) {
240 send_to[i].push(j);
241 }
242 }
243 }
244 }
245 AlgoDir::In => {
246 for &nbr in topo.neighbors(et, Direction::In, id).as_ref() {
248 if let Some(&j) = id_to_idx.get(&nbr) {
249 if !send_to[i].contains(&j) {
250 send_to[i].push(j);
251 }
252 }
253 }
254 }
255 AlgoDir::Both => {
256 for dir in [Direction::Out, Direction::In] {
258 for &nbr in topo.neighbors(et, dir, id).as_ref() {
259 if let Some(&j) = id_to_idx.get(&nbr) {
260 if !send_to[i].contains(&j) {
261 send_to[i].push(j);
262 }
263 }
264 }
265 }
266 }
267 }
268 }
269 }
270
271 let mut receive_from: Vec<Vec<(usize, f64)>> = vec![Vec::new(); n];
274 let mut dangling: Vec<usize> = Vec::new();
275
276 for (i, send) in send_to.iter().enumerate() {
277 let out_deg = send.len();
278 if out_deg == 0 {
279 dangling.push(i);
280 } else {
281 let w = 1.0 / out_deg as f64;
282 for &j in send {
283 receive_from[j].push((i, w));
284 }
285 }
286 }
287
288 let nf = n as f64;
290 let d = config.damping;
291 let teleport = (1.0 - d) / nf;
292 let mut pr: Vec<f64> = vec![1.0 / nf; n];
293 let mut converged = false;
294
295 for _iter in 0..config.max_iters {
296 if let Some(dl) = deadline {
298 if Instant::now() >= dl {
299 break;
300 }
301 }
302
303 let dangling_sum: f64 = dangling.iter().map(|&i| pr[i]).sum::<f64>() * d / nf;
305
306 let mut new_pr = vec![teleport + dangling_sum; n];
307 for j in 0..n {
308 let received: f64 = receive_from[j].iter().map(|&(i, w)| pr[i] * w).sum();
309 new_pr[j] += d * received;
310 }
311
312 let delta: f64 = pr
314 .iter()
315 .zip(new_pr.iter())
316 .map(|(a, b)| (a - b).abs())
317 .sum();
318 pr = new_pr;
319
320 if delta < config.tol {
321 converged = true;
322 break;
323 }
324 }
325
326 let mut scores: Vec<(String, f64)> = node_keys.into_iter().zip(pr).collect();
328 scores.sort_by(|(ka, sa), (kb, sb)| {
329 sb.partial_cmp(sa)
330 .unwrap_or(std::cmp::Ordering::Equal)
331 .then(ka.cmp(kb))
332 });
333
334 PageRankReport { scores, converged }
335}
336
337#[derive(Debug, Clone, Serialize, Deserialize)]
343#[serde(default)]
344pub struct WccConfig {
345 pub edge_type: Option<String>,
347 pub budget_ms: u64,
349}
350
351impl Default for WccConfig {
352 fn default() -> Self {
353 Self {
354 edge_type: None,
355 budget_ms: 5_000,
356 }
357 }
358}
359
360#[derive(Debug, Clone, Serialize, Deserialize)]
362pub struct WccReport {
363 pub components: Vec<(String, String)>,
366 pub truncated: bool,
368}
369
370struct UnionFind {
372 parent: Vec<usize>,
373 rank: Vec<u8>,
374}
375
376impl UnionFind {
377 fn new(n: usize) -> Self {
378 Self {
379 parent: (0..n).collect(),
380 rank: vec![0; n],
381 }
382 }
383
384 fn find(&mut self, mut x: usize) -> usize {
385 while self.parent[x] != x {
386 self.parent[x] = self.parent[self.parent[x]]; x = self.parent[x];
388 }
389 x
390 }
391
392 fn union(&mut self, a: usize, b: usize) {
393 let ra = self.find(a);
394 let rb = self.find(b);
395 if ra == rb {
396 return;
397 }
398 match self.rank[ra].cmp(&self.rank[rb]) {
399 std::cmp::Ordering::Less => self.parent[ra] = rb,
400 std::cmp::Ordering::Greater => self.parent[rb] = ra,
401 std::cmp::Ordering::Equal => {
402 self.parent[rb] = ra;
403 self.rank[ra] += 1;
404 }
405 }
406 }
407}
408
409pub(crate) fn wcc(
413 topo: &TopologyView,
414 idmap: &IdMap,
415 syms: &Interner,
416 labels: &[u32],
417 config: &WccConfig,
418) -> WccReport {
419 let deadline = if config.budget_ms > 0 {
420 Some(Instant::now() + Duration::from_millis(config.budget_ms))
421 } else {
422 None
423 };
424
425 let (node_ids, node_keys) = live_nodes(idmap, labels);
426 let n = node_ids.len();
427
428 if n == 0 {
429 return WccReport {
430 components: Vec::new(),
431 truncated: false,
432 };
433 }
434
435 let mut id_to_idx: BTreeMap<u32, usize> = BTreeMap::new();
437 for (i, &id) in node_ids.iter().enumerate() {
438 id_to_idx.insert(id, i);
439 }
440
441 let etype_filter = match resolve_etype(syms, config.edge_type.as_deref()) {
443 None => {
444 let mut components: Vec<(String, String)> =
446 node_keys.iter().map(|k| (k.clone(), k.clone())).collect();
447 components.sort();
448 return WccReport {
449 components,
450 truncated: false,
451 };
452 }
453 Some(f) => f,
454 };
455
456 let etypes = etypes_filtered(topo, etype_filter);
457
458 let mut uf = UnionFind::new(n);
459 let mut truncated = false;
460
461 'outer: for &et in &etypes {
463 for (i, &id) in node_ids.iter().enumerate() {
464 if let Some(dl) = deadline {
465 if Instant::now() >= dl {
466 truncated = true;
467 break 'outer;
468 }
469 }
470 for &nbr in topo.neighbors(et, Direction::Out, id).as_ref() {
472 if let Some(&j) = id_to_idx.get(&nbr) {
473 uf.union(i, j);
474 }
475 }
476 for &nbr in topo.neighbors(et, Direction::In, id).as_ref() {
480 if let Some(&j) = id_to_idx.get(&nbr) {
481 uf.union(i, j);
482 }
483 }
484 }
485 }
486
487 let mut root_min_key: BTreeMap<usize, &str> = BTreeMap::new();
489 for (i, key_str) in node_keys.iter().enumerate() {
490 let root = uf.find(i);
491 let key = key_str.as_str();
492 let entry = root_min_key.entry(root).or_insert(key);
493 if key < *entry {
494 *entry = key;
495 }
496 }
497
498 let mut components: Vec<(String, String)> = node_keys
499 .iter()
500 .enumerate()
501 .map(|(i, key_str)| {
502 let root = uf.find(i);
503 let comp_id = root_min_key[&root].to_string();
504 (key_str.clone(), comp_id)
505 })
506 .collect();
507 components.sort_by(|(ka, ca), (kb, cb)| ca.cmp(cb).then(ka.cmp(kb)));
508
509 WccReport {
510 components,
511 truncated,
512 }
513}
514
515#[derive(Debug, Clone, Serialize, Deserialize)]
521#[serde(default)]
522pub struct DegreeConfig {
523 pub edge_type: Option<String>,
525 pub direction: AlgoDir,
527 pub budget_ms: u64,
529}
530
531impl Default for DegreeConfig {
532 fn default() -> Self {
533 Self {
534 edge_type: None,
535 direction: AlgoDir::Both,
536 budget_ms: 5_000,
537 }
538 }
539}
540
541#[derive(Debug, Clone, Serialize, Deserialize)]
543pub struct DegreeReport {
544 pub scores: Vec<(String, u64)>,
546 pub truncated: bool,
548}
549
550pub(crate) fn degree_centrality(
552 topo: &TopologyView,
553 idmap: &IdMap,
554 syms: &Interner,
555 labels: &[u32],
556 config: &DegreeConfig,
557) -> DegreeReport {
558 let deadline = if config.budget_ms > 0 {
559 Some(Instant::now() + Duration::from_millis(config.budget_ms))
560 } else {
561 None
562 };
563
564 let (node_ids, node_keys) = live_nodes(idmap, labels);
565 let n = node_ids.len();
566
567 if n == 0 {
568 return DegreeReport {
569 scores: Vec::new(),
570 truncated: false,
571 };
572 }
573
574 let etype_filter = match resolve_etype(syms, config.edge_type.as_deref()) {
576 None => {
577 let scores = node_keys.iter().map(|k| (k.clone(), 0u64)).collect();
579 return DegreeReport {
580 scores,
581 truncated: false,
582 };
583 }
584 Some(f) => f,
585 };
586
587 let etypes = etypes_filtered(topo, etype_filter);
588 let mut degrees: Vec<u64> = vec![0u64; n];
589 let mut truncated = false;
590
591 for (i, &id) in node_ids.iter().enumerate() {
592 if let Some(dl) = deadline {
593 if Instant::now() >= dl {
594 truncated = true;
595 break;
596 }
597 }
598 for &et in &etypes {
599 match config.direction {
600 AlgoDir::Out => {
601 degrees[i] += topo.neighbors(et, Direction::Out, id).len() as u64;
602 }
603 AlgoDir::In => {
604 degrees[i] += topo.neighbors(et, Direction::In, id).len() as u64;
605 }
606 AlgoDir::Both => {
607 degrees[i] += topo.neighbors(et, Direction::Out, id).len() as u64;
608 degrees[i] += topo.neighbors(et, Direction::In, id).len() as u64;
609 }
610 }
611 }
612 }
613
614 let mut scores: Vec<(String, u64)> = node_keys.into_iter().zip(degrees).collect();
615 scores.sort_by(|(ka, da), (kb, db)| db.cmp(da).then(ka.cmp(kb)));
616
617 DegreeReport { scores, truncated }
618}