Type/to search

TradFi 矩阵统计套利

Common strategy
Created: 2026-07-29 11:24:02
Last modified: in 8 hours
2
Follow
502
Followers

TradFi 矩阵统计套利 v1.0 — 策略描述

一句话概括

在交易所的 TradFi 永续合约(美股/指数类,Gate.io contract_type=stocks)上,用低秩矩阵分解 + 零空间投影寻找由共同因子驱动的一篮子标的,实时监测其偏离稳定关系的程度,在"结构错位而非板块异动"时做均值回复交易。默认 paper(模拟盘),不真实下单。


核心思想

把 N 个标的的对数价格排成时间×标的矩阵 X (T×N)。若这些标的由 k 个共同趋势驱动,则 X 近似低秩:X = UΣV'rank(X) ≈ k

  • V 的前 k 个右奇异向量 → 共同因子空间(板块整体在动)
  • 后 N−k 个右奇异向量 → 近似零空间,其组合 Xw 近似平稳,即可交易的套利篮子

一个篮子同时给出 r = N−k 个正交的均值回复方向 —— 这就是"矩阵套利",而非传统配对套利。

稳定性判据(Davis-Kahan / Wedin sinθ 定理):子空间在扰动下的旋转 ‖sinΘ‖ ≤ ‖E‖ / (σ_k − σ_{k+1})。奇异值间隙 gap = σ_k/σ_{k+1} 越大,零空间越抗噪、越可交易。策略把 gap滚动子空间主夹角作为篮子准入的硬门槛。


流水线(5 个阶段)

阶段做什么关键闸门
0 合约池拉全市场行情,筛 stocks 类永续排除 pre-IPO 合成合约、OI/用户数下限、持仓量分位截断、去杠杆版本(相关高+波动比大)
1 候选篮子收益率相关矩阵 → Ward 层次聚类,切出 [6,14] 规模的簇候选不足时启用 kNN 兜底,按平均相关性排序截断
2 稳定空间检验对每个候选做四道检验① MP 定因子数 k ② 奇异值间隙 ≥1.3 ③ 子空间主夹角 ≤20° ④ 零空间组合 ADF 平稳 + OU 半衰期 ≤200bar
3 实时监测每轮拉一次全市场行情算两个统计量Q-score(零空间投影的 Mahalanobis 距离,度量结构错位)+ F-shock(因子得分一阶差分距离,度量板块异动)
4 样本外验证报警后追踪 Q 是否在 TRACK_BARS 内回落回归率 / 回归耗时 / 最大不利偏移 —— 判断策略前提是否成立的唯一依据

两个监测统计量(关键设计)

经典 MSPC 的 T²/Q 假设变量平稳,而对数价格是 I(1) 的,不能直接套用,因此做了改造:

  • Q-score ~ χ²(N−k):当前价格在零空间的投影,用训练期残差协方差做 Mahalanobis 距离,并按日内季节性(UTC 小时分槽)归一化。度量"结构错位"——这是要交易的方向
  • F-shock ~ χ²(k):因子得分差分的 Mahalanobis 距离,度量"板块整体剧烈异动"。

开仓条件:Q-score 超限 且 F-shock 正常。 若 F-shock 同时超限,说明新信息正在定价(不是错位),此时做均值回复是"接飞刀"。

此外还有事件跳空过滤:单轮 Q 跳变过大判为跳空,进入冷却期禁止开仓。


异常归因(LOO)

残差是组合属性,Q 超限只说明"这堆错位了"。定位用 leave-one-out:用其余 N−1 个标的对标的 i 做 OLS 重构,重构残差相对训练期标准差的 z 值 即 i 的特有偏离(对应 MSPC 的 contribution plot)。开仓方向由归因最强的标的 z 符号决定。


模拟盘执行

  • 按 LOO 权重构造多空腿,预算 = 权益 × 杠杆 × 单次占比,整数张取整
  • 对冲保真度检查:取整后实际权重与目标权重的余弦相似度须 ≥ 门槛,否则放弃(避免掉腿失衡)
  • 成本建模:按各合约 Info 读取的 taker 费率 + 滑点系数;资金费率按秒累计(多空组合中不自动抵消,是 spread 的确定性漂移)
  • 离场:Q 回落至离场线 / 超过 N 倍半衰期 / 保证金亏损止损 / 模型失效
  • 持仓用头寸自带的模型快照评估 Q,与篮子重建解耦

已知风险(策略作者自陈)

  1. TradFi 永续无法与真实美股实物交割,缺少强制收敛机制——约束价格的只是资金费率、做市商库存与报价习惯,弱于真套利。
  2. 资金费率在多空组合中不抵消,构成 spread 的确定性漂移。
  3. 合约上线时间短,有效样本可能不足一个完整财报周期。

使用方式

  • TRADE_MODE 默认 paper(live 仅占位未实现)
  • 命令交互:refresh:universe(刷新合约池)、rebuild:baskets(重建篮子)、flat:all(全平)、reset:paper(重置账户)、clear:store(清空状态)
  • 仪表板首屏即为样本外验证统计——这是判断"均值回复前提是否成立"的唯一客观依据;其余为账户、稳定子空间、实时监测、持仓详情、交易流水、去重与淘汰记录。

这个策略本质是监测/研究工具 + 模拟盘验证器,设计哲学非常克制:把"结论正确"让位给"前提可证伪",用样本外回归率而非样本内 z 回归来自我检验。

Source
Python
# -*- coding: utf-8 -*-
"""
TradFi 矩阵统计套利 v1.0(监测 + 模拟盘)
================================================================================

【核心思想】
    横轴时间、纵轴标的,构造对数价格矩阵 X (T×N)。
    若这 N 个标的由 k 个共同趋势驱动,则 X 近似低秩:

        X = U Σ V'      rank(X) ≈ k

    V 的后 N−k 个右奇异向量张成的子空间 = 近似零空间。
    任取其中向量 w,组合 Xw 近似平稳 —— 这就是可交易的套利组合。
    (Stock-Watson 共同趋势表示的几何形式)

    一个篮子给出 r = N−k 个相互正交的均值回复方向,
    因此是"矩阵套利"而非配对套利:同时监测 r 维残差空间。

【为什么零空间是稳定的】
    Davis-Kahan / Wedin sinθ 定理:子空间在扰动 E 下的旋转满足

        ||sinΘ|| ≤ ||E|| / (σ_k − σ_{k+1})

    奇异值间隙就是稳定性的定量度量。gap 大 → 子空间抗噪 → 关系可交易。
    gap 小 → 特征向量在方向间翻转 → 权重每次重估都乱跳。
    本策略把 gap 与滚动子空间主夹角作为篮子准入的硬门槛。

【两个监测统计量(关键设计)】
    经典 MSPC 的 T²/Q 假设变量平稳,而对数价格是 I(1) 的,不能直接套用。
    因此本策略做了改造:

    Q-score  : 当前价格向量在零空间的投影 r_t,用训练窗口的残差协方差
               做 Mahalanobis 距离。r_t 平稳(若协整成立),~ χ²(N−k)。
               → 度量"结构错位",这是要交易的方向。

    F-shock  : 因子得分的一阶差分 d_t 的 Mahalanobis 距离,~ χ²(k)。
               → 度量"板块整体剧烈异动"。

    只在 Q-score 超限 且 F-shock 正常 时开仓。
    F-shock 同时超限意味着新信息正在定价,不是错位,此时做均值回复是接刀。

【异常归因】
    残差是组合的属性,不是单个标的的属性。Q 超限只说明"这堆错位了"。
    定位用 leave-one-out:用其余 N−1 个标的对标的 i 做 OLS 重构,
    重构残差相对训练期标准差的 z 值即 i 的特有偏离。
    (对应 MSPC 的 contribution plot)

【样本外验证(一等公民)】
    z-score 在样本内必然回归 —— Σz_t = 0 是代数恒等式,随机游走也满足。
    唯一有意义的检验是样本外:用训练窗口估的参数,看之后的偏离会不会回归。
    仪表板首屏即为该统计:报警数、回归率、回归耗时分布、最大不利偏移。
    这张表是判断策略前提是否成立的唯一依据。

【已知风险】
    - 币安 TradFi 永续无法与真实美股实物交割,缺少强制收敛机制。
      约束这些价格的是资金费率、做市商库存与相对报价习惯,弱于真套利。
    - 资金费率在多空组合中不会自动抵消,构成 spread 的确定性漂移。
    - 合约上线时间短,有效样本可能不足以覆盖一个完整财报周期。

TRADE_MODE 默认 paper,不会真实下单。
================================================================================
"""

import json
import math
import time

import numpy as np
from scipy import stats as sps
from scipy.cluster.hierarchy import linkage, fcluster
from scipy.spatial.distance import squareform


