run_task_executor.py 6.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168
  1. """Entry point to start the local Motor-CAD task executor (P4-M2 / P5-M2).
  2. Usage:
  3. python scripts/run_task_executor.py [--config path] [--instances N]
  4. [--interval N] [--mock]
  5. [--log-dir dir] [--log-level LVL]
  6. Configuration is loaded from executor_config.json (see executor_config.py
  7. for the lookup order and env-var overrides). This single entry point can
  8. launch N parallel executor instances via config["instances"] or --instances.
  9. NOTE: All strings must be ASCII only.
  10. """
  11. import argparse
  12. import logging
  13. import os
  14. import sys
  15. import time
  16. from datetime import datetime
  17. _SCRIPTS_DIR = os.path.dirname(os.path.abspath(__file__))
  18. if _SCRIPTS_DIR not in sys.path:
  19. sys.path.insert(0, _SCRIPTS_DIR)
  20. import executor_config # noqa: E402
  21. from task_executor import MotorCADTaskExecutor # noqa: E402
  22. def setup_logging(log_dir, level_name):
  23. """Configure a file + console logger for the executor."""
  24. os.makedirs(log_dir, exist_ok=True)
  25. log_file = os.path.join(
  26. log_dir, "executor_%s.log" % datetime.now().strftime("%Y%m%d_%H%M%S")
  27. )
  28. logging.basicConfig(
  29. level=getattr(logging, level_name, logging.INFO),
  30. format="%(asctime)s %(levelname)s %(message)s",
  31. handlers=[
  32. logging.FileHandler(log_file, encoding="utf-8"),
  33. logging.StreamHandler(),
  34. ],
  35. )
  36. return log_file
  37. def make_log_callback(kind):
  38. """Return an executor callback that logs and echoes to console."""
  39. logger = logging.getLogger("executor")
  40. def _cb(*args):
  41. msg = " ".join(str(a) for a in args)
  42. if kind == "progress":
  43. logger.info("%s", msg)
  44. elif kind == "complete":
  45. logger.info("%s", msg)
  46. else:
  47. logger.error("%s", msg)
  48. print(msg, flush=True)
  49. return _cb
  50. def build_executors(cfg):
  51. """Create cfg['instances'] MotorCADTaskExecutor objects."""
  52. multi = cfg["instances"] > 1
  53. executors = []
  54. for i in range(cfg["instances"]):
  55. ex = MotorCADTaskExecutor(
  56. web_base_url=cfg["web_base_url"],
  57. model_path=cfg["model_path"],
  58. enable_mock=cfg["enable_mock"],
  59. tool=cfg["tool"],
  60. enable_thermal=cfg["enable_thermal"],
  61. executor_id=("motorcad-executor-%s" % i) if multi else None,
  62. on_progress=make_log_callback("progress"),
  63. on_complete=make_log_callback("complete"),
  64. on_error=make_log_callback("error"),
  65. )
  66. executors.append(ex)
  67. return executors
  68. def main():
  69. parser = argparse.ArgumentParser(prog="pcb-afm-executor")
  70. parser.add_argument("--version", action="store_true",
  71. help="print version and exit")
  72. parser.add_argument("--self-test", action="store_true",
  73. help="run a mock single-point execution and exit "
  74. "(no web backend, no Motor-CAD)")
  75. parser.add_argument("--config", default=None,
  76. help="path to executor_config.json")
  77. parser.add_argument("--instances", type=int, default=None,
  78. help="parallel executor count (overrides config)")
  79. parser.add_argument("--interval", type=int, default=None,
  80. help="poll interval seconds (overrides config)")
  81. parser.add_argument("--mock", action="store_true",
  82. help="use mock solver (no Motor-CAD)")
  83. parser.add_argument("--log-dir", default=None,
  84. help="log output dir (overrides config)")
  85. parser.add_argument("--log-level", default=None,
  86. help="log level (overrides config)")
  87. args = parser.parse_args()
  88. if args.version:
  89. print("PCB-AFM Executor 1.1.0 (platform P5-M2, configurable)")
  90. return
  91. if args.self_test:
  92. import tempfile
  93. from task_executor import TaskExecutor
  94. _ex = TaskExecutor(task_dir=tempfile.mkdtemp(prefix="selftest_"),
  95. enable_mock=True)
  96. _ex.dispatch_task = lambda tid: True
  97. _captured = {}
  98. _ex.on_complete = lambda tid, res, met: _captured.update(
  99. {tid: (len(res), res[0].get("status") if res else None)})
  100. _ex.execute_task({
  101. "task_id": "selftest-1",
  102. "parameters": [{"airgap_mm": 1.0, "point_id": 1}],
  103. })
  104. _n, _st = _captured.get("selftest-1", (0, None))
  105. print("self-test OK: %d point(s), status=%s" % (_n, _st))
  106. return
  107. # Load configuration (CLI / env / config file / defaults).
  108. cfg = executor_config.load_config(cli_path=args.config)
  109. if args.instances is not None:
  110. cfg["instances"] = args.instances
  111. if args.interval is not None:
  112. cfg["poll_interval"] = args.interval
  113. if args.mock:
  114. cfg["enable_mock"] = True
  115. if args.log_dir:
  116. cfg["log_dir"] = args.log_dir
  117. if args.log_level:
  118. cfg["log_level"] = args.log_level
  119. executor_config.validate_config(cfg)
  120. log_file = setup_logging(cfg["log_dir"], cfg["log_level"])
  121. logger = logging.getLogger("executor")
  122. logger.info("config source: %s", cfg["config_source"])
  123. logger.info("model=%s", cfg["model_path"])
  124. logger.info("web_base_url=%s", cfg["web_base_url"])
  125. logger.info("instances=%s poll_interval=%s tool=%s mock=%s log=%s",
  126. cfg["instances"], cfg["poll_interval"], cfg["tool"],
  127. cfg["enable_mock"], log_file)
  128. executors = build_executors(cfg)
  129. threads = []
  130. for ex in executors:
  131. threads.append(ex.start_polling(interval=int(cfg["poll_interval"])))
  132. logger.info("executor started: %s", ex.executor_id)
  133. print("Task executor started. Model=%s" % cfg["model_path"], flush=True)
  134. print("Web base URL: %s" % cfg["web_base_url"], flush=True)
  135. print("Ctrl+C to stop.", flush=True)
  136. try:
  137. while any(t.is_alive() for t in threads):
  138. time.sleep(1)
  139. except KeyboardInterrupt:
  140. for ex in executors:
  141. ex.stop()
  142. ex.cleanup()
  143. logger.info("All executors stopped.")
  144. print("Executor stopped.", flush=True)
  145. if __name__ == "__main__":
  146. main()