|
| 1 | +"""Undirected hop builds both edge orientations without deduplicating the doubled frame. |
| 2 | +
|
| 3 | +Every consumer of the doubled frame dedups on the edge id, so the whole-frame dedup was |
| 4 | +redundant. Pins: parity with an independent oracle on a graph with self-loops, parallel |
| 5 | +edges and hub seeds on every engine; self-loops kept once through the multi-hop path; and |
| 6 | +categorical endpoint columns with differing category sets keep working. |
| 7 | +""" |
| 8 | +import numpy as np |
| 9 | +import pandas as pd |
| 10 | +import pytest |
| 11 | + |
| 12 | +import graphistry |
| 13 | +from graphistry.compute.ast import e_undirected, n |
| 14 | +from graphistry.compute.predicates.is_in import is_in |
| 15 | + |
| 16 | +ENGINES = ["pandas", "polars", "cudf"] |
| 17 | + |
| 18 | + |
| 19 | +def _frames(seed=0): |
| 20 | + rng = np.random.default_rng(seed) |
| 21 | + N, E = 3000, 12000 |
| 22 | + edges = pd.DataFrame({"s": rng.integers(0, N, E), "d": rng.integers(0, N, E)}) |
| 23 | + edges.loc[:40, "d"] = edges.loc[:40, "s"] # self-loops, some on hub seeds |
| 24 | + edges = pd.concat([edges, edges.iloc[:30]], ignore_index=True) # parallel edges |
| 25 | + nodes = pd.DataFrame({"id": np.arange(N)}) |
| 26 | + return nodes, edges |
| 27 | + |
| 28 | + |
| 29 | +def _graph(engine): |
| 30 | + nodes, edges = _frames() |
| 31 | + if engine == "polars": |
| 32 | + pl = pytest.importorskip("polars") |
| 33 | + nodes, edges = pl.from_pandas(nodes), pl.from_pandas(edges) |
| 34 | + elif engine == "cudf": |
| 35 | + cudf = pytest.importorskip("cudf") |
| 36 | + nodes, edges = cudf.from_pandas(nodes), cudf.from_pandas(edges) |
| 37 | + return graphistry.nodes(nodes, "id").edges(edges, "s", "d") |
| 38 | + |
| 39 | + |
| 40 | +def _oracle_ball(edges, seeds, hops): |
| 41 | + """Node set within `hops` undirected steps of the seeds, by plain set expansion.""" |
| 42 | + adj = {} |
| 43 | + for s, d in zip(edges["s"].tolist(), edges["d"].tolist()): |
| 44 | + adj.setdefault(s, set()).add(d) |
| 45 | + adj.setdefault(d, set()).add(s) |
| 46 | + frontier, seen = set(seeds), set(seeds) |
| 47 | + for _ in range(hops): |
| 48 | + nxt = set() |
| 49 | + for u in frontier: |
| 50 | + nxt |= adj.get(u, set()) |
| 51 | + frontier = nxt - seen |
| 52 | + seen |= nxt |
| 53 | + return seen |
| 54 | + |
| 55 | + |
| 56 | +def _ids(res): |
| 57 | + nodes = res._nodes |
| 58 | + df = nodes.to_pandas() if hasattr(nodes, "to_pandas") else pd.DataFrame(nodes) |
| 59 | + return set(int(v) for v in df["id"].tolist()) |
| 60 | + |
| 61 | + |
| 62 | +@pytest.mark.parametrize("engine", ENGINES) |
| 63 | +@pytest.mark.parametrize("hops", [1, 2]) |
| 64 | +def test_undirected_hop_matches_set_oracle_with_self_loops_and_parallel_edges(engine, hops): |
| 65 | + nodes, edges = _frames() |
| 66 | + deg = pd.concat([edges["s"], edges["d"]]).value_counts() |
| 67 | + seeds = [int(v) for v in deg.index[:5]] |
| 68 | + seeds.append(int(edges.loc[0, "s"])) # a seed carrying a self-loop |
| 69 | + g = _graph(engine) |
| 70 | + got = g.gfql([n({"id": is_in(seeds)}), e_undirected(hops=hops), n()], engine=engine) |
| 71 | + expected = _oracle_ball(edges, seeds, hops) |
| 72 | + # the wavefront keeps a seed only when an edge reaches it; the oracle includes every seed |
| 73 | + assert _ids(got) - set(seeds) == expected - set(seeds) |
| 74 | + assert _ids(got) <= expected |
| 75 | + |
| 76 | + |
| 77 | +@pytest.mark.parametrize("engine", ENGINES) |
| 78 | +@pytest.mark.parametrize("hops", [1, 2]) |
| 79 | +def test_undirected_edges_keep_every_self_loop_once(engine, hops): |
| 80 | + """hops=2 goes through the doubled-frame build; hops=1 is served by the seeded lane.""" |
| 81 | + nodes, edges = _frames() |
| 82 | + loop_seed = int(edges.loc[0, "s"]) |
| 83 | + g = _graph(engine) |
| 84 | + out = g.gfql([n({"id": is_in([loop_seed])}), e_undirected(hops=hops), n()], engine=engine) |
| 85 | + e = out._edges |
| 86 | + e = e.to_pandas() if hasattr(e, "to_pandas") else pd.DataFrame(e) |
| 87 | + loops = e[(e["s"] == loop_seed) & (e["d"] == loop_seed)] |
| 88 | + expected = edges[(edges["s"] == loop_seed) & (edges["d"] == loop_seed)] |
| 89 | + assert len(loops) == len(expected) |
| 90 | + # every returned edge row is one input edge row (parallel edges included): the |
| 91 | + # reverse orientation never surfaces as an extra row |
| 92 | + counts = e.groupby(["s", "d"]).size() |
| 93 | + source = edges.groupby(["s", "d"]).size() |
| 94 | + assert all(counts[k] <= source[k] for k in counts.index) |
| 95 | + |
| 96 | + |
| 97 | +@pytest.mark.parametrize("engine", ["pandas"]) # cuDF cannot concat categoricals of differing category sets |
| 98 | +def test_categorical_endpoints_with_different_category_sets(engine): |
| 99 | + edges = pd.DataFrame({"s": ["a", "b", "b", "a"], "d": ["b", "c", "b", "d"]}).astype( |
| 100 | + {"s": "category", "d": "category"}) |
| 101 | + nodes = pd.DataFrame({"id": list("abcd")}) |
| 102 | + if engine == "cudf": |
| 103 | + cudf = pytest.importorskip("cudf") |
| 104 | + nodes, edges = cudf.from_pandas(nodes), cudf.from_pandas(edges) |
| 105 | + g = graphistry.nodes(nodes, "id").edges(edges, "s", "d") |
| 106 | + seeds = pd.DataFrame({"id": ["a"]}) |
| 107 | + if engine == "cudf": |
| 108 | + import cudf |
| 109 | + seeds = cudf.from_pandas(seeds) |
| 110 | + out = g.hop(nodes=seeds, hops=2, direction="undirected") |
| 111 | + assert _ids_str(out) == {"a", "b", "c", "d"} |
| 112 | + assert len(out._edges) == 4 |
| 113 | + |
| 114 | + |
| 115 | +def _ids_str(res): |
| 116 | + nodes = res._nodes |
| 117 | + df = nodes.to_pandas() if hasattr(nodes, "to_pandas") else pd.DataFrame(nodes) |
| 118 | + return set(str(v) for v in df["id"].tolist()) |
0 commit comments