#!/usr/bin/env python3
"""update_daily.py — 每日增量拉四池行情, 落盘 12 列(raw + 复权因子)。

落盘列(定稿 2026-09-12 迁移后, 顺序即文件列顺序, 勿改):
  ts_code,trade_date,raw_open,raw_high,raw_low,raw_close,adj_factor,
  pre_close,change,pct_chg,vol,amount

- raw_*      : tushare pro.daily 未复权原始价 —— 撮合/估值口径(**唯一价格事实源**)
- adj_factor : 当日复权因子(全市场一次拉取, 同时追加到 adj_factor_all.csv)
- qfq_*/hfq_*: **不落盘**, 由 `grid_or._derive_split_prices()` 在内存派生
               hfq = raw × factor(判定链口径) ; qfq = raw × factor / 末行factor(图表展示)

注: tushare pro.daily 的字段名是 open/high/low/close(未复权), 不是 raw_* ——
取行时必须用 row['close'] 等原字段名(2026-09-12 曾误写成 row['qfq_close'] 导致 KeyError)。

因子来源(优先级): ①当日 pro.adj_factor 全市场 → ②该股文件末行因子(adj_factor 列,
或旧的 hfq_close/raw_close 反推) → ③写不出因子则跳过该股并计数, **绝不写 NaN、绝不用 1.0 兜底**。
除权守卫(2026-09-13 用户批准, **B-only 判据**): ②兜底 × 当日真除权 = 静默错(因子没更新 →
少一次除权调整, 该股此后所有派生价继承) → 当日判定除权则**拒写并告警**(不猜因子)。
判据单点在 `code/exdiv_judge.py`: |(raw_close/该股上一有效交易日raw_close-1)*100 - pct_chg| ≥ 0.05pp。
**不要用 A=(raw_close/pre_close-1)*100 参与判定** —— `pre_close` 是脏列(全库 4.9%/34.3万行既不
等于前日 raw_close 也不自洽, 推算残留), 且大比例分红票 A/B 双侧同时失配会被漏判(000338.SZ 20241018)。
实测(真值=tushare 因子变动): B-only 命中 87.3% vs A∧B 83.2%(漏的 16 条跳幅≥0.5%, 最大 2.81%);
B-only 漏判 56 条跳幅中位 0.02%/max 0.10%(材料性可忽略); 非除权日误报 0.27~0.63%, 但只在②兜底
路径生效、后果是拒写等接口恢复, 属安全侧。实测 20260911 命中 1 只真除权(600309.SH 因子
42.74168→43.195301); 拦截后需 adj_factor 接口恢复时 `--force` 重算该日。
回归测试: `code/preflight_incremental_update.py` ⑦ 段(合成除权→期望拒写) + `code/exdiv_judge.py` 自检。

用法: update_daily.py [YYYYMMDD] [--dry] [--force] [--check]
  --dry    只算不写(校验用)
  --force  即使文件已有该日行也重算(校验用)
  --check  只做四池末行列完整性体检(不拉行情、不写盘)
环境变量: ZT_NO_FAC=1 跳过 adj_factor 拉取(模拟因子表不可用, 验证②兜底路径) ; ZT_BASE 覆写数据根
"""
import tushare as ts
import pandas as pd
import os, sys, datetime

sys.path.insert(0, os.path.dirname(os.path.abspath(__file__)))
from exdiv_judge import is_exdiv as _is_exdiv    # B-only 除权判据(单点定稿, 见该文件)

BASE = os.environ.get('ZT_BASE', '/Users/xpresso/zt_app/backtest_zt_full')
FAC_TABLE = f'{BASE}/adj_factor_all.csv'
COLS = ['ts_code', 'trade_date',
        'raw_open', 'raw_high', 'raw_low', 'raw_close',
        'adj_factor', 'pre_close', 'change', 'pct_chg', 'vol', 'amount']
PRICE_COLS = [c for c in COLS if c not in ('ts_code', 'trade_date')]

