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
0026
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
0105 from pandaserver.workflow.workflow_base import (
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
0117
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
0203
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
0250
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
0310
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
0336
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
0346
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())