From 1da3fb8b3abf72da8903d999e483de56a03e2c1f Mon Sep 17 00:00:00 2001 From: B-Deprez Date: Fri, 18 Sep 2026 14:12:42 +0200 Subject: [PATCH] Phase 3: PEP 8 names, type hints and documented public API --- CHANGELOG.md | 31 +- .../0008-node-identity-is-preserved.md | 29 ++ docs/decisions/index.md | 1 + pyproject.toml | 19 +- src/garg_aml/_blocks.py | 223 +++++----- src/garg_aml/_ordering.py | 93 ++-- src/garg_aml/features.py | 274 +++++++----- src/garg_aml/measures.py | 271 +++++++----- src/garg_aml/preprocess.py | 171 +++++--- src/garg_aml/scores.py | 397 +++++++++--------- tests/_pipeline.py | 39 +- tests/test_equivalence.py | 27 +- tests/test_golden.py | 12 +- 13 files changed, 920 insertions(+), 667 deletions(-) create mode 100644 docs/decisions/0008-node-identity-is-preserved.md diff --git a/CHANGELOG.md b/CHANGELOG.md index facefed..10b2312 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -24,10 +24,35 @@ is `0.x` the public API may change with a minor bump, always with an entry here. latter needs the research repository present and is removed once the extraction is complete. +### Changed + +- PEP 8 names throughout, with the two per-node measure functions merged into + one `block_measures(graph, node, directed=False)`: + + | was | is | + |---|---| + | `GARG_AML_node_{un,}directed_measures` | `block_measures` | + | `calculate_score_{un,}directed` | `score_from_measures` | + | `define_gargaml_scores` | `scores_from_measures` | + | `graph_community` | `reduce_graph` | + | `graph_degree` | `drop_hubs` | + | `summaries_neighbourhoors_node` | `neighbour_score_stats` | + | `degree_neighbours_node` | `neighbour_degree_stats` | + | `summarise_gargaml_scores` | `build_features` | + | `measure_NN_function` | private `_block_NN` | + +- **`score_type` now defaults to `"weighted_average"`**, the aggregation every + published experiment uses. The previous default, `"basic"`, gives different + numbers. +- `reduce_graph` takes `seed` (default 1997) rather than hard-coding it. +- `build_features` returns every feature column by default rather than only the + score. +- Node identity is preserved: the old directed path cast integer node ids to + float. Cosmetic, but the package no longer does it — see `docs/decisions/0008`. +- Type hints on every function and NumPy-style docstrings with paper references + on every public one, each carrying a doctest that runs in CI. + ### Notes -- Public names are still the research repository's (`GARG_AML_node_*_measures`, - `define_gargaml_scores`, ...). Renaming to PEP 8, type hints and NumPy-style - docstrings follow in the next release step, with the fixtures green throughout. - The original's bare `except:` around the neighbour statistics is written as the explicit empty check it always was. Same result, verified by both test layers. diff --git a/docs/decisions/0008-node-identity-is-preserved.md b/docs/decisions/0008-node-identity-is-preserved.md new file mode 100644 index 0000000..45c818a --- /dev/null +++ b/docs/decisions/0008-node-identity-is-preserved.md @@ -0,0 +1,29 @@ +# 0008 — Node identity is preserved + +**Decision.** Whatever the caller's node ids are — `str`, `int`, tuple — they +come back as the index of every returned frame, unchanged. + +**Context.** The pre-extraction code did not do this consistently. +`define_gargaml_scores_undirected` collected node ids with +`measures["node"].tolist()`, preserving their dtype; +`define_gargaml_scores_directed` collected them inside a `DataFrame.iterrows()` +loop, which coerces each row to a single dtype. With float measure columns in the +frame, that silently turned integer node ids into floats. The frozen fixtures +record it: `scores_directed_raw.csv` is indexed `0.0, 1.0` where +`scores_undirected_raw.csv` is indexed `0, 1`. + +**Why this is not a correction to published results.** It is cosmetic. Python +hashes `0.0` and `0` identically, so the downstream dictionary and networkx +lookups resolve either way, and pandas joins a float64 index against an int64 +index by value. The IBM account ids are strings in any case, so `iterrows()` +left them alone. Nothing downstream was reading a wrong number. + +**Why change it.** A library that renames its caller's keys is surprising, and +node ids are the one thing a user matches results back to their own data with. +The package uses `.tolist()` on both paths. + +**Consequence.** The directed score and feature fixtures carry a float index +that the package deliberately does not reproduce. Both test layers therefore +compare indices **by value** rather than by dtype, and say so inline. That is +strictly stronger than what came before: the golden comparison used to drop the +index entirely, so node alignment was never checked at all. diff --git a/docs/decisions/index.md b/docs/decisions/index.md index 0400d83..9108dcd 100644 --- a/docs/decisions/index.md +++ b/docs/decisions/index.md @@ -17,3 +17,4 @@ values produced the results published in | [0005](0005-synthetic-generator-differs.md) | The synthetic generator does not reproduce the paper's datasets | | [0006](0006-no-torch-dependency.md) | No PyTorch, ever | | [0007](0007-directed-score-is-equation-14.md) | The directed score is Eq. 14, without the transpose-max | +| [0008](0008-node-identity-is-preserved.md) | Node ids come back unchanged; the old directed path cast them to float | diff --git a/pyproject.toml b/pyproject.toml index 78fbfa5..cdf8a0e 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -76,13 +76,6 @@ max-complexity = 10 # the condition that triggers it; SIM108 would bury them at the end of a # ternary, in the one file where those values most need to be obvious. "src/garg_aml/_blocks.py" = ["SIM108"] -# TEMPORARY, remove in Phase 3. Phase 2 moves the implementation across with -# its function bodies unchanged, so that a golden-fixture mismatch can only -# mean the move was wrong. These three rules all ask for a rewrite of a body -# (ternaries, unpacking instead of list concatenation, no inplace=True) -- -# correct requests, but ones that belong to the rename-and-modernise pass, -# where the fixtures are already green and can vouch for each change. -"src/garg_aml/*.py" = ["SIM108", "RUF005", "PD002"] # --- Types ------------------------------------------------------------------ [tool.mypy] @@ -96,17 +89,11 @@ warn_unused_ignores = true module = ["networkx.*", "scipy.*"] ignore_missing_imports = true -# Phase 2 moves the implementation across verbatim so that any fixture mismatch -# points at the move rather than at a rewrite; annotations are Phase 3. Delete -# this override then, and the global disallow_untyped_defs takes effect. -[[tool.mypy.overrides]] -module = ["garg_aml.*"] -disallow_untyped_defs = false - # --- Tests ------------------------------------------------------------------ [tool.pytest.ini_options] -testpaths = ["tests"] -addopts = "--strict-markers --cov=garg_aml --cov-report=term-missing --cov-fail-under=90" +testpaths = ["tests", "src/garg_aml"] +addopts = """--strict-markers --doctest-modules \ + --cov=garg_aml --cov-report=term-missing --cov-fail-under=90""" markers = [ "requires_data: needs the IBM AMLworld dataset, which is not redistributable", ] diff --git a/src/garg_aml/_blocks.py b/src/garg_aml/_blocks.py index 422fa21..290a315 100644 --- a/src/garg_aml/_blocks.py +++ b/src/garg_aml/_blocks.py @@ -3,192 +3,185 @@ Each function takes the ordered adjacency matrix and the block sizes, and returns the block's density over its *free* entries together with the count of -those entries. Entries that are structurally fixed -- the diagonal, the ego -node's own row and column -- are excluded from both, which is why the returned -size is often smaller than the block. +those entries. Structurally fixed entries -- the diagonal, the ego node's own +row and column -- are excluded from both, which is why the returned size is +often smaller than the block. -Three functions for the undirected analysis (paper section 3.2) and nine for the -directed one (section 3.3). Copied unchanged from the research repository; the -degenerate-case constants are deliberate, see docs/decisions/0003. +Three blocks for the undirected analysis (paper section 3.2) and nine for the +directed one (section 3.3). The fallback values in the empty-block branches are +frozen behaviour: see docs/decisions/0003. """ +import numpy as np -def measure_1_function(piece_1_dim, adj_full): + +def _block_1(dim_1: list[int], adjacency: np.ndarray) -> tuple[float, int]: """Density among the ego node and its second-order neighbours.""" - piece_1 = adj_full[: piece_1_dim[0], : piece_1_dim[1]] - total_sum_1 = piece_1.sum() - total_size_1 = piece_1.size - reduced_size_1 = total_size_1 - (3 * piece_1_dim[0]) + 2 + block = adjacency[: dim_1[0], : dim_1[1]] + free = block.size - (3 * dim_1[0]) + 2 - if reduced_size_1 > 0: - rel_1 = total_sum_1 / reduced_size_1 + if free > 0: + density = block.sum() / free else: - rel_1 = 0 + density = 0 - return rel_1, reduced_size_1 + return density, free -def measure_2_function(piece_1_dim, piece_2_dim, adj_full): +def _block_2( + dim_1: list[int], dim_2: list[int], adjacency: np.ndarray +) -> tuple[float, int]: """Density between first-order and second-order neighbours.""" - piece_2 = adj_full[ - piece_1_dim[0] : piece_1_dim[0] + piece_2_dim[0], : piece_2_dim[1] - ] - total_sum_2 = piece_2.sum() - reduced_sum_2 = total_sum_2 - piece_2_dim[0] - total_size_2 = piece_2.size - reduced_size_2 = total_size_2 - piece_2_dim[0] - - if reduced_size_2 > 0: - rel_2 = reduced_sum_2 / reduced_size_2 + block = adjacency[dim_1[0] : dim_1[0] + dim_2[0], : dim_2[1]] + # The ego node's own edge to each first-order neighbour is structural. + total = block.sum() - dim_2[0] + free = block.size - dim_2[0] + + if free > 0: + density = total / free else: - rel_2 = 1 + density = 1 - return rel_2, reduced_size_2 + return density, free -def measure_3_function(piece_1_dim, piece_2_dim, piece_3_dim, adj_full): +def _block_3( + dim_1: list[int], dim_2: list[int], dim_3: list[int], adjacency: np.ndarray +) -> tuple[float, int]: """Density among the first-order neighbours.""" - piece_3 = adj_full[piece_1_dim[0] :, piece_2_dim[1] :] - total_sum_3 = piece_3.sum() - total_size_3 = piece_3.size - reduced_size_3 = total_size_3 - piece_3_dim[0] + block = adjacency[dim_1[0] :, dim_2[1] :] + free = block.size - dim_3[0] - if reduced_size_3 > 0: - rel_3 = total_sum_3 / reduced_size_3 + if free > 0: + density = block.sum() / free else: - rel_3 = 0 + density = 0 - return rel_3, reduced_size_3 + return density, free -def measure_00_function(adj_full, size_0): +def _block_00(adjacency: np.ndarray, size_0: int) -> tuple[float, int]: """Density among level-0 nodes (senders).""" - piece_00 = adj_full[:size_0, :size_0] - total_sum_00 = piece_00.sum() - total_size_00 = piece_00.size - reduced_size_00 = total_size_00 - (3 * size_0) + 2 + block = adjacency[:size_0, :size_0] + free = block.size - (3 * size_0) + 2 - if reduced_size_00 > 0: - rel_00 = total_sum_00 / reduced_size_00 + if free > 0: + density = block.sum() / free else: - rel_00 = 0 + density = 0 - return rel_00, reduced_size_00 + return density, free -def measure_01_function(adj_full, size_0, size_1): +def _block_01(adjacency: np.ndarray, size_0: int, size_1: int) -> tuple[float, int]: """Density from level 0 (senders) to level 1 (mules).""" - piece_01 = adj_full[:size_0, size_0 : size_0 + size_1] - total_sum_01 = piece_01.sum() - total_size_01 = piece_01.size + block = adjacency[:size_0, size_0 : size_0 + size_1] + free = block.size - if total_size_01 > 0: - rel_01 = total_sum_01 / total_size_01 + if free > 0: + density = block.sum() / free else: - rel_01 = 1 # Since block only contains sure connections => full sum + density = 1 # Block holds only sure connections => treat as full - return rel_01, total_size_01 + return density, free -def measure_02_function(adj_full, size_0, size_1, size_2): +def _block_02( + adjacency: np.ndarray, size_0: int, size_1: int, size_2: int +) -> tuple[float, int]: """Density from level 0 (senders) to level 2 (receivers).""" - piece_02 = adj_full[:size_0, size_0 + size_1 :] - total_sum_02 = piece_02.sum() - total_size_02 = piece_02.size - reduced_size_02 = total_size_02 - size_2 + block = adjacency[:size_0, size_0 + size_1 :] + free = block.size - size_2 - if reduced_size_02 > 0: - rel_02 = total_sum_02 / reduced_size_02 + if free > 0: + density = block.sum() / free else: - rel_02 = 0 + density = 0 - return rel_02, reduced_size_02 + return density, free -def measure_10_function(adj_full, size_0, size_1): +def _block_10(adjacency: np.ndarray, size_0: int, size_1: int) -> tuple[float, int]: """Density from level 1 (mules) back to level 0 (senders).""" - piece_10 = adj_full[size_0 : size_0 + size_1, :size_0] - total_sum_10 = piece_10.sum() - total_size_10 = piece_10.size + block = adjacency[size_0 : size_0 + size_1, :size_0] + free = block.size - if total_size_10 > 0: - rel_10 = total_sum_10 / total_size_10 + if free > 0: + density = block.sum() / free else: - rel_10 = 0 + density = 0 - return rel_10, total_size_10 + return density, free -def measure_11_function(adj_full, size_0, size_1): +def _block_11(adjacency: np.ndarray, size_0: int, size_1: int) -> tuple[float, int]: """Density among level-1 nodes (mules).""" - piece_11 = adj_full[size_0 : size_0 + size_1, size_0 : size_0 + size_1] - total_sum_11 = piece_11.sum() - total_size_11 = piece_11.size - reduced_size_11 = total_size_11 - size_1 + block = adjacency[size_0 : size_0 + size_1, size_0 : size_0 + size_1] + free = block.size - size_1 - if reduced_size_11 > 0: - rel_11 = total_sum_11 / reduced_size_11 + if free > 0: + density = block.sum() / free else: - rel_11 = 0 + density = 0 - return rel_11, reduced_size_11 + return density, free -def measure_12_function(adj_full, size_0, size_1, size_2): +def _block_12( + adjacency: np.ndarray, size_0: int, size_1: int, size_2: int +) -> tuple[float, int]: """Density from level 1 (mules) to level 2 (receivers).""" - piece_12 = adj_full[size_0 : size_0 + size_1, size_0 + size_1 :] - total_sum_12 = piece_12.sum() - total_size_12 = piece_12.size + block = adjacency[size_0 : size_0 + size_1, size_0 + size_1 :] + free = block.size - if total_size_12 > 0: - rel_12 = total_sum_12 / total_size_12 + if free > 0: + density = block.sum() / free elif size_2 > 0: - rel_12 = 1 # Since block only contains sure connections => full sum + density = 1 # Block holds only sure connections => treat as full else: - rel_12 = 0 # No connections at all + density = 0 # No connections at all - return rel_12, total_size_12 + return density, free -def measure_20_function(adj_full, size_0, size_1, size_2): +def _block_20( + adjacency: np.ndarray, size_0: int, size_1: int, size_2: int +) -> tuple[float, int]: """Density from level 2 (receivers) back to level 0 (senders).""" - piece_20 = adj_full[size_0 + size_1 :, :size_0] - total_sum_20 = piece_20.sum() - total_size_20 = piece_20.size - reduced_size_20 = total_size_20 - size_2 + block = adjacency[size_0 + size_1 :, :size_0] + free = block.size - size_2 - if reduced_size_20 > 0: - rel_20 = total_sum_20 / reduced_size_20 + if free > 0: + density = block.sum() / free else: - rel_20 = 0 + density = 0 - return rel_20, reduced_size_20 + return density, free -def measure_21_function(adj_full, size_0, size_1): +def _block_21(adjacency: np.ndarray, size_0: int, size_1: int) -> tuple[float, int]: """Density from level 2 (receivers) back to level 1 (mules).""" - piece_21 = adj_full[size_0 + size_1 :, size_0 : size_0 + size_1] - total_sum_21 = piece_21.sum() - total_size_21 = piece_21.size + block = adjacency[size_0 + size_1 :, size_0 : size_0 + size_1] + free = block.size - if total_size_21 > 0: - rel_21 = total_sum_21 / total_size_21 + if free > 0: + density = block.sum() / free else: - rel_21 = 0 + density = 0 - return rel_21, total_size_21 + return density, free -def measure_22_function(adj_full, size_0, size_1, size_2): +def _block_22( + adjacency: np.ndarray, size_0: int, size_1: int, size_2: int +) -> tuple[float, int]: """Density among level-2 nodes (receivers).""" - piece_22 = adj_full[size_0 + size_1 :, size_0 + size_1 :] - total_sum_22 = piece_22.sum() - total_size_22 = piece_22.size - reduced_size_22 = total_size_22 - size_2 + block = adjacency[size_0 + size_1 :, size_0 + size_1 :] + free = block.size - size_2 - if reduced_size_22 > 0: - rel_22 = total_sum_22 / reduced_size_22 + if free > 0: + density = block.sum() / free else: - rel_22 = 0 + density = 0 - return rel_22, reduced_size_22 + return density, free diff --git a/src/garg_aml/_ordering.py b/src/garg_aml/_ordering.py index 59c2d59..a3dfa09 100644 --- a/src/garg_aml/_ordering.py +++ b/src/garg_aml/_ordering.py @@ -2,72 +2,69 @@ Ordering of a second-order ego graph's nodes into its analysis blocks. The block structure GARG-AML measures only appears under a specific node order. -Undirected (section 3.2): ``[ego, second-order neighbours, first-order +Undirected (paper section 3.2): ``[ego, second-order neighbours, first-order neighbours]``. Directed (section 3.3): ``[level 0, level 1, level 2]``, where a node at distance two sits at level 0 when no directed path of length two reaches it in either direction -- Eq. (11) applied to the graph and to its reverse. - -Copied unchanged from the research repository. """ +from collections.abc import Hashable + import networkx as nx -def GARG_AML_nodeselection_undirected(G_ego_second, node): +def _undirected_order(ego: nx.Graph, node: Hashable) -> tuple[list, list, list]: """Order as ego, second-order neighbours, first-order neighbours.""" - nodes_1 = list(nx.ego_graph(G_ego_second, node).nodes) - nodes_1.remove(node) - nodes_2 = list(G_ego_second.nodes) - nodes_2.remove(node) - for n in nodes_1: - nodes_2.remove(n) + first = list(nx.ego_graph(ego, node).nodes) + first.remove(node) + second = list(ego.nodes) + second.remove(node) + for n in first: + second.remove(n) - # For undirected networks, specific order to obtain scores - # (group node with second order neighbours) - nodes_ordered = [node] + nodes_2 + nodes_1 + # The ego node is grouped with its second-order neighbours: that is the + # partition whose on-diagonal blocks a smurfing pattern leaves empty. + return first, second, [node, *second, *first] - return nodes_1, nodes_2, nodes_ordered - -def GARG_AML_nodeselection_directed( - G_ego_second, G_ego_second_und, G_ego_second_rev, node -): +def _directed_order( + ego: nx.DiGraph, ego_undirected: nx.Graph, ego_reverse: nx.DiGraph, node: Hashable +) -> tuple[list, list, list, list]: """Order by level: senders, mules, receivers.""" - nodes_1 = list(nx.ego_graph(G_ego_second_und, node).nodes) - nodes_1.remove(node) - nodes_2 = list(G_ego_second.nodes) - - nodes_2_s = list(nx.ego_graph(G_ego_second, node, radius=2).nodes) - - nodes_2_rs = list(nx.ego_graph(G_ego_second_rev, node, radius=2).nodes) - - nodes_0 = list( - set(nodes_2) - .difference(set(nodes_2_s)) - .difference(set(nodes_2_rs)) - .difference(set(nodes_1)) + level_1 = list(nx.ego_graph(ego_undirected, node).nodes) + level_1.remove(node) + level_2 = list(ego.nodes) + + reachable = list(nx.ego_graph(ego, node, radius=2).nodes) + reaching = list(nx.ego_graph(ego_reverse, node, radius=2).nodes) + + # Eq. (11): a node at distance two with no directed path of length two in + # either direction is a sender. + level_0 = list( + set(level_2) + .difference(set(reachable)) + .difference(set(reaching)) + .difference(set(level_1)) ) - nodes_0 = [node] + nodes_0 - - for n in nodes_0: - nodes_2.remove(n) - for n in nodes_1: - nodes_2.remove(n) + level_0 = [node, *level_0] - # For directed network, specific order to obtain scores (in order of "group") - nodes_ordered = nodes_0 + nodes_1 + nodes_2 + for n in level_0: + level_2.remove(n) + for n in level_1: + level_2.remove(n) - return nodes_0, nodes_1, nodes_2, nodes_ordered + return level_0, level_1, level_2, level_0 + level_1 + level_2 -def GARG_AML_nodeselection( - G_ego_second, node, directed, G_ego_second_und=None, G_ego_second_rev=None -): +def node_order( + ego: nx.Graph, + node: Hashable, + directed: bool, + ego_undirected: nx.Graph | None = None, + ego_reverse: nx.DiGraph | None = None, +) -> tuple: """Dispatch to the directed or undirected ordering.""" if directed: - return GARG_AML_nodeselection_directed( - G_ego_second, G_ego_second_und, G_ego_second_rev, node - ) - else: - return GARG_AML_nodeselection_undirected(G_ego_second, node) + return _directed_order(ego, ego_undirected, ego_reverse, node) + return _undirected_order(ego, node) diff --git a/src/garg_aml/features.py b/src/garg_aml/features.py index d8a0430..96e77ac 100644 --- a/src/garg_aml/features.py +++ b/src/garg_aml/features.py @@ -1,124 +1,188 @@ """ Neighbourhood summary statistics built on top of the GARG-AML score. -For every node: the min/mean/max/std of its neighbours' scores and of their -degrees, plus its own degree. These are the features the paper's tree and -boosting models receive alongside the score itself, and they are what makes the -score usable as an input to a model of your own. +For every node: the min, mean, max and standard deviation of its neighbours' +scores and of their degrees, plus its own degree. These are the features the +paper's tree and boosting models receive alongside the score itself, and they +are what makes the score usable as an input to a model of your own. A node with no neighbours returns zero for all eight statistics. That is -deliberate and it is load-bearing -- Louvain reduction strands many nodes, so -the branch is common rather than exotic. See docs/decisions/0003. - -Copied unchanged from the research repository, except that the original's bare -``except:`` around the numpy reductions is written here as the explicit empty -check it always was. +deliberate, and it is common rather than exotic: Louvain reduction strands many +nodes. See ``docs/decisions/0003``. """ +from collections.abc import Hashable, Mapping, Sequence + import networkx as nx import numpy as np import pandas as pd -from .scores import define_gargaml_scores - - -def summaries_neighbourhoors_node(node, G_copy, measures): - """Min, mean, max and std of the neighbours' scores.""" - G_ego = nx.ego_graph(G_copy, node) - G_ego.remove_node(node) - ego_list = list(G_ego.nodes) - - list_ego_measures = [] - for n in ego_list: - list_ego_measures.append(measures[n]) - - if not list_ego_measures: - return [0, 0, 0, 0] - - return [ - np.min(list_ego_measures), - np.mean(list_ego_measures), - np.max(list_ego_measures), - np.std(list_ego_measures), - ] - - -def degree_neighbours_node(node, G_copy, G_degree_dict): - """Min, mean, max and std of the neighbours' degrees.""" - G_ego = nx.ego_graph(G_copy, node) - G_ego.remove_node(node) - ego_list = list(G_ego.nodes) - - list_ego_degree = [] - for n in ego_list: - list_ego_degree.append(G_degree_dict[n]) +__all__ = [ + "FEATURE_COLUMNS", + "build_features", + "neighbour_degree_stats", + "neighbour_score_stats", +] - if not list_ego_degree: - return [0, 0, 0, 0] +#: Every column :func:`build_features` can return, in order. +FEATURE_COLUMNS = [ + "GARGAML", + "GARGAML_min", + "GARGAML_mean", + "GARGAML_max", + "GARGAML_std", + "degree", + "degree_min", + "degree_mean", + "degree_max", + "degree_std", +] - return [ - np.min(list_ego_degree), - np.mean(list_ego_degree), - np.max(list_ego_degree), - np.std(list_ego_degree), +_EMPTY_STATS = [0, 0, 0, 0] + + +def _stats(values: list) -> list: + """Min, mean, max and std, or zeros when there is nothing to summarise.""" + if not values: + return list(_EMPTY_STATS) + return [np.min(values), np.mean(values), np.max(values), np.std(values)] + + +def _neighbours(graph: nx.Graph, node: Hashable) -> list: + ego = nx.ego_graph(graph, node) + ego.remove_node(node) + return list(ego.nodes) + + +def neighbour_score_stats(graph: nx.Graph, node: Hashable, scores: Mapping) -> list: + """ + Min, mean, max and standard deviation of the neighbours' scores. + + Parameters + ---------- + graph : networkx.Graph or networkx.DiGraph + The graph the scores were computed on. + node : hashable + The node whose neighbourhood to summarise. + scores : mapping + Score per node. Must cover every neighbour of ``node``. + + Returns + ------- + list + ``[min, mean, max, std]``, or ``[0, 0, 0, 0]`` when ``node`` has no + neighbours. + + Examples + -------- + >>> import networkx as nx + >>> graph = nx.Graph([("a", "b"), ("a", "c")]) + >>> [float(v) for v in neighbour_score_stats(graph, "a", {"b": 0.0, "c": 1.0})] + [0.0, 0.5, 1.0, 0.5] + """ + return _stats([scores[n] for n in _neighbours(graph, node)]) + + +def neighbour_degree_stats(graph: nx.Graph, node: Hashable, degrees: Mapping) -> list: + """ + Min, mean, max and standard deviation of the neighbours' degrees. + + Parameters + ---------- + graph : networkx.Graph or networkx.DiGraph + The graph the degrees were taken from. + node : hashable + The node whose neighbourhood to summarise. + degrees : mapping + Degree per node, as ``dict(graph.degree())`` gives it. + + Returns + ------- + list + ``[min, mean, max, std]``, or ``[0, 0, 0, 0]`` when ``node`` has no + neighbours. + + Examples + -------- + >>> import networkx as nx + >>> graph = nx.Graph([("a", "b"), ("a", "c"), ("c", "d")]) + >>> [float(v) for v in neighbour_degree_stats(graph, "a", dict(graph.degree()))] + [1.0, 1.5, 2.0, 0.5] + """ + return _stats([degrees[n] for n in _neighbours(graph, node)]) + + +def _assemble( + graph: nx.Graph, scores: Mapping, score_stats: Mapping, degree_stats: Mapping +) -> pd.DataFrame: + """Join scores, neighbour statistics and degrees into one frame.""" + frames = [ + pd.DataFrame(scores, index=["GARGAML"]).transpose(), + pd.DataFrame( + score_stats, + index=["GARGAML_min", "GARGAML_mean", "GARGAML_max", "GARGAML_std"], + ).transpose(), + pd.DataFrame(dict(graph.degree()), index=["degree"]).transpose(), + pd.DataFrame( + degree_stats, + index=["degree_min", "degree_mean", "degree_max", "degree_std"], + ).transpose(), ] - -def combine_GARG_AML(G_selection, measures_dict, summary_dict, neigh_degree_dict): - """Assemble scores, neighbour statistics and degrees into one frame.""" - degree_df = pd.DataFrame(dict(G_selection.degree()), index=["degree"]).transpose() - - measures_df = pd.DataFrame(measures_dict, index=["GARGAML"]).transpose() - - summary_df = pd.DataFrame( - summary_dict, - index=["GARGAML_min", "GARGAML_mean", "GARGAML_max", "GARGAML_std"], - ).transpose() - - neigh_degree_df = pd.DataFrame( - neigh_degree_dict, - index=["degree_min", "degree_mean", "degree_max", "degree_std"], - ).transpose() - - GARG_AML_df = ( - measures_df.merge(summary_df, left_index=True, right_index=True) - .merge(degree_df, left_index=True, right_index=True) - .merge(neigh_degree_df, left_index=True, right_index=True) - ) - - return GARG_AML_df - - -def summarise_gargaml_scores(G_reduced, df_results, columns=None): - """Neighbourhood statistics for a frame of scores.""" + assembled = frames[0] + for frame in frames[1:]: + assembled = assembled.merge(frame, left_index=True, right_index=True) + return assembled + + +def build_features( + graph: nx.Graph, scores: pd.DataFrame, columns: Sequence[str] | None = None +) -> pd.DataFrame: + """ + Neighbourhood features for every scored node. + + Parameters + ---------- + graph : networkx.Graph or networkx.DiGraph + The graph the scores were computed on -- the *reduced* graph, if the + scores came from one, since the neighbourhoods must match. + scores : pandas.DataFrame + Indexed by node, with a ``GARGAML`` column, as + :func:`garg_aml.scores.scores_from_measures` returns. + columns : sequence of str, optional + Which of :data:`FEATURE_COLUMNS` to return. All of them by default. + + Returns + ------- + pandas.DataFrame + Indexed by node, one column per requested feature. + + Notes + ----- + Degrees are taken from ``graph``, so on a Louvain-reduced graph they are + intra-community degrees, not raw transaction counts. + + Examples + -------- + >>> import networkx as nx + >>> import pandas as pd + >>> graph = nx.Graph([("a", "m1"), ("a", "m2"), ("m1", "b"), ("m2", "b")]) + >>> scores = pd.DataFrame({"GARGAML": [1.0, -0.5, -0.5, 1.0]}, + ... index=["a", "m1", "m2", "b"]) + >>> features = build_features(graph, scores, columns=["degree", "GARGAML_mean"]) + >>> [float(v) for v in features.loc["a"]] + [2.0, -0.5] + """ if columns is None: - columns = ["GARGAML"] + columns = FEATURE_COLUMNS - G_degree_dict = dict(G_reduced.degree()) - nodes = list(df_results.index) + degrees = dict(graph.degree()) + score_map = dict(zip(scores.index, scores["GARGAML"], strict=True)) - gargaml_values = dict(zip(df_results.index, df_results["GARGAML"], strict=False)) + score_stats = {} + degree_stats = {} + for node in scores.index: + score_stats[node] = neighbour_score_stats(graph, node, score_map) + degree_stats[node] = neighbour_degree_stats(graph, node, degrees) - summaries_neighbourhood = {} - summaries_degree = {} - - for node in nodes: - summaries_neighbourhood[node] = summaries_neighbourhoors_node( - node, G_reduced, gargaml_values - ) - summaries_degree[node] = degree_neighbours_node(node, G_reduced, G_degree_dict) - - GARGAML_df = combine_GARG_AML( - G_reduced, gargaml_values, summaries_neighbourhood, summaries_degree - ) - - return GARGAML_df[columns] - - -__all__ = [ - "combine_GARG_AML", - "define_gargaml_scores", - "degree_neighbours_node", - "summaries_neighbourhoors_node", - "summarise_gargaml_scores", -] + return _assemble(graph, score_map, score_stats, degree_stats)[list(columns)] diff --git a/src/garg_aml/measures.py b/src/garg_aml/measures.py index 91f481a..7122a98 100644 --- a/src/garg_aml/measures.py +++ b/src/garg_aml/measures.py @@ -1,130 +1,189 @@ """ Per-node block measures: the expensive first stage of GARG-AML. -For one node, build its second-order ego graph, order the nodes into blocks -(:mod:`garg_aml._ordering`) and take each block's density over its free entries -(:mod:`garg_aml._blocks`). Three blocks undirected, nine directed. +For one node, build its second-order ego graph, order the nodes into blocks and +take each block's density over its free entries. Three blocks undirected, nine +directed. -This stage is deliberately separate from scoring: the measures are what is worth -persisting, and any score variant is a cheap re-aggregation of them. See +This stage is deliberately separate from scoring. The measures are what is worth +persisting, and any score variant is a cheap re-aggregation of them -- see docs/decisions/0004. - -Copied unchanged from the research repository. """ +from collections.abc import Hashable + import networkx as nx from ._blocks import ( - measure_00_function, - measure_01_function, - measure_02_function, - measure_1_function, - measure_2_function, - measure_3_function, - measure_10_function, - measure_11_function, - measure_12_function, - measure_20_function, - measure_21_function, - measure_22_function, + _block_00, + _block_01, + _block_02, + _block_1, + _block_2, + _block_3, + _block_10, + _block_11, + _block_12, + _block_20, + _block_21, + _block_22, ) -from ._ordering import GARG_AML_nodeselection +from ._ordering import node_order +__all__ = ["block_measures"] -def GARG_AML_node_undirected_measures(node, G_copy, include_size=False): - """Undirected block densities for one node (paper section 3.2).""" - G_ego_second = nx.ego_graph(G_copy, node, 2) - # nodes_ordered are the nodes ordered as node, 2nd order and 1st order neighbours - nodes_1, nodes_2, nodes_ordered = GARG_AML_nodeselection( - G_ego_second, node, directed=False - ) +def _undirected_measures( + graph: nx.Graph, node: Hashable, include_sizes: bool +) -> tuple[float, ...]: + """Three block densities for one node of an undirected graph.""" + ego = nx.ego_graph(graph, node, 2) - adj_full = nx.adjacency_matrix(G_ego_second, nodelist=nodes_ordered).toarray() + first, second, ordered = node_order(ego, node, directed=False) - size_second = len(nodes_2) - size_first = len(nodes_1) + adjacency = nx.adjacency_matrix(ego, nodelist=ordered).toarray() - piece_1_dim = [size_second + 1, size_second + 1] - piece_2_dim = [size_first, size_second + 1] - piece_3_dim = [size_first, size_first] + n_second = len(second) + n_first = len(first) - measure_1, size_1 = measure_1_function(piece_1_dim, adj_full) - measure_2, size_2 = measure_2_function(piece_1_dim, piece_2_dim, adj_full) - measure_3, size_3 = measure_3_function( - piece_1_dim, piece_2_dim, piece_3_dim, adj_full - ) + dim_1 = [n_second + 1, n_second + 1] + dim_2 = [n_first, n_second + 1] + dim_3 = [n_first, n_first] + + measure_1, size_1 = _block_1(dim_1, adjacency) + measure_2, size_2 = _block_2(dim_1, dim_2, adjacency) + measure_3, size_3 = _block_3(dim_1, dim_2, dim_3, adjacency) - if include_size: + if include_sizes: return (measure_1, measure_2, measure_3, size_1, size_2, size_3) - else: - return (measure_1, measure_2, measure_3) - - -def GARG_AML_node_directed_measures( - node, G_copy, G_copy_und, G_copy_rev, include_size=False -): - """Directed block densities for one node (paper section 3.3).""" - # Use both incoming and outgoing edges - G_ego_second_und = nx.ego_graph(G_copy_und, node, 2) - G_ego_second = nx.subgraph(G_copy, G_ego_second_und.nodes) - # Look at the reverse graph to get the incoming edges - G_ego_second_rev = nx.ego_graph(G_copy_rev, node, 2) - - nodes_0, nodes_1, nodes_2, nodes_ordered = GARG_AML_nodeselection( - G_ego_second, - node, - directed=True, - G_ego_second_und=G_ego_second_und, - G_ego_second_rev=G_ego_second_rev, + return (measure_1, measure_2, measure_3) + + +def _directed_measures( + graph: nx.DiGraph, + node: Hashable, + undirected: nx.Graph, + reverse: nx.DiGraph, + include_sizes: bool, +) -> tuple[float, ...]: + """Nine block densities for one node of a directed graph.""" + # The neighbourhood is taken on the undirected view, so that a node is a + # neighbour regardless of which way its edges point; the levels are then + # assigned from the directed and reversed views. + ego_undirected = nx.ego_graph(undirected, node, 2) + ego = nx.subgraph(graph, ego_undirected.nodes) + ego_reverse = nx.ego_graph(reverse, node, 2) + + level_0, level_1, level_2, ordered = node_order( + ego, node, directed=True, ego_undirected=ego_undirected, ego_reverse=ego_reverse ) - adj_full = nx.adjacency_matrix(G_ego_second, nodelist=nodes_ordered).toarray() - - size_0 = len(nodes_0) - size_1 = len(nodes_1) - size_2 = len(nodes_2) - - measure_00, size_00 = measure_00_function(adj_full, size_0) - measure_01, size_01 = measure_01_function(adj_full, size_0, size_1) - measure_02, size_02 = measure_02_function(adj_full, size_0, size_1, size_2) - measure_10, size_10 = measure_10_function(adj_full, size_0, size_1) - measure_11, size_11 = measure_11_function(adj_full, size_0, size_1) - measure_12, size_12 = measure_12_function(adj_full, size_0, size_1, size_2) - measure_20, size_20 = measure_20_function(adj_full, size_0, size_1, size_2) - measure_21, size_21 = measure_21_function(adj_full, size_0, size_1) - measure_22, size_22 = measure_22_function(adj_full, size_0, size_1, size_2) - - if include_size: - return ( - measure_00, - measure_01, - measure_02, - measure_10, - measure_11, - measure_12, - measure_20, - measure_21, - measure_22, - size_00, - size_01, - size_02, - size_10, - size_11, - size_12, - size_20, - size_21, - size_22, - ) - else: + adjacency = nx.adjacency_matrix(ego, nodelist=ordered).toarray() + + size_0 = len(level_0) + size_1 = len(level_1) + size_2 = len(level_2) + + measure_00, block_00 = _block_00(adjacency, size_0) + measure_01, block_01 = _block_01(adjacency, size_0, size_1) + measure_02, block_02 = _block_02(adjacency, size_0, size_1, size_2) + measure_10, block_10 = _block_10(adjacency, size_0, size_1) + measure_11, block_11 = _block_11(adjacency, size_0, size_1) + measure_12, block_12 = _block_12(adjacency, size_0, size_1, size_2) + measure_20, block_20 = _block_20(adjacency, size_0, size_1, size_2) + measure_21, block_21 = _block_21(adjacency, size_0, size_1) + measure_22, block_22 = _block_22(adjacency, size_0, size_1, size_2) + + measures = ( + measure_00, + measure_01, + measure_02, + measure_10, + measure_11, + measure_12, + measure_20, + measure_21, + measure_22, + ) + if include_sizes: return ( - measure_00, - measure_01, - measure_02, - measure_10, - measure_11, - measure_12, - measure_20, - measure_21, - measure_22, + *measures, + block_00, + block_01, + block_02, + block_10, + block_11, + block_12, + block_20, + block_21, + block_22, ) + return measures + + +def block_measures( + graph: nx.Graph, + node: Hashable, + directed: bool = False, + *, + undirected: nx.Graph | None = None, + reverse: nx.DiGraph | None = None, + include_sizes: bool = False, +) -> tuple[float, ...]: + """ + Block densities of one node's second-order neighbourhood. + + Parameters + ---------- + graph : networkx.Graph or networkx.DiGraph + The transaction graph. Usually reduced first with + :func:`garg_aml.preprocess.reduce_graph`. + node : hashable + The node to measure. Must be in ``graph``. + directed : bool, default False + Use the nine-block directed analysis rather than the three-block + undirected one. + undirected, reverse : networkx.Graph, optional + Precomputed ``graph.to_undirected()`` and ``graph.reverse()``, used only + when ``directed`` is True. Scoring many nodes of one graph is much + faster if these are built once and passed in; they are derived on demand + when omitted. + include_sizes : bool, default False + Also return each block's number of free entries, which is what the + ``weighted_average`` score type weights by. + + Returns + ------- + tuple of float + Three densities undirected, nine directed, in row-major block order. + With ``include_sizes``, the matching free-entry counts follow. + + Notes + ----- + A block with no free entries does not yield a missing value: the + off-diagonal blocks fall back to 1 and the rest to 0. Those constants are + frozen behaviour -- see ``docs/decisions/0003``. + + References + ---------- + Deprez et al. (2025), sections 3.2 and 3.3. + + Examples + -------- + One source paying two mules, which both pay one target -- a pure smurfing + pattern, so the off-diagonal block is full and the others are empty. + + >>> import networkx as nx + >>> graph = nx.Graph([("a", "m1"), ("a", "m2"), ("m1", "b"), ("m2", "b")]) + >>> [round(float(m), 3) for m in block_measures(graph, "a")] + [0.0, 1.0, 0.0] + """ + if not directed: + return _undirected_measures(graph, node, include_sizes) + + if undirected is None: + undirected = graph.to_undirected() + if reverse is None: + reverse = graph.reverse(copy=True) + + return _directed_measures(graph, node, undirected, reverse, include_sizes) diff --git a/src/garg_aml/preprocess.py b/src/garg_aml/preprocess.py index 07af4b6..27a40b1 100644 --- a/src/garg_aml/preprocess.py +++ b/src/garg_aml/preprocess.py @@ -1,71 +1,132 @@ """ Optional graph preprocessing applied before scoring. -``graph_community`` partitions the graph with Louvain and keeps only -intra-community edges (paper Algorithm 1). It is **lossy**: on the IBM data it +:func:`reduce_graph` partitions the graph with Louvain and keeps only +intra-community edges (paper Algorithm 1). It is **lossy** -- on the IBM data it removes the large majority of edges, and it changes which neighbourhoods exist -at all. That is why it stays a separate, explicit step rather than something -``score`` does behind the caller's back -- see docs/decisions/0002. +at all. That is why it is a separate, explicit step rather than something +scoring does behind the caller's back; see ``docs/decisions/0002``. -``graph_degree`` removes the highest-degree nodes, which is not part of the -published pipeline but is useful on graphs with hubs. - -Copied unchanged from the research repository. +:func:`drop_hubs` removes the highest-degree nodes. It is not part of the +published pipeline, but it is useful on graphs with heavy hubs. """ import networkx as nx import pandas as pd - -def graph_degree(G, degree_cutoff=0.01): - """Remove the top ``degree_cutoff`` quantile of nodes by degree.""" - # Delete the hubs - # The cut-off is defined as a relative number - G_copy = G.copy() - - degree_df = pd.DataFrame(dict(G_copy.degree()), index=["Degree"]).transpose() - - degree_threshold = degree_df["Degree"].quantile(1 - degree_cutoff) - hub_criteria = degree_df["Degree"] >= degree_threshold - - hubs_deleted = list(degree_df[hub_criteria].reset_index()["index"]) - - G_copy.remove_nodes_from(hubs_deleted) - - return G_copy - - -def graph_community(G, resolution=10): - """Keep only intra-community edges of a Louvain partition.""" - # large resolution to have smaller communities - directed = nx.is_directed(G) - - if directed: - G_undirected = G.copy().to_undirected() - else: - G_undirected = G.copy() - - community_list = nx.community.louvain_communities( - G_undirected, resolution=resolution, seed=1997 +__all__ = ["drop_hubs", "reduce_graph"] + +#: The seed used throughout the paper. +SEED = 1997 + + +def drop_hubs(graph: nx.Graph, quantile: float = 0.01) -> nx.Graph: + """ + Remove the highest-degree nodes. + + Parameters + ---------- + graph : networkx.Graph or networkx.DiGraph + The graph to prune. It is not modified. + quantile : float, default 0.01 + Fraction of nodes to treat as hubs. Every node whose degree is at or + above the ``1 - quantile`` degree quantile is removed, so ties at the + threshold can push the number removed above the nominal fraction. + + Returns + ------- + networkx.Graph + A copy without the hub nodes or their edges. + + Examples + -------- + >>> import networkx as nx + >>> graph = nx.barabasi_albert_graph(30, 2, seed=1) + >>> drop_hubs(graph, quantile=0.1).number_of_nodes() + 27 + """ + pruned = graph.copy() + + degrees = pd.DataFrame(dict(pruned.degree()), index=["degree"]).transpose() + + threshold = degrees["degree"].quantile(1 - quantile) + hubs = list(degrees[degrees["degree"] >= threshold].reset_index()["index"]) + + pruned.remove_nodes_from(hubs) + + return pruned + + +def reduce_graph(graph: nx.Graph, resolution: float = 10, seed: int = SEED) -> nx.Graph: + """ + Keep only the intra-community edges of a Louvain partition. + + Parameters + ---------- + graph : networkx.Graph or networkx.DiGraph + The transaction graph. It is not modified, and a directed graph stays + directed: the partition is found on the undirected view, but the edges + that survive keep their direction. + resolution : float, default 10 + Louvain resolution. Higher values give smaller communities and so cut + more edges. The published results use 10. + seed : int, default 1997 + Seed for the Louvain partition, which is otherwise not deterministic. + + Returns + ------- + networkx.Graph + A graph with every node of the input but only the edges whose endpoints + share a community. Node and edge attributes are preserved. + + Notes + ----- + This step is lossy by design, and nodes that lose all their edges are kept + rather than dropped -- they go on to score -1. See ``docs/decisions/0002`` + and ``docs/decisions/0003``. + + References + ---------- + Deprez et al. (2025), Algorithm 1. + + Examples + -------- + Two cliques joined by a single edge. At resolution 1 the bridge is cut and + both cliques survive intact: + + >>> import networkx as nx + >>> graph = nx.disjoint_union(nx.complete_graph(4), nx.complete_graph(4)) + >>> graph.add_edge(0, 4) + >>> reduce_graph(graph, resolution=1).number_of_edges() + 12 + + At the published resolution of 10 the communities are smaller than a clique, + so every edge goes and all eight nodes are left isolated. Nothing is dropped + -- they simply have no neighbourhood left to score: + + >>> reduce_graph(graph).number_of_edges() + 0 + >>> reduce_graph(graph).number_of_nodes() + 8 + """ + directed = nx.is_directed(graph) + + undirected = graph.to_undirected() if directed else graph.copy() + + communities = nx.community.louvain_communities( + undirected, resolution=resolution, seed=seed ) - # Create a dictionary to map nodes to their community - node_community = {} - for idx, community in enumerate(community_list): + membership = {} + for index, community in enumerate(communities): for node in community: - node_community[node] = idx - - # Create a new graph with only intra-community edges - if directed: - H = nx.DiGraph() - else: - H = nx.Graph() + membership[node] = index - H.add_nodes_from(G.nodes(data=True)) # Add all nodes with their attributes + reduced = nx.DiGraph() if directed else nx.Graph() + reduced.add_nodes_from(graph.nodes(data=True)) - # Add only edges that connect nodes within the same community - for u, v in G.edges(): - if node_community[u] == node_community[v]: - H.add_edge(u, v, **G[u][v]) + for u, v in graph.edges(): + if membership[u] == membership[v]: + reduced.add_edge(u, v, **graph[u][v]) - return H + return reduced diff --git a/src/garg_aml/scores.py b/src/garg_aml/scores.py index 545c6e6..965a074 100644 --- a/src/garg_aml/scores.py +++ b/src/garg_aml/scores.py @@ -1,221 +1,232 @@ """ Aggregation of block measures into the GARG-AML score: the cheap second stage. -Two aggregations. ``basic`` takes an unweighted mean over the penalty blocks; +Two aggregations. ``basic`` takes an unweighted mean over the blocks; ``weighted_average`` weights each block by its number of free entries, which is -Eq. (8) undirected and the size-weighted form of Eq. (14) directed. The package -defaults to ``weighted_average`` because that is what the paper reports -- see -docs/decisions/0001. +Eq. (8) undirected and the size-weighted form of Eq. (14) directed. +``weighted_average`` is the default because it is what the paper reports -- see +``docs/decisions/0001``. For a directed graph the score is Eq. (14) on the graph's given orientation. -``GARGAML_transposed`` and ``GARGAML_max`` are reported for reference; the -transpose-max is not the paper's reverse-flow handling and is not the score -- -see docs/decisions/0007. - -Copied unchanged from the research repository. +:func:`scores_from_measures` also reports the same quantity computed on the +transpose, and the larger of the two, for reference. Neither is the score, and +the transpose-max is **not** the paper's reverse-flow handling -- that lives in +the level assignment. See ``docs/decisions/0007``. """ +from collections.abc import Mapping +from typing import Any + import numpy as np import pandas as pd +__all__ = ["score_from_measures", "scores_from_measures"] -def calculate_score_directed(line, score_type="basic"): - """Directed score and its transpose for one row of block measures.""" - measure_00 = line["measure_00"] - measure_01 = line["measure_01"] - measure_02 = line["measure_02"] - measure_10 = line["measure_10"] - measure_11 = line["measure_11"] - measure_12 = line["measure_12"] - measure_20 = line["measure_20"] - measure_21 = line["measure_21"] - measure_22 = line["measure_22"] +#: One node's block measures: a plain mapping, or a row of a measures frame. +MeasureRow = Mapping[str, Any] | pd.Series - if score_type == "basic": - measure_high = np.mean([measure_01, measure_12]) - measure_low = np.mean( - [ - measure_10, - measure_21, - measure_00, - measure_02, - measure_11, - measure_20, - measure_22, - ] - ) - measure = measure_high - measure_low - - measure_high_transpose = np.mean([measure_10, measure_21]) - measure_low_transpose = np.mean( - [ - measure_01, - measure_12, - measure_00, - measure_20, - measure_11, - measure_02, - measure_22, - ] - ) - measure_transpose = measure_high_transpose - measure_low_transpose - - elif score_type == "weighted_average": - size_00 = line["size_00"] - size_01 = line["size_01"] - size_02 = line["size_02"] - size_10 = line["size_10"] - size_11 = line["size_11"] - size_12 = line["size_12"] - size_20 = line["size_20"] - size_21 = line["size_21"] - size_22 = line["size_22"] - - if size_01 + size_12 > 0: - measure_high = (size_01 * measure_01 + size_12 * measure_12) / ( - size_01 + size_12 - ) - else: - # both sizes are 0, revert to basic measure - measure_high = np.mean([measure_01, measure_12]) - - if size_10 + size_21 + size_00 + size_02 + size_11 + size_20 + size_22 > 0: - measure_low = ( - size_10 * measure_10 - + size_21 * measure_21 - + size_00 * measure_00 - + size_02 * measure_02 - + size_11 * measure_11 - + size_20 * measure_20 - + size_22 * measure_22 - ) / (size_10 + size_21 + size_00 + size_02 + size_11 + size_20 + size_22) - else: - # all sizes are 0, revert to basic measure - measure_low = np.mean( - [ - measure_10, - measure_21, - measure_00, - measure_02, - measure_11, - measure_20, - measure_22, - ] - ) - - measure = measure_high - measure_low - - if size_10 + size_21 > 0: - measure_high_transpose = (size_10 * measure_10 + size_21 * measure_21) / ( - size_10 + size_21 - ) - else: - # both sizes are 0, revert to basic measure - measure_high_transpose = np.mean([measure_10, measure_21]) - - if size_01 + size_12 + size_00 + size_20 + size_11 + size_02 + size_22 > 0: - measure_low_transpose = ( - size_01 * measure_01 - + size_12 * measure_12 - + size_00 * measure_00 - + size_20 * measure_20 - + size_11 * measure_11 - + size_02 * measure_02 - + size_22 * measure_22 - ) / (size_01 + size_12 + size_00 + size_20 + size_11 + size_02 + size_22) - else: - # all sizes are 0, revert to basic measure - measure_low_transpose = np.mean( - [ - measure_01, - measure_12, - measure_00, - measure_20, - measure_11, - measure_02, - measure_22, - ] - ) - - measure_transpose = measure_high_transpose - measure_low_transpose - - return measure, measure_transpose - - -def define_gargaml_scores_directed(results_df_measures, score_type="basic"): - """Directed scores for a frame of block measures.""" - nodes = [] - gargaml = [] - transposed_gargaml = [] - max_gargaml = [] - - for _i, line in results_df_measures.iterrows(): - measure, measure_transpose = calculate_score_directed( - line, score_type=score_type - ) +#: What the published experiments use. +DEFAULT_SCORE_TYPE = "weighted_average" - nodes.append(line["node"]) - gargaml.append(measure) - transposed_gargaml.append(measure_transpose) - max_gargaml.append(max(measure, measure_transpose)) +SCORE_TYPES = ("basic", "weighted_average") - dict_gargaml = { - "node": nodes, - "GARGAML": gargaml, - "GARGAML_transposed": transposed_gargaml, - "GARGAML_max": max_gargaml, - } - results_df = pd.DataFrame(dict_gargaml) - results_df.set_index("node", inplace=True) - return results_df +def _check_score_type(score_type: str) -> None: + if score_type not in SCORE_TYPES: + raise ValueError( + f"unknown score_type {score_type!r}; expected one of {SCORE_TYPES}" + ) -def calculate_score_undirected(line, score_type="basic"): - """Undirected score for one row of block measures (Eq. 8).""" - measure_1 = line["measure_1"] - measure_2 = line["measure_2"] - measure_3 = line["measure_3"] +def _directed_score(row: MeasureRow, score_type: str) -> tuple[float, float]: + """Eq. (14) on the given orientation, and the same on the transpose.""" + _check_score_type(score_type) + + measure_00 = row["measure_00"] + measure_01 = row["measure_01"] + measure_02 = row["measure_02"] + measure_10 = row["measure_10"] + measure_11 = row["measure_11"] + measure_12 = row["measure_12"] + measure_20 = row["measure_20"] + measure_21 = row["measure_21"] + measure_22 = row["measure_22"] + + dense = [measure_01, measure_12] + sparse = [ + measure_10, + measure_21, + measure_00, + measure_02, + measure_11, + measure_20, + measure_22, + ] + dense_t = [measure_10, measure_21] + sparse_t = [ + measure_01, + measure_12, + measure_00, + measure_20, + measure_11, + measure_02, + measure_22, + ] if score_type == "basic": - measure = measure_2 - (measure_1 + measure_3) / 2 - - elif score_type == "weighted_average": - size_1 = line["size_1"] - size_2 = line["size_2"] - size_3 = line["size_3"] - total_size = size_1 + size_3 - if total_size > 0: - measure = measure_2 - (size_1 * measure_1 + size_3 * measure_3) / total_size - elif size_2 > 0: - measure = measure_2 # both sizes are 0, so only measure_2 is relevant - else: - measure = -1 # Far away from smurfing - return measure - - -def define_gargaml_scores_undirected(results_df_measures, score_type="basic"): - """Undirected scores for a frame of block measures.""" - nodes = results_df_measures["node"].tolist() - gargaml = [ - calculate_score_undirected(line, score_type=score_type) - for _, line in results_df_measures.iterrows() - ] + return ( + np.mean(dense) - np.mean(sparse), + np.mean(dense_t) - np.mean(sparse_t), + ) + + size_00 = row["size_00"] + size_01 = row["size_01"] + size_02 = row["size_02"] + size_10 = row["size_10"] + size_11 = row["size_11"] + size_12 = row["size_12"] + size_20 = row["size_20"] + size_21 = row["size_21"] + size_22 = row["size_22"] + + dense_sizes = [size_01, size_12] + sparse_sizes = [size_10, size_21, size_00, size_02, size_11, size_20, size_22] + dense_sizes_t = [size_10, size_21] + sparse_sizes_t = [size_01, size_12, size_00, size_20, size_11, size_02, size_22] + + return ( + _weighted(dense, dense_sizes) - _weighted(sparse, sparse_sizes), + _weighted(dense_t, dense_sizes_t) - _weighted(sparse_t, sparse_sizes_t), + ) + - dict_gargaml = {"node": nodes, "GARGAML": gargaml} +def _weighted(measures: list[float], sizes: list[int]) -> float: + """Size-weighted mean, falling back to the plain mean when every size is 0.""" + total = sum(sizes) + if total > 0: + return sum(m * s for m, s in zip(measures, sizes, strict=True)) / total + return float(np.mean(measures)) - results_df = pd.DataFrame(dict_gargaml) - results_df.set_index("node", inplace=True) - return results_df +def _undirected_score(row: MeasureRow, score_type: str) -> float: + """Eq. (8).""" + _check_score_type(score_type) -def define_gargaml_scores(results_df_measures, directed, score_type="basic"): - """Dispatch to the directed or undirected aggregation.""" + measure_1 = row["measure_1"] + measure_2 = row["measure_2"] + measure_3 = row["measure_3"] + + if score_type == "basic": + return measure_2 - (measure_1 + measure_3) / 2 + + size_1 = row["size_1"] + size_2 = row["size_2"] + size_3 = row["size_3"] + + total = size_1 + size_3 + if total > 0: + return measure_2 - (size_1 * measure_1 + size_3 * measure_3) / total + if size_2 > 0: + return measure_2 # only the off-diagonal block carries any information + return -1 # no neighbourhood at all: as far from smurfing as the range goes + + +def score_from_measures( + row: MeasureRow, directed: bool = False, score_type: str = DEFAULT_SCORE_TYPE +) -> float: + """ + Return the GARG-AML score for one node's block measures. + + Parameters + ---------- + row : mapping + Block measures for one node, keyed as :func:`block_measures` names them: + ``measure_1`` to ``measure_3`` undirected, ``measure_00`` to + ``measure_22`` directed, plus the matching ``size_*`` entries when + ``score_type`` is ``"weighted_average"``. A row of the frame written by + :func:`scores_from_measures`' input works directly. + directed : bool, default False + Use Eq. (14) rather than Eq. (8). + score_type : {"weighted_average", "basic"}, default "weighted_average" + Weight each block by its free entries, or take an unweighted mean. + + Returns + ------- + float + A score in [-1, 1]. Higher is more smurfing-like; 1 is a pure pattern. + + Raises + ------ + ValueError + If ``score_type`` is not one of the two known values. + + References + ---------- + Deprez et al. (2025), Eq. (8) and Eq. (14). + + Examples + -------- + >>> row = {"measure_1": 0.0, "measure_2": 1.0, "measure_3": 0.0, + ... "size_1": 0, "size_2": 2, "size_3": 2} + >>> score_from_measures(row) + 1.0 + """ if directed: - return define_gargaml_scores_directed( - results_df_measures, score_type=score_type + score, _transposed = _directed_score(row, score_type) + return score + return _undirected_score(row, score_type) + + +def scores_from_measures( + measures: pd.DataFrame, directed: bool = False, score_type: str = DEFAULT_SCORE_TYPE +) -> pd.DataFrame: + """ + Scores for a frame of block measures. + + Parameters + ---------- + measures : pandas.DataFrame + One row per node, with a ``node`` column and the measure and size + columns :func:`block_measures` produces. + directed : bool, default False + Use Eq. (14) rather than Eq. (8). + score_type : {"weighted_average", "basic"}, default "weighted_average" + Weight each block by its free entries, or take an unweighted mean. + + Returns + ------- + pandas.DataFrame + Indexed by node. Undirected: a ``GARGAML`` column. Directed: also + ``GARGAML_transposed`` and ``GARGAML_max``, which are reported for + reference and are not the score -- see ``docs/decisions/0007``. + + Examples + -------- + >>> import pandas as pd + >>> measures = pd.DataFrame({"node": ["a"], "measure_1": [0.0], + ... "measure_2": [1.0], "measure_3": [0.0], + ... "size_1": [0], "size_2": [2], "size_3": [2]}) + >>> float(scores_from_measures(measures).loc["a", "GARGAML"]) + 1.0 + """ + if directed: + pairs = [_directed_score(row, score_type) for _, row in measures.iterrows()] + frame = pd.DataFrame( + { + "node": measures["node"].tolist(), + "GARGAML": [score for score, _ in pairs], + "GARGAML_transposed": [transposed for _, transposed in pairs], + "GARGAML_max": [max(pair) for pair in pairs], + } ) else: - return define_gargaml_scores_undirected( - results_df_measures, score_type=score_type + frame = pd.DataFrame( + { + "node": measures["node"].tolist(), + "GARGAML": [ + _undirected_score(row, score_type) for _, row in measures.iterrows() + ], + } ) + + return frame.set_index("node") diff --git a/tests/_pipeline.py b/tests/_pipeline.py index b051d21..e53e838 100644 --- a/tests/_pipeline.py +++ b/tests/_pipeline.py @@ -5,15 +5,14 @@ and that one would make the comparison meaningless, so keep them in step. """ +import numbers + import networkx as nx import pandas as pd -from garg_aml.features import summarise_gargaml_scores -from garg_aml.measures import ( - GARG_AML_node_directed_measures, - GARG_AML_node_undirected_measures, -) -from garg_aml.scores import define_gargaml_scores +from garg_aml.features import build_features +from garg_aml.measures import block_measures +from garg_aml.scores import scores_from_measures SCORE_TYPES = ["basic", "weighted_average"] FEATURE_SCORE_TYPE = "weighted_average" @@ -44,6 +43,20 @@ ] +def index_values(frame): + """ + A frame's index as comparable values, ignoring numeric dtype. + + The old directed path collected node ids through DataFrame.iterrows(), + which cast integer ids to float; the package preserves them. Comparing + numeric ids as floats lets both test layers check that the right nodes are + present and aligned without demanding the old artefact. Non-numeric ids -- + the string account names of the smurfing fixtures -- compare as themselves. + See docs/decisions/0008. + """ + return [float(i) if isinstance(i, numbers.Number) else i for i in frame.index] + + def load_graph(edges_path, directed, nodes_path=None): """Rebuild a graph from a frozen node list and edge list.""" G = nx.DiGraph() if directed else nx.Graph() @@ -76,14 +89,14 @@ def measures_frame(G, directed): G_rev = G.reverse(copy=True) columns = DIRECTED_COLUMNS rows = [ - GARG_AML_node_directed_measures(n, G, G_und, G_rev, include_size=True) + block_measures( + G, n, directed=True, undirected=G_und, reverse=G_rev, include_sizes=True + ) for n in nodes ] else: columns = UNDIRECTED_COLUMNS - rows = [ - GARG_AML_node_undirected_measures(n, G, include_size=True) for n in nodes - ] + rows = [block_measures(G, n, include_sizes=True) for n in nodes] frame = pd.DataFrame(rows, columns=columns) frame.insert(0, "node", nodes) @@ -93,7 +106,7 @@ def measures_frame(G, directed): def scores_frame(measures, directed): """Stage 2a: measures to score, for every score type, side by side.""" parts = [ - define_gargaml_scores(measures, directed, score_type=score_type).add_prefix( + scores_from_measures(measures, directed, score_type=score_type).add_prefix( score_type + "__" ) for score_type in SCORE_TYPES @@ -103,5 +116,5 @@ def scores_frame(measures, directed): def features_frame(G, measures, directed): """Stage 2b: score plus neighbourhood summary statistics.""" - scores = define_gargaml_scores(measures, directed, score_type=FEATURE_SCORE_TYPE) - return summarise_gargaml_scores(G, scores, columns=FEATURE_COLUMNS).sort_index() + scores = scores_from_measures(measures, directed, score_type=FEATURE_SCORE_TYPE) + return build_features(G, scores, columns=FEATURE_COLUMNS).sort_index() diff --git a/tests/test_equivalence.py b/tests/test_equivalence.py index 795fb4f..70d4e6a 100644 --- a/tests/test_equivalence.py +++ b/tests/test_equivalence.py @@ -21,9 +21,10 @@ import pandas as pd import pytest -from garg_aml.preprocess import graph_community +from garg_aml._ordering import node_order +from garg_aml.preprocess import reduce_graph -from ._pipeline import features_frame, measures_frame, scores_frame +from ._pipeline import features_frame, index_values, measures_frame, scores_frame RESEARCH_REPO = Path( os.environ.get("GARGAML_RESEARCH_REPO", Path(__file__).parents[3] / "GARG-AML") @@ -114,18 +115,14 @@ def test_node_ordering_matches(name, directed): expected = old["nodeselection"]( ego, node, True, G_ego_second_und=ego_und, G_ego_second_rev=ego_rev ) - from garg_aml._ordering import GARG_AML_nodeselection - - produced = GARG_AML_nodeselection( - ego, node, True, G_ego_second_und=ego_und, G_ego_second_rev=ego_rev + produced = node_order( + ego, node, True, ego_undirected=ego_und, ego_reverse=ego_rev ) else: ego = nx.ego_graph(G, node, 2) expected = old["nodeselection"](ego, node, False) - from garg_aml._ordering import GARG_AML_nodeselection - - produced = GARG_AML_nodeselection(ego, node, False) + produced = node_order(ego, node, False) assert produced == expected, f"ordering differs at node {node}" @@ -168,11 +165,15 @@ def test_scores_and_features_match(name, directed): for score_type in ("basic", "weighted_average"): expected = old["scores"](measures, directed, score_type=score_type) produced = scores_frame(measures, directed) + # check_index=False: the old directed path returned a float node + # index, an artefact of collecting ids through iterrows(). The values + # are what must agree. See docs/decisions/0008. for column in expected.columns: pd.testing.assert_series_equal( produced[f"{score_type}__{column}"].rename(column), expected[column].sort_index(), rtol=0, + check_index=False, ) expected_features = old["summarise"]( @@ -191,8 +192,12 @@ def test_scores_and_features_match(name, directed): "degree_std", ], ).sort_index() + produced_features = features_frame(G, measures, directed) + assert index_values(produced_features) == index_values(expected_features) pd.testing.assert_frame_equal( - features_frame(G, measures, directed), expected_features, rtol=0 + produced_features.reset_index(drop=True), + expected_features.reset_index(drop=True), + rtol=0, ) @@ -202,7 +207,7 @@ def test_louvain_reduction_matches(name, directed): old = _old() G = _build(name, directed) - produced = graph_community(G, resolution=10) + produced = reduce_graph(G, resolution=10) expected = old["community"](G, resolution=10) assert sorted(produced.nodes, key=str) == sorted(expected.nodes, key=str) diff --git a/tests/test_golden.py b/tests/test_golden.py index 65beb1c..fc54e96 100644 --- a/tests/test_golden.py +++ b/tests/test_golden.py @@ -9,11 +9,12 @@ import pandas as pd import pytest -from garg_aml.preprocess import graph_community +from garg_aml.preprocess import reduce_graph from ._pipeline import ( edge_frame, features_frame, + index_values, load_graph, measures_frame, node_frame, @@ -38,6 +39,13 @@ def _graph(golden, case, direction, variant): def _same(produced, expected): + # Indices are compared by value, not by dtype. The directed score and + # feature fixtures carry a *float* node index because the old + # implementation collected node ids through DataFrame.iterrows(), which + # coerces a row to one dtype; the package preserves whatever the caller's + # node ids are. See docs/decisions/0008. + assert index_values(produced) == index_values(expected) + pd.testing.assert_frame_equal( produced.reset_index(drop=True), expected.reset_index(drop=True), @@ -52,7 +60,7 @@ def test_reduce_graph_reproduces_the_frozen_partition(golden, case, direction): # The version-sensitive fixture: Louvain's partition is networkx's to # change. Kept separate so an upgrade breaks this and nothing else. G = load_graph(golden / case / "edges.csv", direction == "directed") - reduced = graph_community(G, resolution=RESOLUTION) + reduced = reduce_graph(G, resolution=RESOLUTION) _same( node_frame(reduced),