File indexing completed on 2026-09-28 09:36:49
0001 """
0002 Offline check of workflow step task submission.
0003
0004 Covers the pieces that turn a parsed step into a queued task: resolving input dataset references
0005 against the datasets actually produced upstream, refusing a production label without the role,
0006 resolving the late-bound task ID into the step's output dataset names, and the task status mapping.
0007
0008 Run from the repository root: python3 pandaserver/workflow/examples/step_submission_test.py
0009 """
0010
0011 import copy
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 LOGGED: dict[str, list[str]] = {"warning": [], "error": []}
0059
0060
0061 class Log:
0062 def __init__(self, *a, **k):
0063 pass
0064
0065 def info(self, m):
0066 pass
0067
0068 def debug(self, m):
0069 pass
0070
0071 def warning(self, m):
0072 LOGGED["warning"].append(m)
0073
0074 def error(self, m):
0075 LOGGED["error"].append(m)
0076
0077
0078 stub("pandacommon")
0079 pandautils = stub("pandacommon.pandautils")
0080 pandautils.__path__ = []
0081 stub("pandacommon.pandautils.base", SpecBase=object)
0082 pandalogger = stub("pandacommon.pandalogger")
0083 pandalogger.__path__ = []
0084 stub("pandacommon.pandalogger.LogWrapper", LogWrapper=Log)
0085 stub("pandacommon.pandalogger.PandaLogger", PandaLogger=lambda: types.SimpleNamespace(getLogger=lambda n: None))
0086 stub("pandaserver.config", panda_config=types.SimpleNamespace(schemaJEDI="ATLAS_PANDA", schemaDEFT="ATLAS_DEFT"))
0087
0088 from pandaserver.workflow.step_handler_plugins.panda_task_step_handler import (
0089 PandaTaskStepHandler,
0090 )
0091 from pandaserver.workflow.workflow_base import (
0092 PARENT_TASKID_PLACEHOLDER,
0093 TASKID_PLACEHOLDER,
0094 WFStepSpec,
0095 WFStepStatus,
0096 )
0097
0098
0099 class FakeData:
0100 def __init__(self, name, target_id, source_step_id=None):
0101 self.name = name
0102 self.target_id = target_id
0103 self.workflow_id = 1
0104 self.data_id = abs(hash(name)) % 1000
0105
0106 self.source_step_id = source_step_id
0107
0108
0109 class FakeStep(WFStepSpec):
0110 def __init__(self, definition, parameters=None):
0111 self.workflow_id = 1
0112 self.step_id = 7
0113 self.name = "probe"
0114 self.flavor = "panda_task"
0115 self.target_id = None
0116 self.status = WFStepStatus.ready
0117
0118
0119 self.definition_json = json.dumps(definition)
0120 self.parameters = json.dumps(parameters or {})
0121
0122
0123 class FakeTaskBuffer:
0124 def __init__(self, data_by_name=None, task_id=49900001, error="", deft_status=None, steps_by_id=None):
0125 self.deft_status = deft_status
0126 self.data = data_by_name or {}
0127 self.steps = steps_by_id or {}
0128 self.task_id = task_id
0129 self.error = error
0130 self.inserted = []
0131 self.inserted_parent_tids = []
0132 self.updated_data = []
0133
0134 def get_workflow_step(self, step_id):
0135 return self.steps.get(step_id)
0136
0137 def get_steps_of_workflow(self, workflow_id, status_filter_list=None, status_exclusion_list=None):
0138 return list(self.steps.values())
0139
0140 def get_workflow_data_by_name(self, name, workflow_id):
0141 return self.data.get(name)
0142
0143 def update_workflow_data(self, data_spec):
0144 self.updated_data.append((data_spec.name, data_spec.target_id))
0145
0146 def update_workflow_step(self, step_spec):
0147 pass
0148
0149 def insert_step_task(self, task_params_map, user_dn, parent_tid=None):
0150 self.inserted.append(copy.deepcopy(task_params_map))
0151 self.inserted_parent_tids.append(parent_tid)
0152 if self.task_id is None:
0153 return None, self.error
0154 return self.task_id, ""
0155
0156 def getTaskStatusSuperstatus(self, task_id):
0157 return self._status
0158
0159 def get_deft_task_status(self, task_id):
0160 return self.deft_status
0161
0162 def set_status(self, status, superstatus=None):
0163
0164 self._status = None if status is None else (status, superstatus or status)
0165
0166 def getTaskWithID_JEDI(self, task_id, *args, **kwargs):
0167
0168 return (False, None) if self._status is None else (True, None)
0169
0170
0171 def check(label, condition, detail=""):
0172 print(f" {'PASS' if condition else 'FAIL'} {label}{' ' + str(detail) if detail and not condition else ''}")
0173 return condition
0174
0175
0176 def main():
0177 failures = 0
0178 wfd = json.load(open(os.path.join(os.path.dirname(__file__), "production_chain_wfd.json")))
0179
0180 simul_params = copy.deepcopy(wfd["steps"]["simul"]["task_params"])
0181
0182 simul_params = json.loads(json.dumps(simul_params).replace("${WFID}", "12345"))
0183
0184 def make_step(prod_role=True, all_inputs_complete=True, params=None):
0185 return FakeStep(
0186 {
0187 "task_params": copy.deepcopy(params if params is not None else simul_params),
0188 "user_dn": "/DC=ch/CN=test",
0189 "prod_role": prod_role,
0190 "output_data_list": ["simul/HITS"],
0191 },
0192 {"all_inputs_complete": all_inputs_complete},
0193 )
0194
0195 produced_evnt = "mc23_13p6TeV.526140.x.merge.EVNT.e8590_e8586_wfid12345_tid48810699_00"
0196
0197 def make_tbif(**kw):
0198 data = {
0199 "merge_evnt/EVNT": FakeData("merge_evnt/EVNT", produced_evnt),
0200 "simul/HITS": FakeData("simul/HITS", f"mc23_13p6TeV.526140.x.simul.HITS.e8590_e8586_a934_wfid12345_tid{TASKID_PLACEHOLDER}_00"),
0201 }
0202 return FakeTaskBuffer(data_by_name=data, **kw)
0203
0204 print("\n=== submit_target: the happy path ===")
0205 tbif = make_tbif()
0206 handler = PandaTaskStepHandler(tbif)
0207 step = make_step()
0208 res = handler.submit_target(step)
0209 failures += not check("submitted", res.success is True, res.message)
0210 failures += not check("target_id is the task id", res.target_id == "49900001", res.target_id)
0211 submitted = tbif.inserted[0]
0212 inputs = [p["dataset"] for p in submitted["jobParameters"] if p.get("param_type") == "input"]
0213 failures += not check("input reference resolved to the produced dataset", inputs == [produced_evnt], inputs)
0214 failures += not check("no brace reference left in any input dataset", not any("{" in d for d in inputs), inputs)
0215
0216
0217 outputs = [p["dataset"] for p in submitted["jobParameters"] if p.get("param_type") == "output"]
0218 failures += not check("output datasets still carry the placeholder for the DB layer to resolve", all(TASKID_PLACEHOLDER in d for d in outputs), outputs)
0219 failures += not check("pseudo_input left untouched", any(p.get("dataset") == "seq_number" for p in submitted["jobParameters"]))
0220 failures += not check("workflowHoldup not set when inputs are complete", "workflowHoldup" not in submitted)
0221 failures += not check(
0222 "output dataset name resolved from the task id",
0223 tbif.updated_data == [("simul/HITS", "mc23_13p6TeV.526140.x.simul.HITS.e8590_e8586_a934_wfid12345_tid49900001_00")],
0224 tbif.updated_data,
0225 )
0226 failures += not check("submission attempt recorded", step.get_parameter("submit_attempt_task_name") == simul_params["taskName"])
0227
0228 print("\n=== workflowHoldup is set while inputs are incomplete ===")
0229 tbif = make_tbif()
0230 handler = PandaTaskStepHandler(tbif)
0231 handler.submit_target(make_step(all_inputs_complete=False))
0232 failures += not check("workflowHoldup set", tbif.inserted[0].get("workflowHoldup") is True)
0233
0234 print("\n=== a production label without the role is refused ===")
0235
0236 for label in ["managed"]:
0237 params = copy.deepcopy(simul_params)
0238 params["prodSourceLabel"] = label
0239 tbif = make_tbif()
0240 handler = PandaTaskStepHandler(tbif)
0241 res = handler.submit_target(make_step(prod_role=False, params=params))
0242 failures += not check(f"{label} refused without the role", res.success is not True)
0243 failures += not check(f"{label} reason mentions the production role", "production role" in res.message, res.message)
0244 failures += not check(f"{label} submitted nothing", tbif.inserted == [])
0245
0246 tbif = make_tbif()
0247 handler = PandaTaskStepHandler(tbif)
0248 res = handler.submit_target(make_step(prod_role=True, params=params))
0249 failures += not check(f"{label} accepted with the role", res.success is True, res.message)
0250
0251
0252
0253 for label in ["user", "test", "ptest"]:
0254 params = copy.deepcopy(simul_params)
0255 params["prodSourceLabel"] = label
0256 tbif = make_tbif()
0257 handler = PandaTaskStepHandler(tbif)
0258 res = handler.submit_target(make_step(prod_role=False, params=params))
0259 failures += not check(f"{label} needs no production role", res.success is True, res.message)
0260
0261 print("\n=== an unresolved upstream output blocks submission ===")
0262 tbif = make_tbif()
0263 tbif.data["merge_evnt/EVNT"].target_id = f"mc23...merge.EVNT.e8590_wfid12345_tid{TASKID_PLACEHOLDER}_00"
0264 handler = PandaTaskStepHandler(tbif)
0265 res = handler.submit_target(make_step())
0266 failures += not check("refused", res.success is not True)
0267 failures += not check("reason mentions it is not resolved yet", "not resolved yet" in res.message, res.message)
0268 failures += not check("nothing submitted", tbif.inserted == [])
0269
0270 print("\n=== a missing upstream output blocks submission ===")
0271 tbif = make_tbif()
0272 del tbif.data["merge_evnt/EVNT"]
0273 handler = PandaTaskStepHandler(tbif)
0274 res = handler.submit_target(make_step())
0275 failures += not check("refused", res.success is not True)
0276 failures += not check("nothing submitted", tbif.inserted == [])
0277
0278 print("\n=== a second attempt for the same taskName is refused ===")
0279 tbif = make_tbif()
0280 handler = PandaTaskStepHandler(tbif)
0281 step = make_step()
0282 handler.submit_target(step)
0283 res2 = handler.submit_target(step)
0284 failures += not check("second attempt refused", res2.success is not True)
0285 failures += not check("reason mentions the previous attempt", "previous attempt" in res2.message, res2.message)
0286 failures += not check("submitted only once", len(tbif.inserted) == 1, len(tbif.inserted))
0287
0288 print("\n=== a literal (external) dataset is passed through ===")
0289 recon_params = json.loads(json.dumps(copy.deepcopy(wfd["steps"]["recon"]["task_params"])).replace("${WFID}", "12345"))
0290 tbif = make_tbif()
0291 tbif.data["merge_hits/HITS"] = FakeData("merge_hits/HITS", "mc23...merge.HITS..._tid48810713_00")
0292 tbif.data["rdo_bkg"] = FakeData("rdo_bkg", wfd["inputs"]["rdo_bkg"])
0293 handler = PandaTaskStepHandler(tbif)
0294 step = FakeStep({"task_params": recon_params, "user_dn": "/DC=ch/CN=t", "prod_role": True, "output_data_list": []}, {"all_inputs_complete": True})
0295 res = handler.submit_target(step)
0296 failures += not check("submitted", res.success is True, res.message)
0297 datasets = [p["dataset"] for p in tbif.inserted[0]["jobParameters"] if p.get("param_type") == "input"]
0298 failures += not check("both inputs resolved", set(datasets) == {"mc23...merge.HITS..._tid48810713_00", wfd["inputs"]["rdo_bkg"]}, datasets)
0299
0300 print("\n=== a task submitted but not yet refined into JEDI ===")
0301
0302
0303 del LOGGED["warning"][:], LOGGED["error"][:]
0304 tbif = make_tbif(deft_status="waiting")
0305 tbif.set_status(None)
0306 handler = PandaTaskStepHandler(tbif)
0307 step = make_step()
0308 step.status = WFStepStatus.running
0309 step.target_id = "49900001"
0310 result = handler.check_target(step)
0311 failures += not check("check_target does not advance the step", not result.success)
0312
0313 failures += not check("success is None, not False", result.success is None, result.success)
0314 failures += not check("logged as a warning, not an error", LOGGED["warning"] and not LOGGED["error"], (LOGGED["warning"], LOGGED["error"]))
0315 failures += not check("message explains it is queued in DEFT", "queued in DEFT" in result.message and "not yet refined" in result.message, result.message)
0316
0317 del LOGGED["warning"][:], LOGGED["error"][:]
0318 handler.on_all_inputs_done(step)
0319 failures += not check("on_all_inputs_done also warns rather than errors", LOGGED["warning"] and not LOGGED["error"], (LOGGED["warning"], LOGGED["error"]))
0320
0321 print("\n=== a task missing from both JEDI and DEFT is still an error ===")
0322 del LOGGED["warning"][:], LOGGED["error"][:]
0323 tbif = make_tbif(deft_status=None)
0324 tbif.set_status(None)
0325 handler = PandaTaskStepHandler(tbif)
0326 step = make_step()
0327 step.status = WFStepStatus.running
0328 step.target_id = "49900001"
0329 result = handler.check_target(step)
0330 failures += not check("logged as an error", LOGGED["error"] and not LOGGED["warning"], (LOGGED["warning"], LOGGED["error"]))
0331 failures += not check("message says neither JEDI nor DEFT", "not found in JEDI or DEFT" in result.message, result.message)
0332 failures += not check("success is False, so the caller reports a failure", result.success is False, result.success)
0333
0334 print("\n=== the tri-state is honoured by every check_target return path ===")
0335
0336 for label, setup, expected in [
0337 ("step not in a checkable status", lambda st, tb: setattr(st, "status", WFStepStatus.pending), None),
0338
0339 ("wrong flavor for this handler", lambda st, tb: setattr(st, "flavor", "something_else"), False),
0340 ("target not submitted yet", lambda st, tb: setattr(st, "target_id", None), None),
0341 ("unrecognised native status", lambda st, tb: tb.set_status("nonsense_status"), False),
0342 ]:
0343 tbif = make_tbif(deft_status="waiting")
0344 tbif.set_status("running")
0345 step = make_step()
0346 step.status = WFStepStatus.running
0347 step.target_id = "49900001"
0348 setup(step, tbif)
0349 got = PandaTaskStepHandler(tbif).check_target(step).success
0350 failures += not check(f"{label} -> success {expected}", got is expected, f"got {got}")
0351
0352 print("\n=== a flavor mismatch fails loudly on every entry point ===")
0353 for name in ("submit_target", "check_target", "cancel_target", "on_all_inputs_done"):
0354 del LOGGED["warning"][:], LOGGED["error"][:]
0355 tbif = make_tbif()
0356 tbif.set_status("running")
0357 step = make_step()
0358 step.status = WFStepStatus.running
0359 step.target_id = "49900001"
0360 step.flavor = "something_else"
0361 result = getattr(PandaTaskStepHandler(tbif), name)(step)
0362 failures += not check(f"{name} logs an error, not a warning", LOGGED["error"] and not LOGGED["warning"], (LOGGED["warning"], LOGGED["error"]))
0363 failures += not check(f"{name} names the wrong handler", any("wrong step handler" in m for m in LOGGED["error"]), LOGGED["error"])
0364 if result is not None:
0365 failures += not check(f"{name} reports success False", result.success is False, result.success)
0366 failures += not check(f"{name} did nothing", tbif.inserted == [] and tbif.updated_data == [])
0367
0368 print("\n=== ${PARENT_TASKID}: resolved from the step that produced the input ===")
0369
0370 def make_parent_step(parent_tid_value, noWaitParent=True, inputs=("merge_evnt/EVNT",)):
0371 params = copy.deepcopy(simul_params)
0372 params["parent_tid"] = parent_tid_value
0373 if noWaitParent:
0374 params["noWaitParent"] = True
0375 else:
0376 params.pop("noWaitParent", None)
0377 step = make_step(params=params)
0378 definition = step.definition_json_map
0379 definition["input_data_dict"] = {name: {} for name in inputs}
0380 step.definition_json = json.dumps(definition)
0381 return step
0382
0383
0384 def parent_tbif(**kw):
0385 data = {
0386 "merge_evnt/EVNT": FakeData("merge_evnt/EVNT", produced_evnt, source_step_id=3),
0387 "simul/HITS": FakeData("simul/HITS", f"...tid{TASKID_PLACEHOLDER}_00"),
0388 }
0389 data.update(kw.pop("extra_data", {}))
0390 feeding = FakeStep({}, None)
0391 feeding.step_id, feeding.name = 3, "merge_evnt"
0392 return FakeTaskBuffer(data_by_name=data, steps_by_id={3: feeding}, **kw)
0393
0394 tbif = parent_tbif()
0395 tbif.steps[3].target_id = "48810699"
0396 res = PandaTaskStepHandler(tbif).submit_target(make_parent_step(PARENT_TASKID_PLACEHOLDER))
0397 failures += not check("the task is submitted", res.success is True, res.message)
0398 failures += not check("with the producing step's task as parent", tbif.inserted_parent_tids == [48810699], tbif.inserted_parent_tids)
0399 failures += not check("and parent_tid is taken out of the task params", "parent_tid" not in tbif.inserted[0], sorted(tbif.inserted[0])[:5])
0400
0401 print("\n=== a step with no parent inside the workflow ===")
0402 tbif = parent_tbif(extra_data={"rdo_bkg": FakeData("rdo_bkg", "external.dataset", source_step_id=None)})
0403 res = PandaTaskStepHandler(tbif).submit_target(make_parent_step(PARENT_TASKID_PLACEHOLDER, inputs=("rdo_bkg",)))
0404 failures += not check("is submitted", res.success is True, res.message)
0405 failures += not check("with no parent, so JEDI makes the task its own", tbif.inserted_parent_tids == [None], tbif.inserted_parent_tids)
0406
0407 print("\n=== a step that does not ask for it is unaffected ===")
0408 tbif = parent_tbif()
0409 res = PandaTaskStepHandler(tbif).submit_target(make_step())
0410 failures += not check("submitted with no parent", res.success is True and tbif.inserted_parent_tids == [None], tbif.inserted_parent_tids)
0411
0412 print("\n=== a literal task ID is passed through ===")
0413 tbif = parent_tbif()
0414 tbif.steps[3].target_id = "48810699"
0415 PandaTaskStepHandler(tbif).submit_target(make_parent_step(52397622))
0416 failures += not check("as given, without consulting the graph", tbif.inserted_parent_tids == [52397622], tbif.inserted_parent_tids)
0417
0418 print("\n=== a joining step is refused rather than guessed at ===")
0419 tbif = parent_tbif(
0420 extra_data={
0421 "left/OUT": FakeData("left/OUT", "a", source_step_id=3),
0422 "right/OUT": FakeData("right/OUT", "b", source_step_id=4),
0423 }
0424 )
0425 tbif.steps[3].target_id = "52000003"
0426 tbif.steps[4] = FakeStep({}, None)
0427 tbif.steps[4].step_id, tbif.steps[4].name = 4, "right"
0428 tbif.steps[4].target_id = "52000004"
0429 res = PandaTaskStepHandler(tbif).submit_target(make_parent_step(PARENT_TASKID_PLACEHOLDER, inputs=("left/OUT", "right/OUT")))
0430 failures += not check("not submitted", res.success is not True, res.success)
0431 failures += not check("and says it is ambiguous", "ambiguous" in res.message, res.message)
0432 failures += not check("and offers the named form", "PARENT_TASKID:<step name>" in res.message, res.message)
0433 failures += not check("nothing was inserted", tbif.inserted == [], tbif.inserted)
0434
0435 print("\n=== a value that is neither is refused ===")
0436 tbif = parent_tbif()
0437 res = PandaTaskStepHandler(tbif).submit_target(make_parent_step("not-a-task"))
0438 failures += not check("not submitted", res.success is not True, res.success)
0439 failures += not check("and says why", "neither" in res.message, res.message)
0440
0441 print("\n=== a parent whose step has not submitted is an inconsistency ===")
0442 tbif = parent_tbif()
0443 res = PandaTaskStepHandler(tbif).submit_target(make_parent_step(PARENT_TASKID_PLACEHOLDER))
0444 failures += not check("not submitted", res.success is not True, res.success)
0445 failures += not check("and says the parent has no task", "no task yet" in res.message, res.message)
0446
0447 print("\n=== setting a parent without noWaitParent is warned about ===")
0448 del LOGGED["warning"][:]
0449 tbif = parent_tbif()
0450 tbif.steps[3].target_id = "48810699"
0451 res = PandaTaskStepHandler(tbif).submit_target(make_parent_step(PARENT_TASKID_PLACEHOLDER, noWaitParent=False))
0452 failures += not check("still submitted", res.success is True, res.message)
0453 failures += not check("with a warning naming noWaitParent", any("noWaitParent" in m for m in LOGGED["warning"]), LOGGED["warning"])
0454
0455 print("\n=== ${PARENT_TASKID:<step name>} on a joining step ===")
0456
0457 def joining_tbif():
0458 data = {
0459 "left/OUT": FakeData("left/OUT", "a", source_step_id=3),
0460 "right/OUT": FakeData("right/OUT", "b", source_step_id=4),
0461
0462
0463 "merge_evnt/EVNT": FakeData("merge_evnt/EVNT", produced_evnt, source_step_id=None),
0464 "simul/HITS": FakeData("simul/HITS", f"...tid{TASKID_PLACEHOLDER}_00"),
0465 }
0466 left, right = FakeStep({}, None), FakeStep({}, None)
0467 left.name, left.step_id, left.target_id = "left", 3, "52000003"
0468 right.name, right.step_id, right.target_id = "right", 4, "52000004"
0469 return FakeTaskBuffer(data_by_name=data, steps_by_id={3: left, 4: right})
0470
0471
0472 for named, expected_task in (("left", 52000003), ("right", 52000004)):
0473 tbif = joining_tbif()
0474 res = PandaTaskStepHandler(tbif).submit_target(make_parent_step(f"${{PARENT_TASKID:{named}}}", inputs=("left/OUT", "right/OUT")))
0475 failures += not check(f"naming {named} submits", res.success is True, res.message)
0476 failures += not check(f"...with its task {expected_task}", tbif.inserted_parent_tids == [expected_task], tbif.inserted_parent_tids)
0477
0478 print("\n=== naming a step that does not feed this one ===")
0479 tbif = joining_tbif()
0480 res = PandaTaskStepHandler(tbif).submit_target(make_parent_step("${PARENT_TASKID:elsewhere}", inputs=("left/OUT", "right/OUT")))
0481 failures += not check("is refused", res.success is not True, res.success)
0482 failures += not check("naming what does feed it", "left" in res.message and "right" in res.message, res.message)
0483 failures += not check("nothing inserted", tbif.inserted == [], tbif.inserted)
0484
0485 print("\n=== the bare form on a join now points at the named form ===")
0486 tbif = joining_tbif()
0487 res = PandaTaskStepHandler(tbif).submit_target(make_parent_step(PARENT_TASKID_PLACEHOLDER, inputs=("left/OUT", "right/OUT")))
0488 failures += not check("still refused", res.success is not True, res.success)
0489 failures += not check("and says how to disambiguate", "PARENT_TASKID:<step name>" in res.message, res.message)
0490
0491 print("\n=== naming the only feeding step is the same as the bare form ===")
0492 tbif = parent_tbif()
0493 tbif.steps[3].target_id = "48810699"
0494 res = PandaTaskStepHandler(tbif).submit_target(make_parent_step("${PARENT_TASKID:merge_evnt}"))
0495 failures += not check("submits with that task", res.success is True and tbif.inserted_parent_tids == [48810699], tbif.inserted_parent_tids)
0496
0497 print("\n=== check_target status mapping ===")
0498 expectations = {
0499 WFStepStatus.running: ["running", "scouting", "scouted", "throttled", "prepared", "finishing", "passed", "merging", "toretry", "toincexec", "paused"],
0500 WFStepStatus.starting: [
0501 "registered",
0502 "defined",
0503 "assigned",
0504 "activated",
0505 "starting",
0506 "ready",
0507 "topreprocess",
0508 "preprocessing",
0509 "staging",
0510 "staged",
0511 "rerefine",
0512 ],
0513 WFStepStatus.done: ["done", "finished"],
0514 WFStepStatus.failed: ["failed", "exhausted", "aborted", "toabort", "aborting", "broken", "tobroken"],
0515 }
0516 for expected_status, statuses in expectations.items():
0517 for status in statuses:
0518 tbif = make_tbif()
0519 tbif.set_status(status)
0520 handler = PandaTaskStepHandler(tbif)
0521 step = make_step()
0522 step.status = WFStepStatus.running
0523 step.target_id = "49900001"
0524 result = handler.check_target(step)
0525 if not (result.success and result.step_status == expected_status):
0526 failures += not check(f"{status} -> {expected_status}", False, f"got success={result.success} status={result.step_status}")
0527 failures += not check(f"all {sum(len(v) for v in expectations.values())} task statuses map without error", True)
0528
0529 tbif = make_tbif()
0530 tbif.set_status("nonsense_status")
0531 handler = PandaTaskStepHandler(tbif)
0532 step = make_step()
0533 step.status = WFStepStatus.running
0534 step.target_id = "49900001"
0535 result = handler.check_target(step)
0536 failures += not check("an unknown status is still an error", result.success is False and "unknown" in result.message, result.message)
0537
0538 print(f"\n{'ALL CHECKS PASSED' if not failures else f'{failures} CHECK(S) FAILED'}")
0539 return 1 if failures else 0
0540
0541
0542 if __name__ == "__main__":
0543 sys.exit(main())