bbqddt2/Antigravity
0
1# -*- coding: utf-8 -*-2"""3Antigravity 云端守护进程 V2.0 (腾讯云 CVM 专用)4=================================================5利用 32核/123GB 资源进行大规模并行计算。6 7调度计划:8 每30分钟: 数据采集 + 预测引擎9 每6小时: 非随机性检测10 每天: 滚动回测11 每周: 公式进化12 13用法:14 python cloud_runner_v2.py --daemon # 守护模式(默认)15 python cloud_runner_v2.py --once # 单轮运行16 python cloud_runner_v2.py --backtest # 仅回测17 python cloud_runner_v2.py --evolve # 仅进化18 python cloud_runner_v2.py --status # 显示状态19"""20import sys21import os22import time23import logging24import json25import signal26import traceback27from pathlib import Path28from datetime import datetime, timedelta29from typing import Optional, Dict, Any30 31# Fix encoding32sys.stdout.reconfigure(encoding="utf-8") if hasattr(sys.stdout, 'reconfigure') else None33 34_PROJECT_ROOT = Path(__file__).resolve().parent35sys.path.insert(0, str(_PROJECT_ROOT))36 37# ─── 日志 ─────────────────────────────────────────────38LOG_DIR = _PROJECT_ROOT / "logs"39LOG_DIR.mkdir(parents=True, exist_ok=True)40 41logger = logging.getLogger("CloudRunnerV2")42logger.setLevel(logging.INFO)43 44# 文件处理器45fh = logging.FileHandler(LOG_DIR / "cloud_runner_v2.log", encoding="utf-8")46fh.setFormatter(logging.Formatter("%(asctime)s [%(levelname)s] %(message)s"))47logger.addHandler(fh)48 49# 控制台处理器50sh = logging.StreamHandler(sys.stdout)51sh.setFormatter(logging.Formatter("%(asctime)s [%(levelname)s] %(message)s"))52logger.addHandler(sh)53 54# ─── 状态管理 ──────────────────────────────────────────55STATE_FILE = _PROJECT_ROOT / "cloud_runner_state.json"56 57 58def load_state() -> Dict[str, Any]:59 if STATE_FILE.exists():60 with open(STATE_FILE, "r", encoding="utf-8") as f:61 return json.load(f)62 return {"runs": [], "last_run": None, "total_cycles": 0, "errors": []}63 64 65def save_state(state: Dict[str, Any]):66 with open(STATE_FILE, "w", encoding="utf-8") as f:67 json.dump(state, f, indent=2, ensure_ascii=False)68 69 70def record_run(task: str, status: str, elapsed: float, error: str = None):71 state = load_state()72 entry = {73 "time": datetime.now().isoformat(),74 "task": task,75 "status": status,76 "elapsed_seconds": round(elapsed, 2),77 }78 if error:79 entry["error"] = error80 state["runs"].append(entry)81 # 只保留最近100条记录82 if len(state["runs"]) > 100:83 state["runs"] = state["runs"][-100:]84 state["last_run"] = entry["time"]85 state["total_cycles"] = state.get("total_cycles", 0) + 186 if status == "error":87 state.setdefault("errors", []).append(error)88 if len(state["errors"]) > 50:89 state["errors"] = state["errors"][-50:]90 save_state(state)91 92 93# ─── 单轮任务执行器 ─────────────────────────────────────94 95def run_data_collection() -> bool:96 """数据采集"""97 logger.info(" [1/4] 数据采集...")98 try:99 from data_updater_v2 import main as updater_main100 updater_main()101 logger.info(" ✅ 数据采集完成")102 return True103 except Exception as e:104 logger.warning(f" ⚠️ 数据采集失败: {e}")105 return False106 107 108def run_prediction() -> bool:109 """预测引擎"""110 logger.info(" [2/4] 预测引擎...")111 try:112 from orchestrate import run_full_pipeline113 result = run_full_pipeline(top_k=5)114 target = result.get("target_period", "?")115 logger.info(f" ✅ 预测完成: 目标 #{target}")116 return True117 except Exception as e:118 logger.warning(f" ⚠️ 预测引擎失败: {e}")119 return False120 121 122def run_nonrandomness() -> bool:123 """非随机性检测"""124 logger.info(" [3/4] 非随机性检测...")125 try:126 from nonrandomness_detector import run_detection127 run_detection()128 logger.info(" ✅ 非随机性检测完成")129 return True130 except Exception as e:131 logger.warning(f" ⚠️ 非随机性检测失败: {e}")132 return False133 134 135def run_backtest() -> bool:136 """滚动回测"""137 logger.info(" [4/4] 滚动回测...")138 try:139 from walkforward_backtest_v2 import run as backtest_run140 backtest_run()141 logger.info(" ✅ 滚动回测完成")142 return True143 except Exception as e:144 logger.warning(f" ⚠️ 滚动回测失败: {e}")145 return False146 147 148def run_full_cycle():149 """运行完整调度周期"""150 logger.info("=" * 60)151 logger.info(" 云端守护进程 V2.0 — 完整调度周期")152 logger.info(f" 主机: {os.popen('hostname').read().strip()}")153 logger.info(f" CPU: {os.popen('nproc').read().strip()} 核 | "154 f"内存: {os.popen(\"free -g | awk '/Mem:/{print $2}'\").read().strip()}G")155 logger.info("=" * 60)156 157 start = time.time()158 results = {}159 160 # 数据采集161 results["data_collection"] = run_data_collection()162 163 # 预测引擎(依赖数据)164 if results["data_collection"]:165 results["prediction"] = run_prediction()166 else:167 results["prediction"] = False168 logger.warning(" 跳过预测引擎(数据未更新)")169 170 elapsed = time.time() - start171 logger.info(f"\n 📊 周期完成: {elapsed:.1f}s | "172 f"成功: {sum(1 for v in results.values() if v)}/{len(results)}")173 174 record_run("full_cycle", "ok" if all(results.values()) else "partial", elapsed)175 176 return results177 178 179def run_single_task(task_name: str) -> bool:180 """运行单个任务"""181 logger.info(f"运行单个任务: {task_name}")182 start = time.time()183 try:184 if task_name == "data_collection":185 result = run_data_collection()186 elif task_name == "prediction":187 result = run_prediction()188 elif task_name == "nonrandomness":189 result = run_nonrandomness()190 elif task_name == "backtest":191 result = run_backtest()192 else:193 logger.error(f"未知任务: {task_name}")194 return False195 elapsed = time.time() - start196 record_run(task_name, "ok" if result else "error", elapsed)197 return result198 except Exception as e:199 elapsed = time.time() - start200 record_run(task_name, "error", elapsed, str(e))201 raise202 203 204def show_status():205 """显示云端状态"""206 state = load_state()207 print("\n" + "=" * 60)208 print(" Antigravity 云端守护进程状态")209 print("=" * 60)210 print(f"\n 总运行周期: {state.get('total_cycles', 0)}")211 print(f" 上次运行: {state.get('last_run', '从未')}")212 213 if state.get("runs"):214 print(f"\n 最近5次运行:")215 for run in state["runs"][-5:]:216 status_icon = "✅" if run["status"] == "ok" else "❌"217 print(f" {status_icon} [{run['time']}] {run['task']} "218 f"({run['elapsed_seconds']}s)")219 220 if state.get("errors"):221 print(f"\n 最近错误:")222 for err in state["errors"][-3:]:223 print(f" ⚠️ {err[:100]}")224 225 # 系统信息226 print(f"\n 系统信息:")227 print(f" 主机: {os.popen('hostname').read().strip()}")228 print(f" CPU: {os.popen('nproc').read().strip()} 核")229 print(f" 内存: {os.popen('free -g | awk \'/Mem:/{{print $2}}\'').read().strip()}G")230 print(f" 负载: {os.popen('cat /proc/loadavg').read().strip()}")231 232 # 最新预测233 pred_file = _PROJECT_ROOT / "latest_prediction.json"234 if pred_file.exists():235 with open(pred_file, "r", encoding="utf-8") as f:236 pred = json.load(f)237 target = pred.get("target_period", "?")238 fusion = pred.get("fusion", [])239 if fusion:240 best = fusion[0]241 reds = ", ".join(f"{r:02d}" for r in best.get("reds", []))242 blue = f"{best.get('blue', 0):02d}"243 print(f"\n 最新预测 (#{target}):")244 print(f" 🔴[{reds}] 🔵{blue}")245 print(f" 跨引擎支持: {best.get('cross_engine_support', '?')}")246 print("=" * 60)247 248 249# ─── 守护模式 ──────────────────────────────────────────250 251class CloudRunnerDaemon:252 """云端守护进程调度器"""253 254 def __init__(self, interval_minutes: int = 30):255 self.interval = interval_minutes256 self.running = True257 self.cycle_count = 0258 259 # 信号处理260 def signal_handler(sig, frame):261 logger.info("收到停止信号,优雅退出...")262 self.running = False263 264 signal.signal(signal.SIGINT, signal_handler)265 signal.signal(signal.SIGTERM, signal_handler)266 267 def should_run_nonrandomness(self) -> bool:268 """判断是否应该运行非随机性检测(每6小时)"""269 hour = datetime.now().hour270 return hour % 6 == 0271 272 def should_run_backtest(self) -> bool:273 """判断是否应该运行回测(每天周日)"""274 return datetime.now().weekday() == 6 # Sunday275 276 def run_loop(self):277 """主循环"""278 logger.info("=" * 60)279 logger.info(" 云端守护进程 V2.0 — 守护模式启动")280 logger.info(f" 调度间隔: {self.interval} 分钟")281 logger.info(f" 主机: {os.popen('hostname').read().strip()}")282 logger.info(f" 资源: {os.popen('nproc').read().strip()}核 / "283 f"{os.popen('free -g | awk \'/Mem:/{{print $2}}\'').read().strip()}G")284 logger.info("=" * 60)285 286 while self.running:287 self.cycle_count += 1288 logger.info(f"\n>>> 第 {self.cycle_count} 轮 <<<")289 290 # 标准周期:数据 + 预测291 run_full_cycle()292 293 # 特殊周期:非随机性检测294 if self.should_run_nonrandomness():295 logger.info("\n 📅 运行非随机性检测...")296 run_nonrandomness()297 298 # 特殊周期:回测299 if self.should_run_backtest():300 logger.info("\n 📅 运行滚动回测...")301 run_backtest()302 303 # 等待下一轮304 logger.info(f"\n 💤 等待 {self.interval} 分钟...")305 for _ in range(self.interval * 60):306 if not self.running:307 break308 time.sleep(1)309 310 # 优雅退出311 state = load_state()312 state["stopped_at"] = datetime.now().isoformat()313 save_state(state)314 logger.info(f"\n守护进程退出,共运行 {self.cycle_count} 轮")315 316 317# ─── CLI ──────────────────────────────────────────────318 319if __name__ == "__main__":320 import argparse321 parser = argparse.ArgumentParser(description="Antigravity 云端守护进程 V2.0")322 parser.add_argument("--daemon", action="store_true", help="守护模式(持续运行)")323 parser.add_argument("--once", action="store_true", help="单轮运行")324 parser.add_argument("--task", type=str, help="运行单个任务 (data_collection/prediction/nonrandomness/backtest)")325 parser.add_argument("--status", action="store_true", help="显示状态")326 parser.add_argument("--interval", type=int, default=30, help="守护模式间隔(分钟)")327 328 args = parser.parse_args()329 330 if args.status:331 show_status()332 elif args.task:333 run_single_task(args.task)334 elif args.daemon:335 daemon = CloudRunnerDaemon(interval_minutes=args.interval)336 daemon.run_loop()337 elif args.once:338 run_full_cycle()339 else:340 # 默认:守护模式341 daemon = CloudRunnerDaemon(interval_minutes=args.interval)342 daemon.run_loop()343 