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 (
0084 TASKID_PLACEHOLDER,
0085 WFID_PLACEHOLDER,
0086 has_placeholder,
0087 substitute_placeholder,
0088 )
0089 from pandaserver.workflow.workflow_parser import (
0090 INLINE_DESCRIPTION_KEY,
0091 parse_raw_request,
0092 )
0093
0094
0095
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())