# ══════════════════════════════════════════════════════════════════════════════
#  参数(平台注入优先,否则用默认值)
# ══════════════════════════════════════════════════════════════════════════════

def _p(name, default):
    v = globals().get(name, None)
    return default if v is None or v == "" else v


TRADE_MODE      = _p("TRADE_MODE", "paper")      # paper | live(live 未实现,仅占位)
INIT_CAPITAL    = float(_p("INIT_CAPITAL", 1000))

# —— 合约筛选(Gate.io 格式:Info.contract_type)——
# Gate.io 实测品类分布:stocks=284 indices=16 metals=12 commodities=3 forex=3
CONTRACT_TYPES  = str(_p("CONTRACT_TYPES", "stocks"))   # 逗号分隔;空 contract_type 一律排除
SKIP_PREMARKET  = bool(_p("SKIP_PREMARKET", True))      # 排除 pre-IPO 合成合约(ANTHROPIC/ANDURIL等)
MIN_OI_USDT     = float(_p("MIN_OI_USDT", 100000))      # 持仓量绝对下限(名义USDT)
MIN_USERS       = int(_p("MIN_USERS", 40))              # 多空用户数合计下限
OI_PCTL         = float(_p("OI_PCTL", 35))              # 再按持仓量分位数截断(0=不启用)
MAX_UNIVERSE    = int(_p("MAX_UNIVERSE", 120))          # 池子上限,按持仓量取前N(控制K线请求量)
DEDUP_RHO       = float(_p("DEDUP_RHO", 0.97))          # 近乎同一标的的相关系数门槛
DEDUP_VOL_RATIO = float(_p("DEDUP_VOL_RATIO", 1.40))    # 波动比超此值 → 判为杠杆版本,剔除
VOL_OUTLIER_MULT = float(_p("VOL_OUTLIER_MULT", 1.80)) # 篮子内波动/中位数超此值 → 疑杠杆品,剔除
MIN_HEDGE_FID   = float(_p("MIN_HEDGE_FID", 0.97))     # 整数张取整后实际权重与目标的余弦相似度下限
MIN_LEG_CONTRACTS = int(_p("MIN_LEG_CONTRACTS", 1))    # 实测:设高反而因掉腿降低保真度,交由保真度检查把关
FID_TARGET_CT   = int(_p("FID_TARGET_CT", 10))          # 每腿目标张数(10张≈5%取整误差),用于估算最小可交易敞口

# —— 成本 ——
TAKER_FEE       = float(_p("TAKER_FEE", 0.00075))       # 缺省值;实际按各合约 Info 读取
COST_MULT       = float(_p("COST_MULT", 1.0))           # 滑点放大系数

# —— 数据层 ——
KLINE_PERIOD    = PERIOD_M15
BAR_MINUTES     = int(_p("BAR_MINUTES", 15))
TRAIN_BARS      = int(_p("TRAIN_BARS", 2400))     # 期望训练窗口;不足时自动下调(见 effective_train)
MIN_TRAIN_BARS  = int(_p("MIN_TRAIN_BARS", 600))  # 训练窗口硬下限,低于此值不做回归
TN_RATIO        = float(_p("TN_RATIO", 20))      # 要求 T/N ≥ 此值,否则协方差估计不可靠
MAX_BAR_LEN     = int(_p("MAX_BAR_LEN", 3000))   # 向平台申请的K线条数上限
REF_PRESENCE    = float(_p("REF_PRESENCE", 0.60))  # 时间点被多少比例的标的覆盖才进参考网格
MIN_COVERAGE    = float(_p("MIN_COVERAGE", 0.95))  # 标的对参考网格的覆盖率下限,不足则剔除
FREEZE_FRAC     = float(_p("FREEZE_FRAC", 0.98))  # 仅剔除"全员冻结"的bar(标记价格未更新),不剔除清淡时段
VOL_SLOTS       = int(_p("VOL_SLOTS", 24))        # 波动率日内季节性分槽数(UTC小时)
MIN_SLOT_BARS   = int(_p("MIN_SLOT_BARS", 12))    # 每槽最少样本,不足则退回全局尺度
JUMP_MULT       = float(_p("JUMP_MULT", 0.60))    # 单轮Q跳变 > 报警线×此值 → 判为事件跳空
JUMP_COOL_S     = int(_p("JUMP_COOL_S", 4 * 3600))  # 判为跳空后的禁开仓冷却

# —— 篮子构建 ——
MIN_BASKET      = int(_p("MIN_BASKET", 6))
MAX_BASKET      = int(_p("MAX_BASKET", 14))
MAX_BASKETS     = int(_p("MAX_BASKETS", 6))      # 最多保留几个通过检验的篮子
LINKAGE_METHOD  = str(_p("LINKAGE_METHOD", "ward"))   # ward 比 average 更均衡,缓解链式效应
KNN_TRIGGER     = int(_p("KNN_TRIGGER", 12))     # 层次聚类候选少于此数 → 启用kNN兜底
MAX_CANDIDATES  = int(_p("MAX_CANDIDATES", 150)) # 候选上限(按平均相关性排序后截断)

# —— 稳定性门槛(阶段2 准入) ——
MIN_GAP         = float(_p("MIN_GAP", 1.30))     # σ_k / σ_{k+1}
MAX_ANGLE_DEG   = float(_p("MAX_ANGLE_DEG", 20.0))  # 子空间最大主夹角
ADF_CRIT        = float(_p("ADF_CRIT", -2.86))   # 5% 临界值(含常数项)
MAX_HALFLIFE    = float(_p("MAX_HALFLIFE", 200)) # 半衰期上限(bar)
MIN_STABILITY   = float(_p("MIN_STABILITY", 0.5))

# —— 监测阈值 ——
ENTRY_Q_P       = float(_p("ENTRY_Q_P", 0.995))  # Q-score 报警分位
EXIT_Q_P        = float(_p("EXIT_Q_P", 0.80))    # Q-score 回落分位(离场)
FSHOCK_P        = float(_p("FSHOCK_P", 0.99))    # F-shock 超此分位 → 禁止开仓
MIN_LOO_Z       = float(_p("MIN_LOO_Z", 2.0))    # 归因 z 门槛

# —— 交易 ——
LEVERAGE        = float(_p("LEVERAGE", 3))
GROSS_PER_TRADE = float(_p("GROSS_PER_TRADE", 0.25))  # 单次篮子交易占 权益×杠杆 的比例
MAX_OPEN        = int(_p("MAX_OPEN", 3))         # 最多同时持有几个篮子头寸
MIN_LEG_USDT    = float(_p("MIN_LEG_USDT", 3))
HOLD_HL_MULT    = float(_p("HOLD_HL_MULT", 3.0)) # 最大持仓 = 半衰期 × 该倍数
STOP_LOSS_PCT   = float(_p("STOP_LOSS_PCT", 18.0))  # 篮子头寸保证金亏损止损%

# —— 节奏 ——
POLL_S          = int(_p("POLL_S", 60))
HEARTBEAT_CYCLES = int(_p("HEARTBEAT_CYCLES", 10))  # 每N轮打一次心跳,避免"静默"与"卡死"无法区分
UNIVERSE_S      = int(_p("UNIVERSE_S", 6 * 3600))
REBUILD_S       = int(_p("REBUILD_S", 12 * 3600))
TRACK_BARS      = int(_p("TRACK_BARS", 192))     # 样本外验证:报警后追踪多少bar

STORE_KEY = "tradfi_matrix_statarb_v1"

# 内存缓存(不进 _G,避免存储膨胀)
_REC_CACHE = {}
_chart = None


# ══════════════════════════════════════════════════════════════════════════════
#  Store
# ══════════════════════════════════════════════════════════════════════════════

def new_store():
    return {
        "markets": {},          # symbol -> {contract, ...}
        "baskets": {},          # bid -> 通过检验的篮子(含模型参数)
        "rejected": [],         # 被淘汰的候选及原因
        "monitor": {},          # bid -> 最近一次监测结果
        "alerts": [],           # 报警记录(含样本外追踪)
        "positions": {},        # bid -> 模拟持仓
        "trades": [],
        "logs": [],
        "cash": INIT_CAPITAL,
        "init_capital": INIT_CAPITAL,
        "realized": 0.0,
        "last_universe": 0,
        "last_rebuild": 0,
    }


def load_store():
    s = _G(STORE_KEY)
    if not s:
        return new_store()
    for k, v in new_store().items():
        if k not in s:
            s[k] = v
    return s


def save_store(s):
    _G(STORE_KEY, s)


def record(store, step, data, color="#999999"):
    store["logs"].insert(0, {"time": int(time.time()), "step": step, "data": data})
    store["logs"] = store["logs"][:300]
    Log("STEP | {} | {}".format(step, json.dumps(data, ensure_ascii=False)[:400]), color)


