File indexing completed on 2026-09-28 09:36:49
0001 """
0002 Offline check of raw-task-params ("task") workflow step parsing.
0003
0004 Stubs the runtime dependencies which are not installed in a bare checkout (pandaclient,
0005 pandacommon, ruamel, pandaserver.config) so that parse_workflow_data / resolve_nodes /
0006 resolve_params / verify can be exercised directly on the example description.
0007
0008 Run from the repository root: python3 pandaserver/workflow/examples/parse_task_steps_test.py
0009 """
0010
0011 import json
0012 import os
0013 import sys
0014 import types
0015 from typing import Any
0016
0017 REPO_ROOT = os.path.abspath(os.path.join(os.path.dirname(__file__), "..", "..", ".."))
0018 sys.path.insert(0, REPO_ROOT)
0019
0020
0021
0022 def _stub(name, **attrs):
0023 mod = types.ModuleType(name)
0024 for k, v in attrs.items():
0025 setattr(mod, k, v)
0026 sys.modules[name] = mod
0027 return mod
0028
0029
0030 _stub("pandaclient", PhpoScript=types.SimpleNamespace(main=None), PrunScript=types.SimpleNamespace(main=None))
0031 _stub("pandacommon")
0032 _stub("pandacommon.pandautils")
0033 _stub("pandacommon.pandautils.base", SpecBase=object)
0034 _stub("pandaserver.config", panda_config=types.SimpleNamespace(schemaJEDI="ATLAS_PANDA", schemaDEFT="ATLAS_DEFT"))
0035
0036 from pandaserver.workflow import workflow_native_utils as wnu
0037 from pandaserver.workflow.workflow_base import (
0038 TASKID_PLACEHOLDER,
0039 WFID_PLACEHOLDER,
0040 has_placeholder,
0041 substitute_placeholder,
0042 )
0043
0044
0045 class Log:
0046 def info(self, m):
0047 pass
0048
0049 def debug(self, m):
0050 pass
0051
0052 def warning(self, m):
0053 print(f" WARNING: {m}")
0054
0055 def error(self, m):
0056 print(f" ERROR: {m}")
0057
0058
0059 def check(label, cond, detail=""):
0060 print(f" {'PASS' if cond else 'FAIL'} {label}{' ' + detail if detail and not cond else ''}")
0061 return cond
0062
0063
0064 def main():
0065 log = Log()
0066 failures = 0
0067 wfd = json.load(open(os.path.join(os.path.dirname(__file__), "production_chain_wfd.json")))
0068
0069 print("\n=== parse_workflow_data ===")
0070 nodes, root_in = wnu.parse_workflow_data(wfd, log)
0071 by_name = {n.name: n for n in nodes}
0072 failures += not check("all 9 steps parsed", len(nodes) == 9, f"got {len(nodes)}")
0073 failures += not check("all steps are leaves", all(n.is_leaf for n in nodes))
0074 failures += not check("all steps typed 'task'", all(n.type == "task" for n in nodes))
0075 failures += not check(
0076 "tails are merge_ntup + deriv_phys",
0077 {n.name for n in nodes if n.is_tail} == {"merge_ntup", "deriv_phys"},
0078 str({n.name for n in nodes if n.is_tail}),
0079 )
0080
0081 print("\n--- derived outputs (one workflow datum per output job parameter)")
0082 for name in ["evgen", "merge_hits", "deriv_phys"]:
0083 print(f" {name}: {sorted(by_name[name].outputs)}")
0084 failures += not check(
0085 "deriv_phys exposes both DAODs independently",
0086 set(by_name["deriv_phys"].outputs) == {"deriv_phys/DAOD_PHYS", "deriv_phys/DAOD_PHYSLITE"},
0087 str(set(by_name["deriv_phys"].outputs)),
0088 )
0089 failures += not check(
0090 "merge_hits honours the outputs override (HITS, not HITS_MRG)",
0091 set(by_name["merge_hits"].outputs) == {"merge_hits/HITS"},
0092 str(set(by_name["merge_hits"].outputs)),
0093 )
0094 failures += not check(
0095 "the 'number' job parameter is not mistaken for an output",
0096 all("DAOD" in k or k.endswith("PHYSLITE") for k in by_name["deriv_phys"].outputs),
0097 )
0098 failures += not check("no outDS alias is invented", not any(k.endswith("/outDS") for n in nodes for k in n.outputs))
0099
0100 print("\n--- derived inputs and dependency edges")
0101 failures += not check("evgen has no inputs (noInput head step)", by_name["evgen"].inputs == {})
0102 recon_sources = {v["source"] for v in by_name["recon"].inputs.values()}
0103 print(f" recon sources: {sorted(recon_sources)}")
0104
0105
0106
0107 failures += not check(
0108 "recon takes the upstream HITS and the external pileup input",
0109 recon_sources == {"merge_hits/HITS", "rdo_bkg"},
0110 str(recon_sources),
0111 )
0112 failures += not check(
0113 "pseudo_input (seq_number) never becomes an input",
0114 not any("seq_number" in json.dumps(n.inputs) for n in nodes),
0115 )
0116
0117 print("\n=== resolve_nodes ===")
0118 serial, tails, nodes = wnu.resolve_nodes(nodes, root_in, wfd.get("inputs", {}), 0, set(), None, log)
0119 by_name = {n.name: n for n in nodes}
0120 failures += not check("all 9 steps survive resolve_nodes", len(nodes) == 9, str(len(nodes)))
0121 failures += not check("member_id assigned to every step", all(n.member_id for n in nodes))
0122
0123
0124
0125 failures += not check(
0126 "is_tail preserved through resolve_nodes",
0127 {n.name for n in nodes if n.is_tail} == {"merge_ntup", "deriv_phys"},
0128 str({n.name for n in nodes if n.is_tail}),
0129 )
0130
0131 evgen_out = by_name["evgen"].outputs["evgen/EVNT"]["value"]
0132 print(f" evgen/EVNT -> {evgen_out}")
0133 failures += not check("author-supplied output name survives resolve_nodes", evgen_out.startswith("mc23_13p6TeV.526140"))
0134 failures += not check("no generated name was substituted in", "_001_evgen" not in evgen_out)
0135
0136
0137
0138 failures += not check("${TASKID} still outstanding", has_placeholder(evgen_out, TASKID_PLACEHOLDER))
0139 failures += not check("no workflow ID in the dataset name, per the ATLAS convention", not has_placeholder(evgen_out, WFID_PLACEHOLDER), evgen_out)
0140
0141 print("\n--- parent edges")
0142 id_to_name = {n.id: n.name for n in nodes}
0143 for name in ["merge_evnt", "recon", "deriv_phys"]:
0144 print(f" {name} parents: {sorted(id_to_name[p] for p in by_name[name].parents)}")
0145 failures += not check("deriv_phys branches off merge_aod", {id_to_name[p] for p in by_name["deriv_phys"].parents} == {"merge_aod"})
0146 failures += not check("recon depends only on merge_hits", {id_to_name[p] for p in by_name["recon"].parents} == {"merge_hits"})
0147
0148 print("\n--- input values resolved from the producing step's outputs")
0149 merge_evnt_in = [v.get("value") for v in by_name["merge_evnt"].inputs.values()]
0150 print(f" merge_evnt input value: {merge_evnt_in}")
0151 failures += not check("merge_evnt input resolved to evgen's output dataset", merge_evnt_in == [evgen_out], str(merge_evnt_in))
0152 rdo = [v.get("value") for v in by_name["recon"].inputs.values() if v["source"] == "rdo_bkg"]
0153 print(f" recon workflow-input value: {rdo}")
0154 failures += not check("workflow input reference resolved", rdo and rdo[0] == wfd["inputs"]["rdo_bkg"], str(rdo))
0155
0156 print("\n=== resolve_params / verify ===")
0157 id_map = wnu.get_node_id_map(nodes)
0158 for n in nodes:
0159 n.resolve_params(None, id_map)
0160 failures += not check("task_params built without a CLI template", all(n.task_params for n in nodes))
0161 for n in nodes:
0162 ok, msg = n.verify()
0163 if not ok:
0164 failures += not check(f"verify {n.name}", False, msg)
0165 failures += not check("all steps verify", all(n.verify()[0] for n in nodes))
0166
0167 print("\n--- verify rejects bad task params")
0168 for label, mutate in [
0169 ("missing taskName", lambda tp: tp.pop("taskName")),
0170 ("parentTaskName present", lambda tp: tp.update({"parentTaskName": "some.parent"})),
0171 ("${TASKID} in taskName", lambda tp: tp.update({"taskName": "x_tid${TASKID}"})),
0172 ("output without dataset", lambda tp: [j.pop("dataset") for j in tp["jobParameters"] if j.get("param_type") == "output"]),
0173 ("no output at all", lambda tp: tp.update({"jobParameters": [j for j in tp["jobParameters"] if j.get("param_type") != "output"]})),
0174 ]:
0175 node = wnu.Node(1, "task", None, True, "probe")
0176 node.task_params = json.loads(json.dumps(wfd["steps"]["evgen"]["task_params"]))
0177 mutate(node.task_params)
0178 ok, msg = node.verify()
0179 failures += not check(f"rejects {label}", not ok, "was accepted")
0180 if not ok:
0181 print(f" -> {msg}")
0182
0183 print("\n=== ${TASKID} substitution ===")
0184 resolved = substitute_placeholder(evgen_out, TASKID_PLACEHOLDER, 49900001)
0185 print(f" authored {evgen_out}")
0186 print(f" + TASKID {resolved}")
0187 failures += not check("resolved to the task id", not has_placeholder(resolved, TASKID_PLACEHOLDER) and "tid49900001_00" in resolved, resolved)
0188 failures += not check("matches the original production name shape", resolved.endswith("_tid49900001_00"), resolved)
0189 failures += not check("JEDI per-job templates untouched", "${SN}" in substitute_placeholder("log.${TASKID}._${SN}.tgz", TASKID_PLACEHOLDER, 7))
0190
0191 print("\n=== ${PARENT_TASKID} is only allowed where it can be resolved ===")
0192 base = json.loads(json.dumps(wfd["steps"]["simul"]["task_params"]))
0193 ok_params = {**base, "parent_tid": "${PARENT_TASKID}"}
0194 node = wnu.Node(1, "task", {}, True, "probe")
0195 node.task_params = ok_params
0196 good, reason = node.verify_task_params()
0197 failures += not check("accepted as the value of parent_tid", good, reason)
0198 bad_params = {**base, "taskType": "${PARENT_TASKID}"}
0199 node.task_params = bad_params
0200 good, reason = node.verify_task_params()
0201 failures += not check("rejected anywhere else", not good, reason)
0202 failures += not check("...naming the parameter it belongs in", "parent_tid" in reason, reason)
0203
0204 print("\n=== a joining step asking for ${PARENT_TASKID} is refused on submission ===")
0205
0206
0207 join_wfd: dict[str, Any] = {
0208 "name": "join_probe",
0209 "inputs": {},
0210 "outputs": {"out": {"from": "join/AOD"}},
0211 "steps": {},
0212 }
0213
0214 def task_step(name: str, inputs: list[str], output: str) -> dict[str, Any]:
0215 return {
0216 "type": "task",
0217 "task_params": {
0218 "taskName": f"probe.{name}",
0219 "vo": "atlas",
0220 "prodSourceLabel": "managed",
0221 "transPath": "Reco_tf.py",
0222 "log": {"type": "template", "param_type": "log", "value": "log.tgz", "dataset": f"probe.{name}.log"},
0223 "jobParameters": [{"type": "template", "param_type": "input", "value": "--inputFile=x", "dataset": ds} for ds in inputs]
0224 + [{"type": "template", "param_type": "output", "value": f"--output{output}File=x", "dataset": f"probe.{name}.{output}"}],
0225 },
0226 }
0227
0228 join_wfd["steps"]["left"] = task_step("left", [], "AOD")
0229 join_wfd["steps"]["right"] = task_step("right", [], "AOD")
0230 join_wfd["steps"]["join"] = task_step("join", ["{left/AOD}", "{right/AOD}"], "AOD")
0231 is_valid, errors = wnu.validate_workflow_description(join_wfd)
0232 failures += not check("valid while no step asks for a parent", is_valid, errors)
0233 join_wfd["steps"]["join"]["task_params"]["parent_tid"] = "${PARENT_TASKID}"
0234 is_valid, errors = wnu.validate_workflow_description(join_wfd)
0235 failures += not check("rejected once the joining step asks", not is_valid, errors)
0236 failures += not check("naming both feeding steps", errors and "left" in errors[0] and "right" in errors[0], errors)
0237 failures += not check("offering the step-name form, which an author can know", errors and "PARENT_TASKID:<step name>" in errors[0], errors)
0238
0239 join_wfd["steps"]["join"]["task_params"]["parent_tid"] = "${PARENT_TASKID:left}"
0240 is_valid, errors = wnu.validate_workflow_description(join_wfd)
0241 failures += not check("naming a feeding step is accepted", is_valid, errors)
0242 join_wfd["steps"]["join"]["task_params"]["parent_tid"] = "${PARENT_TASKID:elsewhere}"
0243 is_valid, errors = wnu.validate_workflow_description(join_wfd)
0244 failures += not check("naming a step that does not feed it is refused", not is_valid, errors)
0245 join_wfd["steps"]["join"]["task_params"]["parent_tid"] = "${PARENT_TASKID}"
0246
0247 join_wfd["steps"]["join"]["task_params"]["jobParameters"][1]["dataset"] = "some.external.dataset"
0248 is_valid, errors = wnu.validate_workflow_description(join_wfd)
0249 failures += not check("one step plus a literal dataset is unambiguous", is_valid, errors)
0250
0251 print("\n=== regression: prun steps are unaffected ===")
0252 prun_wfd: dict[str, Any] = {
0253 "name": "prun_probe",
0254 "inputs": {"sig": "some:dataset"},
0255 "outputs": {"out": {"from": "b/outDS", "output_types": ["aaa.root"]}},
0256 "steps": {
0257 "a": {"type": "prun", "inDS": "{sig}", "args": "--outputs x.root", "exec": "echo"},
0258 "b": {"type": "prun", "inDS": "a/outDS", "args": "--outputs aaa.root", "exec": "echo"},
0259 },
0260 }
0261 p_nodes, p_root_in = wnu.parse_workflow_data(prun_wfd, log)
0262 failures += not check("prun outputs still keyed on outDS", all(set(n.outputs) == {f"{n.name}/outDS"} for n in p_nodes))
0263 failures += not check("prun outputs start empty (no pre-set value)", all(v == {} for n in p_nodes for v in n.outputs.values()))
0264 _, _, p_nodes = wnu.resolve_nodes(p_nodes, p_root_in, prun_wfd["inputs"], 0, set(), "user.me.myOut", log)
0265 p_by = {n.name: n for n in p_nodes}
0266 gen = p_by["b"].outputs["b/outDS"]["value"]
0267 print(f" generated prun name: {gen}")
0268 failures += not check("prun name generation unchanged", gen == "user.me.myOut_002_b", gen)
0269
0270 print(f"\n{'ALL CHECKS PASSED' if not failures else f'{failures} CHECK(S) FAILED'}")
0271 return 1 if failures else 0
0272
0273
0274 if __name__ == "__main__":
0275 sys.exit(main())