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
0025
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
0105 from pandaserver.workflow.workflow_base import (
0106 WFDataSpec,
0107 WFDataType,
0108 WFStepSpec,
0109 WFStepStatus,
0110 WFStepType,
0111 )
0112
0113
0114
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
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
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
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
0296
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
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
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
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
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
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())