Back to home page

EIC code displayed by LXR

 
 

    


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

0001 """
0002 Check that every task buffer method the workflow code calls actually exists on TaskBuffer.
0003 
0004 The other harnesses drive the workflow against a fake task buffer, which implements whatever the
0005 code happens to call -- so a method that exists on the DB proxy but was never wrapped on TaskBuffer
0006 is invisible to them. That is exactly how on_all_inputs_done shipped calling updateTask_JEDI,
0007 release_task_on_hold and push_task_trigger_message, none of which TaskBuffer exposed: the hook raised
0008 AttributeError inside its own except clause on every cycle, so a workflow step's holdup was never
0009 released and the chain stalled with the error only visible in the handler log.
0010 
0011 TaskBuffer has no __getattr__, so an unwrapped proxy method is an AttributeError at runtime rather
0012 than a delegated call. This test parses the sources (no DB, no imports) and compares.
0013 
0014 Run from the repository root:  python3 pandaserver/workflow/examples/taskbuffer_interface_test.py
0015 """
0016 
0017 import ast
0018 import os
0019 import pathlib
0020 import re
0021 import sys
0022 
0023 REPO_ROOT = pathlib.Path(os.path.abspath(os.path.join(os.path.dirname(__file__), "..", "..", "..")))
0024 TASK_BUFFER = REPO_ROOT / "pandaserver/taskbuffer/TaskBuffer.py"
0025 # the workflow sources whose self.tbif calls must be satisfied by TaskBuffer
0026 WORKFLOW_SOURCES = [
0027     "pandaserver/workflow/workflow_core.py",
0028     "pandaserver/workflow/step_handler_plugins/panda_task_step_handler.py",
0029     "pandaserver/workflow/step_handler_plugins/base_step_handler.py",
0030     "pandaserver/workflow/data_handler_plugins/panda_task_data_handler.py",
0031     "pandaserver/workflow/data_handler_plugins/ddm_collection_data_handler.py",
0032     "pandaserver/workflow/data_handler_plugins/base_data_handler.py",
0033 ]
0034 
0035 
0036 def task_buffer_methods() -> set[str]:
0037     tree = ast.parse(TASK_BUFFER.read_text())
0038     cls = next(n for n in tree.body if isinstance(n, ast.ClassDef) and n.name == "TaskBuffer")
0039     return {f.name for f in cls.body if isinstance(f, ast.FunctionDef)}
0040 
0041 
0042 def has_getattr_delegation() -> bool:
0043     tree = ast.parse(TASK_BUFFER.read_text())
0044     cls = next(n for n in tree.body if isinstance(n, ast.ClassDef) and n.name == "TaskBuffer")
0045     return any(isinstance(f, ast.FunctionDef) and f.name in ("__getattr__", "__getattribute__") for f in cls.body)
0046 
0047 
0048 def called_methods(path: pathlib.Path) -> dict[str, list[int]]:
0049     """Map method name -> sorted line numbers, for every self.tbif.<name>( call."""
0050     calls: dict[str, list[int]] = {}
0051     for lineno, line in enumerate(path.read_text().splitlines(), start=1):
0052         for name in re.findall(r"self\.tbif\.([A-Za-z_][A-Za-z0-9_]*)\s*\(", line):
0053             calls.setdefault(name, []).append(lineno)
0054     return calls
0055 
0056 
0057 def main():
0058     failures = 0
0059     available = task_buffer_methods()
0060     print(f"\nTaskBuffer exposes {len(available)} methods; __getattr__ delegation: {has_getattr_delegation()}")
0061     if has_getattr_delegation():
0062         print("  NOTE: TaskBuffer now delegates unknown attributes, so this test is advisory only")
0063 
0064     total = 0
0065     for rel in WORKFLOW_SOURCES:
0066         path = REPO_ROOT / rel
0067         if not path.exists():
0068             print(f"  FAIL  {rel} does not exist")
0069             failures += 1
0070             continue
0071         calls = called_methods(path)
0072         total += len(calls)
0073         missing = {n: ls for n, ls in calls.items() if n not in available}
0074         status = "PASS" if not missing else "FAIL"
0075         print(f"  {status}  {rel}  ({len(calls)} distinct calls)")
0076         for name, lines in sorted(missing.items()):
0077             print(f"          MISSING on TaskBuffer: {name}  (called at line{'s' if len(lines) > 1 else ''} {', '.join(map(str, lines))})")
0078             failures += 1
0079 
0080     print(f"\nchecked {total} distinct self.tbif call sites across {len(WORKFLOW_SOURCES)} modules")
0081 
0082     # the three that were missing, called by on_all_inputs_done -- named explicitly so a revert is caught
0083     print("\nthe methods whose absence stalled workflow 130:")
0084     for name in ("updateTask_JEDI", "release_task_on_hold", "push_task_trigger_message"):
0085         present = name in available
0086         print(f"  {'PASS' if present else 'FAIL'}  {name}")
0087         if not present:
0088             failures += 1
0089 
0090     print(f"\n{'ALL CHECKS PASSED' if not failures else f'{failures} CHECK(S) FAILED'}")
0091     return 1 if failures else 0
0092 
0093 
0094 if __name__ == "__main__":
0095     sys.exit(main())