DRY = '--dry' in sys.argv
FORCE = '--force' in sys.argv
CHECK = '--check' in sys.argv
NO_FAC = os.environ.get('ZT_NO_FAC') == '1'
args = [a for a in sys.argv[1:] if not a.startswith('--')]
pools = ['hs300', 'zz500', 'zz1000', 'zz2000']


def pool_codes(pool):
    """增量更新的代码全集 = 旧成员表并集 ∪ 目录已有日线文件 (2026-09-16)

    为何并入目录: zz2000 时点成分并集 2934 > 旧表 2000。补数后若仍只按旧表并集跑,
    新增的 934 只「有文件但永不被增量更新」→ 自回填次日起缺行累积, 且池级统计
    (按目录当日有行情文件计) 静默偏小, 无任何告警。取并集后该类代码自动纳入增量。
    注意: 目录里的代码必有文件, 故下方 os.path.exists 分支只会命中「并集有、目录无」者。
    """
    members = pd.read_pickle(f'{BASE}/{pool}_members.pkl')
    codes = set()
    for mset in members.values():
        codes |= set(mset)
    ddir = f'{BASE}/daily_{pool}'
    on_disk = {f[:-4] for f in os.listdir(ddir) if f.endswith('.csv')} if os.path.isdir(ddir) else set()
    extra = on_disk - codes
    if extra:
        print(f'⚠️ {pool}: 目录比成员表并集多 {len(extra)} 只 → 已纳入本次增量 '
              f'(成员表并集 {len(codes)} / 目录 {len(on_disk)}), 例: {sorted(extra)[:6]}')
    return sorted(codes | on_disk)


def do_check():
    """四池末行完整性体检: 列齐 + 价格列无 NaN"""
    bad, total, missing_file = [], 0, 0
    for pool in pools:
        for code in pool_codes(pool):
            fpath = f'{BASE}/daily_{pool}/{code}.csv'
            if not os.path.exists(fpath):
                missing_file += 1
                continue
            total += 1
            try:
                d = pd.read_csv(fpath, dtype={'trade_date': str})
            except Exception as e:
                bad.append((pool, code, f'读取失败 {e}'))
                continue
            lack = [c for c in COLS if c not in d.columns]
            if lack:
                bad.append((pool, code, f'缺列 {lack}'))
                continue
            last = d.iloc[-1]
            nulls = [c for c in PRICE_COLS if pd.isna(last[c])]
            if nulls:
                bad.append((pool, code, f'末行 {last["trade_date"]} 空值 {nulls}'))
    print(f'[--check] 体检 {total} 个文件; 无日线文件 {missing_file} 个; 不完整 {len(bad)} 个')
    for pool, code, why in bad[:30]:
        print(f'   ✗ {pool}/{code}: {why}')
    print('结论:', '全部完整 ✅' if not bad else f'存在 {len(bad)} 个问题 ❌')
    return 1 if bad else 0


if CHECK:
    sys.exit(do_check())

ts.set_token('edf6739fe1a4de0d747600cc753a8b4bf335cf27ef0f5aea2d2aa64c')
pro = ts.pro_api()

if args:
    NEW_DATE = str(args[0])
else:
    NEW_DATE = '20260831'  # 兜底
    cal = pro.trade_cal(exchange='SSE', start_date='20260101',
                        end_date=datetime.date.today().strftime('%Y%m%d'), is_open='1')
    if cal is not None and len(cal) > 0:
        NEW_DATE = str(cal['cal_date'].max())
print(f'目标交易日: {NEW_DATE}   dry={DRY} force={FORCE} no_fac={NO_FAC}')

df = pro.daily(trade_date=NEW_DATE, fields='ts_code,trade_date,open,high,low,close,pre_close,change,pct_chg,vol,amount')
if df is None or len(df) == 0:
    print(f'{NEW_DATE} 行情尚未发布, 无法更新')
    sys.exit(1)
print(f'全市场 {NEW_DATE}: {len(df)}只')
df = df.set_index('ts_code')

