File indexing completed on 2026-09-28 09:36:48
0001 """
0002 Offline check of PandaTaskDataHandler.check_target.
0003
0004 The status this returns drives the workflow data state machine, so a wrong answer here stalls the
0005 step that consumes or produces the data, and with it the whole workflow. In particular a dataset
0006 name still holding ${TASKID} must be reported as non-existent rather than as "no change": reporting
0007 no change leaves the data in checking forever, which leaves the step in checking and the workflow in
0008 starting.
0009
0010 Run from the repository root: python3 pandaserver/workflow/examples/data_handler_test.py
0011 """
0012
0013 import os
0014 import sys
0015 import types
0016
0017 REPO_ROOT = os.path.abspath(os.path.join(os.path.dirname(__file__), "..", "..", ".."))
0018 sys.path.insert(0, REPO_ROOT)
0019
0020 WARNINGS: list[str] = []
0021
0022
0023 def stub(name, **attrs):
0024 module = types.ModuleType(name)
0025 for key, value in attrs.items():
0026 setattr(module, key, value)
0027 sys.modules[name] = module
0028 return module
0029
0030
0031 class Log:
0032 def __init__(self, *a, **k):
0033 pass
0034
0035 def info(self, m):
0036 pass
0037
0038 def debug(self, m):
0039 pass
0040
0041 def warning(self, m):
0042 WARNINGS.append(m)
0043
0044 def error(self, m):
0045 pass
0046
0047
0048 stub("pandacommon")
0049 stub("pandacommon.pandautils").__path__ = []
0050 stub("pandacommon.pandautils.base", SpecBase=object)
0051 stub("pandacommon.pandalogger").__path__ = []
0052 stub("pandacommon.pandalogger.LogWrapper", LogWrapper=Log)
0053 stub("pandacommon.pandalogger.PandaLogger", PandaLogger=lambda: types.SimpleNamespace(getLogger=lambda n: None))
0054 stub("pandaserver.config", panda_config=types.SimpleNamespace(schemaJEDI="ATLAS_PANDA", schemaDEFT="ATLAS_DEFT"))
0055
0056 from pandaserver.workflow.data_handler_plugins.base_data_handler import (
0057 BaseDataHandler,
0058 )
0059 from pandaserver.workflow.data_handler_plugins.panda_task_data_handler import (
0060 PandaTaskDataHandler,
0061 )
0062 from pandaserver.workflow.workflow_base import (
0063 TASKID_PLACEHOLDER,
0064 WFDataSpec,
0065 WFDataStatus,
0066 WFDataTargetCheckStatus,
0067 WFStepStatus,
0068 )
0069
0070
0071
0072
0073 class FakeData(WFDataSpec):
0074 def __init__(self, target_id, output_types, source_step_id=None):
0075 self.target_id = target_id
0076 self.output_types = output_types
0077 self.source_step_id = source_step_id
0078 self.flavor = "panda_task"
0079 self.workflow_id = 10
0080 self.data_id = 2
0081
0082 def get_parameter(self, param):
0083 return self.output_types if param == "output_types" else None
0084
0085
0086 class FakeStep:
0087 def __init__(self, status):
0088 self.status = status
0089 self.step_id = 27
0090
0091
0092 class FakeDDM:
0093 def __init__(self, metadata):
0094 self.metadata = metadata
0095 self.queried = []
0096
0097 def get_dataset_metadata(self, name, **kwargs):
0098 self.queried.append(name)
0099 return self.metadata
0100
0101
0102 class FakeTaskBuffer:
0103 def __init__(self, step=None):
0104 self.step = step
0105
0106 def get_workflow(self, workflow_id):
0107 return None
0108
0109 def get_workflow_step(self, step_id):
0110 return self.step
0111
0112
0113 CLOSED = {"state": "closed", "content_state": "closed", "length": 42}
0114 MISSING = {"state": "missing"}
0115
0116
0117 EMPTY_OPEN = {"state": "open", "content_state": "open", "length": None}
0118 EMPTY_CLOSED = {"state": "closed", "content_state": "closed", "length": None}
0119 UNRESOLVED = f"mc23_13p6TeV.526140.x.evgen.EVNT.e8590_wfid10_tid{TASKID_PLACEHOLDER}_00"
0120 RESOLVED = "mc23_13p6TeV.526140.x.evgen.EVNT.e8590_wfid10_tid49900001_00"
0121
0122
0123 def check(label, condition, detail=""):
0124 print(f" {'PASS' if condition else 'FAIL'} {label}{' ' + str(detail) if detail and not condition else ''}")
0125 return condition
0126
0127
0128 def main():
0129 failures = 0
0130
0131 print("\n=== a name still holding ${TASKID} ===")
0132 ddm = FakeDDM(CLOSED)
0133 handler: BaseDataHandler = PandaTaskDataHandler(FakeTaskBuffer(), ddm)
0134 result = handler.check_target(FakeData(UNRESOLVED, []))
0135 failures += not check(
0136 "reported as non-existent", result.success is True and result.check_status == WFDataTargetCheckStatus.nonexist, (result.success, result.check_status)
0137 )
0138
0139
0140 failures += not check("nonexist is a checked status the state machine acts on", WFDataStatus.checked_nonexist in WFDataStatus.checked_statuses)
0141 failures += not check("no DDM query made with an unresolved name", ddm.queried == [], ddm.queried)
0142
0143 print("\n=== an unresolved name on a finished step is surfaced ===")
0144 del WARNINGS[:]
0145 handler = PandaTaskDataHandler(FakeTaskBuffer(step=FakeStep(WFStepStatus.done)), FakeDDM(CLOSED))
0146 result = handler.check_target(FakeData(UNRESOLVED, [], source_step_id=27))
0147 failures += not check("still reported as non-existent", result.check_status == WFDataTargetCheckStatus.nonexist)
0148 failures += not check("a warning names the situation", any("still unresolved" in m for m in WARNINGS), WARNINGS)
0149 del WARNINGS[:]
0150 handler = PandaTaskDataHandler(FakeTaskBuffer(step=FakeStep(WFStepStatus.running)), FakeDDM(CLOSED))
0151 handler.check_target(FakeData(UNRESOLVED, [], source_step_id=27))
0152 failures += not check("no warning while the step is still running", not any("still unresolved" in m for m in WARNINGS), WARNINGS)
0153
0154 print("\n=== a resolved name with no output types (production shape) ===")
0155 ddm = FakeDDM(CLOSED)
0156 handler = PandaTaskDataHandler(FakeTaskBuffer(), ddm)
0157 result = handler.check_target(FakeData(RESOLVED, []))
0158 failures += not check("closed and non-empty -> complete", result.check_status == WFDataTargetCheckStatus.complete, result.check_status)
0159 failures += not check("queried the bare target_id", ddm.queried == [RESOLVED], ddm.queried)
0160 ddm = FakeDDM(MISSING)
0161 handler = PandaTaskDataHandler(FakeTaskBuffer(), ddm)
0162 failures += not check("missing -> nonexist", handler.check_target(FakeData(RESOLVED, [])).check_status == WFDataTargetCheckStatus.nonexist)
0163
0164 print("\n=== a collection reporting length None (no files yet) ===")
0165 ddm = FakeDDM(EMPTY_OPEN)
0166 handler = PandaTaskDataHandler(FakeTaskBuffer(), ddm)
0167 result = handler.check_target(FakeData(RESOLVED, []))
0168 failures += not check("does not raise on a None length", result.success is True, result.message)
0169 failures += not check("an open empty collection is not complete", result.check_status != WFDataTargetCheckStatus.complete, result.check_status)
0170 ddm = FakeDDM(EMPTY_CLOSED)
0171 handler = PandaTaskDataHandler(FakeTaskBuffer(), ddm)
0172 result = handler.check_target(FakeData(RESOLVED, []))
0173 failures += not check("does not raise when closed but empty", result.success is True, result.message)
0174 failures += not check(
0175 "a closed but empty collection is not complete",
0176 result.check_status != WFDataTargetCheckStatus.complete,
0177 result.check_status,
0178 )
0179
0180 print("\n=== an analysis output with output types is unaffected ===")
0181 ddm = FakeDDM(CLOSED)
0182 handler = PandaTaskDataHandler(FakeTaskBuffer(), ddm)
0183 result = handler.check_target(FakeData("user.me.myOut_002_b", ["aaa.root", "bbb.root"]))
0184 failures += not check(
0185 "expanded per output type",
0186 ddm.queried == ["user.me.myOut_002_b_aaa.root", "user.me.myOut_002_b_bbb.root"],
0187 ddm.queried,
0188 )
0189 failures += not check("complete", result.check_status == WFDataTargetCheckStatus.complete)
0190
0191 print("\n=== a done source step short-circuits to complete ===")
0192 handler = PandaTaskDataHandler(FakeTaskBuffer(step=FakeStep(WFStepStatus.done)), FakeDDM(MISSING))
0193 result = handler.check_target(FakeData(RESOLVED, [], source_step_id=27))
0194 failures += not check("complete regardless of DDM", result.check_status == WFDataTargetCheckStatus.complete)
0195
0196 print("\n=== ddm_collection handler: an empty input collection ===")
0197 from pandaserver.workflow.data_handler_plugins.ddm_collection_data_handler import (
0198 DDMCollectionDataHandler,
0199 )
0200
0201 class InputData(FakeData):
0202 def __init__(self, target_id):
0203 super().__init__(target_id, [])
0204 self.flavor = "ddm_collection"
0205
0206
0207
0208 handler = DDMCollectionDataHandler(FakeTaskBuffer(), FakeDDM(EMPTY_OPEN))
0209 result = handler.check_target(InputData("mc23_13p6TeV.some.input"))
0210 failures += not check("open and empty -> insuffi", result.check_status == WFDataTargetCheckStatus.insuffi, result.check_status)
0211 handler = DDMCollectionDataHandler(FakeTaskBuffer(), FakeDDM({"state": "open", "length": 0}))
0212 failures += not check(
0213 "open with an explicit zero length -> insuffi",
0214 handler.check_target(InputData("mc23_13p6TeV.some.input")).check_status == WFDataTargetCheckStatus.insuffi,
0215 )
0216 handler = DDMCollectionDataHandler(FakeTaskBuffer(), FakeDDM({"state": "open", "length": 7}))
0217 failures += not check(
0218 "open with files -> suffice",
0219 handler.check_target(InputData("mc23_13p6TeV.some.input")).check_status == WFDataTargetCheckStatus.suffice,
0220 )
0221 handler = DDMCollectionDataHandler(FakeTaskBuffer(), FakeDDM(CLOSED))
0222 failures += not check("closed -> complete", handler.check_target(InputData("mc23_13p6TeV.some.input")).check_status == WFDataTargetCheckStatus.complete)
0223 handler = DDMCollectionDataHandler(FakeTaskBuffer(), FakeDDM(MISSING))
0224 failures += not check("missing -> nonexist", handler.check_target(InputData("mc23_13p6TeV.some.input")).check_status == WFDataTargetCheckStatus.nonexist)
0225
0226 print(f"\n{'ALL CHECKS PASSED' if not failures else f'{failures} CHECK(S) FAILED'}")
0227 return 1 if failures else 0
0228
0229
0230 if __name__ == "__main__":
0231 sys.exit(main())