Back to home page

EIC code displayed by LXR

 
 

    


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 (  # noqa: E402
0089     PandaTaskStepHandler,
0090 )
0091 from pandaserver.workflow.workflow_base import (  # noqa: E402
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         # which step produced it, which is what ${PARENT_TASKID} resolves through
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         # definition_json and parameters are the real backing attributes, so the spec's own
0118         # definition_json_map and get/set_parameter do the work and are exercised as written
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         # a real getTaskStatusSuperstatus returns a falsy value when the task is not in JEDI
0164         self._status = None if status is None else (status, superstatus or status)
0165 
0166     def getTaskWithID_JEDI(self, task_id, *args, **kwargs):
0167         # mirrors the real signature: (found, task_spec); None when the task is not in JEDI yet
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     # the simul step consumes {merge_evnt/EVNT} and produces one HITS dataset
0180     simul_params = copy.deepcopy(wfd["steps"]["simul"]["task_params"])
0181     # resolve ${WFID} the way registration would, so this works whether or not the description uses it
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     # Output dataset names still carry ${TASKID} at this point on purpose: the ID does not exist
0216     # until the insert allocates it, so insert_step_task resolves it inside the same transaction.
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     # only the task-level production label is guarded; see PRODUCTION_SOURCE_LABELS
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         # ... and is accepted once the submitter holds it
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     # non-production task labels need no role. "test" is included deliberately: a JEDI instance may
0251     # route it to the production refiner, but which task labels it accepts is its own configuration,
0252     # so the server does not gate it. "prod_test" is a job-level label and never reaches here.
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     # JEDI has no row yet, but DEFT does: expected for up to a refiner cycle, so this must be a
0302     # warning naming the situation, not an error suggesting the task was lost.
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     # the tri-state contract: None means "not checkable yet", which the caller treats as a wait
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     # None for the cases that are waits or skips, False only for genuine failures
0336     for label, setup, expected in [
0337         ("step not in a checkable status", lambda st, tb: setattr(st, "status", WFStepStatus.pending), None),
0338         # a flavor mismatch is a routing bug, not a wait: False so the caller reports it
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     # merge_evnt/EVNT was produced by step 3, whose task is 48810699
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()  # step 3 keeps target_id None
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             # the step's job parameters still reference this one, which input_data_dict does not
0462             # list, so it takes no part in working out the parent
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     # `expected` already names a bool | None earlier in this function
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     # an unknown status must still be reported rather than silently mapped
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())