from __future__ import annotations import argparse import csv import hashlib import json import re from datetime import datetime, timedelta, timezone from pathlib import Path import pandas as pd RUN_ID = "RUN-ANA-WUJI-V1-BUY-POINT-SECOND-REVIEW-20260615-001" SOURCE_RUN_ID = "RUN-ANA-WUJI-V1-STRICT-NOTE-FULL-RERUN-20260614-001" TZ = timezone(timedelta(hours=8)) ROOT = Path(__file__).resolve().parents[1] PACKET_ROOT = ROOT / "p2_step_review_packets" INDEX_PATH = PACKET_ROOT / "p2_step_review_index.csv" LEDGER_PATH = PACKET_ROOT / "p2_step_review_resolution_ledger.csv" SUMMARY_PATH = PACKET_ROOT / "p2_step_review_resolution_summary.md" ROLLUP_PATH = ROOT / "buy_point_resolution_rollup.csv" def now_date() -> str: return datetime.now(TZ).date().isoformat() def now_iso() -> str: return datetime.now(TZ).isoformat(timespec="seconds") def sha256_file(path: Path) -> str: h = hashlib.sha256() with path.open("rb") as f: for chunk in iter(lambda: f.read(1024 * 1024), b""): h.update(chunk) return h.hexdigest() def write_csv(df: pd.DataFrame, path: Path) -> None: path.parent.mkdir(parents=True, exist_ok=True) df.to_csv(path, index=False, encoding="utf-8-sig") def write_text(path: Path, text: str) -> None: path.parent.mkdir(parents=True, exist_ok=True) path.write_text(text, encoding="utf-8-sig") def status_label(action: str) -> str: if action == "BUY": return "改为 BUY" if action == "REVIEW_HELD": return "维持 REVIEW_HELD" return "数据不足,继续待审" def impact_note(original_action: str, final_action: str, candidate_id: str) -> str: if original_action == "REVIEW_HELD" and final_action == "BUY": return ( f"源结果包当前未生成该 candidate 的 BUY/order/lot;正式返修时需要新增 {candidate_id} " "对应 BUY,并重跑后续卖点、lot、case 汇总和收益口径。" ) if original_action == "BUY" and final_action == "REVIEW_HELD": return ( f"正式返修时需要撤销 {candidate_id} 对应 BUY/order/lot;如已有后续 SELL,需同步重算 lot、case 和收益口径。" ) return "维持源裁决;不触发交易链路返修。" def packet_decision_block(row: pd.Series, final_action: str, reason: str) -> str: return "\n".join( [ "## 人工最终裁决", "", f"- 裁决日期:{now_date()}", "- 裁决来源:用户人工确认", f"- 最终结论:`{final_action}`", f"- 人工理由:{reason}", f"- 处理含义:{impact_note(row['original_action'], final_action, row['candidate_id'])}", "", ] ) def update_packet(row: pd.Series, final_action: str, reason: str) -> None: packet_path = Path(row["packet_path"]) text = packet_path.read_text(encoding="utf-8-sig") if row["original_action"] == "REVIEW_HELD": checked = { "维持 REVIEW_HELD": final_action == "REVIEW_HELD", "改为 BUY": final_action == "BUY", "数据不足,继续待审": final_action == "DATA_INSUFFICIENT", } replacements = [ f"- [{'x' if checked['维持 REVIEW_HELD'] else ' '}] 维持 REVIEW_HELD", f"- [{'x' if checked['改为 BUY'] else ' '}] 改为 BUY", f"- [{'x' if checked['数据不足,继续待审'] else ' '}] 数据不足,继续待审", ] else: checked = { "维持 BUY": final_action == "BUY", "改为 REVIEW_HELD": final_action == "REVIEW_HELD", "数据不足,继续待审": final_action == "DATA_INSUFFICIENT", } replacements = [ f"- [{'x' if checked['维持 BUY'] else ' '}] 维持 BUY", f"- [{'x' if checked['改为 REVIEW_HELD'] else ' '}] 改为 REVIEW_HELD", f"- [{'x' if checked['数据不足,继续待审'] else ' '}] 数据不足,继续待审", ] text = re.sub( r"(?ms)(## 你要裁决\n\n)(?:- \[[ x]\] .+\n){3}", "\\1" + "\n".join(replacements) + "\n", text, count=1, ) block = packet_decision_block(row, final_action, reason) if "## 人工最终裁决" in text: text = re.sub(r"(?ms)## 人工最终裁决\n\n.*?(?=\n## 原人工裁决和二审异议)", block, text, count=1) else: text = text.replace("## 原人工裁决和二审异议", block + "\n## 原人工裁决和二审异议", 1) packet_path.write_text(text, encoding="utf-8-sig") def update_ledger(row: pd.Series, final_action: str, reason: str) -> None: columns = [ "resolution_time", "packet_order", "case_id", "candidate_id", "symbol", "entry_trade_date", "original_action", "final_action", "resolution_source", "resolution_reason_cn", "impact_note", ] if LEDGER_PATH.exists(): ledger = pd.read_csv(LEDGER_PATH, dtype=str, encoding="utf-8-sig").fillna("") else: ledger = pd.DataFrame(columns=columns) ledger = ledger[~ledger["candidate_id"].eq(row["candidate_id"])] if not ledger.empty else ledger new_row = { "resolution_time": now_date(), "packet_order": str(row["order"]), "case_id": row["case_id"], "candidate_id": row["candidate_id"], "symbol": row["symbol"], "entry_trade_date": row["entry_trade_date"], "original_action": row["original_action"], "final_action": final_action, "resolution_source": "USER_MANUAL_P2_CONFIRMATION", "resolution_reason_cn": reason, "impact_note": impact_note(row["original_action"], final_action, row["candidate_id"]), } ledger = pd.concat([ledger, pd.DataFrame([new_row])], ignore_index=True) ledger["packet_order_num"] = pd.to_numeric(ledger["packet_order"], errors="coerce") ledger = ledger.sort_values("packet_order_num").drop(columns=["packet_order_num"]) write_csv(ledger[columns], LEDGER_PATH) def formal_repair_action(original_action: str, final_action: str) -> str: if original_action == "REVIEW_HELD" and final_action == "BUY": return "ADD_BUY_ORDER_LOT_AND_RECALC" if original_action == "BUY" and final_action == "REVIEW_HELD": return "REVOKE_BUY_ORDER_LOT_AND_RECALC" return "KEEP_SOURCE_DECISION_NO_REPAIR" def update_rollup(row: pd.Series, final_action: str, reason: str) -> pd.DataFrame: rollup = pd.read_csv(ROLLUP_PATH, dtype=str, encoding="utf-8-sig").fillna("") mask = rollup["candidate_id"].eq(row["candidate_id"]) if not mask.any(): raise ValueError(f"candidate not found in rollup: {row['candidate_id']}") rollup.loc[mask, "final_action"] = final_action rollup.loc[mask, "resolution_bucket"] = "USER_P2_FINAL" rollup.loc[mask, "resolution_source"] = "USER_MANUAL_P2_CONFIRMATION" rollup.loc[mask, "resolution_status"] = "FINAL_CONFIRMED" rollup.loc[mask, "direct_resolution_flag"] = "1" rollup.loc[mask, "codex_direct_resolution_flag"] = "0" rollup.loc[mask, "formal_repair_action"] = formal_repair_action(row["original_action"], final_action) rollup.loc[mask, "resolution_reason_cn"] = reason write_csv(rollup, ROLLUP_PATH) return rollup def counts_table(series: pd.Series) -> str: lines = ["| item | count |", "|---|---:|"] for item, count in series.value_counts().sort_index().items(): lines.append(f"| `{item}` | {int(count)} |") return "\n".join(lines) def update_p2_summary() -> None: ledger = pd.read_csv(LEDGER_PATH, dtype=str, encoding="utf-8-sig").fillna("") lines = [ "# P2 买点逐条复核裁决汇总", "", f"- updated_at: {now_iso()}", f"- resolved_count: {len(ledger)}", f"- 改为 BUY: {int((ledger['final_action'] == 'BUY').sum()) if not ledger.empty else 0}", f"- 维持 REVIEW_HELD: {int((ledger['final_action'] == 'REVIEW_HELD').sum()) if not ledger.empty else 0}", f"- 数据不足继续待审: {int((ledger['final_action'] == 'DATA_INSUFFICIENT').sum()) if not ledger.empty else 0}", "", ] if ledger.empty: lines.append("- 状态:等待逐条人工裁决。") else: lines += [ "| # | case | candidate | symbol | date | original | final | reason |", "|---:|---|---|---|---|---|---|---|", ] for _, r in ledger.iterrows(): lines.append( f"| {r['packet_order']} | {r['case_id']} | {r['candidate_id']} | {r['symbol']} | {r['entry_trade_date']} | " f"{r['original_action']} | {r['final_action']} | {r['resolution_reason_cn']} |" ) lines += [ "", "## 待处理影响", "", "- `REVIEW_HELD -> BUY` 的条目需要在正式返修时新增 BUY,并重跑后续卖点、lot、case 汇总和收益口径。", "- `BUY -> REVIEW_HELD` 的条目需要在正式返修时撤销 BUY/order/lot,并重算后续链路。", ] write_text(SUMMARY_PATH, "\n".join(lines) + "\n") def update_resolution_summary(rollup: pd.DataFrame) -> None: confirmed = rollup[rollup["direct_resolution_flag"].eq("1")] codex = rollup[rollup["codex_direct_resolution_flag"].eq("1")] p1 = rollup[rollup["resolution_bucket"].eq("USER_P1_FINAL")] p2 = rollup[rollup["resolution_bucket"].eq("USER_P2_FINAL")] revoke = rollup[rollup["formal_repair_action"].eq("REVOKE_BUY_ORDER_LOT_AND_RECALC")] add = rollup[rollup["formal_repair_action"].eq("ADD_BUY_ORDER_LOT_AND_RECALC")] pending = rollup[rollup["resolution_status"].str.contains("PENDING", na=False)] lines = [ "# 买点裁决总览(含 Codex 数据直裁)", "", f"- run_id: {RUN_ID}", f"- source_run_id: {SOURCE_RUN_ID}", f"- updated_at: {now_iso()}", f"- source_rows: {len(rollup)}", "", "## 当前结论", "", f"- 已形成明确裁决:{len(confirmed)} 条", f"- 其中用户 P1 最终裁决:{len(p1)} 条", f"- 其中用户 P2 最终裁决:{len(p2)} 条", f"- 其中 Codex 数据二审直接同意源裁决:{len(codex)} 条", f"- 已确认 BUY:{len(confirmed[confirmed['final_action'].eq('BUY')])} 条", f"- 已确认 REVIEW_HELD:{len(confirmed[confirmed['final_action'].eq('REVIEW_HELD')])} 条", f"- 需要正式返修撤销 BUY/order/lot 并重算:{len(revoke)} 条", f"- 需要正式返修新增 BUY 并重算:{len(add)} 条", f"- 仍需 P2/P3 人工复核或抽查:{len(pending)} 条", f"- 无本地分钟数据覆盖、仅沿用源人工裁决且未二审确认:{len(rollup[rollup['resolution_bucket'].eq('SOURCE_MANUAL_CARRIED_NO_DATA_RECHECK')])} 条", "", "## 裁决来源分布", "", counts_table(rollup["resolution_bucket"]), "", "## 最终动作分布", "", counts_table(rollup["final_action"]), "", "## 输出文件", "", "- `buy_point_resolution_rollup.csv`: 1141 条总台账,合并用户 P1/P2 裁决、Codex 数据直裁、待复核和无数据覆盖状态。", "- `buy_point_codex_direct_resolution.csv`: Codex 数据二审直接同意源裁决的明细。", "- `buy_point_formal_repair_required.csv`: 需要正式返修撤销 BUY/order/lot 并重算的明细。", "- `buy_point_formal_buy_add_required.csv`: 需要正式返修新增 BUY 并重算的明细。", "", "## 口径边界", "", "- `CODEX_DATA_DIRECT_AGREE` 表示本机分钟数据复算与源人工裁决一致,因此列入“数据二审直接同意源裁决”;它不是新的买卖规则,也不替代后续正式结果包返修。", "- `SOURCE_MANUAL_CARRIED_NO_DATA_RECHECK` 表示本机分钟数据没有覆盖,当前没有形成二审证据,不把它算作 Codex 直接裁决。", "- 正式结果包返修时,不能只改人工裁决表;凡新增或撤销 BUY 的条目,都需要同步重算 order、lot、case 汇总和收益口径。", ] write_text(ROOT / "buy_point_resolution_summary.md", "\n".join(lines) + "\n") def update_repair_files(rollup: pd.DataFrame) -> None: write_csv( rollup[rollup["formal_repair_action"].eq("REVOKE_BUY_ORDER_LOT_AND_RECALC")], ROOT / "buy_point_formal_repair_required.csv", ) write_csv( rollup[rollup["formal_repair_action"].eq("ADD_BUY_ORDER_LOT_AND_RECALC")], ROOT / "buy_point_formal_buy_add_required.csv", ) def update_manifest() -> None: rows = [] for path in sorted(ROOT.rglob("*")): if path.is_file() and "__pycache__" not in path.parts: rows.append({"path": path.relative_to(ROOT).as_posix(), "size": path.stat().st_size, "sha256": sha256_file(path)}) with (ROOT / "manifest.csv").open("w", encoding="utf-8-sig", newline="") as f: writer = csv.DictWriter(f, fieldnames=["path", "size", "sha256"]) writer.writeheader() writer.writerows(rows) (ROOT / "manifest.json").write_text( json.dumps({"run_id": RUN_ID, "updated_at": now_iso(), "file_count": len(rows), "files": rows}, ensure_ascii=False, indent=2), encoding="utf-8", ) def parse_args() -> argparse.Namespace: parser = argparse.ArgumentParser() parser.add_argument("--order", type=int, required=True) parser.add_argument("--final-action", choices=["BUY", "REVIEW_HELD", "DATA_INSUFFICIENT"], required=True) parser.add_argument("--reason", default="") return parser.parse_args() def main() -> None: args = parse_args() index = pd.read_csv(INDEX_PATH, dtype=str, encoding="utf-8-sig").fillna("") row_match = index[index["order"].astype(int).eq(args.order)] if row_match.empty: raise ValueError(f"P2 order not found: {args.order}") row = row_match.iloc[0] reason = args.reason or f"{status_label(args.final_action)}。" update_packet(row, args.final_action, reason) update_ledger(row, args.final_action, reason) rollup = update_rollup(row, args.final_action, reason) update_p2_summary() update_resolution_summary(rollup) update_repair_files(rollup) update_manifest() if __name__ == "__main__": main()