CoolFace
Apppublic

bbqddt2/Antigravity

sourceHugging Faceupdated 2mo agoView on Hugging Face
0likes
cloud_runner_v2.py343 linesDownload Raw Back to root
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