# ══════════════════════════════════════════════════════════════════════════════
#  数学工具
# ══════════════════════════════════════════════════════════════════════════════

def adf_tstat(y, lags=4):
    """ADF 检验(含常数项)的 t 统计量。返回 None 表示样本不足。

    Δy_t = a + γ·y_{t-1} + Σ δ_i·Δy_{t-i} + ε
    检验 γ = 0(单位根)。t 统计量越负 → 越拒绝单位根 → 越平稳。
    """
    y = np.asarray(y, dtype=float)
    n = len(y)
    if n < lags + 20:
        return None
    dy = np.diff(y)
    rows = n - 1 - lags
    if rows < 10:
        return None
    Y = dy[lags:]
    X = [np.ones(rows), y[lags:n - 1]]
    for i in range(1, lags + 1):
        X.append(dy[lags - i:n - 1 - i])
    X = np.column_stack(X)
    try:
        beta, _, _, _ = np.linalg.lstsq(X, Y, rcond=None)
    except np.linalg.LinAlgError:
        return None
    resid = Y - X @ beta
    dof = rows - X.shape[1]
    if dof <= 0:
        return None
    s2 = float(resid @ resid) / dof
    try:
        xtx_inv = np.linalg.inv(X.T @ X)
    except np.linalg.LinAlgError:
        return None
    se = math.sqrt(max(s2 * xtx_inv[1, 1], 1e-18))
    return float(beta[1] / se)


def ou_halflife(y):
    """OU 过程半衰期(单位:bar)。dX = θ(μ−X)dt + σdW → 半衰期 = ln2/θ。"""
    y = np.asarray(y, dtype=float)
    if len(y) < 30:
        return None
    lag = y[:-1]
    dy = np.diff(y)
    X = np.column_stack([np.ones(len(lag)), lag])
    try:
        beta, _, _, _ = np.linalg.lstsq(X, dy, rcond=None)
    except np.linalg.LinAlgError:
        return None
    theta = -float(beta[1])
    if theta <= 1e-9:
        return None                      # 非回复(发散或随机游走)
    return math.log(2.0) / theta


def principal_angles_deg(Q1, Q2):
    """两个子空间(列正交基)之间的最大主夹角,单位度。"""
    if Q1.shape[1] == 0 or Q2.shape[1] == 0:
        return 90.0
    s = np.linalg.svd(Q1.T @ Q2, compute_uv=False)
    s = np.clip(s, -1.0, 1.0)
    return float(np.degrees(np.arccos(s.min())))


