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()
|