from __future__ import annotations import hashlib import json from datetime import datetime from pathlib import Path import pandas as pd RUN_ID = "RUN-ANA-WUJI-BASELINE-PILOT-20260607-001" ROOT = Path(__file__).resolve().parents[1] REQUIRED_EVIDENCE_COLUMNS = [ "amount", "prev60_high_ref_date", "prev60_high_ref_policy", ] 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 select_candidates(case_index: pd.DataFrame, candidate_ledger: pd.DataFrame) -> pd.DataFrame: selected_rows = [] for _, case in case_index.iterrows(): rows = candidate_ledger[candidate_ledger["entry_trade_date"] == case["entry_trade_date"]].copy() if rows.empty: continue if case["selection_bucket"] == "PREV_HIGH_REVIEW_RISK": review = rows[rows["candidate_status"] != "PASS"].sort_values("candidate_rank").head(2) strict = rows[rows["candidate_status"] == "PASS"].sort_values("candidate_rank").head(3) chosen = pd.concat([strict, review], ignore_index=True).sort_values("candidate_rank").head(5) else: strict = rows[rows["candidate_status"] == "PASS"].sort_values("candidate_rank").head(5) chosen = strict if len(strict) >= 5 else rows.sort_values("candidate_rank").head(5) chosen = chosen.copy() chosen["case_id"] = case["case_id"] chosen["case_status"] = case["case_status"] chosen["selection_bucket"] = case["selection_bucket"] selected_rows.append(chosen) return pd.concat(selected_rows, ignore_index=True) if selected_rows else pd.DataFrame() def update_case_manifest(case_dir: Path) -> None: case_files = [ "candidate_ledger.csv", "image_manifest.csv", "case_image_board.md", "case_story_board.md", ] existing_manifest = {} manifest_path = case_dir / "manifest.json" if manifest_path.exists(): existing_manifest = json.loads(manifest_path.read_text(encoding="utf-8")) manifest = { "case_id": case_dir.name, "run_id": RUN_ID, "stage": existing_manifest.get("stage", "CANDIDATE_DAILY_IMAGE_PACKAGE_READY"), "files": [ { "path": name, "size": (case_dir / name).stat().st_size, "sha256": sha256_file(case_dir / name), } for name in case_files if (case_dir / name).exists() ], "image_count": existing_manifest.get("image_count", 0), } manifest_path.write_text(json.dumps(manifest, ensure_ascii=False, indent=2) + "\n", encoding="utf-8") def update_candidate_image_summary(selected_path: Path, image_manifest_path: Path, root_board_path: Path) -> None: summary_path = ROOT / "candidate_image_generation_summary.json" if not summary_path.exists(): return summary = json.loads(summary_path.read_text(encoding="utf-8")) summary["generated_at"] = datetime.now().astimezone().isoformat(timespec="seconds") summary["repair_note"] = ( "Candidate evidence columns amount / prev60_high_ref_date / prev60_high_ref_policy " "were propagated without regenerating images." ) summary.setdefault("artifacts", {}) summary["artifacts"]["selected_candidate_ledger.csv"] = { "size": selected_path.stat().st_size, "sha256": sha256_file(selected_path), } if image_manifest_path.exists(): summary["artifacts"]["image_manifest.csv"] = { "size": image_manifest_path.stat().st_size, "sha256": sha256_file(image_manifest_path), } if root_board_path.exists(): summary["artifacts"]["case_image_board.md"] = { "size": root_board_path.stat().st_size, "sha256": sha256_file(root_board_path), } summary_path.write_text(json.dumps(summary, ensure_ascii=False, indent=2) + "\n", encoding="utf-8") def main() -> None: candidate_ledger = pd.read_csv(ROOT / "candidate_ledger.csv", encoding="utf-8-sig") case_index = pd.read_csv(ROOT / "case_index.csv", encoding="utf-8-sig") missing = [col for col in REQUIRED_EVIDENCE_COLUMNS if col not in candidate_ledger.columns] if missing: raise RuntimeError(f"candidate_ledger.csv missing evidence columns: {missing}") old_selected_path = ROOT / "selected_candidate_ledger.csv" old_selected_ids: list[str] = [] if old_selected_path.exists(): old_selected = pd.read_csv(old_selected_path, encoding="utf-8-sig") old_selected_ids = old_selected["candidate_id"].astype(str).tolist() selected = select_candidates(case_index, candidate_ledger) if selected.empty: raise RuntimeError("No selected candidates after evidence repair.") new_selected_ids = selected["candidate_id"].astype(str).tolist() selected_id_changed = bool(old_selected_ids and old_selected_ids != new_selected_ids) if selected_id_changed: raise RuntimeError("Selected candidate IDs changed during evidence-only repair.") selected.to_csv(old_selected_path, index=False, encoding="utf-8-sig") case_rows = [] for _, case in case_index.iterrows(): case_id = case["case_id"] case_dir = ROOT / "cases" / case_id case_dir.mkdir(parents=True, exist_ok=True) case_candidates = selected[selected["case_id"] == case_id].copy() case_candidates.to_csv(case_dir / "candidate_ledger.csv", index=False, encoding="utf-8-sig") update_case_manifest(case_dir) case_rows.append( { "case_id": case_id, "candidate_rows": int(len(case_candidates)), "candidate_ledger_size": (case_dir / "candidate_ledger.csv").stat().st_size, "candidate_ledger_sha256": sha256_file(case_dir / "candidate_ledger.csv"), } ) update_candidate_image_summary( old_selected_path, ROOT / "image_manifest.csv", ROOT / "case_image_board.md", ) summary = { "schema_version": "1.0", "run_id": RUN_ID, "generated_at": datetime.now().astimezone().isoformat(timespec="seconds"), "stage": "CANDIDATE_POOL_EVIDENCE_REPAIR_DONE", "repair_scope": "evidence_only_no_candidate_selection_change_no_image_regeneration", "required_evidence_columns": REQUIRED_EVIDENCE_COLUMNS, "candidate_rows": int(len(candidate_ledger)), "selected_candidate_rows": int(len(selected)), "selected_candidate_ids_changed": selected_id_changed, "case_count": int(case_index["case_id"].nunique()), "case_candidate_ledgers": case_rows, "artifacts": { "candidate_ledger.csv": { "size": (ROOT / "candidate_ledger.csv").stat().st_size, "sha256": sha256_file(ROOT / "candidate_ledger.csv"), }, "selected_candidate_ledger.csv": { "size": old_selected_path.stat().st_size, "sha256": sha256_file(old_selected_path), }, "candidate_generation_summary.json": { "size": (ROOT / "candidate_generation_summary.json").stat().st_size, "sha256": sha256_file(ROOT / "candidate_generation_summary.json"), }, }, "boundary": ( "This repair only exposes amount sorting evidence and previous-high reference volume policy. " "It does not change candidate IDs, selected pilot cases, buy/sell decisions, return stats, or images." ), } (ROOT / "candidate_pool_evidence_repair_summary.json").write_text( json.dumps(summary, ensure_ascii=False, indent=2) + "\n", encoding="utf-8", ) (ROOT / "candidate_pool_evidence_repair_summary.md").write_text( "\n".join( [ "# candidate_pool_evidence_repair_summary", "", f"run_id:`{RUN_ID}`", "阶段:`CANDIDATE_POOL_EVIDENCE_REPAIR_DONE`", "", "## 修复范围", "", "- `candidate_ledger.csv` 保留 `amount` 作为候选排名兜底排序证据。", "- `candidate_ledger.csv` 保留 `prev60_high_ref_date` 和 `prev60_high_ref_policy`。", "- 前高参考成交量口径为 `FIRST_PREVIOUS_HIGH_IN_60D_WINDOW`。", "- 同步 `selected_candidate_ledger.csv` 和各 case 的 `candidate_ledger.csv`。", "- 不重新生成图片,不修改买卖裁决,不修改收益准备包。", "", "## 校验", "", f"- 候选行数:{summary['candidate_rows']}", f"- selected candidate 行数:{summary['selected_candidate_rows']}", f"- selected candidate IDs 是否变化:{summary['selected_candidate_ids_changed']}", "", ] ), encoding="utf-8", ) if __name__ == "__main__": main()