# --- 当日因子(全市场一次拉取) + 落盘因子表 ---
latest_fac, fac_src = {}, {}
if not NO_FAC:
    try:
        _af = pro.adj_factor(trade_date=NEW_DATE)
        latest_fac = {r['ts_code']: float(r['adj_factor']) for _, r in _af.iterrows()}
        for c, v in latest_fac.items():
            fac_src[c] = 'today'
        print(f'当日复权因子: {len(latest_fac)}只')
        try:
            tbl = _af[['ts_code', 'trade_date', 'adj_factor']].copy()
            if os.path.exists(FAC_TABLE):
                old = pd.read_csv(FAC_TABLE, dtype={'trade_date': str})
                tbl = pd.concat([old, tbl.assign(trade_date=str(NEW_DATE))], ignore_index=True)
                tbl = tbl.drop_duplicates(subset=['ts_code', 'trade_date'], keep='last')
            if not DRY:
                tbl.to_csv(FAC_TABLE, index=False)
            print(f'因子表{"(DRY跳过写盘)" if DRY else "已更新"}: {FAC_TABLE} ({len(tbl)}行)')
        except Exception as e:
            print(f'因子表落盘失败(不影响当日追加): {e}')
    except Exception as e:
        print(f'adj_factor 拉取失败, 将回退到文件末行隐含因子: {e}')

