From 16849e200a3c7e94374aed685ca1c3583f81ddde Mon Sep 17 00:00:00 2001 From: B-Deprez Date: Fri, 18 Sep 2026 13:44:21 +0200 Subject: [PATCH] Phase 2: move the GARG-AML scoring core across verbatim --- CHANGELOG.md | 18 +++ pyproject.toml | 22 +++- src/garg_aml/_blocks.py | 194 ++++++++++++++++++++++++++++++++ src/garg_aml/_ordering.py | 73 ++++++++++++ src/garg_aml/features.py | 124 +++++++++++++++++++++ src/garg_aml/measures.py | 130 ++++++++++++++++++++++ src/garg_aml/preprocess.py | 71 ++++++++++++ src/garg_aml/scores.py | 221 +++++++++++++++++++++++++++++++++++++ tests/_pipeline.py | 107 ++++++++++++++++++ tests/test_equivalence.py | 209 +++++++++++++++++++++++++++++++++++ tests/test_golden.py | 103 +++++++++++++++++ 11 files changed, 1271 insertions(+), 1 deletion(-) create mode 100644 src/garg_aml/_blocks.py create mode 100644 src/garg_aml/_ordering.py create mode 100644 src/garg_aml/features.py create mode 100644 src/garg_aml/measures.py create mode 100644 src/garg_aml/preprocess.py create mode 100644 src/garg_aml/scores.py create mode 100644 tests/_pipeline.py create mode 100644 tests/test_equivalence.py create mode 100644 tests/test_golden.py diff --git a/CHANGELOG.md b/CHANGELOG.md index f584038..facefed 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -13,3 +13,21 @@ is `0.x` the public API may change with a minor bump, always with an entry here. - Frozen-behaviour fixtures under `tests/golden/`, generated from the pre-extraction implementation, with the tooling that built them in `tools/`. - Initial architecture decision records under `docs/decisions/`. +- The GARG-AML scoring core, moved across from the research repository with its + function bodies unchanged: `_blocks` (the 3 undirected and 9 directed block + densities), `_ordering` (the level assignment), `measures` (per-node block + measures), `scores` (aggregation into the score), `preprocess` (Louvain + reduction, hub removal) and `features` (neighbourhood summary statistics). +- `tests/test_golden.py`, comparing every stage against the frozen fixtures, and + `tests/test_equivalence.py`, comparing the moved code against the + implementation it came from over 22 graph shapes in both directions. The + latter needs the research repository present and is removed once the + extraction is complete. + +### 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/pyproject.toml b/pyproject.toml index cb3a249..78fbfa5 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -31,7 +31,7 @@ dependencies = ["numpy", "pandas", "networkx>=3.0", "scipy"] progress = ["tqdm"] parallel = ["joblib"] sklearn = ["scikit-learn"] -dev = ["ruff", "mypy", "pytest", "pytest-cov", "pre-commit"] +dev = ["ruff", "mypy", "pytest", "pytest-cov", "pre-commit", "pandas-stubs"] docs = ["mkdocs-material", "mkdocstrings[python]"] [project.urls] @@ -70,6 +70,19 @@ max-complexity = 10 [tool.ruff.lint.per-file-ignores] "tests/*" = ["D"] # tests document themselves by name "tools/*" = ["D"] +# The else-branches in _blocks.py encode the frozen degenerate-case constants +# (docs/decisions/0003): an empty off-diagonal block has density 1, everything +# else falls back to 0. An explicit if/else keeps each fallback visible beside +# 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] @@ -83,6 +96,13 @@ 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"] diff --git a/src/garg_aml/_blocks.py b/src/garg_aml/_blocks.py new file mode 100644 index 0000000..422fa21 --- /dev/null +++ b/src/garg_aml/_blocks.py @@ -0,0 +1,194 @@ +""" +Block densities of a second-order ego graph's adjacency matrix. + +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. + +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. +""" + + +def measure_1_function(piece_1_dim, adj_full): + """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 + + if reduced_size_1 > 0: + rel_1 = total_sum_1 / reduced_size_1 + else: + rel_1 = 0 + + return rel_1, reduced_size_1 + + +def measure_2_function(piece_1_dim, piece_2_dim, adj_full): + """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 + else: + rel_2 = 1 + + return rel_2, reduced_size_2 + + +def measure_3_function(piece_1_dim, piece_2_dim, piece_3_dim, adj_full): + """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] + + if reduced_size_3 > 0: + rel_3 = total_sum_3 / reduced_size_3 + else: + rel_3 = 0 + + return rel_3, reduced_size_3 + + +def measure_00_function(adj_full, size_0): + """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 + + if reduced_size_00 > 0: + rel_00 = total_sum_00 / reduced_size_00 + else: + rel_00 = 0 + + return rel_00, reduced_size_00 + + +def measure_01_function(adj_full, size_0, size_1): + """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 + + if total_size_01 > 0: + rel_01 = total_sum_01 / total_size_01 + else: + rel_01 = 1 # Since block only contains sure connections => full sum + + return rel_01, total_size_01 + + +def measure_02_function(adj_full, size_0, size_1, size_2): + """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 + + if reduced_size_02 > 0: + rel_02 = total_sum_02 / reduced_size_02 + else: + rel_02 = 0 + + return rel_02, reduced_size_02 + + +def measure_10_function(adj_full, size_0, size_1): + """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 + + if total_size_10 > 0: + rel_10 = total_sum_10 / total_size_10 + else: + rel_10 = 0 + + return rel_10, total_size_10 + + +def measure_11_function(adj_full, size_0, size_1): + """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 + + if reduced_size_11 > 0: + rel_11 = total_sum_11 / reduced_size_11 + else: + rel_11 = 0 + + return rel_11, reduced_size_11 + + +def measure_12_function(adj_full, size_0, size_1, size_2): + """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 + + if total_size_12 > 0: + rel_12 = total_sum_12 / total_size_12 + elif size_2 > 0: + rel_12 = 1 # Since block only contains sure connections => full sum + else: + rel_12 = 0 # No connections at all + + return rel_12, total_size_12 + + +def measure_20_function(adj_full, size_0, size_1, size_2): + """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 + + if reduced_size_20 > 0: + rel_20 = total_sum_20 / reduced_size_20 + else: + rel_20 = 0 + + return rel_20, reduced_size_20 + + +def measure_21_function(adj_full, size_0, size_1): + """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 + + if total_size_21 > 0: + rel_21 = total_sum_21 / total_size_21 + else: + rel_21 = 0 + + return rel_21, total_size_21 + + +def measure_22_function(adj_full, size_0, size_1, size_2): + """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 + + if reduced_size_22 > 0: + rel_22 = total_sum_22 / reduced_size_22 + else: + rel_22 = 0 + + return rel_22, reduced_size_22 diff --git a/src/garg_aml/_ordering.py b/src/garg_aml/_ordering.py new file mode 100644 index 0000000..59c2d59 --- /dev/null +++ b/src/garg_aml/_ordering.py @@ -0,0 +1,73 @@ +""" +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 +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. +""" + +import networkx as nx + + +def GARG_AML_nodeselection_undirected(G_ego_second, node): + """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) + + # For undirected networks, specific order to obtain scores + # (group node with second order neighbours) + nodes_ordered = [node] + nodes_2 + nodes_1 + + return nodes_1, nodes_2, nodes_ordered + + +def GARG_AML_nodeselection_directed( + G_ego_second, G_ego_second_und, G_ego_second_rev, node +): + """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)) + ) + + nodes_0 = [node] + nodes_0 + + for n in nodes_0: + nodes_2.remove(n) + for n in nodes_1: + nodes_2.remove(n) + + # For directed network, specific order to obtain scores (in order of "group") + nodes_ordered = nodes_0 + nodes_1 + nodes_2 + + return nodes_0, nodes_1, nodes_2, nodes_ordered + + +def GARG_AML_nodeselection( + G_ego_second, node, directed, G_ego_second_und=None, G_ego_second_rev=None +): + """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) diff --git a/src/garg_aml/features.py b/src/garg_aml/features.py new file mode 100644 index 0000000..d8a0430 --- /dev/null +++ b/src/garg_aml/features.py @@ -0,0 +1,124 @@ +""" +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. + +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. +""" + +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]) + + if not list_ego_degree: + return [0, 0, 0, 0] + + return [ + np.min(list_ego_degree), + np.mean(list_ego_degree), + np.max(list_ego_degree), + np.std(list_ego_degree), + ] + + +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.""" + if columns is None: + columns = ["GARGAML"] + + G_degree_dict = dict(G_reduced.degree()) + nodes = list(df_results.index) + + gargaml_values = dict(zip(df_results.index, df_results["GARGAML"], strict=False)) + + 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", +] diff --git a/src/garg_aml/measures.py b/src/garg_aml/measures.py new file mode 100644 index 0000000..91f481a --- /dev/null +++ b/src/garg_aml/measures.py @@ -0,0 +1,130 @@ +""" +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. + +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. +""" + +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, +) +from ._ordering import GARG_AML_nodeselection + + +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 + ) + + adj_full = nx.adjacency_matrix(G_ego_second, nodelist=nodes_ordered).toarray() + + size_second = len(nodes_2) + size_first = len(nodes_1) + + piece_1_dim = [size_second + 1, size_second + 1] + piece_2_dim = [size_first, size_second + 1] + piece_3_dim = [size_first, size_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 + ) + + if include_size: + 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, + ) + + 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: + return ( + measure_00, + measure_01, + measure_02, + measure_10, + measure_11, + measure_12, + measure_20, + measure_21, + measure_22, + ) diff --git a/src/garg_aml/preprocess.py b/src/garg_aml/preprocess.py new file mode 100644 index 0000000..07af4b6 --- /dev/null +++ b/src/garg_aml/preprocess.py @@ -0,0 +1,71 @@ +""" +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 +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. + +``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. +""" + +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 + ) + + # Create a dictionary to map nodes to their community + node_community = {} + for idx, community in enumerate(community_list): + 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() + + H.add_nodes_from(G.nodes(data=True)) # Add all nodes with their attributes + + # 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]) + + return H diff --git a/src/garg_aml/scores.py b/src/garg_aml/scores.py new file mode 100644 index 0000000..545c6e6 --- /dev/null +++ b/src/garg_aml/scores.py @@ -0,0 +1,221 @@ +""" +Aggregation of block measures into the GARG-AML score: the cheap second stage. + +Two aggregations. ``basic`` takes an unweighted mean over the penalty 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. + +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. +""" + +import numpy as np +import pandas as pd + + +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"] + + 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 + ) + + nodes.append(line["node"]) + gargaml.append(measure) + transposed_gargaml.append(measure_transpose) + max_gargaml.append(max(measure, measure_transpose)) + + 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 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"] + + 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() + ] + + dict_gargaml = {"node": nodes, "GARGAML": gargaml} + + results_df = pd.DataFrame(dict_gargaml) + results_df.set_index("node", inplace=True) + return results_df + + +def define_gargaml_scores(results_df_measures, directed, score_type="basic"): + """Dispatch to the directed or undirected aggregation.""" + if directed: + return define_gargaml_scores_directed( + results_df_measures, score_type=score_type + ) + else: + return define_gargaml_scores_undirected( + results_df_measures, score_type=score_type + ) diff --git a/tests/_pipeline.py b/tests/_pipeline.py new file mode 100644 index 0000000..b051d21 --- /dev/null +++ b/tests/_pipeline.py @@ -0,0 +1,107 @@ +""" +Each pipeline stage rebuilt from garg_aml, for comparison against the fixtures. + +Mirrors tools/make_golden.py stage for stage. Any divergence between this file +and that one would make the comparison meaningless, so keep them in step. +""" + +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 + +SCORE_TYPES = ["basic", "weighted_average"] +FEATURE_SCORE_TYPE = "weighted_average" + +FEATURE_COLUMNS = [ + "GARGAML", + "GARGAML_min", + "GARGAML_mean", + "GARGAML_max", + "GARGAML_std", + "degree", + "degree_min", + "degree_mean", + "degree_max", + "degree_std", +] + +DIRECTED_COLUMNS = [f"measure_{i}{j}" for i in range(3) for j in range(3)] + [ + f"size_{i}{j}" for i in range(3) for j in range(3) +] +UNDIRECTED_COLUMNS = [ + "measure_1", + "measure_2", + "measure_3", + "size_1", + "size_2", + "size_3", +] + + +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() + if nodes_path is not None: + G.add_nodes_from(pd.read_csv(nodes_path)["node"]) + data = pd.read_csv(edges_path) + G.add_edges_from(zip(data["source"], data["target"], strict=True)) + G.remove_edges_from([(n, n) for n in G.nodes() if G.has_edge(n, n)]) + return G + + +def edge_frame(G): + """Sorted edge list of a graph.""" + return pd.DataFrame( + sorted((int(u), int(v)) for u, v in G.edges()), columns=["source", "target"] + ) + + +def node_frame(G): + """Sorted node list of a graph.""" + return pd.DataFrame({"node": sorted(G.nodes)}) + + +def measures_frame(G, directed): + """Stage 1: per-node block measures and block sizes.""" + nodes = sorted(G.nodes) + + if directed: + G_und = G.to_undirected() + G_rev = G.reverse(copy=True) + columns = DIRECTED_COLUMNS + rows = [ + GARG_AML_node_directed_measures(n, G, G_und, G_rev, include_size=True) + for n in nodes + ] + else: + columns = UNDIRECTED_COLUMNS + rows = [ + GARG_AML_node_undirected_measures(n, G, include_size=True) for n in nodes + ] + + frame = pd.DataFrame(rows, columns=columns) + frame.insert(0, "node", nodes) + return frame + + +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( + score_type + "__" + ) + for score_type in SCORE_TYPES + ] + return pd.concat(parts, axis=1).sort_index() + + +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() diff --git a/tests/test_equivalence.py b/tests/test_equivalence.py new file mode 100644 index 0000000..795fb4f --- /dev/null +++ b/tests/test_equivalence.py @@ -0,0 +1,209 @@ +""" +The L2 test: the extracted code must agree with the implementation it came from. + +Where test_golden.py pins four fixed graphs, this sweeps a spread of shapes -- +including the degenerate ones no real dataset would contain -- and compares the +new implementation against the old one node by node. + +It needs the research repository beside this one, so it skips in CI. Set +GARGAML_RESEARCH_REPO to point elsewhere. + +**This file is deleted at Phase 5.** Its job is to prove the extraction, and +keeping a dependency on the old tree afterwards is the coupling the migration +exists to remove. +""" + +import os +import sys +from pathlib import Path + +import networkx as nx +import pandas as pd +import pytest + +from garg_aml.preprocess import graph_community + +from ._pipeline import features_frame, measures_frame, scores_frame + +RESEARCH_REPO = Path( + os.environ.get("GARGAML_RESEARCH_REPO", Path(__file__).parents[3] / "GARG-AML") +) + +pytestmark = pytest.mark.skipif( + not (RESEARCH_REPO / "src" / "methods" / "GARGAML.py").exists(), + reason=f"research repository not found at {RESEARCH_REPO}", +) + + +def _old(): + """Import the pre-extraction implementation.""" + if str(RESEARCH_REPO) not in sys.path: + sys.path.insert(0, str(RESEARCH_REPO)) + + from src.methods.GARGAML import ( + GARG_AML_node_directed_measures, + GARG_AML_node_undirected_measures, + ) + from src.methods.gargaml_scores import ( + define_gargaml_scores, + summarise_gargaml_scores, + ) + from src.methods.utils.neighbourhood_functions import GARG_AML_nodeselection + from src.utils.graph_processing import graph_community as old_community + + return { + "undirected_measures": GARG_AML_node_undirected_measures, + "directed_measures": GARG_AML_node_directed_measures, + "nodeselection": GARG_AML_nodeselection, + "scores": define_gargaml_scores, + "summarise": summarise_gargaml_scores, + "community": old_community, + } + + +def _smurfing_edges(n_mules): + """One source paying n mules, which all pay one target.""" + source, target = "source", "target" + mules = [f"mule_{i}" for i in range(n_mules)] + return [(source, m) for m in mules] + [(m, target) for m in mules] + + +# Shapes chosen to exercise the degenerate paths as hard as the ordinary ones: +# empty and single-node graphs, isolated nodes, a clique (no block structure at +# all), and a perfect smurfing pattern (maximal block structure). +GRAPH_SPECS = { + "empty": nx.empty_graph(0), + "single_node": nx.empty_graph(1), + "isolated_pair": nx.empty_graph(2), + "one_edge": nx.path_graph(2), + "star_8": nx.star_graph(8), + "clique_6": nx.complete_graph(6), + "path_7": nx.path_graph(7), + "cycle_7": nx.cycle_graph(7), + "smurf_3": nx.Graph(_smurfing_edges(3)), + "smurf_8": nx.Graph(_smurfing_edges(8)), + "two_components": nx.disjoint_union(nx.star_graph(4), nx.complete_graph(4)), + **{f"ba_50_{s}": nx.barabasi_albert_graph(50, 2, seed=s) for s in (1, 2, 3)}, + **{f"er_50_{s}": nx.erdos_renyi_graph(50, 0.08, seed=s) for s in (1, 2, 3)}, + **{f"ws_50_{s}": nx.watts_strogatz_graph(50, 4, 0.1, seed=s) for s in (1, 2, 3)}, +} + +CASES = sorted(GRAPH_SPECS) + + +def _build(name, directed): + """Same edge set as a Graph or a DiGraph, with node labels normalised.""" + source = GRAPH_SPECS[name] + G = nx.DiGraph() if directed else nx.Graph() + G.add_nodes_from(source.nodes) + G.add_edges_from(source.edges) + return G + + +@pytest.mark.parametrize("name", CASES) +@pytest.mark.parametrize("directed", [False, True], ids=["undirected", "directed"]) +def test_node_ordering_matches(name, directed): + old = _old() + G = _build(name, directed) + + for node in sorted(G.nodes, key=str): + if directed: + ego_und = nx.ego_graph(G.to_undirected(), node, 2) + ego = nx.subgraph(G, ego_und.nodes) + ego_rev = nx.ego_graph(G.reverse(copy=True), node, 2) + 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 + ) + 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) + + assert produced == expected, f"ordering differs at node {node}" + + +@pytest.mark.parametrize("name", CASES) +@pytest.mark.parametrize("directed", [False, True], ids=["undirected", "directed"]) +def test_block_measures_match(name, directed): + old = _old() + G = _build(name, directed) + + if directed: + G_und, G_rev = G.to_undirected(), G.reverse(copy=True) + + for node in sorted(G.nodes, key=str): + if directed: + expected = old["directed_measures"]( + node, G, G_und, G_rev, include_size=True + ) + else: + expected = old["undirected_measures"](node, G, include_size=True) + + produced = measures_frame(G, directed) + row = produced[produced["node"] == node].iloc[0].drop("node").tolist() + + assert row == pytest.approx(list(expected), rel=0, abs=0), ( + f"measures differ at node {node}" + ) + + +@pytest.mark.parametrize("name", CASES) +@pytest.mark.parametrize("directed", [False, True], ids=["undirected", "directed"]) +def test_scores_and_features_match(name, directed): + old = _old() + G = _build(name, directed) + if G.number_of_nodes() == 0: + pytest.skip("no nodes to score") + + measures = measures_frame(G, directed) + + for score_type in ("basic", "weighted_average"): + expected = old["scores"](measures, directed, score_type=score_type) + produced = scores_frame(measures, directed) + for column in expected.columns: + pd.testing.assert_series_equal( + produced[f"{score_type}__{column}"].rename(column), + expected[column].sort_index(), + rtol=0, + ) + + expected_features = old["summarise"]( + G, + old["scores"](measures, directed, score_type="weighted_average"), + columns=[ + "GARGAML", + "GARGAML_min", + "GARGAML_mean", + "GARGAML_max", + "GARGAML_std", + "degree", + "degree_min", + "degree_mean", + "degree_max", + "degree_std", + ], + ).sort_index() + pd.testing.assert_frame_equal( + features_frame(G, measures, directed), expected_features, rtol=0 + ) + + +@pytest.mark.parametrize("name", CASES) +@pytest.mark.parametrize("directed", [False, True], ids=["undirected", "directed"]) +def test_louvain_reduction_matches(name, directed): + old = _old() + G = _build(name, directed) + + produced = graph_community(G, resolution=10) + expected = old["community"](G, resolution=10) + + assert sorted(produced.nodes, key=str) == sorted(expected.nodes, key=str) + assert sorted(map(str, produced.edges)) == sorted(map(str, expected.edges)) diff --git a/tests/test_golden.py b/tests/test_golden.py new file mode 100644 index 0000000..65beb1c --- /dev/null +++ b/tests/test_golden.py @@ -0,0 +1,103 @@ +""" +The L1 test: garg_aml must reproduce the frozen fixtures exactly. + +These numbers came from the implementation that produced the published paper +results. A failure here means the code changed, never that a fixture needs +updating -- see tests/golden/README.md. +""" + +import pandas as pd +import pytest + +from garg_aml.preprocess import graph_community + +from ._pipeline import ( + edge_frame, + features_frame, + load_graph, + measures_frame, + node_frame, + scores_frame, +) +from .conftest import CASES, DIRECTIONS, VARIANTS + +RTOL = 1e-12 +RESOLUTION = 10 + + +def _graph(golden, case, direction, variant): + case_dir = golden / case + directed = direction == "directed" + if variant == "raw": + return load_graph(case_dir / "edges.csv", directed) + return load_graph( + case_dir / f"reduced_edges_{direction}.csv", + directed, + nodes_path=case_dir / f"reduced_nodes_{direction}.csv", + ) + + +def _same(produced, expected): + pd.testing.assert_frame_equal( + produced.reset_index(drop=True), + expected.reset_index(drop=True), + rtol=RTOL, + check_dtype=False, + ) + + +@pytest.mark.parametrize("case", CASES) +@pytest.mark.parametrize("direction", DIRECTIONS) +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) + + _same( + node_frame(reduced), + pd.read_csv(golden / case / f"reduced_nodes_{direction}.csv"), + ) + _same( + edge_frame(reduced), + pd.read_csv(golden / case / f"reduced_edges_{direction}.csv"), + ) + + +@pytest.mark.parametrize("case", CASES) +@pytest.mark.parametrize("direction", DIRECTIONS) +@pytest.mark.parametrize("variant", VARIANTS) +def test_measures_reproduce(golden, case, direction, variant): + graph = _graph(golden, case, direction, variant) + produced = measures_frame(graph, direction == "directed") + _same(produced, pd.read_csv(golden / case / f"measures_{direction}_{variant}.csv")) + + +@pytest.mark.parametrize("case", CASES) +@pytest.mark.parametrize("direction", DIRECTIONS) +@pytest.mark.parametrize("variant", VARIANTS) +def test_scores_reproduce(golden, case, direction, variant): + directed = direction == "directed" + graph = _graph(golden, case, direction, variant) + produced = scores_frame(measures_frame(graph, directed), directed) + _same( + produced, + pd.read_csv( + golden / case / f"scores_{direction}_{variant}.csv", index_col="node" + ), + ) + + +@pytest.mark.parametrize("case", CASES) +@pytest.mark.parametrize("direction", DIRECTIONS) +@pytest.mark.parametrize("variant", VARIANTS) +def test_features_reproduce(golden, case, direction, variant): + directed = direction == "directed" + graph = _graph(golden, case, direction, variant) + produced = features_frame(graph, measures_frame(graph, directed), directed) + _same( + produced, + pd.read_csv( + golden / case / f"features_{direction}_{variant}.csv", index_col="node" + ), + )