Back to home page

EIC code displayed by LXR

 
 

    


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

0001 """
0002 Offline check of the late-bound ID placeholders in a workflow description.
0003 
0004 ${TASKID} is used by the production example, since a dataset name embeds the ID of the task that
0005 produced it. ${WFID} is supported but deliberately unused there: ATLAS dataset names follow a
0006 convention that carries no workflow ID. This test therefore exercises ${WFID} against a synthetic
0007 description, so the capability keeps working should a future convention want it.
0008 
0009 Run from the repository root:  python3 pandaserver/workflow/examples/placeholder_resolution_test.py
0010 """
0011 
0012 import importlib.abc
0013 import importlib.machinery
0014 import json
0015 import os
0016 import sys
0017 import types
0018 
0019 REPO_ROOT = os.path.abspath(os.path.join(os.path.dirname(__file__), "..", "..", ".."))
0020 sys.path.insert(0, REPO_ROOT)
0021 
0022 AUTO_STUB_ROOTS = ("idds", "pandaclient", "ruamel", "requests", "snakemake")
0023 
0024 
0025 class AutoStubFinder(importlib.abc.MetaPathFinder, importlib.abc.Loader):
0026     def find_spec(self, name, path=None, target=None):
0027         if name.split(".")[0] in AUTO_STUB_ROOTS:
0028             return importlib.machinery.ModuleSpec(name, self, is_package=True)
0029         return None
0030 
0031     def create_module(self, spec):
0032         module = types.ModuleType(spec.name)
0033         module.__path__ = []
0034         return module
0035 
0036     def exec_module(self, module):
0037         class Anything:
0038             def __init__(self, *a, **k):
0039                 pass
0040 
0041             def __call__(self, *a, **k):
0042                 return self
0043 
0044         module.__getattr__ = lambda name: Anything
0045 
0046 
0047 sys.meta_path.insert(0, AutoStubFinder())
0048 
0049 
0050 def stub(name, **attrs):
0051     module = types.ModuleType(name)
0052     for key, value in attrs.items():
0053         setattr(module, key, value)
0054     sys.modules[name] = module
0055     return module
0056 
0057 
0058 class Log:
0059     def __init__(self, *a, **k):
0060         pass
0061 
0062     def info(self, m):
0063         pass
0064 
0065     def debug(self, m):
0066         pass
0067 
0068     def warning(self, m):
0069         print(f"  WARN: {m}")
0070 
0071     def error(self, m):
0072         print(f"  ERROR: {str(m)[:200]}")
0073 
0074 
0075 stub("pandacommon")
0076 stub("pandacommon.pandautils").__path__ = []
0077 stub("pandacommon.pandautils.base", SpecBase=object)
0078 stub("pandacommon.pandalogger").__path__ = []
0079 stub("pandacommon.pandalogger.LogWrapper", LogWrapper=Log)
0080 stub("pandacommon.pandalogger.PandaLogger", PandaLogger=lambda: types.SimpleNamespace(getLogger=lambda n: None))
0081 stub("pandaserver.config", panda_config=types.SimpleNamespace(schemaJEDI="ATLAS_PANDA", schemaDEFT="ATLAS_DEFT"))
0082 
0083 from pandaserver.workflow.workflow_base import (  # noqa: E402
0084     TASKID_PLACEHOLDER,
0085     WFID_PLACEHOLDER,
0086     has_placeholder,
0087     substitute_placeholder,
0088 )
0089 from pandaserver.workflow.workflow_parser import (  # noqa: E402
0090     INLINE_DESCRIPTION_KEY,
0091     parse_raw_request,
0092 )
0093 
0094 # A synthetic description using BOTH placeholders, including ${WFID} which the production example
0095 # deliberately does not use. Kept minimal: one step, one output.
0096 DESCRIPTION = {
0097     "name": "placeholder_probe",
0098     "outputs": {"out": {"from": "only_step/EVNT"}},
0099     "steps": {
0100         "only_step": {
0101             "type": "task",
0102             "task_params": {
0103                 "taskName": "probe.evgen.e0001",
0104                 "taskType": "prod",
0105                 "prodSourceLabel": "managed",
0106                 "vo": "atlas",
0107                 "userName": "prober",
0108                 "transPath": "Gen_tf.py",
0109                 "noInput": True,
0110                 "log": {
0111                     "type": "template",
0112                     "param_type": "log",
0113                     "value": "log.${TASKID}._${SN}.job.log.tgz",
0114                     "dataset": "probe.evgen.log.e0001_wfid${WFID}_tid${TASKID}_00",
0115                 },
0116                 "jobParameters": [
0117                     {
0118                         "type": "template",
0119                         "param_type": "output",
0120                         "value": "--outputEVNTFile=EVNT.${TASKID}._${SN}.pool.root",
0121                         "dataset": "probe.evgen.EVNT.e0001_wfid${WFID}_tid${TASKID}_00",
0122                     }
0123                 ],
0124             },
0125         }
0126     },
0127 }
0128 
0129 
0130 def check(label, condition, detail=""):
0131     print(f"  {'PASS' if condition else 'FAIL'}  {label}{'  ' + str(detail) if detail and not condition else ''}")
0132     return condition
0133 
0134 
0135 def main():
0136     failures = 0
0137 
0138     print("\n=== the helpers ===")
0139     both = "ds.e0001_wfid${WFID}_tid${TASKID}_00"
0140     after_wfid = substitute_placeholder(both, WFID_PLACEHOLDER, 10)
0141     after_both = substitute_placeholder(after_wfid, TASKID_PLACEHOLDER, 4004198)
0142     failures += not check("WFID resolves, TASKID left alone", after_wfid == "ds.e0001_wfid10_tid${TASKID}_00", after_wfid)
0143     failures += not check("TASKID then resolves", after_both == "ds.e0001_wfid10_tid4004198_00", after_both)
0144     failures += not check("has_placeholder tracks both", has_placeholder(both, WFID_PLACEHOLDER) and not has_placeholder(after_wfid, WFID_PLACEHOLDER))
0145     failures += not check("substitution recurses into nested structures", substitute_placeholder({"a": ["x${WFID}"]}, WFID_PLACEHOLDER, 7) == {"a": ["x7"]})
0146     failures += not check("JEDI per-job templates untouched", substitute_placeholder("log.${TASKID}._${SN}.tgz", TASKID_PLACEHOLDER, 7) == "log.7._${SN}.tgz")
0147 
0148     print("\n=== ${WFID} is resolved when the description is parsed ===")
0149     is_ok, is_fatal, definition = parse_raw_request(None, "<probe>", "prober", {INLINE_DESCRIPTION_KEY: DESCRIPTION}, workflow_id=4242)
0150     failures += not check("parsed", is_ok and not is_fatal)
0151     blob = json.dumps(definition, default=str)
0152     failures += not check("${WFID} resolved to the workflow id", "${WFID}" not in blob and "wfid4242" in blob)
0153     failures += not check("${TASKID} deliberately left for submission time", "${TASKID}" in blob)
0154     node = definition["nodes"][0]
0155     print(f"    output dataset: {node['outputs']['only_step/EVNT']['value']}")
0156     print(f"    log dataset:    {node['task_params']['log']['dataset']}")
0157     failures += not check(
0158         "resolved in the output dataset name",
0159         node["outputs"]["only_step/EVNT"]["value"] == "probe.evgen.EVNT.e0001_wfid4242_tid${TASKID}_00",
0160         node["outputs"]["only_step/EVNT"]["value"],
0161     )
0162     failures += not check(
0163         "resolved in the log dataset name too",
0164         node["task_params"]["log"]["dataset"] == "probe.evgen.log.e0001_wfid4242_tid${TASKID}_00",
0165         node["task_params"]["log"]["dataset"],
0166     )
0167 
0168     print("\n=== without a workflow id the placeholder is left in place ===")
0169     is_ok2, _, definition2 = parse_raw_request(None, "<probe>", "prober", {INLINE_DESCRIPTION_KEY: DESCRIPTION}, workflow_id=None)
0170     failures += not check("parsed", is_ok2)
0171     failures += not check("${WFID} untouched", "${WFID}" in json.dumps(definition2, default=str))
0172 
0173     print(f"\n{'ALL CHECKS PASSED' if not failures else f'{failures} CHECK(S) FAILED'}")
0174     return 1 if failures else 0
0175 
0176 
0177 if __name__ == "__main__":
0178     sys.exit(main())