updated, skipped, no_quote, rejected, fac_fallback = 0, 0, 0, 0, []
bad_schema = []      # 列结构与 COLS 不符(混 schema) → 拒写
exdiv_missing_fac = []   # ②因子兜底 且 当日真除权 → 拒写(记 raw_prev/raw_close/pre_close/pct)
for pool in pools:
    ddir = f'{BASE}/daily_{pool}'
    upd = skp = 0
    for code in pool_codes(pool):
        fpath = f'{ddir}/{code}.csv'
        if not os.path.exists(fpath):
            skp += 1
            continue
        if code not in df.index:
            no_quote += 1
            continue  # 停牌/无行情
        row = df.loc[code]
        old = pd.read_csv(fpath, dtype={'trade_date': str})
        if len(old) > 0 and str(old['trade_date'].iloc[-1]) >= NEW_DATE and not FORCE:
            skp += 1
            continue
        if FORCE and len(old) > 0 and str(old['trade_date'].iloc[-1]) == NEW_DATE:
            old = old.iloc[:-1]        # 去掉待重算的末日行
        # --- schema 守卫(2026-09-13): 混 schema 追加会被 pd.concat 按列名对齐,
        #     静默产出多余列 + 末行 NaN(且"价格列非空"自检查不出来) → 列不符直接拒写 ---
        if list(old.columns) != COLS:
            rejected += 1
            bad_schema.append((pool, code, list(old.columns)))
            continue
        raw_open, raw_high, raw_low = float(row['open']), float(row['high']), float(row['low'])
        raw_close = float(row['close'])
        pct = float(row['pct_chg'])
        # --- 因子三级优先: ①当日全市场 adj_factor ②文件末行因子(新schema: adj_factor 列 /
        #     旧schema: hfq_close÷raw_close 反推) ③写不出则拒写该行并计数(绝不写 NaN/1.0 兜底) ---
        fac = latest_fac.get(code)
        fac_from_fallback = False
        if fac is None and len(old):
            lc = old.iloc[-1]
            try:
                if 'adj_factor' in old.columns and pd.notna(lc['adj_factor']) and float(lc['adj_factor']) > 0:
                    fac = float(lc['adj_factor'])                      # 新 12 列 schema
                    fac_fallback.append(code)
                    fac_from_fallback = True
                elif {'hfq_close', 'raw_close'} <= set(old.columns):
                    if float(lc['raw_close']) > 0 and pd.notna(lc['hfq_close']):
                        fac = float(lc['hfq_close']) / float(lc['raw_close'])   # 旧 19 列 schema
                        fac_fallback.append(code)
                        fac_from_fallback = True
            except Exception:
                fac = None
        if fac is None or fac <= 0:
            rejected += 1
            continue
        # --- 除权守卫(2026-09-13, B-only 判据定稿见 code/exdiv_judge.py): ②末行因子兜底 ×
        #     当日真除权 = 静默错(因子没更新 → 少一次除权调整, 该股此后所有派生价继承)。
        #     判据只比 B=(raw_close/该股上一有效交易日raw_close-1)*100 与 pct_chg(除权日
        #     走除权参考价口径 → 必不吻合); 不用 A(pre_close 是脏列, 且大比例分红票 A/B 双侧
        #     同时失配会漏判)。判定除权则拒写并告警, 绝不猜因子。
        exdiv_reject = False
        pc0 = None
        if fac_from_fallback and len(old):
            raw_prev = old.iloc[-1].get('raw_close')   # 该股自身上一有效交易日(非全市场日历前一日)
            pc0 = float(row['pre_close'])
            if _is_exdiv(raw_close, raw_prev, pct):
                exdiv_reject = True
        if exdiv_reject:
            rejected += 1
            exdiv_missing_fac.append((pool, code, round(float(old.iloc[-1]['raw_close']), 3),
                                      raw_close, pc0, pct))
            continue
        new_row = {
            'ts_code': code, 'trade_date': NEW_DATE,
            'raw_open': raw_open, 'raw_high': raw_high, 'raw_low': raw_low, 'raw_close': raw_close,
            'adj_factor': round(fac, 6),
            'pre_close': row['pre_close'], 'change': row['change'],
            'pct_chg': pct, 'vol': row['vol'], 'amount': row['amount'],
        }
        # --- 自校验: 任一价格列为空则拒写(绝不落 NaN) ---
        nulls = [c for c in PRICE_COLS if pd.isna(new_row.get(c))]
        if nulls:
            rejected += 1
            print(f'  ✗ 拒写 {pool}/{code}: 空值 {nulls}')
            continue
        merged = pd.concat([old, pd.DataFrame([new_row])[COLS]], ignore_index=True)
        if DRY:
            prev = pd.read_csv(fpath, dtype={'trade_date': str})
            p = prev[prev['trade_date'] == NEW_DATE]
            if len(p):
                p = p.iloc[0]
                print(f"[DRY] {pool}/{code} 重算 vs 已存: "
                      f"qfq_close {new_row['qfq_close']} vs {p['qfq_close']} | "
                      f"raw_close {new_row['raw_close']} vs {p['raw_close']} | "
                      f"hfq_close {new_row['hfq_close']} vs {p['hfq_close']} | "
                      f"hfq_open {new_row['hfq_open']} vs {p.get('hfq_open')} | "
                      f"fac {fac:.6f}")
        else:
            merged.to_csv(fpath, index=False)
        upd += 1
    print(f'{pool}: 更新{upd}只 (跳过{skp})')
    updated += upd
    skipped += skp

print(f'汇总: 更新{updated} 跳过{skipped} 无行情{no_quote} 拒写{rejected}')
if bad_schema:
    print(f'⚠️ schema 不符被拒写(列结构 != 12列定稿): {len(bad_schema)} 个, 例: '
          f'{[f"{p}/{c}({len(cols)}列)" for p, c, cols in bad_schema[:6]]}')
if fac_fallback:
    print(f'因子回退(用文件末行隐含因子)的股票数: {len(fac_fallback)} 例: {fac_fallback[:8]}')
if exdiv_missing_fac:
    print(f'⛔ 除权守卫拦截({len(exdiv_missing_fac)}只): 因子走②末行兜底 但当日为除权日 → 拒写, '
          f'需等 adj_factor 接口恢复后 --force 重算该日。明细 [池/代码/前日raw_close/当日raw_close/pre_close/pct_chg]: {exdiv_missing_fac[:6]}')
print('完成(三档价)' if not DRY else '完成(DRY, 未写盘)')
sys.exit(1 if rejected else 0)
