Back to home page

EIC code displayed by LXR

 
 

    


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

0001 """
0002 Offline check of when a running workflow is allowed to become done.
0003 
0004 A step reaches done as soon as its task does, but the task's output dataset is closed in DDM
0005 slightly later, so the data pass that moves an output from generating_suffice to done_generated
0006 necessarily runs a cycle behind the step transition. Finishing the workflow on the step transition
0007 alone leaves that output non-terminal for good, since a done workflow is no longer in
0008 active_statuses and is never processed again. This was seen for real: workflow 133 finished with
0009 merge_ntup/NTUP_PILEUP frozen in generating_suffice.
0010 
0011 So the workflow waits for its outputs, and the wait is bounded, so that an output whose dataset is
0012 never closed cannot keep the workflow running forever.
0013 
0014 Run from the repository root:  python3 pandaserver/workflow/examples/workflow_done_transition_test.py
0015 """
0016 
0017 import importlib.abc
0018 import importlib.machinery
0019 import os
0020 import sys
0021 import types
0022 import warnings
0023 from datetime import datetime, timedelta
0024 
0025 # Unrelated modules in the import graph still carry unescaped regex literals; their SyntaxWarnings
0026 # say nothing about this check.
0027 warnings.filterwarnings("ignore", category=SyntaxWarning)
0028 
0029 REPO_ROOT = os.path.abspath(os.path.join(os.path.dirname(__file__), "..", "..", ".."))
0030 sys.path.insert(0, REPO_ROOT)
0031 
0032 AUTO_STUB_ROOTS = ("idds", "pandaclient", "ruamel", "requests", "rucio")
0033 
0034 
0035 class AutoStubFinder(importlib.abc.MetaPathFinder, importlib.abc.Loader):
0036     def find_spec(self, name, path=None, target=None):
0037         if name.split(".")[0] in AUTO_STUB_ROOTS:
0038             return importlib.machinery.ModuleSpec(name, self, is_package=True)
0039         return None
0040 
0041     def create_module(self, spec):
0042         module = types.ModuleType(spec.name)
0043         module.__path__ = []
0044         return module
0045 
0046     def exec_module(self, module):
0047         class Anything:
0048             def __init__(self, *a, **k):
0049                 pass
0050 
0051             def __call__(self, *a, **k):
0052                 return self
0053 
0054         module.__getattr__ = lambda name: Anything
0055 
0056 
0057 sys.meta_path.insert(0, AutoStubFinder())
0058 
0059 
0060 def stub(name, **attrs):
0061     module = types.ModuleType(name)
0062     for key, value in attrs.items():
0063         setattr(module, key, value)
0064     sys.modules[name] = module
0065     return module
0066 
0067 
0068 class Log:
0069     """Collects what the code under test logs, so the warnings it emits can be asserted on"""
0070 
0071     def __init__(self, *a, **k):
0072         self.messages = []
0073 
0074     def info(self, m):
0075         self.messages.append(("info", str(m)))
0076 
0077     def debug(self, m):
0078         self.messages.append(("debug", str(m)))
0079 
0080     def warning(self, m):
0081         self.messages.append(("warning", str(m)))
0082 
0083     def error(self, m):
0084         self.messages.append(("error", str(m)))
0085 
0086 
0087 class SpecBase:
0088     """Enough of pandacommon's SpecBase for the real spec classes to be instantiated offline"""
0089 
0090     def __init__(self):
0091         for attribute in getattr(self, "attributes", ()):
0092             setattr(self, attribute, None)
0093 
0094 
0095 stub("pandacommon")
0096 stub("pandacommon.pandautils").__path__ = []
0097 stub("pandacommon.pandautils.base", SpecBase=SpecBase)
0098 stub("pandacommon.pandautils.PandaUtils", naive_utcnow=datetime.utcnow, get_sql_IN_bind_variables=lambda *a, **k: (None, None))
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     WFDataStatus,
0107     WFDataTargetCheckStatus,
0108     WFDataType,
0109     WFStepStatus,
0110     WorkflowSpec,
0111     WorkflowStatus,
0112 )
0113 
0114 NOW = datetime(2026, 9, 10, 12, 0, 0)
0115 
0116 # process_workflow_running stamps its own "now" from naive_utcnow, which the bounded-wait case has
0117 # to be able to place relative to the recorded wait, so the clock is pinned here.
0118 workflow_core.naive_utcnow = lambda: NOW
0119 
0120 
0121 class FakeData:
0122     def __init__(self, name, status, data_type=WFDataType.output, data_id=1, flavor="ddm_collection"):
0123         self.name = name
0124         self.status = status
0125         self.type = data_type
0126         self.data_id = data_id
0127         self.workflow_id = 133
0128         self.flavor = flavor
0129         self.end_time = None
0130         self.check_time = None
0131 
0132 
0133 class FakeStep:
0134     def __init__(self, name, status=WFStepStatus.done):
0135         self.name = name
0136         self.status = status
0137         self.flavor = "panda_task"
0138         self.target_id = "1"
0139         self.member_id = 1
0140 
0141 
0142 class FakeTaskBuffer:
0143     def __init__(self, data_specs, step_specs):
0144         self.data_specs = data_specs
0145         self.step_specs = step_specs
0146         self.updated_workflows = []
0147         self.updated_data = []
0148 
0149     def get_data_of_workflow(self, workflow_id, status_exclusion_list=None, type_filter_list=None):
0150         return list(self.data_specs)
0151 
0152     def get_steps_of_workflow(self, workflow_id, status_filter_list=None):
0153         return list(self.step_specs)
0154 
0155     def update_workflow(self, workflow_spec):
0156         self.updated_workflows.append((workflow_spec.status, workflow_spec.parameters))
0157         return True
0158 
0159     def update_workflow_data(self, data_spec):
0160         self.updated_data.append((data_spec.name, data_spec.status))
0161         return True
0162 
0163 
0164 def make_workflow(steps_final_time=None):
0165     workflow_spec = WorkflowSpec()
0166     workflow_spec.workflow_id = 133
0167     workflow_spec.status = WorkflowStatus.running
0168     if steps_final_time is not None:
0169         workflow_spec.set_parameter(workflow_core.STEPS_FINAL_TIME_PARAM, steps_final_time)
0170     return workflow_spec
0171 
0172 
0173 class FakeDataHandler:
0174     """Reports whatever DDM state a case wants, so the waiting transitions can be driven"""
0175 
0176     def __init__(self, check_status):
0177         self.check_status = check_status
0178 
0179     def check_target(self, data_spec, **kwargs):
0180         result = types.SimpleNamespace(success=True, check_status=self.check_status, message="", metadata={})
0181         return result
0182 
0183 
0184 class StubbedInterface(workflow_core.WorkflowInterface):
0185     """The real workflow transitions, with the step and data passes replaced
0186 
0187     Subclassed rather than monkey-patched so the stubs are checked against the methods they stand
0188     in for. __init__ is bypassed because the real one opens a message broker and a DDM client.
0189     """
0190 
0191     def __init__(self, data_specs, step_specs, all_steps_final=True, check_status=None):
0192         self.tbif = FakeTaskBuffer(data_specs, step_specs)
0193         self.full_pid = "test-0-0"
0194         self.plugin_map = {}
0195         self.mb_proxy = None
0196         self._all_steps_final = all_steps_final
0197         self._data_handler = FakeDataHandler(check_status if check_status is not None else WFDataTargetCheckStatus.suffice)
0198 
0199     def get_plugin(self, plugin_type, flavor):
0200         return self._data_handler
0201 
0202     # The data pass is what would advance an output on a later cycle; here it changes nothing, so
0203     # each case controls the output statuses directly.
0204     def process_datas(self, data_specs, by="dog"):
0205         return {"n_processed": len(data_specs), "processed": {}, "changed": {}}
0206 
0207     def process_steps(self, step_specs, data_spec_map=None, by="dog"):
0208         status = WFStepStatus.done if self._all_steps_final else WFStepStatus.running
0209         return {"n_processed": len(step_specs), "processed": {status: len(step_specs)}, "changed": {}}
0210 
0211 
0212 def make_interface(data_specs, step_specs, all_steps_final=True, check_status=None):
0213     return StubbedInterface(data_specs, step_specs, all_steps_final, check_status)
0214 
0215 
0216 def check(label, condition, detail=""):
0217     print(f"  {'PASS' if condition else 'FAIL'}  {label}{'  ' + str(detail) if detail and not condition else ''}")
0218     return condition
0219 
0220 
0221 def main():
0222     failures = 0
0223     interface = make_interface([], [])
0224 
0225     print("\n=== are_all_outputs_good is three-valued ===")
0226     failures += not check("no output at all -> None", interface.are_all_outputs_good({}) is None)
0227     good = {"a": FakeData("a", WFDataStatus.done_generated), "b": FakeData("b", WFDataStatus.done_waited)}
0228     failures += not check("every output terminal -> True", interface.are_all_outputs_good(good) is True)
0229     mixed = {"a": FakeData("a", WFDataStatus.done_generated), "b": FakeData("b", WFDataStatus.generating_suffice)}
0230     failures += not check("one output still generating -> False", interface.are_all_outputs_good(mixed) is False)
0231 
0232     print("\n=== record_steps_final_time measures the wait ===")
0233     workflow_spec = make_workflow()
0234     failures += not check("first cycle reports 0", interface.record_steps_final_time(workflow_spec, NOW) == 0)
0235     failures += not check(
0236         "first cycle records the time",
0237         workflow_spec.get_parameter(workflow_core.STEPS_FINAL_TIME_PARAM) == NOW.isoformat(),
0238         workflow_spec.parameters,
0239     )
0240     failures += not check(
0241         "a later cycle reports the elapsed seconds",
0242         interface.record_steps_final_time(workflow_spec, NOW + timedelta(seconds=125)) == 125,
0243     )
0244     broken = make_workflow(steps_final_time="not-a-timestamp")
0245     failures += not check("an unparsable value restarts the wait instead of hanging", interface.record_steps_final_time(broken, NOW) == 0)
0246     failures += not check("...and is replaced", broken.get_parameter(workflow_core.STEPS_FINAL_TIME_PARAM) == NOW.isoformat())
0247 
0248     print("\n=== a workflow whose steps are all final waits for its outputs ===")
0249     # This is workflow 133: every step done, the last output still generating_suffice because its
0250     # dataset is closed in DDM only after the task reached done.
0251     lagging = FakeData("merge_ntup/NTUP_PILEUP", WFDataStatus.generating_suffice)
0252     data_specs = [FakeData("deriv_phys/DAOD_PHYS", WFDataStatus.done_generated), lagging]
0253     interface = make_interface(data_specs, [FakeStep("merge_ntup"), FakeStep("deriv_phys")])
0254     workflow_spec = make_workflow()
0255     result = interface.process_workflow_running(workflow_spec)
0256     failures += not check("stays running", workflow_spec.status == WorkflowStatus.running, workflow_spec.status)
0257     failures += not check("reports success rather than an error", result.success is True, result.message)
0258     failures += not check("does not announce a new status", result.new_status is None, result.new_status)
0259     failures += not check("the lagging output is untouched", lagging.status == WFDataStatus.generating_suffice, lagging.status)
0260     failures += not check("the wait is recorded on the workflow", workflow_spec.get_parameter(workflow_core.STEPS_FINAL_TIME_PARAM) is not None)
0261     failures += not check("the recorded wait is persisted", interface.tbif.updated_workflows and interface.tbif.updated_workflows[-1][1])
0262 
0263     print("\n=== the output catching up finishes the workflow ===")
0264     data_specs = [FakeData("deriv_phys/DAOD_PHYS", WFDataStatus.done_generated), FakeData("merge_ntup/NTUP_PILEUP", WFDataStatus.done_generated)]
0265     interface = make_interface(data_specs, [FakeStep("merge_ntup"), FakeStep("deriv_phys")])
0266     workflow_spec = make_workflow(steps_final_time=NOW.isoformat())
0267     result = interface.process_workflow_running(workflow_spec)
0268     failures += not check("becomes done", workflow_spec.status == WorkflowStatus.done, workflow_spec.status)
0269     failures += not check("nothing had to be settled", not interface.tbif.updated_data, interface.tbif.updated_data)
0270 
0271     print("\n=== the wait is bounded ===")
0272     stuck = FakeData("merge_ntup/NTUP_PILEUP", WFDataStatus.generating_suffice)
0273     interface = make_interface([stuck], [FakeStep("merge_ntup")])
0274     expired = (NOW - timedelta(seconds=workflow_core.OUTPUT_SETTLE_GRACE_SEC + 1)).isoformat()
0275     workflow_spec = make_workflow(steps_final_time=expired)
0276     result = interface.process_workflow_running(workflow_spec)
0277     failures += not check("becomes done once the grace period is over", workflow_spec.status == WorkflowStatus.done, workflow_spec.status)
0278     failures += not check("the outstanding output is settled", stuck.status == WFDataStatus.done_generated, stuck.status)
0279     failures += not check("the settled output is persisted", ("merge_ntup/NTUP_PILEUP", WFDataStatus.done_generated) in interface.tbif.updated_data)
0280 
0281     print("\n=== settle_pending_outputs only settles what it can ===")
0282     log = Log()
0283     generating = FakeData("out/generating", WFDataStatus.generating_insuffi)
0284     waiting = FakeData("out/waiting", WFDataStatus.waiting_suffice)
0285     already = FakeData("out/already", WFDataStatus.done_waited)
0286     never_bound = FakeData("out/never_bound", WFDataStatus.checking)
0287     interface = make_interface([], [])
0288     interface.settle_pending_outputs(log, {d.name: d for d in [generating, waiting, already, never_bound]}, NOW)
0289     failures += not check("generating -> done_generated", generating.status == WFDataStatus.done_generated, generating.status)
0290     failures += not check("waiting -> done_waited", waiting.status == WFDataStatus.done_waited, waiting.status)
0291     failures += not check("an already terminal output is left alone", already.status == WFDataStatus.done_waited)
0292     failures += not check("one that was never bound is not given a final status", never_bound.status == WFDataStatus.checking, never_bound.status)
0293     failures += not check("every settled output is reported as a warning", sum(1 for level, _ in log.messages if level == "warning") == 3, log.messages)
0294 
0295     print("\n=== a workflow with no output data still finishes ===")
0296     interface = make_interface([], [FakeStep("only_step")])
0297     workflow_spec = make_workflow()
0298     result = interface.process_workflow_running(workflow_spec)
0299     failures += not check("becomes done rather than waiting for outputs it does not have", workflow_spec.status == WorkflowStatus.done, workflow_spec.status)
0300 
0301     print("\n=== steps that are not all final are unaffected ===")
0302     interface = make_interface([FakeData("out/a", WFDataStatus.generating_suffice)], [FakeStep("a", WFStepStatus.running)], all_steps_final=False)
0303     workflow_spec = make_workflow()
0304     result = interface.process_workflow_running(workflow_spec)
0305     failures += not check("stays running", workflow_spec.status == WorkflowStatus.running, workflow_spec.status)
0306     failures += not check("no wait is recorded", workflow_spec.get_parameter(workflow_core.STEPS_FINAL_TIME_PARAM) is None)
0307 
0308     print("\n=== a root input already parked in waiting_suffice ===")
0309     # Workflow 133's rdo_bkg: it exists and has files, but DDM never closes it. A datum parked
0310     # there before the rule existed is let out here; a new one never gets parked at all.
0311     root_input = FakeData("rdo_bkg", WFDataStatus.waiting_suffice, WFDataType.input)
0312     interface = make_interface([root_input], [])
0313     result = interface.process_data_waiting(root_input)
0314     failures += not check("it becomes done_waited", root_input.status == WFDataStatus.done_waited, root_input.status)
0315     failures += not check("the transition is reported", result.new_status == WFDataStatus.done_waited, result.new_status)
0316     failures += not check("an end time is stamped", root_input.end_time is not None)
0317 
0318     print("\n=== but data produced inside the workflow still waits ===")
0319     for data_type in (WFDataType.mid, WFDataType.output):
0320         produced = FakeData(f"step/{data_type}", WFDataStatus.waiting_suffice, data_type)
0321         make_interface([produced], []).process_data_waiting(produced)
0322         failures += not check(f"{data_type} stays waiting_suffice", produced.status == WFDataStatus.waiting_suffice, produced.status)
0323 
0324     print("\n=== a root input that is not sufficient yet still waits ===")
0325     not_enough = FakeData("rdo_bkg", WFDataStatus.waiting_insuffi, WFDataType.input)
0326     make_interface([not_enough], [], check_status=WFDataTargetCheckStatus.insuffi).process_data_waiting(not_enough)
0327     failures += not check("it stays waiting_insuffi", not_enough.status == WFDataStatus.waiting_insuffi, not_enough.status)
0328 
0329     print("\n=== a closed collection is unchanged by this ===")
0330     closed = FakeData("rdo_bkg", WFDataStatus.waiting_suffice, WFDataType.input)
0331     make_interface([closed], [], check_status=WFDataTargetCheckStatus.complete).process_data_waiting(closed)
0332     failures += not check("complete still means done_waited", closed.status == WFDataStatus.done_waited, closed.status)
0333 
0334     print("\n=== and the step can now finish ===")
0335     # The point of the change: a done root input makes all_inputs_complete true, which is what
0336     # releases workflowHoldup so the task is allowed to finish.
0337     stats = make_interface([], [])._check_all_inputs_of_step(Log(), ["rdo_bkg"], {"rdo_bkg": root_input})
0338     failures += not check("all_inputs_complete is true once it is done", stats["all_inputs_complete"] is True, stats)
0339     still_waiting = FakeData("rdo_bkg", WFDataStatus.waiting_suffice, WFDataType.input)
0340     stats = make_interface([], [])._check_all_inputs_of_step(Log(), ["rdo_bkg"], {"rdo_bkg": still_waiting})
0341     failures += not check("...and was false while it waited", stats["all_inputs_complete"] is False, stats)
0342     failures += not check("...though it was already good enough to start on", stats["all_inputs_sufficient"] is True, stats)
0343 
0344     print("\n=== a root input never enters the waiting path in the first place ===")
0345     # The first check is where it is settled: an open but sufficient collection is called complete,
0346     # so the datum goes checked_complete -> done_skipped and is terminal straight away.
0347     fresh = FakeData("rdo_bkg", WFDataStatus.checking, WFDataType.input)
0348     interface = make_interface([fresh], [], check_status=WFDataTargetCheckStatus.suffice)
0349     interface.process_data_checking(fresh)
0350     failures += not check("checked_complete, not checked_suffice", fresh.status == WFDataStatus.checked_complete, fresh.status)
0351     interface.process_data_checked(fresh)
0352     failures += not check("and then done_skipped", fresh.status == WFDataStatus.done_skipped, fresh.status)
0353 
0354     print("\n=== data produced inside the workflow is unaffected at that check ===")
0355     for data_type in (WFDataType.mid, WFDataType.output):
0356         produced = FakeData(f"step/{data_type}", WFDataStatus.checking, data_type)
0357         make_interface([produced], [], check_status=WFDataTargetCheckStatus.suffice).process_data_checking(produced)
0358         failures += not check(f"{data_type} is still checked_suffice", produced.status == WFDataStatus.checked_suffice, produced.status)
0359 
0360     print("\n=== an insufficient root input is still not complete ===")
0361     thin = FakeData("rdo_bkg", WFDataStatus.checking, WFDataType.input)
0362     make_interface([thin], [], check_status=WFDataTargetCheckStatus.insuffi).process_data_checking(thin)
0363     failures += not check("checked_insuffi is untouched", thin.status == WFDataStatus.checked_insuffi, thin.status)
0364 
0365     print("\n=== a missing root input is still missing ===")
0366     gone = FakeData("rdo_bkg", WFDataStatus.checking, WFDataType.input)
0367     make_interface([gone], [], check_status=WFDataTargetCheckStatus.nonexist).process_data_checking(gone)
0368     failures += not check("checked_nonexist is untouched", gone.status == WFDataStatus.checked_nonexist, gone.status)
0369 
0370     print(f"\n{'ALL CHECKS PASSED' if not failures else f'{failures} CHECK(S) FAILED'}")
0371     return 1 if failures else 0
0372 
0373 
0374 if __name__ == "__main__":
0375     sys.exit(main())