Back to home page

EIC code displayed by LXR

 
 

    


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 (  # noqa: E402
0057     BaseDataHandler,
0058 )
0059 from pandaserver.workflow.data_handler_plugins.panda_task_data_handler import (  # noqa: E402
0060     PandaTaskDataHandler,
0061 )
0062 from pandaserver.workflow.workflow_base import (  # noqa: E402
0063     TASKID_PLACEHOLDER,
0064     WFDataSpec,
0065     WFDataStatus,
0066     WFDataTargetCheckStatus,
0067     WFStepStatus,
0068 )
0069 
0070 
0071 # Subclassing the real spec rather than standing in for it: the handlers declare WFDataSpec, and a
0072 # fake that has drifted from it should be reported here rather than pass as any object would.
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 # Rucio reports a collection with no files as length None, not 0 -- the shape that crashed
0116 # check_target in production the first time a freshly created output dataset was checked
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     # nonexist is the only check status that advances an output to binding, which is what lets the
0139     # step generating it leave checking. Anything else stalls the workflow.
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     # An open collection with no files must be insufficient, not sufficient: reporting suffice would
0207     # let a step start with zero input files.
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())