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
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
0087 CASES = {
0088
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
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
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
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
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())