Back to home page

EIC code displayed by LXR

 
 

    


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 # ---- stub the dependencies which are absent in a bare checkout ----------------------------
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  # noqa: E402
0037 from pandaserver.workflow.workflow_base import (  # noqa: E402
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     # The parent-resolution pass strips the braces from a workflow-input reference so that the
0105     # source matches the root_inputs key, while a "step/output" reference is left as it is and
0106     # becomes a parent edge. Both forms are expected here in their post-resolution shape.
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     # resolve_nodes returns every leaf in its tail list (pre-existing behaviour, also true for a
0123     # prun-only workflow). The workflow tails that matter are the is_tail flags, which is what
0124     # workflow_parser keys root_outputs off, so assert those survive instead.
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     # Dataset names follow the standard ATLAS convention, which carries no workflow ID, so only
0136     # ${TASKID} is expected here. ${WFID} remains supported and is covered by
0137     # placeholder_resolution_test.py.
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     # parent_tid holds one task, so a step fed by two steps has to say which. Caught here rather
0206     # than when the workflow is running and the step simply fails to start.
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     # naming which feeding step is meant resolves it
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     # a step fed by one step plus an external input is fine
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())