Back to home page

EIC code displayed by LXR

 
 

    


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

0001 """
0002 Offline check of WorkflowInterface.get_step_relations.
0003 
0004 The engine is data-driven: a step starts because its inputs are good, never because a parent step
0005 finished, so it holds no step-to-step edge of its own. Consumers that model a chain as related
0006 tasks -- DEFT above all -- need those edges, and this derives them from what is recorded: each
0007 datum's source_step_id, and each step's input_data_dict.
0008 
0009 The shape checked here is the production chain PanDA ran as workflow 133, including the branch
0010 where deriv_prw and deriv_phys both consume merge_aod/AOD, and the external background dataset
0011 that must contribute no parent.
0012 
0013 Run from the repository root:  python3 pandaserver/workflow/examples/step_relations_test.py
0014 """
0015 
0016 import importlib.abc
0017 import importlib.machinery
0018 import json
0019 import os
0020 import sys
0021 import types
0022 import warnings
0023 
0024 # Unrelated modules in the import graph still carry unescaped regex literals; their SyntaxWarnings
0025 # say nothing about this check.
0026 warnings.filterwarnings("ignore", category=SyntaxWarning)
0027 
0028 REPO_ROOT = os.path.abspath(os.path.join(os.path.dirname(__file__), "..", "..", ".."))
0029 sys.path.insert(0, REPO_ROOT)
0030 
0031 AUTO_STUB_ROOTS = ("idds", "pandaclient", "ruamel", "requests", "rucio")
0032 
0033 
0034 class AutoStubFinder(importlib.abc.MetaPathFinder, importlib.abc.Loader):
0035     def find_spec(self, name, path=None, target=None):
0036         if name.split(".")[0] in AUTO_STUB_ROOTS:
0037             return importlib.machinery.ModuleSpec(name, self, is_package=True)
0038         return None
0039 
0040     def create_module(self, spec):
0041         module = types.ModuleType(spec.name)
0042         module.__path__ = []
0043         return module
0044 
0045     def exec_module(self, module):
0046         class Anything:
0047             def __init__(self, *a, **k):
0048                 pass
0049 
0050             def __call__(self, *a, **k):
0051                 return self
0052 
0053         module.__getattr__ = lambda name: Anything
0054 
0055 
0056 sys.meta_path.insert(0, AutoStubFinder())
0057 
0058 LOGGED: dict[str, list[str]] = {"warning": [], "debug": []}
0059 
0060 
0061 class Log:
0062     def __init__(self, *a, **k):
0063         pass
0064 
0065     def info(self, m):
0066         pass
0067 
0068     def debug(self, m):
0069         LOGGED["debug"].append(str(m))
0070 
0071     def warning(self, m):
0072         LOGGED["warning"].append(str(m))
0073 
0074     def error(self, m):
0075         LOGGED["warning"].append(str(m))
0076 
0077 
0078 class SpecBase:
0079     """Enough of pandacommon's SpecBase for the real spec classes to be instantiated offline"""
0080 
0081     def __init__(self):
0082         for attribute in getattr(self, "attributes", ()):
0083             setattr(self, attribute, None)
0084 
0085 
0086 def stub(name, **attrs):
0087     module = types.ModuleType(name)
0088     for key, value in attrs.items():
0089         setattr(module, key, value)
0090     sys.modules[name] = module
0091     return module
0092 
0093 
0094 stub("pandacommon")
0095 stub("pandacommon.pandautils").__path__ = []
0096 stub("pandacommon.pandautils.base", SpecBase=SpecBase)
0097 stub("pandacommon.pandautils.PandaUtils", naive_utcnow=lambda: None, get_sql_IN_bind_variables=lambda *a, **k: (None, None))
0098 stub("pandacommon.pandautils.thread_utils", GenericThread=object)
0099 stub("pandacommon.pandalogger").__path__ = []
0100 stub("pandacommon.pandalogger.LogWrapper", LogWrapper=Log)
0101 stub("pandacommon.pandalogger.PandaLogger", PandaLogger=lambda: types.SimpleNamespace(getLogger=lambda n: None))
0102 stub("pandaserver.config", panda_config=types.SimpleNamespace(schemaJEDI="ATLAS_PANDA", schemaDEFT="ATLAS_DEFT"))
0103 
0104 from pandaserver.workflow import workflow_core  # noqa: E402
0105 from pandaserver.workflow.workflow_base import (  # noqa: E402
0106     WFDataSpec,
0107     WFDataType,
0108     WFStepSpec,
0109     WFStepStatus,
0110     WFStepType,
0111 )
0112 
0113 # The production chain of workflow 133: step name -> (inputs it consumes, outputs it produces).
0114 # recon also reads rdo_bkg, which the workflow does not produce.
0115 CHAIN = [
0116     ("evgen", [], ["evgen/EVNT"]),
0117     ("merge_evnt", ["evgen/EVNT"], ["merge_evnt/EVNT"]),
0118     ("simul", ["merge_evnt/EVNT"], ["simul/HITS"]),
0119     ("merge_hits", ["simul/HITS"], ["merge_hits/HITS"]),
0120     ("recon", ["merge_hits/HITS", "rdo_bkg"], ["recon/AOD"]),
0121     ("merge_aod", ["recon/AOD"], ["merge_aod/AOD"]),
0122     ("deriv_prw", ["merge_aod/AOD"], ["deriv_prw/NTUP_PILEUP"]),
0123     ("merge_ntup", ["deriv_prw/NTUP_PILEUP"], ["merge_ntup/NTUP_PILEUP"]),
0124     ("deriv_phys", ["merge_aod/AOD"], ["deriv_phys/DAOD_PHYS", "deriv_phys/DAOD_PHYSLITE"]),
0125 ]
0126 
0127 TASK_IDS = {
0128     "evgen": "52382519",
0129     "merge_evnt": "52382898",
0130     "simul": "52383720",
0131     "merge_hits": "52397622",
0132     "recon": "52401216",
0133     "merge_aod": "52416493",
0134     "deriv_prw": "52421202",
0135     "merge_ntup": "52421461",
0136     "deriv_phys": "52421203",
0137 }
0138 
0139 
0140 def make_step(step_id, name, inputs, flavor="panda_task", status=WFStepStatus.done, target_id=None, workflow_id=133):
0141     step_spec = WFStepSpec()
0142     step_spec.step_id = step_id
0143     step_spec.name = name
0144     step_spec.workflow_id = workflow_id
0145     step_spec.type = WFStepType.sub_workflow if flavor != "panda_task" else WFStepType.ordinary
0146     step_spec.flavor = flavor
0147     step_spec.status = status
0148     step_spec.target_id = target_id
0149     step_spec.definition_json = json.dumps({"input_data_dict": {name: {} for name in inputs}})
0150     return step_spec
0151 
0152 
0153 def make_data(name, source_step_id, data_type=WFDataType.mid, workflow_id=133):
0154     data_spec = WFDataSpec()
0155     data_spec.name = name
0156     data_spec.workflow_id = workflow_id
0157     data_spec.source_step_id = source_step_id
0158     data_spec.type = data_type
0159     return data_spec
0160 
0161 
0162 def build_chain():
0163     """The workflow-133 steps and data, with step ids 1..9 in chain order"""
0164     step_specs = []
0165     data_specs = [make_data("rdo_bkg", None, WFDataType.input)]
0166     for index, (name, inputs, outputs) in enumerate(CHAIN, start=1):
0167         step_specs.append(make_step(index, name, inputs, target_id=TASK_IDS[name]))
0168         for output_name in outputs:
0169             data_specs.append(make_data(output_name, index))
0170     return step_specs, data_specs
0171 
0172 
0173 class FakeTaskBuffer:
0174     """Serves one or more workflows, so the nested cases can be built"""
0175 
0176     def __init__(self, step_specs, data_specs, by_workflow=None):
0177         self.by_workflow = dict(by_workflow or {})
0178         self.by_workflow.setdefault(133, (step_specs, data_specs))
0179 
0180     def get_steps_of_workflow(self, workflow_id, status_filter_list=None, status_exclusion_list=None):
0181         return list(self.by_workflow.get(workflow_id, ([], []))[0])
0182 
0183     def get_data_of_workflow(self, workflow_id, status_filter_list=None, status_exclusion_list=None, type_filter_list=None):
0184         return list(self.by_workflow.get(workflow_id, ([], []))[1])
0185 
0186     def get_steps_by_target_id(self, target_id, flavor_filter_list=None):
0187         found = []
0188         for step_specs, _ in self.by_workflow.values():
0189             for step_spec in step_specs:
0190                 if step_spec.target_id == target_id and (not flavor_filter_list or step_spec.flavor in flavor_filter_list):
0191                     found.append(step_spec)
0192         return found
0193 
0194 
0195 def make_interface(step_specs, data_specs, by_workflow=None):
0196     """A WorkflowInterface with no message broker or DDM behind it"""
0197     interface = workflow_core.WorkflowInterface.__new__(workflow_core.WorkflowInterface)
0198     interface.tbif = FakeTaskBuffer(step_specs, data_specs, by_workflow)
0199     interface.full_pid = "test-0-0"
0200     interface.plugin_map = {}
0201     interface.mb_proxy = None
0202     return interface
0203 
0204 
0205 def check(label, condition, detail=""):
0206     print(f"  {'PASS' if condition else 'FAIL'}  {label}{'  ' + str(detail) if detail and not condition else ''}")
0207     return condition
0208 
0209 
0210 def by_name(result, name):
0211     """The task entry with that step name, or None"""
0212     return next((task for task in result["tasks"] if task["name"] == name), None)
0213 
0214 
0215 def task_by_key_name(result, key):
0216     """The step name behind a task key, so expectations read as names rather than ids"""
0217     return next(task["name"] for task in result["tasks"] if task["key"] == key)
0218 
0219 
0220 def parents_by_name(result):
0221     by_id = {step["step_id"]: step["name"] for step in result["steps"]}
0222     return {step["name"]: sorted(by_id[parent] for parent in step["parent_step_ids"]) for step in result["steps"]}
0223 
0224 
0225 def main():
0226     failures = 0
0227 
0228     print("\n=== the production chain of workflow 133 ===")
0229     step_specs, data_specs = build_chain()
0230     result = make_interface(step_specs, data_specs).get_step_relations(133)
0231     failures += not check("a result is returned", result is not None)
0232     if result is None:
0233         print("\n1 CHECK(S) FAILED")
0234         return 1
0235     failures += not check("the workflow id is echoed back", result["workflow_id"] == 133, result["workflow_id"])
0236     failures += not check("every step is reported", len(result["steps"]) == 9, len(result["steps"]))
0237     failures += not check("steps come back in step_id order", [s["step_id"] for s in result["steps"]] == list(range(1, 10)))
0238 
0239     relations = parents_by_name(result)
0240     expected = {
0241         "evgen": [],
0242         "merge_evnt": ["evgen"],
0243         "simul": ["merge_evnt"],
0244         "merge_hits": ["simul"],
0245         "recon": ["merge_hits"],
0246         "merge_aod": ["recon"],
0247         "deriv_prw": ["merge_aod"],
0248         "merge_ntup": ["deriv_prw"],
0249         "deriv_phys": ["merge_aod"],
0250     }
0251     for name, expected_parents in expected.items():
0252         failures += not check(f"{name} <- {expected_parents or 'nothing'}", relations[name] == expected_parents, relations[name])
0253 
0254     print("\n=== what the edges do and do not come from ===")
0255     failures += not check("an entry step has no parent rather than itself", relations["evgen"] == [])
0256     failures += not check("the external rdo_bkg contributes no parent to recon", relations["recon"] == ["merge_hits"], relations["recon"])
0257     branch = [name for name, parents in relations.items() if parents == ["merge_aod"]]
0258     failures += not check("merge_aod is the parent of two steps", sorted(branch) == ["deriv_phys", "deriv_prw"], branch)
0259 
0260     print("\n=== each step carries what a task view needs ===")
0261     evgen = next(s for s in result["steps"] if s["name"] == "evgen")
0262     failures += not check("the step's target is reported", evgen["target_id"] == "52382519", evgen["target_id"])
0263     failures += not check("so is its flavor", evgen["flavor"] == "panda_task", evgen["flavor"])
0264     failures += not check("and its status", evgen["status"] == WFStepStatus.done, evgen["status"])
0265 
0266     print("\n=== a step with several parents ===")
0267     # A join: one step consuming the outputs of two independent steps.
0268     step_specs = [
0269         make_step(1, "left", []),
0270         make_step(2, "right", []),
0271         make_step(3, "join", ["left/OUT", "right/OUT"]),
0272     ]
0273     data_specs = [make_data("left/OUT", 1), make_data("right/OUT", 2)]
0274     result = make_interface(step_specs, data_specs).get_step_relations(133)
0275     join = next(s for s in result["steps"] if s["name"] == "join")
0276     failures += not check("both parents are reported", join["parent_step_ids"] == [1, 2], join["parent_step_ids"])
0277 
0278     print("\n=== a running workflow ===")
0279     # simul is running and has a task; everything after it has neither started nor got a task.
0280     step_specs, data_specs = build_chain()
0281     for step_spec in step_specs:
0282         if step_spec.step_id > 3:
0283             step_spec.status = WFStepStatus.registered
0284             step_spec.target_id = None
0285         elif step_spec.step_id == 3:
0286             step_spec.status = WFStepStatus.running
0287     # only the data produced so far has a producer recorded
0288     for data_spec in data_specs:
0289         if data_spec.source_step_id is not None and data_spec.source_step_id > 3:
0290             data_spec.source_step_id = None
0291     result = make_interface(step_specs, data_specs).get_step_relations(133)
0292     relations = parents_by_name(result)
0293     failures += not check("the whole graph is reported, not only the started part", len(result["steps"]) == 9, len(result["steps"]))
0294     failures += not check("edges among started steps are there", relations["simul"] == ["merge_evnt"], relations["simul"])
0295     # simul is running, so its output is already bound to it: the edge into the next step exists
0296     # before that output is complete, which is what makes the graph usable mid-flight.
0297     failures += not check("a running step is already a parent", relations["merge_hits"] == ["simul"], relations["merge_hits"])
0298     failures += not check("a step whose producer has not started has no parent yet", relations["recon"] == [], relations["recon"])
0299     not_started = next(s for s in result["steps"] if s["name"] == "recon")
0300     failures += not check("a step that has not started reports no target", not_started["target_id"] is None, not_started["target_id"])
0301     failures += not check("...but is still in the answer with its status", not_started["status"] == WFStepStatus.registered, not_started["status"])
0302 
0303     print("\n=== a nested sub-workflow step ===")
0304     step_specs = [
0305         make_step(1, "prepare", [], target_id="52400001"),
0306         make_step(2, "child", ["prepare/OUT"], flavor="sub_workflow", target_id="277"),
0307     ]
0308     data_specs = [make_data("prepare/OUT", 1)]
0309     result = make_interface(step_specs, data_specs).get_step_relations(133)
0310     child = next(s for s in result["steps"] if s["name"] == "child")
0311     failures += not check("the sub-workflow step is reported as itself", child["parent_step_ids"] == [1], child["parent_step_ids"])
0312     failures += not check("its target is the child workflow id", child["target_id"] == "277", child["target_id"])
0313     failures += not check("its flavor says a caller can recurse", child["flavor"] == "sub_workflow", child["flavor"])
0314 
0315     print("\n=== data the engine cannot relate ===")
0316     del LOGGED["warning"][:]
0317     step_specs = [make_step(1, "producer", []), make_step(2, "consumer", ["producer/OUT", "never/heard/of"])]
0318     data_specs = [make_data("producer/OUT", 1)]
0319     result = make_interface(step_specs, data_specs).get_step_relations(133)
0320     consumer = next(s for s in result["steps"] if s["name"] == "consumer")
0321     failures += not check("an unknown input is skipped, not fatal", consumer["parent_step_ids"] == [1], consumer["parent_step_ids"])
0322 
0323     # A datum naming a producer outside this workflow must not become a dangling edge.
0324     del LOGGED["warning"][:]
0325     data_specs = [make_data("producer/OUT", 1), make_data("never/heard/of", 99)]
0326     result = make_interface(step_specs, data_specs).get_step_relations(133)
0327     consumer = next(s for s in result["steps"] if s["name"] == "consumer")
0328     failures += not check("a producer outside the workflow is not reported as a parent", consumer["parent_step_ids"] == [1], consumer["parent_step_ids"])
0329     failures += not check("...and is warned about", any("not in this workflow" in m for m in LOGGED["warning"]), LOGGED["warning"])
0330 
0331     # A step recorded as producing its own input would make a consumer walking parents loop.
0332     del LOGGED["warning"][:]
0333     step_specs = [make_step(1, "selfish", ["selfish/OUT"])]
0334     data_specs = [make_data("selfish/OUT", 1)]
0335     result = make_interface(step_specs, data_specs).get_step_relations(133)
0336     failures += not check("a step is never its own parent", result["steps"][0]["parent_step_ids"] == [], result["steps"][0]["parent_step_ids"])
0337     failures += not check("...and is warned about", any("its own input" in m for m in LOGGED["warning"]), LOGGED["warning"])
0338 
0339     print("\n=== a workflow with no steps ===")
0340     failures += not check("reports nothing rather than an empty graph", make_interface([], []).get_step_relations(999) is None)
0341 
0342     print("\n=== the answer is JSON-serializable ===")
0343     step_specs, data_specs = build_chain()
0344     result = make_interface(step_specs, data_specs).get_step_relations(133)
0345     try:
0346         json.dumps(result)
0347         failures += not check("json.dumps accepts it as it stands", True)
0348     except Exception as exc:
0349         failures += not check("json.dumps accepts it as it stands", False, exc)
0350 
0351     print("\n=== the task view of workflow 133 ===")
0352     step_specs, data_specs = build_chain()
0353     result = make_interface(step_specs, data_specs).get_task_relations(133)
0354     failures += not check("a result is returned", result is not None)
0355     task_by_key = {task["key"]: task for task in result["tasks"]}
0356     failures += not check("every step became a task", len(result["tasks"]) == 9, len(result["tasks"]))
0357     failures += not check("task ids come from the step targets", task_by_key["133:5"]["task_id"] == 52401216, task_by_key["133:5"]["task_id"])
0358     task_parents = {task["name"]: sorted(task_by_key[key]["name"] for key in task["parents"]) for task in result["tasks"]}
0359     failures += not check("the chain is unchanged by the projection", task_parents == expected, task_parents)
0360 
0361     print("\n=== a step that is not a task is collapsed, not dropped ===")
0362     # task A -> a step with some other target -> task B must report A as B's parent.
0363     step_specs = [
0364         make_step(1, "A", [], target_id="52400001"),
0365         make_step(2, "middle", ["A/OUT"], flavor="future_thing", target_id="whatever"),
0366         make_step(3, "B", ["middle/OUT"], target_id="52400003"),
0367     ]
0368     data_specs = [make_data("A/OUT", 1), make_data("middle/OUT", 2)]
0369     result = make_interface(step_specs, data_specs).get_task_relations(133)
0370     names = sorted(task["name"] for task in result["tasks"])
0371     failures += not check("the non-task step is not a node", names == ["A", "B"], names)
0372     task_b = next(task for task in result["tasks"] if task["name"] == "B")
0373     failures += not check("the relation passes through it", [task_by_key_name(result, k) for k in task_b["parents"]] == ["A"], task_b["parents"])
0374 
0375     print("\n=== a nested workflow is replaced by the tasks inside it ===")
0376     # outer: prepare -> child(workflow 277) -> after
0377     outer_steps = [
0378         make_step(1, "prepare", [], target_id="52400001"),
0379         make_step(2, "child", ["prepare/OUT"], flavor="sub_workflow", target_id="277"),
0380         make_step(3, "after", ["child/OUT"], target_id="52400009"),
0381     ]
0382     outer_data = [make_data("prepare/OUT", 1), make_data("child/OUT", 2)]
0383     inner_steps = [
0384         make_step(1, "inner_first", [], target_id="52400005", workflow_id=277),
0385         make_step(2, "inner_last", ["inner_first/OUT"], target_id="52400006", workflow_id=277),
0386     ]
0387     inner_data = [make_data("inner_first/OUT", 1, workflow_id=277)]
0388     interface = make_interface(outer_steps, outer_data, by_workflow={133: (outer_steps, outer_data), 277: (inner_steps, inner_data)})
0389     result = interface.get_task_relations(133)
0390     names = sorted(task["name"] for task in result["tasks"])
0391     failures += not check("the inner tasks appear", names == ["after", "inner_first", "inner_last", "prepare"], names)
0392     failures += not check("the sub-workflow step itself is not a node", "child" not in names)
0393     failures += not check(
0394         "the child's entry task takes the outer parent", [task_by_key_name(result, k) for k in by_name(result, "inner_first")["parents"]] == ["prepare"]
0395     )
0396     failures += not check(
0397         "the outer consumer takes the child's tail task", [task_by_key_name(result, k) for k in by_name(result, "after")["parents"]] == ["inner_last"]
0398     )
0399     failures += not check("the inner chain is kept", [task_by_key_name(result, k) for k in by_name(result, "inner_last")["parents"]] == ["inner_first"])
0400     failures += not check("keys are unique across the workflows walked", len({task["key"] for task in result["tasks"]}) == 4)
0401 
0402     print("\n=== a step with no task yet is a placeholder ===")
0403     step_specs = [make_step(1, "done_one", [], target_id="52400001"), make_step(2, "not_yet", ["done_one/OUT"], status=WFStepStatus.registered)]
0404     data_specs = [make_data("done_one/OUT", 1)]
0405     result = make_interface(step_specs, data_specs).get_task_relations(133)
0406     pending = by_name(result, "not_yet")
0407     failures += not check("it is still reported", pending is not None)
0408     failures += not check("with no task id", pending["task_id"] is None, pending["task_id"])
0409     failures += not check("and its parent, so the shape is visible", [task_by_key_name(result, k) for k in pending["parents"]] == ["done_one"])
0410 
0411     print("\n=== steps are resolved parents-first whatever their ids ===")
0412     # ids deliberately against the flow: step 3 feeds step 1.
0413     step_specs = [make_step(3, "first", [], target_id="52400001"), make_step(1, "second", ["first/OUT"], target_id="52400002")]
0414     data_specs = [make_data("first/OUT", 3)]
0415     result = make_interface(step_specs, data_specs).get_task_relations(133)
0416     failures += not check("the edge survives the ordering", [task_by_key_name(result, k) for k in by_name(result, "second")["parents"]] == ["first"])
0417 
0418     print("\n=== entering by task id ===")
0419     step_specs, data_specs = build_chain()
0420     interface = make_interface(step_specs, data_specs)
0421     result = interface.get_task_relations_of_task(52401216)
0422     failures += not check("the workflow of that task is reported", result is not None and result["workflow_id"] == 133)
0423     failures += not check("the whole chain comes back", len(result["tasks"]) == 9, len(result["tasks"]))
0424     failures += not check("the task asked about is named", result["asked_for"] == "133:5", result["asked_for"])
0425     failures += not check("a task no workflow runs reports nothing", interface.get_task_relations_of_task(99999999) is None)
0426 
0427     print(f"\n{'ALL CHECKS PASSED' if not failures else f'{failures} CHECK(S) FAILED'}")
0428     return 1 if failures else 0
0429 
0430 
0431 if __name__ == "__main__":
0432     sys.exit(main())