def mp_factor_count(returns):
    """Marchenko-Pastur 上边界法确定共同因子数 k。

    returns: (T, N) 收益率矩阵。相关矩阵中大于 (1+√q)² 的特征值视为真实因子,
    其余落在 MP 谱内的是纯噪声。
    """
    T, N = returns.shape
    if T <= N + 2:
        return max(1, N // 3)
    sd = returns.std(axis=0, ddof=1)
    sd[sd < 1e-12] = 1e-12
    Z = (returns - returns.mean(axis=0)) / sd
    C = (Z.T @ Z) / (T - 1)
    ev = np.linalg.eigvalsh(C)
    q = float(N) / float(T)
    upper = (1.0 + math.sqrt(q)) ** 2
    k = int(np.sum(ev > upper))
    return int(min(max(k, 1), N - 2))


def orth_basis(M):
    """取矩阵列空间的正交基(用 SVD,数值稳定)。"""
    u, s, _ = np.linalg.svd(M, full_matrices=False)
    tol = max(M.shape) * np.finfo(float).eps * (s[0] if len(s) else 0.0)
    return u[:, s > tol]


# ══════════════════════════════════════════════════════════════════════════════
#  阶段0:合约池 + 数据清洗
# ══════════════════════════════════════════════════════════════════════════════

def refresh_universe(store):
    record(store, "universe_start", {}, "#00AAFF")
    try:
        ms = exchange.GetMarkets()
    except Exception as e:
        record(store, "universe_failed", {"err": str(e)}, "#FF0000")
        return

    want = {t.strip().lower() for t in CONTRACT_TYPES.split(",") if t.strip()}
    markets, seen_types = {}, {}
    drops = {"type": 0, "status": 0, "premarket": 0, "oi": 0, "users": 0}

    for key, m in (ms or {}).items():
        if ".swap" not in key:
            continue
        info = m.get("Info", {}) or {}

        # —— Gate.io:contract_type 标记品类,股票类为 "stocks" ——
        ctype = str(info.get("contract_type", "") or "").lower()
        seen_types[ctype or "(空)"] = seen_types.get(ctype or "(空)", 0) + 1
        if not ctype or ctype not in want:
            drops["type"] += 1
            continue

        # —— 可交易性 ——
        if str(info.get("status", "trading")) != "trading" or info.get("in_delisting", False):
            drops["status"] += 1
            continue

        # —— 排除 pre-IPO 合成合约 ——
        if SKIP_PREMARKET and info.get("is_pre_market", False):
            drops["premarket"] += 1
            continue

        ct_val = float(m.get("CtVal", 1) or 1)
        mark = float(info.get("mark_price", 0) or info.get("last_price", 0) or 0)

        # —— 流动性闸门 ——
        oi_usdt = float(info.get("position_size", 0) or 0) * ct_val * mark
        users = int(info.get("long_users", 0) or 0) + int(info.get("short_users", 0) or 0)
        if oi_usdt < MIN_OI_USDT:
            drops["oi"] += 1
            continue
        if users < MIN_USERS:
            drops["users"] += 1
            continue

        sym = key.replace("_USDT.swap", "").replace(".swap", "").upper()
        markets[sym] = {
            "symbol": sym,
            "contract": key,
            "amountPrecision": int(m.get("AmountPrecision", 0) or 0),
            "pricePrecision": int(m.get("PricePrecision", 2) or 2),
            "ctVal": ct_val,
            "minQty": float(m.get("MinQty", 1) or 1),
            "takerFee": float(info.get("taker_fee_rate", TAKER_FEE) or TAKER_FEE),
            "fundingRate": float(info.get("funding_rate", 0) or 0),
            "fundingInterval": int(info.get("funding_interval", 28800) or 28800),
            "oiUsdt": round(oi_usdt, 0),
            "users": users,
        }

    # —— 相对流动性截断 ——
    ois = sorted(m["oiUsdt"] for m in markets.values())
    pctl_info = {}
    if ois:
        pctl_info = {"p10": round(np.percentile(ois, 10)), "p50": round(np.percentile(ois, 50)),
                     "p90": round(np.percentile(ois, 90)), "max": round(ois[-1])}
        if OI_PCTL > 0 and len(markets) > MIN_BASKET * 2:
            cut = float(np.percentile(ois, OI_PCTL))
            removed = [s for s, m in markets.items() if m["oiUsdt"] < cut]
            for s in removed:
                del markets[s]
            drops["oi_pctl"] = len(removed)

    if len(markets) > MAX_UNIVERSE:
        top = sorted(markets.items(), key=lambda kv: -kv[1]["oiUsdt"])[:MAX_UNIVERSE]
        drops["universe_cap"] = len(markets) - MAX_UNIVERSE
        markets = dict(top)

    store["markets"] = markets
    record(store, "universe_done", {
        "kept": len(markets), "dropped": drops,
        "持仓量分布USDT": pctl_info,
        "contract_type分布": dict(sorted(seen_types.items(), key=lambda x: -x[1])[:8]),
    }, "#00AAFF")


def fetch_records(contract):
    """取K线。优先用三参数形式显式指定条数。"""
    if contract in _REC_CACHE:
        return _REC_CACHE[contract]
    bars = []
    try:
        bars = exchange.GetRecords(contract, KLINE_PERIOD, MAX_BAR_LEN) or []
    except Exception:
        try:
            bars = exchange.GetRecords(contract, KLINE_PERIOD) or []
        except Exception as e:
            Log("GetRecords 失败 | {} | {}".format(contract, str(e)), "#FFAA00")
            bars = []
    _REC_CACHE[contract] = bars
    return bars


def build_price_matrix(store, symbols):
    """对齐 + 清洗,返回 (timestamps, log_price_matrix(T×N), kept_symbols)。

    新逻辑:按最少可用K线数进行对齐,而非硬性要求每个标的都有MIN_TRAIN_BARS根。
    这样即使某些标的历史短,只要池子里有足够的长序列标的,就能凑出长训练窗口。
    """
    series, counts = {}, []
    for s in symbols:
        m = store["markets"].get(s)
        if not m:
            continue
        bars = fetch_records(m["contract"])
        counts.append(len(bars))
        if len(bars) < 50:  # 硬下限:每个标的至少 50 bar(高于此值才能建立基本统计)
            continue
        series[s] = {int(b["Time"]): float(b["Close"]) for b in bars if b.get("Close", 0) > 0}

    if counts:
        cs = sorted(counts)
        Log("K线条数 | 请求上限:{} 实得 min={} 中位={} max={} | 达标(50+根):{}/{}".format(
            MAX_BAR_LEN, cs[0], cs[len(cs) // 2], cs[-1],
            len(series), len(counts)), "#888888")

    if len(series) < MIN_BASKET:
        Log("❌ 达标标的不足({}<{})。若中位数条数明显小于期望,说明平台K线长度受限。".format(
            len(series), MIN_BASKET), "#FF0000")
        return None, None, []

    # —— 时间戳对齐:参考网格 + 覆盖率筛选 ——
    n_sym = len(series)
    counts = {}
    for s in series:
        for t in series[s]:
            counts[t] = counts.get(t, 0) + 1
    ref = {t for t, c in counts.items() if c >= REF_PRESENCE * n_sym}

    # 灵活的样本长度调整
    if len(ref) < 50:
        Log("参考网格不足,放宽条件重试...", "#FFAA00")
        ref = {t for t, c in counts.items() if c >= max(1, 0.4 * n_sym)}

    if len(ref) < 50:
        Log("❌ 对齐失败:无足够重叠时间点", "#FF0000")
        return None, None, []

    kept, thin = [], []
    for s in sorted(series.keys()):
        cov = len(ref & set(series[s].keys())) / float(len(ref))
        (kept if cov >= MIN_COVERAGE else thin).append(s if cov >= MIN_COVERAGE
                                                       else (s, round(cov, 3)))
    if thin:
        Log("覆盖率不足剔除 {} 个(多为新上市): {}".format(
            len(thin), ", ".join("{}({:.0%})".format(a, b) for a, b in thin[:8])), "#999999")
    if len(kept) < MIN_BASKET:
        return None, None, []

    common = ref
    for s in kept:
        common &= set(series[s].keys())
    ts = sorted(common)

    # 根据实际数据长度动态调整下限
    actual_min_bars = min(MIN_TRAIN_BARS, max(100, len(ts) // 2))
    if len(ts) < actual_min_bars:
        Log("⚠️  可用bar({}) 不足目标({}),降低至 {}".format(
            len(ts), MIN_TRAIN_BARS, actual_min_bars), "#FFAA00")

    if len(ts) < 50:
        return None, None, []

    P = np.array([[series[s][t] for s in kept] for t in ts], dtype=float)
    L = np.log(P)

    # —— 只剔除全员冻结的 bar ——
    d = np.diff(L, axis=0)
    zero_frac = (np.abs(d) < 1e-12).mean(axis=1)
    keep_mask = np.concatenate([[True], zero_frac < FREEZE_FRAC])
    L = L[keep_mask]
    ts = [t for t, k in zip(ts, keep_mask) if k]

    if len(ts) < 50:
        return None, None, []
    return np.array(ts), L, kept


# ══════════════════════════════════════════════════════════════════════════════
#  阶段1:候选篮子(层次聚类)
# ══════════════════════════════════════════════════════════════════════════════

def build_candidate_baskets(store):
    """对收益率相关矩阵做平均连接层次聚类,切出规模在 [MIN,MAX] 的簇。"""
    symbols = sorted(store["markets"].keys())
    if len(symbols) < MIN_BASKET:
        record(store, "cluster_skip", {"reason": "universe too small", "n": len(symbols)}, "#FFAA00")
        return []

    ts, L, kept = build_price_matrix(store, symbols)
    if L is None:
        record(store, "cluster_skip", {"reason": "insufficient aligned bars"}, "#FFAA00")
        return []

    R = np.diff(L, axis=0)
    if R.shape[0] < 50:
        record(store, "cluster_skip", {"reason": "too few returns", "rows": R.shape[0]}, "#FFAA00")
        return []

    C = np.corrcoef(R, rowvar=False)
    C = np.nan_to_num(C, nan=0.0)
    np.fill_diagonal(C, 1.0)

    # —— 剔除杠杆版本 ——
    vol = R.std(axis=0, ddof=1)
    drop_idx, dedup_log = set(), []
    for i in range(len(kept)):
        for j in range(i + 1, len(kept)):
            if abs(C[i, j]) < DEDUP_RHO:
                continue
            hi, lo = (i, j) if vol[i] > vol[j] else (j, i)
            ratio = vol[hi] / max(vol[lo], 1e-12)
            if ratio >= DEDUP_VOL_RATIO:
                drop_idx.add(hi)
                dedup_log.append({"drop": kept[hi], "keep": kept[lo],
                                  "rho": round(float(C[i, j]), 4),
                                  "vol_ratio": round(float(ratio), 2)})
    if drop_idx:
        record(store, "leveraged_variants_dropped", {"n": len(drop_idx), "detail": dedup_log[:10]}, "#FFAA00")
        keep_cols = [i for i in range(len(kept)) if i not in drop_idx]
        kept = [kept[i] for i in keep_cols]
        L = L[:, keep_cols]
        R = R[:, keep_cols]
        C = C[np.ix_(keep_cols, keep_cols)]
        if len(kept) < MIN_BASKET:
            record(store, "cluster_skip", {"reason": "剔除杠杆版本后标的不足", "n": len(kept)}, "#FFAA00")
            return []
    store["dedup"] = dedup_log[:20]
    D = np.sqrt(np.maximum(2.0 * (1.0 - C), 0.0))
    np.fill_diagonal(D, 0.0)
    D = (D + D.T) / 2.0

    Z = linkage(squareform(D, checks=False), method=LINKAGE_METHOD)
    n = len(kept)
    pos = {sym: i for i, sym in enumerate(kept)}
    seen, cands, size_hist = set(), [], {}

    def add_cand(members, src):
        mt = tuple(sorted(members))
        if not (MIN_BASKET <= len(mt) <= MAX_BASKET) or mt in seen:
            return
        seen.add(mt)
        idx = [pos[m] for m in mt]
        sub = C[np.ix_(idx, idx)]
        avg_rho = float((sub.sum() - len(idx)) / (len(idx) * (len(idx) - 1)))
        cands.append({"members": list(mt), "avg_rho": round(avg_rho, 4), "src": src})

    for nclust in range(2, min(n, 80) + 1):
        labels = fcluster(Z, t=nclust, criterion="maxclust")
        for lab in set(labels):
            members = [kept[i] for i in range(n) if labels[i] == lab]
            size_hist[len(members)] = size_hist.get(len(members), 0) + 1
            add_cand(members, "hier")

    if len(cands) < KNN_TRIGGER:
        Log("层次聚类候选仅 {} 个(链式效应),启用 kNN 兜底".format(len(cands)), "#FFAA00")
        for i in range(n):
            order = np.argsort(-C[i])
            for size in {MIN_BASKET + 2, (MIN_BASKET + MAX_BASKET) // 2, MAX_BASKET}:
                add_cand([kept[j] for j in order[:size]], "knn")

    cands.sort(key=lambda x: -x["avg_rho"])
    cands = cands[:MAX_CANDIDATES]
    record(store, "candidates_built", {
        "universe": len(kept), "candidates": len(cands),
        "来源": {k: sum(1 for c in cands if c["src"] == k) for k in ("hier", "knn")},
        "簇规模分布": dict(sorted(size_hist.items())[:12]),
        "top": [{"n": len(c["members"]), "rho": c["avg_rho"], "src": c["src"]} for c in cands[:5]],
    }, "#00AAFF")
    return cands


# ══════════════════════════════════════════════════════════════════════════════
#  阶段2:稳定空间检验(准入闸门)
# ══════════════════════════════════════════════════════════════════════════════

def basket_id(members):
    h = 0
    for ch in "|".join(sorted(members)):
        h = (h * 131 + ord(ch)) & 0xFFFFFFFF
    return "B" + format(h, "07x")[:6].upper()


def validate_basket(store, members):
    """对一个候选篮子做四道检验,返回模型或 (None, 淘汰原因)。

    1. MP 定共同因子数 k
    2. 奇异值间隙 σ_k/σ_{k+1}     —— Davis-Kahan 界的分母
    3. 滚动子空间主夹角            —— 零空间是否随时间旋转
    4. 零空间组合的 ADF + OU 半衰期 —— 是否真的均值回复
    """
    ts, L, kept = build_price_matrix(store, members)
    if L is None or len(kept) < MIN_BASKET:
        return None, "数据不足"

    # 杠杆ETF漏网检测
    Rv = np.diff(L, axis=0)
    vols = Rv.std(axis=0, ddof=1)
    med_vol = float(np.median(vols))
    if med_vol > 0:
        out = [i for i in range(len(kept)) if vols[i] > VOL_OUTLIER_MULT * med_vol]
        if out:
            names = [kept[i] for i in out]
            keep_i = [i for i in range(len(kept)) if i not in set(out)]
            if len(keep_i) < MIN_BASKET:
                return None, "剔除波动异常成员{}后规模不足".format(names)
            Log("篮子内剔除波动异常成员(疑杠杆品): {} | 波动/中位={}".format(
                ",".join(names),
                ",".join("{:.1f}x".format(vols[i] / med_vol) for i in out)), "#FFAA00")
            kept = [kept[i] for i in keep_i]
            L = L[:, keep_i]

    T = L.shape[0]
    N = L.shape[1]
    need = max(50, int(TN_RATIO * N))  # 根据 T/N 比动态调整
    if T < need:
        return None, "有效bar不足({}<{}=max({},{}×{}))".format(
            T, need, 50, int(TN_RATIO), N)
    eff = min(TRAIN_BARS, T)

    train = L[-eff:]
    mu = train.mean(axis=0)
    Xc = train - mu

    # 1. 共同因子数
    k = mp_factor_count(np.diff(train, axis=0))
    r = N - k
    if r < 1:
        return None, "零空间维度为0"

    # 2. SVD + 奇异值间隙
    _, S, Vt = np.linalg.svd(Xc, full_matrices=False)
    if len(S) <= k or S[k] < 1e-12:
        return None, "奇异值退化"
    gap = float(S[k - 1] / S[k])
    if gap < MIN_GAP:
        return None, "奇异值间隙不足 {:.2f}<{:.2f}".format(gap, MIN_GAP)

    V = Vt.T
    Vf = V[:, :k]
    Vn = V[:, k:]

    # 3. 子空间稳定性
    half = eff // 2
    angles = []
    for seg in (train[:half], train[half:]):
        Sc = seg - seg.mean(axis=0)
        _, _, vt = np.linalg.svd(Sc, full_matrices=False)
        angles.append(principal_angles_deg(Vn, vt.T[:, k:]))
    max_angle = max(angles)
    if max_angle > MAX_ANGLE_DEG:
        return None, "子空间旋转 {:.1f}°>{:.1f}°".format(max_angle, MAX_ANGLE_DEG)

    # 4. 零空间组合的平稳性与回复速度
    Rn = Xc @ Vn
    adfs, hls = [], []
    for j in range(r):
        a = adf_tstat(Rn[:, j])
        h = ou_halflife(Rn[:, j])
        if a is not None:
            adfs.append(a)
        if h is not None:
            hls.append(h)
    if not adfs or not hls:
        return None, "平稳性检验失败"

    keep_dirs = [j for j in range(r) if (adf_tstat(Rn[:, j]) or 0) < ADF_CRIT]
    pass_adf = len(keep_dirs)
    if pass_adf == 0:
        return None, "无方向通过ADF(最优t={:.2f})".format(min(adfs))
    if pass_adf < r:
        Vn = Vn[:, keep_dirs]
        Rn = Xc @ Vn
        r = pass_adf
    halflife = float(np.median(hls))
    if halflife > MAX_HALFLIFE:
        return None, "半衰期过长 {:.0f}>{:.0f}bar".format(halflife, MAX_HALFLIFE)

    # —— 监测所需的训练期统计量 ——
    resid_cov = np.cov(Rn, rowvar=False)
    if r == 1:
        resid_cov = np.array([[float(resid_cov)]])
    resid_inv = np.linalg.pinv(resid_cov)

    # —— 日内季节性标准化 ——
    Q_train = np.einsum("ij,jk,ik->i", Rn, resid_inv, Rn)
    slots = ((ts[-eff:] // 3600000) % VOL_SLOTS).astype(int)
    q_med_all = float(np.median(Q_train)) or 1.0
    slot_scale = {}
    for h in range(VOL_SLOTS):
        sel = Q_train[slots == h]
        if len(sel) >= MIN_SLOT_BARS:
            slot_scale[str(h)] = round(max(float(np.median(sel)) / q_med_all, 0.05), 4)

    F = Xc @ Vf
    dF = np.diff(F, axis=0)
    f_cov = np.cov(dF, rowvar=False)
    if k == 1:
        f_cov = np.array([[float(f_cov)]])
    f_inv = np.linalg.pinv(f_cov)

    # LOO 归因模型
    loo = []
    for i in range(N):
        others = [j for j in range(N) if j != i]
        A = np.column_stack([np.ones(eff), train[:, others]])
        b, _, _, _ = np.linalg.lstsq(A, train[:, i], rcond=None)
        res = train[:, i] - A @ b
        loo.append({"i": i, "others": others, "beta": b.tolist(),
                    "sd": float(res.std(ddof=1)) or 1e-9})

    # —— 最小可交易敞口 ——
    last_px = np.exp(train[-1])
    gran = np.array([last_px[i] * store["markets"][kept[i]]["ctVal"] for i in range(N)])
    needs = []
    for rec in loo:
        b = np.array(rec["beta"])
        wv = np.zeros(N)
        wv[rec["i"]] = -1.0
        for pj, j in enumerate(rec["others"]):
            wv[j] = float(b[1 + pj])
        aw = np.abs(wv)
        nz = aw > 1e-9
        if nz.sum() < 2:
            continue
        needs.append(float(np.max(FID_TARGET_CT * gran[nz] * aw.sum() / aw[nz])))
    min_gross = float(np.median(needs)) if needs else float("inf")

    stability = float(np.clip(
        0.35 * np.clip((gap - 1.0) / 1.0, 0, 1)
        + 0.30 * np.clip(1.0 - max_angle / MAX_ANGLE_DEG, 0, 1)
        + 0.20 * (pass_adf / r)
        + 0.15 * np.clip(1.0 - halflife / MAX_HALFLIFE, 0, 1), 0, 1))

    model = {
        "members": kept,
        "N": N, "k": k, "r": r,
        "mu": mu.tolist(),
        "Vf": Vf.tolist(), "Vn": Vn.tolist(),
        "resid_inv": resid_inv.tolist(),
        "f_inv": f_inv.tolist(),
        "slot_scale": slot_scale,
        "loo": loo,
        "gap": round(gap, 3),
        "max_angle": round(max_angle, 2),
        "adf_best": round(min(adfs), 3),
        "adf_pass": pass_adf,
        "halflife": round(halflife, 1),
        "min_gross": round(min_gross, 1),
        "stability": round(stability, 3),
        "train_bars": eff,
        "q_entry": float(sps.chi2.ppf(ENTRY_Q_P, r)),
        "q_exit": float(sps.chi2.ppf(EXIT_Q_P, r)),
        "f_limit": float(sps.chi2.ppf(FSHOCK_P, k)),
        "built_at": int(time.time()),
    }
    return model, None


def rebuild_baskets(store):
    record(store, "rebuild_start", {}, "#00AAFF")
    _REC_CACHE.clear()

    cands = build_candidate_baskets(store)
    passed, rejected, used = [], [], set()

    for c in cands:
        members = c["members"]
        if len(set(members) & used) > len(members) // 2:
            continue
        model, reason = validate_basket(store, members)
        if model is None:
            rejected.append({"members": members[:6], "n": len(members), "reason": reason})
            Log("篮子淘汰 | n={} | {} | {}".format(len(members), reason,
                                                 ",".join(members[:6])), "#999999")
            continue
        if model["stability"] < MIN_STABILITY:
            rejected.append({"members": members[:6], "n": len(members),
                             "reason": "稳定性评分不足 {:.2f}".format(model["stability"])})
            continue
        passed.append(model)
        used |= set(model["members"])
        Log("篮子通过 | n={} | gap={} 夹角={}° ADF={} 半衰期={}bar 评分={} | {}".format(
            model["N"], model["gap"], model["max_angle"], model["adf_best"],
            model["halflife"], model["stability"], ",".join(model["members"])), "#00CC66")
        if len(passed) >= MAX_BASKETS:
            break

    passed.sort(key=lambda m: -m["stability"])
    store["baskets"] = {basket_id(m["members"]): m for m in passed}
    store["rejected"] = rejected[:40]
    record(store, "rebuild_done", {
        "passed": len(passed), "rejected": len(rejected),
        "detail": [{"id": bid, "n": m["N"], "k": m["k"], "r": m["r"],
                    "gap": m["gap"], "hl": m["halflife"], "score": m["stability"]}
                   for bid, m in store["baskets"].items()],
    }, "#00AAFF")


# ══════════════════════════════════════════════════════════════════════════════
#  阶段3:实时监测(Q-score / F-shock / LOO 归因)
# ══════════════════════════════════════════════════════════════════════════════

def fetch_all_tickers(store):
    """整轮只拉一次全市场行情,避免逐腿 GetTicker 触发限频。"""
    out = {}
    try:
        ts = exchange.GetTickers() or []
    except Exception as e:
        Log("GetTickers 异常:", str(e), "#FFAA00")
        return out
    by_contract = {}
    for t in ts:
        last = t.get("Last", 0) if isinstance(t, dict) else 0
        sym = t.get("Symbol", "") if isinstance(t, dict) else ""
        if last and last > 0:
            by_contract[sym] = float(last)
    for sym, m in store["markets"].items():
        p = by_contract.get(m["contract"])
        if p:
            out[sym] = p
    return out


def current_log_prices(members, tick):
    out, prices = [], {}
    for s in members:
        p = tick.get(s)
        if not p or p <= 0:
            return None, None
        prices[s] = p
        out.append(math.log(p))
    return np.array(out), prices


def monitor_basket(store, bid, model, prev, tick):
    x, prices = current_log_prices(model["members"], tick)
    if x is None:
        return None

    mu = np.array(model["mu"])
    Vn = np.array(model["Vn"])
    Vf = np.array(model["Vf"])
    xc = x - mu

    # Q-score
    rvec = xc @ Vn
    q_raw = float(rvec @ np.array(model["resid_inv"]) @ rvec)
    slot = str(int((time.time() // 3600) % VOL_SLOTS))
    scale = float(model.get("slot_scale", {}).get(slot, 1.0))
    q = q_raw / scale

    # F-shock
    f = xc @ Vf
    fshock = 0.0
    if prev is not None and prev.get("f") is not None:
        d = f - np.array(prev["f"])
        fshock = float(d @ np.array(model["f_inv"]) @ d)

    # LOO 归因
    contrib = []
    for rec in model["loo"]:
        b = np.array(rec["beta"])
        pred = b[0] + float(b[1:] @ x[rec["others"]])
        z = (x[rec["i"]] - pred) / rec["sd"]
        contrib.append({"symbol": model["members"][rec["i"]], "z": round(float(z), 3)})
    contrib.sort(key=lambda c: -abs(c["z"]))

    # —— 跳空过滤 ——
    now = int(time.time())
    jump_until = (prev or {}).get("jump_until", 0)
    if prev is not None and prev.get("q") is not None:
        if q - prev["q"] > model["q_entry"] * JUMP_MULT:
            jump_until = now + JUMP_COOL_S
            Log("⚡ 判定为事件跳空 | {} | Q {:.1f}→{:.1f} | 冷却 {}h".format(
                bid, prev["q"], q, JUMP_COOL_S // 3600), "#FFAA00")
    in_jump = now < jump_until

    return {
        "bid": bid,
        "q": round(q, 2),
        "q_raw": round(q_raw, 2),
        "slot_scale": round(scale, 3),
        "q_entry": round(model["q_entry"], 2),
        "q_exit": round(model["q_exit"], 2),
        "fshock": round(fshock, 2),
        "f_limit": round(model["f_limit"], 2),
        "alarm": q > model["q_entry"],
        "blocked": fshock > model["f_limit"] or in_jump,
        "block_reason": "因子异动" if fshock > model["f_limit"] else ("事件跳空冷却" if in_jump else ""),
        "jump_until": jump_until,
        "contrib": contrib[:5],
        "f": f.tolist(),
        "prices": prices,
        "at": now,
    }


def run_monitor(store):
    tick = fetch_all_tickers(store)
    results = {}
    for bid, model in store["baskets"].items():
        prev = store["monitor"].get(bid)
        m = monitor_basket(store, bid, model, prev, tick)
        if m:
            results[bid] = m
    store["monitor"] = results
    return results


# ══════════════════════════════════════════════════════════════════════════════
#  阶段4:样本外验证(报警追踪)
# ══════════════════════════════════════════════════════════════════════════════

def track_alerts(store, results):
    """报警后追踪 Q 是否回落。"""
    now = int(time.time())
    horizon = TRACK_BARS * BAR_MINUTES * 60

    for a in store["alerts"]:
        if a.get("resolved") is not None:
            continue
        r = results.get(a["bid"])
        if not r:
            continue
        a["peak_q"] = max(a.get("peak_q", a["q0"]), r["q"])
        if r["q"] <= a["exit_q"]:
            a["resolved"] = True
            a["resolve_s"] = now - a["time"]
        elif now - a["time"] > horizon:
            a["resolved"] = False
            a["resolve_s"] = now - a["time"]

    for bid, r in results.items():
        if not r["alarm"] or r["blocked"]:
            continue
        recent = [a for a in store["alerts"]
                  if a["bid"] == bid and a.get("resolved") is None]
        if recent:
            continue
        model = store["baskets"][bid]
        store["alerts"].insert(0, {
            "time": now, "bid": bid, "q0": r["q"], "peak_q": r["q"],
            "exit_q": model["q_exit"], "resolved": None, "resolve_s": 0,
            "anomaly": r["contrib"][0]["symbol"] if r["contrib"] else "",
            "z": r["contrib"][0]["z"] if r["contrib"] else 0,
        })
        record(store, "alert_raised", {
            "bid": bid, "q": r["q"], "limit": r["q_entry"],
            "anomaly": r["contrib"][0] if r["contrib"] else None,
        }, "#FF0000")

    store["alerts"] = store["alerts"][:200]


def validation_stats(store):
    done = [a for a in store["alerts"] if a.get("resolved") is not None]
    if not done:
        return {"total": len(store["alerts"]), "closed": 0, "rate": None,
                "median_min": None, "mae": None}
    ok = [a for a in done if a["resolved"]]
    times = sorted(a["resolve_s"] / 60.0 for a in ok)
    mae = [a["peak_q"] / max(a["q0"], 1e-9) for a in done]
    return {
        "total": len(store["alerts"]),
        "closed": len(done),
        "rate": round(len(ok) / len(done), 3),
        "median_min": round(times[len(times) // 2], 1) if times else None,
        "mae": round(float(np.mean(mae)), 2),
    }


# ══════════════════════════════════════════════════════════════════════════════
#  模拟盘引擎
# ══════════════════════════════════════════════════════════════════════════════

def leg_pnl(leg, price):
    """合约盈亏。"""
    diff = price - leg["entry"] if leg["side"] == "long" else leg["entry"] - price
    return diff * leg["contracts"] * leg["ctVal"]


def position_pnl(pos, price_map):
    """含手续费与资金费率的净盈亏。"""
    gross = sum(leg_pnl(l, price_map.get(s, l["entry"])) for s, l in pos["legs"].items())
    return gross - pos.get("fee_paid", 0.0) - pos.get("funding_paid", 0.0)


def equity(store, price_map):
    eq = store["cash"]
    for pos in store["positions"].values():
        eq += sum(l["margin"] for l in pos["legs"].values())
        eq += position_pnl(pos, price_map)
    return eq


def all_prices(store):
    pm = {}
    for r in store["monitor"].values():
        pm.update(r.get("prices", {}))
    return pm


def open_basket(store, bid, model, r):
    """按 LOO 归因方向建仓。"""
    if len(store["positions"]) >= MAX_OPEN or bid in store["positions"]:
        return
    top = r["contrib"][0] if r["contrib"] else None
    if not top or abs(top["z"]) < MIN_LOO_Z:
        record(store, "entry_skip_weak_attribution", {"bid": bid, "top": top}, "#FFAA00")
        return

    i = model["members"].index(top["symbol"])
    rec = next(x for x in model["loo"] if x["i"] == i)
    b = np.array(rec["beta"])

    sgn = 1.0 if top["z"] > 0 else -1.0
    w = {top["symbol"]: -sgn}
    for pos_j, j in enumerate(rec["others"]):
        w[model["members"][j]] = sgn * float(b[1 + pos_j])

    gross = sum(abs(v) for v in w.values())
    if gross <= 0:
        return
    budget = equity(store, r["prices"]) * LEVERAGE * GROSS_PER_TRADE

    legs, need, fee, dropped = {}, 0.0, 0.0, []
    for sym, wt in w.items():
        price = r["prices"].get(sym)
        mk = store["markets"].get(sym)
        if not price or not mk:
            continue
        target = budget * abs(wt) / gross

        ct_val = mk["ctVal"]
        per_ct = price * ct_val
        contracts = math.floor(target / per_ct) if per_ct > 0 else 0
        if contracts < max(mk["minQty"], MIN_LEG_CONTRACTS):
            dropped.append({"sym": sym, "want": round(target, 2), "contracts": contracts})
            continue
        notional = contracts * per_ct
        if notional < MIN_LEG_USDT:
            dropped.append({"sym": sym, "notional": round(notional, 2)})
            continue

        margin = notional / LEVERAGE
        legs[sym] = {"side": "long" if wt > 0 else "short", "entry": price,
                     "contracts": contracts, "ctVal": ct_val, "notional": notional,
                     "margin": margin, "w": round(wt, 5)}
        need += margin
        fee += notional * mk["takerFee"] * COST_MULT

    if len(legs) < 2:
        record(store, "entry_skip_legs", {"bid": bid, "legs": len(legs)}, "#FFAA00")
        return

    # —— 对冲保真度检查 ——
    tgt = np.array([w.get(sy, 0.0) for sy in w])
    act = np.array([(legs[sy]["notional"] * (1 if legs[sy]["side"] == "long" else -1))
                    if sy in legs else 0.0 for sy in w])
    nt, na = np.linalg.norm(tgt), np.linalg.norm(act)
    fidelity = float(tgt @ act / (nt * na)) if nt > 0 and na > 0 else 0.0
    if fidelity < MIN_HEDGE_FID:
        record(store, "entry_skip_hedge_fidelity", {
            "bid": bid, "fidelity": round(fidelity, 4), "门槛": MIN_HEDGE_FID}, "#FFAA00")
        return

    if need + fee > store["cash"] - 1:
        record(store, "entry_skip_capital", {"bid": bid, "need": round(need + fee, 2)}, "#FFAA00")
        return

    round_trip = fee * 2
    store["cash"] -= (need + fee)
    store["positions"][bid] = {
        "bid": bid, "legs": legs, "open_time": int(time.time()),
        "last_funding": int(time.time()),
        "entry_q": r["q"], "anomaly": top["symbol"], "z": top["z"],
        "margin": need, "fee_paid": fee, "funding_paid": 0.0,
        "est_round_trip": round(round_trip, 4),
        "fidelity": round(fidelity, 4),
        "max_hold_s": int(model["halflife"] * HOLD_HL_MULT * BAR_MINUTES * 60),
        "model": {k: model[k] for k in
                  ("members", "mu", "Vn", "resid_inv", "q_exit", "slot_scale", "r")},
    }
    trade = {"time": int(time.time()), "type": "open", "bid": bid,
             "anomaly": top["symbol"], "z": top["z"], "q": r["q"],
             "legs": len(legs), "margin": round(need, 2),
             "fee": round(fee, 4), "往返成本": round(round_trip, 4)}
    store["trades"].insert(0, trade)
    store["trades"] = store["trades"][:100]
    record(store, "paper_open", trade, "#00CC66")


def close_basket(store, bid, reason, price_map):
    pos = store["positions"].pop(bid, None)
    if not pos:
        return
    exit_fee = sum(l["notional"] * store["markets"].get(s, {}).get("takerFee", TAKER_FEE)
                   * COST_MULT for s, l in pos["legs"].items())
    pos["fee_paid"] = pos.get("fee_paid", 0.0) + exit_fee
    pnl = position_pnl(pos, price_map)
    for leg in pos["legs"].values():
        store["cash"] += leg["margin"]
    store["cash"] += pnl
    store["realized"] += pnl
    trade = {"time": int(time.time()), "type": "close", "bid": bid,
             "reason": reason, "pnl": round(pnl, 4),
             "fee": round(pos["fee_paid"], 4),
             "funding": round(pos.get("funding_paid", 0.0), 4),
             "hold_min": round((int(time.time()) - pos["open_time"]) / 60.0, 1)}
    store["trades"].insert(0, trade)
    store["trades"] = store["trades"][:100]
    record(store, "paper_close", trade, "#FFAA00" if pnl < 0 else "#00CC66")


def accrue_funding(store, pm):
    """按秒累计资金费率。"""
    now = int(time.time())
    for pos in store["positions"].values():
        dt = now - pos.get("last_funding", now)
        if dt <= 0:
            continue
        pos["last_funding"] = now
        acc = 0.0
        for sym, leg in pos["legs"].items():
            mk = store["markets"].get(sym)
            if not mk or not mk["fundingInterval"]:
                continue
            p = pm.get(sym, leg["entry"])
            notional = leg["contracts"] * leg["ctVal"] * p
            rate = mk["fundingRate"] * dt / mk["fundingInterval"]
            acc += notional * rate if leg["side"] == "long" else -notional * rate
        pos["funding_paid"] = pos.get("funding_paid", 0.0) + acc


def position_q(pos, tick):
    """用头寸自带的模型快照算 Q。"""
    m = pos.get("model")
    if not m:
        return None
    x, _ = current_log_prices(m["members"], tick)
    if x is None:
        return None
    rv = (x - np.array(m["mu"])) @ np.array(m["Vn"])
    q_raw = float(rv @ np.array(m["resid_inv"]) @ rv)
    slot = str(int((time.time() // 3600) % VOL_SLOTS))
    return q_raw / float(m.get("slot_scale", {}).get(slot, 1.0))


def manage_positions(store, results):
    now = int(time.time())
    pm = all_prices(store)
    tick = {s: p for s, p in pm.items()}
    accrue_funding(store, pm)
    for bid in list(store["positions"].keys()):
        pos = store["positions"][bid]

        q_now = position_q(pos, tick)
        if q_now is not None and q_now <= pos["model"]["q_exit"]:
            close_basket(store, bid, "Q回落", pm)
            continue
        if q_now is None and bid not in store["baskets"]:
            close_basket(store, bid, "模型失效(无法评估Q)", pm)
            continue
        if now - pos["open_time"] > pos["max_hold_s"]:
            close_basket(store, bid, "超过{}倍半衰期".format(HOLD_HL_MULT), pm)
            continue
        pnl = position_pnl(pos, pm)
        if pos["margin"] > 0 and pnl / pos["margin"] * 100 <= -STOP_LOSS_PCT:
            close_basket(store, bid, "止损", pm)

    for bid, r in results.items():
        if r["alarm"] and not r["blocked"] and bid in store["baskets"]:
            open_basket(store, bid, store["baskets"][bid], r)


# ══════════════════════════════════════════════════════════════════════════════
#  仪表板
# ══════════════════════════════════════════════════════════════════════════════

def tbl(title, cols, rows):
    return "\n`" + json.dumps({"type": "table", "title": title, "cols": cols,
                               "rows": rows or [["-"] * len(cols)]},
                              ensure_ascii=False) + "`\n"


def short(x, n=40):
    s = str(x or "").replace("\n", " ")
    return s[:n] + ("..." if len(s) > n else "")


def render(store):
    pm = all_prices(store)
    eq = equity(store, pm)
    init = store["init_capital"]
    v = validation_stats(store)

    # 首屏:样本外验证
    t_valid = tbl("🔬 样本外验证(判断前提是否成立的唯一依据)",
                  ["累计报警", "已了结", "回归率", "回归耗时中位数(min)", "最大不利偏移倍数", "结论"],
                  [[v["total"], v["closed"],
                    "—" if v["rate"] is None else "{:.1%}".format(v["rate"]),
                    v["median_min"] if v["median_min"] is not None else "—",
                    v["mae"] if v["mae"] is not None else "—",
                    "样本不足,继续积累" if v["closed"] < 20 else
                    ("✅ 前提成立" if (v["rate"] or 0) >= 0.6 else "⚠️ 回归率偏低,前提存疑")]])

    t_acct = tbl("💰 模拟账户", ["模式", "权益", "总盈亏", "收益率", "现金", "已实现", "持仓篮子"],
                 [[TRADE_MODE, "${:.2f}".format(eq),
                   "{}${:.2f}".format("+" if eq >= init else "-", abs(eq - init)),
                   "{:+.2f}%".format((eq - init) / init * 100),
                   "${:.2f}".format(store["cash"]),
                   "{:+.2f}".format(store["realized"]),
                   "{}/{}".format(len(store["positions"]), MAX_OPEN)]])

    rows = []
    for bid, m in store["baskets"].items():
        budget = eq * LEVERAGE * GROSS_PER_TRADE
        mg = m.get("min_gross", 0)
        rows.append([bid, m["N"], m["k"], m["r"], m["gap"], "{}°".format(m["max_angle"]),
                     m["adf_best"], "{}/{}".format(m["adf_pass"], m["r"]),
                     m["halflife"], m["stability"],
                     "${:,.0f}".format(mg) if mg != float("inf") else "∞",
                     "✅" if budget >= mg else "❌需${:,.0f}权益".format(
                         mg / max(LEVERAGE * GROSS_PER_TRADE, 1e-9)),
                     short(",".join(m["members"]), 50)])
    t_basket = tbl("🧮 通过检验的稳定子空间",
                   ["ID", "N", "k因子", "r零空间", "σ间隙", "子空间夹角", "ADF最优t",
                    "ADF通过", "半衰期(bar)", "稳定评分", "最小敞口", "当前可交易", "成分"], rows)

    rows = []
    for bid, r in store["monitor"].items():
        state = "🔴 报警" if r["alarm"] and not r["blocked"] else \
                ("⏸ " + r.get("block_reason", "禁开") if r["blocked"] else "🟢 正常")
        rows.append([bid, state, r["q"], r.get("q_raw", "-"), r.get("slot_scale", "-"),
                     r["q_entry"], r["q_exit"], r["fshock"], r["f_limit"],
                     " ".join("{}:{}".format(c["symbol"], c["z"]) for c in r["contrib"][:4])])
    t_mon = tbl("📡 实时监测(Q超限 且 无禁开标记 才开仓)",
                ["ID", "状态", "Q(季节归一)", "Q(原始)", "时段因子", "报警线", "离场线",
                 "F-shock", "F上限", "LOO归因(z)"], rows)

    # 新增:每个开仓品种的详细信息
    leg_rows = []
    for bid, pos in store["positions"].items():
        for sym, leg in pos["legs"].items():
            current_price = pm.get(sym, leg["entry"])
            pnl_per_contract = leg_pnl(leg, current_price) / leg["contracts"] if leg["contracts"] > 0 else 0
            roi = (pnl_per_contract / leg["entry"] * 100) if leg["entry"] > 0 else 0
            leg_rows.append([
                bid, sym, leg["side"].upper(),
                "${:.4f}".format(leg["entry"]),
                "${:.4f}".format(current_price),
                "{:+.2f}%".format(roi),
                leg["contracts"],
                "${:.2f}".format(leg_pnl(leg, current_price)),
            ])
    t_legs = tbl("📊 开仓品种详情(开仓价 → 实时价 → 盈利率)",
                 ["篮子ID", "品种", "方向", "开仓价", "实时价", "盈利率%", "张数", "盈亏$"], leg_rows)

    rows = []
    for bid, p in store["positions"].items():
        gross = sum(leg_pnl(l, pm.get(s, l["entry"])) for s, l in p["legs"].items())
        qn = position_q(p, pm)
        rows.append([bid, p["anomaly"], p["z"],
                     "{:.1f}→{:.1f}".format(p["entry_q"], qn) if qn is not None else "—",
                     p.get("fidelity", "—"), len(p["legs"]),
                     "${:.2f}".format(p["margin"]), "{:+.3f}".format(gross),
                     "{:.3f}".format(p.get("fee_paid", 0)),
                     "{:+.4f}".format(p.get("funding_paid", 0)),
                     "{:+.3f}".format(position_pnl(p, pm)),
                     "{:.1f}".format((int(time.time()) - p["open_time"]) / 60.0),
                     " ".join("{}{}x{}".format("+" if l["side"] == "long" else "-", s,
                                               l["contracts"])
                              for s, l in list(p["legs"].items())[:5])])
    t_pos = tbl("📦 持仓篮子汇总", ["ID", "异常标的", "z", "Q(进场→当前)", "对冲保真", "腿数",
                             "保证金", "毛盈亏", "手续费", "资金费", "净盈亏",
                             "持仓(min)", "腿(张数)"], rows)

    rows = [[time.strftime("%m-%d %H:%M", time.localtime(t["time"])), t["type"], t["bid"],
             t.get("anomaly", "-"), t.get("q", t.get("reason", "-")),
             t.get("pnl", "-"), t.get("fee", "-"), t.get("funding", "-"),
             t.get("hold_min", "-")]
            for t in store["trades"][:8]]
    t_trade = tbl("📜 最近交易", ["时间", "类型", "篮子", "异常标的", "Q/原因",
                               "净盈亏$", "手续费", "资金费", "持仓(min)"], rows)

    rows = [[d["drop"], d["keep"], d["rho"], "{}x".format(d["vol_ratio"])]
            for d in store.get("dedup", [])[:8]]
    t_dedup = tbl("🔀 剔除的杠杆版本(波动率衰减会破坏协整)",
                  ["剔除", "保留", "相关系数", "波动比"], rows)

    rows = [[r["n"], short(",".join(r["members"]), 44), r["reason"]] for r in store["rejected"][:8]]
    t_rej = tbl("🚫 淘汰的候选篮子(为什么不交易它)", ["规模", "成分", "淘汰原因"], rows)

    LogStatus("TradFi 矩阵统计套利 v1.0 | {} | 合约池:{} | {}\n".format(
        TRADE_MODE, len(store["markets"]), time.strftime("%Y-%m-%d %H:%M:%S"))
        + t_valid + t_acct + t_basket + t_mon + t_legs + t_pos + t_trade + t_dedup + t_rej)


# ══════════════════════════════════════════════════════════════════════════════
#  命令
# ══════════════════════════════════════════════════════════════════════════════

def handle_command(store):
    cmd = GetCommand()
    if not cmd:
        return
    if cmd == "refresh:universe":
        refresh_universe(store)
        store["last_universe"] = int(time.time())
    elif cmd == "rebuild:baskets":
        rebuild_baskets(store)
        store["last_rebuild"] = int(time.time())
    elif cmd == "flat:all":
        pm = all_prices(store)
        for bid in list(store["positions"].keys()):
            close_basket(store, bid, "手动平仓", pm)
    elif cmd == "reset:paper":
        store["cash"] = INIT_CAPITAL
        store["init_capital"] = INIT_CAPITAL
        store["positions"] = {}
        store["trades"] = []
        store["realized"] = 0.0
        Log("模拟账户已重置", "#FFAA00")
    elif cmd == "clear:store":
        store.clear()
        store.update(new_store())
        Log("store 已清空", "#FFAA00")
    save_store(store)


# ══════════════════════════════════════════════════════════════════════════════
#  主循环
# ══════════════════════════════════════════════════════════════════════════════

def main():
    _G(None)
    global _chart
    LogReset(0)
    Log("TradFi 矩阵统计套利 v1.0 启动 | 模式:", TRADE_MODE,
        "| 初始资金:", INIT_CAPITAL, "#00AAFF")
    Log("注意:币安 TradFi 永续无法与美股实物交割,缺少强制收敛机制。"
        "请以仪表板首屏的样本外验证统计为准。", "#FFAA00")

    _chart = Chart({
        "title": {"text": "模拟账户权益"},
        "xAxis": {"type": "datetime"},
        "yAxis": {"title": {"text": "权益(USDT)"}},
        "series": [{"name": "权益", "type": "line", "data": []}],
    })

    store = load_store()
    Log("载入持久化状态 | 合约池:{} | 篮子:{} | 持仓:{} | 报警:{} | 现金:{:.2f}".format(
        len(store["markets"]), len(store["baskets"]), len(store["positions"]),
        len(store["alerts"]), store["cash"]), "#00AAFF")

    cycle = 0
    while True:
        try:
            handle_command(store)
            now = int(time.time())
            cycle += 1

            need_uni = (not store["markets"]) or (now - store["last_universe"] > UNIVERSE_S)
            if need_uni:
                if not store["markets"] and store["last_universe"] > 0:
                    Log("合约池为空但定时器未到期 → 强制刷新", "#FFAA00")
                refresh_universe(store)
                store["last_universe"] = now
                save_store(store)

            need_reb = store["markets"] and ((not store["baskets"]) or
                                             now - store["last_rebuild"] > REBUILD_S)
            if need_reb:
                if not store["baskets"] and store["last_rebuild"] > 0:
                    Log("无可用篮子但定时器未到期 → 强制重建", "#FFAA00")
                rebuild_baskets(store)
                store["last_rebuild"] = now
                save_store(store)

            if store["baskets"]:
                results = run_monitor(store)
                track_alerts(store, results)
                manage_positions(store, results)

            render(store)
            eq = equity(store, all_prices(store))
            _chart.add(0, [now * 1000, eq])
            LogProfit(eq - store["init_capital"], "&")
            save_store(store)

            if cycle == 1 or cycle % HEARTBEAT_CYCLES == 0:
                alarms = sum(1 for r in store["monitor"].values() if r.get("alarm"))
                Log("♥ 第{}轮 | 池:{} 篮子:{} 监测:{} 报警:{} 持仓:{} | 权益:{:.2f} | "
                    "下次重建:{:.1f}h".format(
                        cycle, len(store["markets"]), len(store["baskets"]),
                        len(store["monitor"]), alarms, len(store["positions"]), eq,
                        max(0, REBUILD_S - (now - store["last_rebuild"])) / 3600.0),
                    "#888888")

        except Exception as e:
            import traceback
            Log("主循环异常:", str(e), "\n", traceback.format_exc()[:1200], "#FF0000")

        Sleep(POLL_S * 1000)
Comment
All comments (0)
No data
No data
  • 1
Forums
PINE Language
Get the app
iPhone Download
© 2015 - ∞ INVENTOR PTE LTD (SG)