Back to home page

EIC code displayed by LXR

 
 

    


File indexing completed on 2026-09-28 09:36:49

0001 """
0002 Differential check: prun / nested / scatter parsing must be identical before and after the
0003 raw-task-params ("task") step support was added.
0004 
0005 Loads the pre-change workflow_native_utils.py out of a pinned baseline commit alongside the
0006 working-tree version, runs both over the same workflow descriptions, and diffs the resulting node
0007 lists.
0008 
0009 The baseline is pinned to an explicit commit rather than HEAD: once the change under test is
0010 committed, HEAD contains it too and the comparison would silently be new-against-new.
0011 
0012 The descriptions below mirror the shapes of pandaserver/workflow/examples/*.yaml. They are given as
0013 dicts rather than loaded from the YAML files because no YAML loader is installed in a bare checkout
0014 -- and because the YAML->dict step is not touched by this change, so parsing from an equivalent
0015 dict exercises exactly the code under test.
0016 
0017 Run from the repository root:  python3 pandaserver/workflow/examples/regression_diff_test.py
0018 """
0019 
0020 import copy
0021 import difflib
0022 import importlib.util
0023 import json
0024 import os
0025 import subprocess
0026 import sys
0027 import tempfile
0028 import types
0029 
0030 REPO_ROOT = os.path.abspath(os.path.join(os.path.dirname(__file__), "..", "..", ".."))
0031 sys.path.insert(0, REPO_ROOT)
0032 MODULE_PATH = "pandaserver/workflow/workflow_native_utils.py"
0033 # Last commit before raw-task-params ("task") step support was added
0034 MODULE_BASE_COMMIT = "83733672"
0035 
0036 
0037 def _stub(name, **attrs):
0038     mod = types.ModuleType(name)
0039     for k, v in attrs.items():
0040         setattr(mod, k, v)
0041     sys.modules[name] = mod
0042 
0043 
0044 _stub("pandaclient", PhpoScript=types.SimpleNamespace(main=None), PrunScript=types.SimpleNamespace(main=None))
0045 _stub("pandacommon")
0046 _stub("pandacommon.pandautils")
0047 _stub("pandacommon.pandautils.base", SpecBase=object)
0048 _stub("pandaserver.config", panda_config=types.SimpleNamespace(schemaJEDI="ATLAS_PANDA", schemaDEFT="ATLAS_DEFT"))
0049 
0050 
0051 class Log:
0052     def info(self, m):
0053         pass
0054 
0055     def debug(self, m):
0056         pass
0057 
0058     def warning(self, m):
0059         pass
0060 
0061     def error(self, m):
0062         pass
0063 
0064 
0065 def load_old_module():
0066     """Load the pre-change module from the pinned baseline commit under its own module name."""
0067     src = subprocess.run(
0068         ["git", "show", f"{MODULE_BASE_COMMIT}:{MODULE_PATH}"],
0069         cwd=REPO_ROOT,
0070         capture_output=True,
0071         text=True,
0072         check=True,
0073     ).stdout
0074     tmp = tempfile.NamedTemporaryFile("w", suffix=".py", delete=False)
0075     tmp.write(src)
0076     tmp.close()
0077     spec = importlib.util.spec_from_file_location("wnu_head", tmp.name)
0078     if spec is None or spec.loader is None:
0079         raise RuntimeError(f"cannot load the baseline module written to {tmp.name}")
0080     mod = importlib.util.module_from_spec(spec)
0081     spec.loader.exec_module(mod)
0082     os.unlink(tmp.name)
0083     return mod
0084 
0085 
0086 # ---- descriptions mirroring the shapes of the checked-in examples -------------------------
0087 CASES = {
0088     # multistep_merge_wfd.yaml -- a plain serial prun chain
0089     "serial_prun_chain": {
0090         "name": "multistep_merge_chain",
0091         "inputs": {"input_to_merge": "user.sgaid:user.sgaid.some.input"},
0092         "outputs": {"final_output": {"from": "third/outDS", "output_types": ["merge.root"]}},
0093         "steps": {
0094             "first": {"type": "prun", "inDS": "{input_to_merge}", "args": "--outputs merge.root --noBuild", "exec": "merge.sh"},
0095             "second": {"type": "prun", "inDS": "first/outDS", "args": "--outputs merge.root --noBuild", "exec": "merge.sh"},
0096             "third": {"type": "prun", "inDS": "second/outDS", "args": "--outputs merge.root --noBuild", "exec": "merge.sh"},
0097         },
0098     },
0099     # signal_background_combine_wfd.yaml -- branches, secondaryDSs, inDsType, containerImage
0100     "branching_prun_with_secondaries": {
0101         "name": "signal_background_combine",
0102         "inputs": {"signal": "mc16_valid:mc16_valid.signal.HITS", "background": "mc16_5TeV.background.HITS/"},
0103         "outputs": {"outDS": {"from": "combine/outDS", "output_types": ["aaa.root"]}},
0104         "steps": {
0105             "make_signal": {
0106                 "type": "prun",
0107                 "inDS": "{signal}",
0108                 "containerImage": "docker://busybox",
0109                 "args": "--outputs abc.dat,def.zip --nFilesPerJob 5",
0110                 "exec": "echo %IN > abc.dat; echo 123 > def.zip",
0111             },
0112             "make_background_1": {"type": "prun", "inDS": "{background}", "args": "--outputs opq.root,xyz.pool", "exec": "echo %IN > opq.root"},
0113             "generate_some": {"type": "prun", "args": "--outputs gen.root --nJobs 10", "exec": "echo %RNDM:10 > gen.root"},
0114             "premix": {
0115                 "type": "prun",
0116                 "inDS": "make_signal/outDS",
0117                 "inDsType": "def.zip",
0118                 "secondaryDSs": ["make_background_1/outDS"],
0119                 "secondaryDsTypes": ["xyz.pool"],
0120                 "args": "--outputs klm.root --secondaryDSs IN2:13:%{SECDS1}",
0121                 "exec": "echo %IN %IN2 > klm.root",
0122             },
0123             "make_background_2": {
0124                 "type": "prun",
0125                 "inDS": "{background}",
0126                 "containerImage": "docker://alpine",
0127                 "secondaryDSs": ["generate_some/outDS"],
0128                 "secondaryDsTypes": ["gen.root"],
0129                 "args": "--outputs ooo.root,jjj.txt --secondaryDSs IN2:10:%{SECDS1}",
0130                 "exec": "echo %IN > ooo.root",
0131             },
0132             "combine": {
0133                 "type": "prun",
0134                 "inDS": "make_signal/outDS",
0135                 "inDsType": "abc.dat",
0136                 "secondaryDSs": ["premix/outDS", "make_background_2/outDS"],
0137                 "secondaryDsTypes": ["klm.root", "ooo.root"],
0138                 "args": "--outputs aaa.root --secondaryDSs IN2:2:%{SECDS1},IN3:5:%{SECDS2}",
0139                 "exec": "echo %IN %IN2 %IN3 > aaa.root",
0140             },
0141         },
0142     },
0143     # nested_workflow_inline_sig_bg_comb_wfd.yaml -- an inline sub-workflow
0144     "inline_sub_workflow": {
0145         "name": "nested_inline",
0146         "inputs": {"signal": "mc16_valid:signal.HITS", "background": "mc16_5TeV:background.HITS/"},
0147         "outputs": {"final": {"from": "merge/outDS", "output_types": ["merged.root"]}},
0148         "steps": {
0149             "sig_bg_comb": {
0150                 "type": "workflow",
0151                 "inputs": {"signal": "{signal}", "background": "{background}"},
0152                 "outputs": {"combined": {"from": "combine/outDS", "output_types": ["aaa.root"]}},
0153                 "steps": {
0154                     "make_signal": {"type": "prun", "inDS": "{signal}", "args": "--outputs abc.dat", "exec": "echo %IN > abc.dat"},
0155                     "combine": {"type": "prun", "inDS": "make_signal/outDS", "args": "--outputs aaa.root", "exec": "echo %IN > aaa.root"},
0156                 },
0157             },
0158             "merge": {"type": "prun", "inDS": "sig_bg_comb/outDS", "args": "--outputs merged.root", "exec": "echo %IN > merged.root"},
0159         },
0160     },
0161     # scatter_sig_bg_comb_wfd.yaml -- a scatter sub-workflow
0162     "scatter_sub_workflow": {
0163         "name": "scatter_nested",
0164         "inputs": {"signals": ["ds.sig.a", "ds.sig.b", "ds.sig.c"]},
0165         "outputs": {"final": {"from": "merge/outDS", "output_types": ["merged.root"]}},
0166         "steps": {
0167             "per_signal": {
0168                 "type": "workflow",
0169                 "scatter_inputs": {"signal": "signals"},
0170                 "scatter_mode": "zip",
0171                 "outputs": {"combined": {"from": "combine/outDS", "output_types": ["aaa.root"]}},
0172                 "steps": {
0173                     "combine": {"type": "prun", "inDS": "{signal}", "args": "--outputs aaa.root", "exec": "echo %IN > aaa.root"},
0174                 },
0175             },
0176             "merge": {"type": "prun", "inDS": "per_signal/outDS", "args": "--outputs merged.root", "exec": "echo %IN > merged.root"},
0177         },
0178     },
0179 }
0180 
0181 
0182 def normalise(nodes):
0183     """Serialise a node list to a stable, comparable form."""
0184 
0185     def default(obj):
0186         if isinstance(obj, set):
0187             return sorted(obj, key=str)
0188         if obj.__class__.__name__ == "Node":
0189             return f"<Node id={obj.id}>"
0190         return str(obj)
0191 
0192     return json.dumps([{k: v for k, v in sorted(vars(n).items())} for n in nodes], indent=2, sort_keys=True, default=default)
0193 
0194 
0195 def run(module, description, out_ds_name):
0196     log = Log()
0197     nodes, root_in = module.parse_workflow_data(copy.deepcopy(description), log)
0198     data = description.get("inputs", {})
0199     serial, tails, nodes = module.resolve_nodes(nodes, dict(root_in), copy.deepcopy(data), 0, set(), out_ds_name, log)
0200     module.set_workflow_outputs(nodes)
0201     id_map = module.get_node_id_map(nodes)
0202     for n in nodes:
0203         # task_template is None so make_task_params (which needs a real PrunScript) is skipped
0204         n.resolve_params(None, id_map)
0205     verdicts = {n.name: n.verify() for n in nodes}
0206     return normalise(nodes), sorted(t.name for t in tails), verdicts
0207 
0208 
0209 def main():
0210     old = load_old_module()
0211     from pandaserver.workflow import workflow_native_utils as new
0212 
0213     print(f"comparing working tree against {MODULE_BASE_COMMIT} ({MODULE_PATH})\n")
0214     failures = 0
0215     for case_name, description in CASES.items():
0216         out_ds_name = "user.me.myOutDS"
0217         old_dump, old_tails, old_verdicts = run(old, description, out_ds_name)
0218         new_dump, new_tails, new_verdicts = run(new, description, out_ds_name)
0219 
0220         identical = old_dump == new_dump and old_tails == new_tails and old_verdicts == new_verdicts
0221         print(f"  {'PASS' if identical else 'FAIL'}  {case_name}  ({len(json.loads(new_dump))} nodes)")
0222         if not identical:
0223             failures += 1
0224             if old_tails != new_tails:
0225                 print(f"        tails differ: {old_tails} -> {new_tails}")
0226             if old_verdicts != new_verdicts:
0227                 print(f"        verify differs: {old_verdicts} -> {new_verdicts}")
0228             for line in list(difflib.unified_diff(old_dump.splitlines(), new_dump.splitlines(), MODULE_BASE_COMMIT, "working-tree", lineterm=""))[:40]:
0229                 print(f"        {line}")
0230 
0231     print(f"\n{'ALL CASES IDENTICAL' if not failures else f'{failures} CASE(S) DIFFER'}")
0232     return 1 if failures else 0
0233 
0234 
0235 if __name__ == "__main__":
0236     sys.exit(main())