| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169 |
- """Entry point to start the local Motor-CAD task executor (P4-M2 / P5-M2).
- Usage:
- python scripts/run_task_executor.py [--config path] [--instances N]
- [--interval N] [--mock]
- [--log-dir dir] [--log-level LVL]
- Configuration is loaded from executor_config.json (see executor_config.py
- for the lookup order and env-var overrides). This single entry point can
- launch N parallel executor instances via config["instances"] or --instances.
- NOTE: All strings must be ASCII only.
- """
- import argparse
- import logging
- import os
- import sys
- import time
- from datetime import datetime
- _SCRIPTS_DIR = os.path.dirname(os.path.abspath(__file__))
- if _SCRIPTS_DIR not in sys.path:
- sys.path.insert(0, _SCRIPTS_DIR)
- import executor_config # noqa: E402
- from task_executor import MotorCADTaskExecutor # noqa: E402
- def setup_logging(log_dir, level_name):
- """Configure a file + console logger for the executor."""
- os.makedirs(log_dir, exist_ok=True)
- log_file = os.path.join(
- log_dir, "executor_%s.log" % datetime.now().strftime("%Y%m%d_%H%M%S")
- )
- logging.basicConfig(
- level=getattr(logging, level_name, logging.INFO),
- format="%(asctime)s %(levelname)s %(message)s",
- handlers=[
- logging.FileHandler(log_file, encoding="utf-8"),
- logging.StreamHandler(),
- ],
- )
- return log_file
- def make_log_callback(kind):
- """Return an executor callback that logs and echoes to console."""
- logger = logging.getLogger("executor")
- def _cb(*args):
- msg = " ".join(str(a) for a in args)
- if kind == "progress":
- logger.info("%s", msg)
- elif kind == "complete":
- logger.info("%s", msg)
- else:
- logger.error("%s", msg)
- print(msg, flush=True)
- return _cb
- def build_executors(cfg):
- """Create cfg['instances'] MotorCADTaskExecutor objects."""
- multi = cfg["instances"] > 1
- executors = []
- for i in range(cfg["instances"]):
- ex = MotorCADTaskExecutor(
- web_base_url=cfg["web_base_url"],
- model_path=cfg["model_path"],
- enable_mock=cfg["enable_mock"],
- tool=cfg["tool"],
- enable_thermal=cfg["enable_thermal"],
- ambient_temperature=cfg["ambient_temperature"],
- executor_id=("motorcad-executor-%s" % i) if multi else None,
- on_progress=make_log_callback("progress"),
- on_complete=make_log_callback("complete"),
- on_error=make_log_callback("error"),
- )
- executors.append(ex)
- return executors
- def main():
- parser = argparse.ArgumentParser(prog="pcb-afm-executor")
- parser.add_argument("--version", action="store_true",
- help="print version and exit")
- parser.add_argument("--self-test", action="store_true",
- help="run a mock single-point execution and exit "
- "(no web backend, no Motor-CAD)")
- parser.add_argument("--config", default=None,
- help="path to executor_config.json")
- parser.add_argument("--instances", type=int, default=None,
- help="parallel executor count (overrides config)")
- parser.add_argument("--interval", type=int, default=None,
- help="poll interval seconds (overrides config)")
- parser.add_argument("--mock", action="store_true",
- help="use mock solver (no Motor-CAD)")
- parser.add_argument("--log-dir", default=None,
- help="log output dir (overrides config)")
- parser.add_argument("--log-level", default=None,
- help="log level (overrides config)")
- args = parser.parse_args()
- if args.version:
- print("PCB-AFM Executor 1.1.0 (platform P5-M2, configurable)")
- return
- if args.self_test:
- import tempfile
- from task_executor import TaskExecutor
- _ex = TaskExecutor(task_dir=tempfile.mkdtemp(prefix="selftest_"),
- enable_mock=True)
- _ex.dispatch_task = lambda tid: True
- _captured = {}
- _ex.on_complete = lambda tid, res, met: _captured.update(
- {tid: (len(res), res[0].get("status") if res else None)})
- _ex.execute_task({
- "task_id": "selftest-1",
- "parameters": [{"airgap_mm": 1.0, "point_id": 1}],
- })
- _n, _st = _captured.get("selftest-1", (0, None))
- print("self-test OK: %d point(s), status=%s" % (_n, _st))
- return
- # Load configuration (CLI / env / config file / defaults).
- cfg = executor_config.load_config(cli_path=args.config)
- if args.instances is not None:
- cfg["instances"] = args.instances
- if args.interval is not None:
- cfg["poll_interval"] = args.interval
- if args.mock:
- cfg["enable_mock"] = True
- if args.log_dir:
- cfg["log_dir"] = args.log_dir
- if args.log_level:
- cfg["log_level"] = args.log_level
- executor_config.validate_config(cfg)
- log_file = setup_logging(cfg["log_dir"], cfg["log_level"])
- logger = logging.getLogger("executor")
- logger.info("config source: %s", cfg["config_source"])
- logger.info("model=%s", cfg["model_path"])
- logger.info("web_base_url=%s", cfg["web_base_url"])
- logger.info("instances=%s poll_interval=%s tool=%s mock=%s log=%s",
- cfg["instances"], cfg["poll_interval"], cfg["tool"],
- cfg["enable_mock"], log_file)
- executors = build_executors(cfg)
- threads = []
- for ex in executors:
- threads.append(ex.start_polling(interval=int(cfg["poll_interval"])))
- logger.info("executor started: %s", ex.executor_id)
- print("Task executor started. Model=%s" % cfg["model_path"], flush=True)
- print("Web base URL: %s" % cfg["web_base_url"], flush=True)
- print("Ctrl+C to stop.", flush=True)
- try:
- while any(t.is_alive() for t in threads):
- time.sleep(1)
- except KeyboardInterrupt:
- for ex in executors:
- ex.stop()
- ex.cleanup()
- logger.info("All executors stopped.")
- print("Executor stopped.", flush=True)
- if __name__ == "__main__":
- main()
|