run_task_executor.py 6.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169
  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. ambient_temperature=cfg["ambient_temperature"],
  62. executor_id=("motorcad-executor-%s" % i) if multi else None,
  63. on_progress=make_log_callback("progress"),
  64. on_complete=make_log_callback("complete"),
  65. on_error=make_log_callback("error"),
  66. )
  67. executors.append(ex)
  68. return executors
  69. def main():
  70. parser = argparse.ArgumentParser(prog="pcb-afm-executor")
  71. parser.add_argument("--version", action="store_true",
  72. help="print version and exit")
  73. parser.add_argument("--self-test", action="store_true",
  74. help="run a mock single-point execution and exit "
  75. "(no web backend, no Motor-CAD)")
  76. parser.add_argument("--config", default=None,
  77. help="path to executor_config.json")
  78. parser.add_argument("--instances", type=int, default=None,
  79. help="parallel executor count (overrides config)")
  80. parser.add_argument("--interval", type=int, default=None,
  81. help="poll interval seconds (overrides config)")
  82. parser.add_argument("--mock", action="store_true",
  83. help="use mock solver (no Motor-CAD)")
  84. parser.add_argument("--log-dir", default=None,
  85. help="log output dir (overrides config)")
  86. parser.add_argument("--log-level", default=None,
  87. help="log level (overrides config)")
  88. args = parser.parse_args()
  89. if args.version:
  90. print("PCB-AFM Executor 1.1.0 (platform P5-M2, configurable)")
  91. return
  92. if args.self_test:
  93. import tempfile
  94. from task_executor import TaskExecutor
  95. _ex = TaskExecutor(task_dir=tempfile.mkdtemp(prefix="selftest_"),
  96. enable_mock=True)
  97. _ex.dispatch_task = lambda tid: True
  98. _captured = {}
  99. _ex.on_complete = lambda tid, res, met: _captured.update(
  100. {tid: (len(res), res[0].get("status") if res else None)})
  101. _ex.execute_task({
  102. "task_id": "selftest-1",
  103. "parameters": [{"airgap_mm": 1.0, "point_id": 1}],
  104. })
  105. _n, _st = _captured.get("selftest-1", (0, None))
  106. print("self-test OK: %d point(s), status=%s" % (_n, _st))
  107. return
  108. # Load configuration (CLI / env / config file / defaults).
  109. cfg = executor_config.load_config(cli_path=args.config)
  110. if args.instances is not None:
  111. cfg["instances"] = args.instances
  112. if args.interval is not None:
  113. cfg["poll_interval"] = args.interval
  114. if args.mock:
  115. cfg["enable_mock"] = True
  116. if args.log_dir:
  117. cfg["log_dir"] = args.log_dir
  118. if args.log_level:
  119. cfg["log_level"] = args.log_level
  120. executor_config.validate_config(cfg)
  121. log_file = setup_logging(cfg["log_dir"], cfg["log_level"])
  122. logger = logging.getLogger("executor")
  123. logger.info("config source: %s", cfg["config_source"])
  124. logger.info("model=%s", cfg["model_path"])
  125. logger.info("web_base_url=%s", cfg["web_base_url"])
  126. logger.info("instances=%s poll_interval=%s tool=%s mock=%s log=%s",
  127. cfg["instances"], cfg["poll_interval"], cfg["tool"],
  128. cfg["enable_mock"], log_file)
  129. executors = build_executors(cfg)
  130. threads = []
  131. for ex in executors:
  132. threads.append(ex.start_polling(interval=int(cfg["poll_interval"])))
  133. logger.info("executor started: %s", ex.executor_id)
  134. print("Task executor started. Model=%s" % cfg["model_path"], flush=True)
  135. print("Web base URL: %s" % cfg["web_base_url"], flush=True)
  136. print("Ctrl+C to stop.", flush=True)
  137. try:
  138. while any(t.is_alive() for t in threads):
  139. time.sleep(1)
  140. except KeyboardInterrupt:
  141. for ex in executors:
  142. ex.stop()
  143. ex.cleanup()
  144. logger.info("All executors stopped.")
  145. print("Executor stopped.", flush=True)
  146. if __name__ == "__main__":
  147. main()