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
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
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())