import subprocess from pathlib import Path import logging import os import json import time import threading from datetime import datetime from typing import List, Optional, Dict logger = logging.getLogger(__name__) PROJECT_ROOT = Path(__file__).resolve().parents[4] if os.getenv('PROJECT_ROOT'): PROJECT_ROOT = Path(os.getenv('PROJECT_ROOT')) _running_processes: Dict[str, dict] = {} _lock = threading.Lock() def _get_cmd(topic_ids: List[str] = None): script_path = PROJECT_ROOT / "scripts" / "compliance_optimizer.py" if not script_path.exists(): raise FileNotFoundError(f"Optimizer script not found: {script_path}") venv_python = PROJECT_ROOT / "platform" / "backend" / "venv" / "bin" / "python" cmd = [str(venv_python), str(script_path)] if venv_python.exists() else ["python3", str(script_path)] if topic_ids: cmd.extend(["--topic-ids", ','.join(topic_ids)]) else: cmd.append("--today-only") return cmd def run_optimizer(topic_ids: List[str] = None): """非阻塞:后台启动合规审查脚本,立即返回""" with _lock: existing = _running_processes.get("optimizer", {}) proc = existing.get("process") if proc and proc.poll() is None: raise RuntimeError("已有合规审查任务正在运行,请等待完成后再试") cmd = _get_cmd(topic_ids) proc = subprocess.Popen(cmd, stdout=subprocess.PIPE, stderr=subprocess.PIPE, text=True, cwd=str(PROJECT_ROOT)) _running_processes["optimizer"] = {"pid": proc.pid, "started_at": time.time(), "process": proc, "topic_ids": topic_ids} logger.info("Optimizer started in background (pid=%s, topic_ids=%s)", proc.pid, topic_ids) return {"ok": True, "pid": proc.pid, "topic_ids": topic_ids, "_proc": proc} def run_optimizer_blocking(topic_ids: List[str] = None, timeout: int = 600): """阻塞版:带超时,给定时任务使用""" cmd = _get_cmd(topic_ids) result = subprocess.run(cmd, cwd=str(PROJECT_ROOT), capture_output=True, text=True, timeout=timeout) if result.returncode != 0: logger.error(f"Optimizer failed: {result.stderr}") return {"ok": False, "error": result.stderr} report_date = datetime.now().strftime("%Y-%m-%d") report_path = PROJECT_ROOT / "automation" / "data" / "drafts" / report_date / "optimization_report.json" if report_path.exists(): report = json.loads(report_path.read_text(encoding='utf-8')) return {"ok": True, "report": report} else: logger.warning(f"Report not found: {report_path}") return {"ok": True, "report": None, "stdout": result.stdout} def get_optimizer_status() -> Optional[Dict]: info = _running_processes.get("optimizer") if not info: return None proc = info["process"] if "finished" in info: return {"status": "completed" if info.get("returncode") == 0 else "failed", "pid": info["pid"], "elapsed": round(time.time() - info["started_at"], 1), "returncode": info.get("returncode")} if proc.poll() is not None: info["returncode"] = proc.returncode info["finished"] = True elapsed = time.time() - info["started_at"] return {"status": "completed" if proc.returncode == 0 else "failed", "pid": info["pid"], "elapsed": round(elapsed, 1), "returncode": proc.returncode} return {"status": "running", "pid": info["pid"], "elapsed": round(time.time() - info["started_at"], 1)}