from __future__ import annotations import hashlib import json import os import re from pathlib import Path import pandas as pd import pymysql RUN_ID = "RUN-ANA-WUJI-BASELINE-PILOT-20260607-001" ROOT = Path(__file__).resolve().parents[1] PROJECT_ROOT = ROOT.parents[2] LOCAL_DB_INDEX = Path( r"D:\strategy_project\s-system-doc\observer\天下模型沉淀\数据库索引数据.md" ) def read_password() -> str: env = os.environ.get("TIANXIA_MYSQL_PASSWORD") or os.environ.get("MYSQL_PWD") if env: return env text = LOCAL_DB_INDEX.read_text(encoding="utf-8") match = re.search(r"^\s*-\s*密码:`([^`]+)`", text, re.MULTILINE) if not match: raise RuntimeError("Unable to read local MySQL credential from approved local index.") return match.group(1) def get_conn(): return pymysql.connect( host="127.0.0.1", port=3306, user="root", password=read_password(), database="tianxia", charset="utf8mb4", connect_timeout=5, read_timeout=120, ) 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, name: str) -> Path: path = ROOT / name df.to_csv(path, index=False, encoding="utf-8-sig") return path def main() -> None: cfg = json.loads((ROOT / "run_config.json").read_text(encoding="utf-8")) minute_start = cfg["data_sources"]["minute_price_source"]["actual_coverage_start"] minute_end = cfg["data_sources"]["minute_price_source"]["actual_coverage_end"] pull_start = "2022-10-01" pull_end = minute_end with get_conn() as conn: daily = pd.read_sql( """ SELECT trade_date, symbol, open_price, high_price, low_price, close_price, volume, amount FROM a_share_daily_price WHERE trade_date BETWEEN %(start)s AND %(end)s ORDER BY symbol, trade_date """, conn, params={"start": pull_start, "end": pull_end}, ) breadth = pd.read_sql( """ SELECT trade_date, stock_count, up_count, flat_count, down_count, run_id, price_source_table FROM ts_market_breadth_daily_cache WHERE trade_date BETWEEN %(start)s AND %(end)s AND scope_type='ALL_A_SHARE' AND scope_value='ALL' ORDER BY trade_date """, conn, params={"start": "2023-01-01", "end": minute_end}, ) calendar = pd.read_sql( """ SELECT calendar_date, is_trading_day, trade_date_rank FROM a_share_trading_calendar WHERE calendar_date BETWEEN %(start)s AND %(end)s ORDER BY calendar_date """, conn, params={"start": "2023-01-01", "end": minute_end}, ) for col in ["trade_date"]: daily[col] = pd.to_datetime(daily[col]) breadth[col] = pd.to_datetime(breadth[col]) calendar["calendar_date"] = pd.to_datetime(calendar["calendar_date"]) price_cols = ["open_price", "high_price", "low_price", "close_price", "volume", "amount"] for col in price_cols: daily[col] = pd.to_numeric(daily[col], errors="coerce") daily = daily.sort_values(["symbol", "trade_date"]).reset_index(drop=True) grouped = daily.groupby("symbol", group_keys=False) daily["prev_close"] = grouped["close_price"].shift(1) daily["prev5_avg_volume"] = grouped["volume"].transform( lambda s: s.shift(1).rolling(5, min_periods=5).mean() ) daily["volume_ratio"] = daily["volume"] / daily["prev5_avg_volume"] daily["daily_return_pct"] = (daily["close_price"] / daily["prev_close"] - 1.0) * 100.0 daily["limit_up_event_flag"] = daily["prev_close"].gt(0) & (daily["high_price"] / daily["prev_close"] - 1.0).ge(0.095) daily["recent_limitup_30_flag"] = grouped["limit_up_event_flag"].transform( lambda s: s.rolling(30, min_periods=1).max().astype(bool) ) daily["upper_shadow_pct"] = ( (daily["high_price"] - daily[["open_price", "close_price"]].max(axis=1)) / daily["prev_close"] ) * 100.0 daily["body_pct"] = ((daily["close_price"] - daily["open_price"]).abs() / daily["prev_close"]) * 100.0 daily["prev60_high"] = grouped["high_price"].transform( lambda s: s.shift(1).rolling(60, min_periods=20).max() ) daily["prev60_low"] = grouped["low_price"].transform( lambda s: s.shift(1).rolling(60, min_periods=20).min() ) daily["flat60_range_pct"] = (daily["prev60_high"] / daily["prev60_low"] - 1.0) * 100.0 daily["flat60_flag"] = daily["flat60_range_pct"].le(35.0) # Previous-high reference volume: use the first occurrence of the 60-day rolling high. # This preserves the audited pilot candidate semantics while exposing the policy in ledgers. prev_high_volume = [] prev_high_ref_date = [] prev_high_ref_policy = [] last_limit_date = [] for _symbol, g in daily.groupby("symbol", sort=False): highs = g["high_price"].to_numpy() vols = g["volume"].to_numpy() dates = g["trade_date"].to_numpy() limit_flags = g["limit_up_event_flag"].to_numpy() for i in range(len(g)): start = max(0, i - 60) if i - start >= 20: window_highs = highs[start:i] max_pos = int(window_highs.argmax()) prev_high_volume.append(vols[start + max_pos]) prev_high_ref_date.append(pd.Timestamp(dates[start + max_pos])) prev_high_ref_policy.append("FIRST_PREVIOUS_HIGH_IN_60D_WINDOW") else: prev_high_volume.append(float("nan")) prev_high_ref_date.append(pd.NaT) prev_high_ref_policy.append("") lim_start = max(0, i - 29) recent_idx = [idx for idx in range(lim_start, i + 1) if limit_flags[idx]] last_limit_date.append(pd.Timestamp(dates[recent_idx[-1]]) if recent_idx else pd.NaT) daily["prev60_high_volume"] = prev_high_volume daily["prev60_high_ref_date"] = prev_high_ref_date daily["prev60_high_ref_policy"] = prev_high_ref_policy daily["last_limitup_date"] = last_limit_date daily["days_since_last_limitup"] = (daily["trade_date"] - daily["last_limitup_date"]).dt.days daily["touch_prev_high_flag"] = daily["prev60_high"].gt(0) & daily["high_price"].ge(daily["prev60_high"] * 0.995) daily["prev_high_volume_pass_flag"] = (~daily["touch_prev_high_flag"]) | daily["volume"].gt(daily["prev60_high_volume"]) trade_dates = sorted(pd.to_datetime(daily["trade_date"].drop_duplicates()).tolist()) entry_links = pd.DataFrame( { "signal_trade_date": trade_dates[:-1], "entry_trade_date": trade_dates[1:], } ) entry_links = entry_links[ (entry_links["entry_trade_date"] >= pd.Timestamp(minute_start)) & (entry_links["entry_trade_date"] <= pd.Timestamp(minute_end)) ].copy() breadth = breadth.rename(columns={"trade_date": "signal_trade_date"}) entry_links = entry_links.merge(breadth, on="signal_trade_date", how="left") entry_links["market_gate_open_flag"] = entry_links["up_count"].fillna(-1).ge(3000) entry_links["market_gate_status"] = entry_links["market_gate_open_flag"].map( {True: "MKT_GATE_OPEN_PREV_DAY_UP_3000", False: "NO_TRADE_MARKET_GATE_CLOSED"} ) signal_dates = set(entry_links["signal_trade_date"]) candidates = daily[daily["trade_date"].isin(signal_dates)].copy() candidates = candidates.merge( entry_links, left_on="trade_date", right_on="signal_trade_date", how="inner", suffixes=("", "_gate"), ) candidates = candidates[ candidates["prev_close"].gt(0) & candidates["prev5_avg_volume"].gt(0) & candidates["recent_limitup_30_flag"] & candidates["volume_ratio"].ge(cfg["candidate_rules"]["volume_ratio_threshold"]) & candidates["upper_shadow_pct"].gt(0) ].copy() candidates["strict_candidate_flag"] = candidates["prev_high_volume_pass_flag"] candidates["candidate_status"] = candidates["strict_candidate_flag"].map( {True: "PASS", False: "FAKE_BREAKOUT_RISK_REVIEW"} ) candidates["candidate_rank"] = ( candidates.sort_values( ["signal_trade_date", "upper_shadow_pct", "volume_ratio", "amount"], ascending=[True, False, False, False], ) .groupby("signal_trade_date") .cumcount() + 1 ) candidates = candidates[candidates["candidate_rank"].le(cfg["candidate_rules"]["upper_shadow_top_n"])].copy() candidates = candidates.sort_values(["entry_trade_date", "candidate_rank", "symbol"]).reset_index(drop=True) candidates["candidate_id"] = [ f"CAND-{RUN_ID}-{i + 1:05d}" for i in range(len(candidates)) ] out_cols = [ "candidate_id", "entry_trade_date", "signal_trade_date", "symbol", "candidate_rank", "candidate_status", "strict_candidate_flag", "market_gate_status", "market_gate_open_flag", "stock_count", "up_count", "flat_count", "down_count", "run_id", "price_source_table", "open_price", "high_price", "low_price", "close_price", "prev_close", "daily_return_pct", "volume", "amount", "prev5_avg_volume", "volume_ratio", "upper_shadow_pct", "body_pct", "recent_limitup_30_flag", "last_limitup_date", "days_since_last_limitup", "touch_prev_high_flag", "prev60_high", "prev60_high_volume", "prev60_high_ref_date", "prev60_high_ref_policy", "prev_high_volume_pass_flag", "flat60_flag", "flat60_range_pct", ] candidate_ledger = candidates[out_cols].copy() for col in ["entry_trade_date", "signal_trade_date", "last_limitup_date", "prev60_high_ref_date"]: candidate_ledger[col] = pd.to_datetime(candidate_ledger[col]).dt.strftime("%Y-%m-%d") entry_summary = entry_links[ [ "entry_trade_date", "signal_trade_date", "market_gate_status", "market_gate_open_flag", "stock_count", "up_count", "flat_count", "down_count", "run_id", "price_source_table", ] ].copy() entry_summary["entry_trade_date"] = entry_summary["entry_trade_date"].dt.strftime("%Y-%m-%d") entry_summary["signal_trade_date"] = entry_summary["signal_trade_date"].dt.strftime("%Y-%m-%d") candidate_counts = candidate_ledger.groupby("entry_trade_date").agg( candidate_count=("candidate_id", "count"), strict_candidate_count=("strict_candidate_flag", "sum"), review_candidate_count=("candidate_status", lambda s: int((s != "PASS").sum())), max_upper_shadow_pct=("upper_shadow_pct", "max"), max_volume_ratio=("volume_ratio", "max"), ).reset_index() summary = entry_summary.merge(candidate_counts, on="entry_trade_date", how="left") for col in ["candidate_count", "strict_candidate_count", "review_candidate_count"]: summary[col] = summary[col].fillna(0).astype(int) candidate_path = write_csv(candidate_ledger, "candidate_ledger.csv") summary_path = write_csv(summary, "candidate_date_summary.csv") open_dates = int(summary["market_gate_open_flag"].sum()) closed_dates = int((~summary["market_gate_open_flag"]).sum()) selected = [] def pick_one(label: str, frame: pd.DataFrame, sort_cols: list[str], ascending: list[bool]) -> None: nonlocal selected used_dates = {row["entry_trade_date"] for row in selected} pool = frame[~frame["entry_trade_date"].isin(used_dates)].copy() if pool.empty: return row = pool.sort_values(sort_cols, ascending=ascending).iloc[0].to_dict() row["selection_bucket"] = label selected.append(row) summary["entry_dt"] = pd.to_datetime(summary["entry_trade_date"]) open_with_candidates = summary[(summary["market_gate_open_flag"]) & (summary["strict_candidate_count"] > 0)] closed_with_candidates = summary[(~summary["market_gate_open_flag"]) & (summary["candidate_count"] > 0)] open_no_candidates = summary[(summary["market_gate_open_flag"]) & (summary["candidate_count"] == 0)] review_risk = summary[(summary["review_candidate_count"] > 0)] near_start = open_with_candidates.sort_values("entry_dt").head(25) near_end = open_with_candidates.sort_values("entry_dt", ascending=False).head(25) pick_one("OPEN_STRONG_CANDIDATE", open_with_candidates, ["strict_candidate_count", "max_upper_shadow_pct"], [False, False]) pick_one("OPEN_HIGH_VOLUME_RATIO", open_with_candidates, ["max_volume_ratio", "strict_candidate_count"], [False, False]) pick_one("MARKET_GATE_CLOSED_WITH_CANDIDATES", closed_with_candidates, ["candidate_count", "max_upper_shadow_pct"], [False, False]) pick_one("OPEN_NO_CANDIDATE", open_no_candidates, ["entry_dt"], [True]) pick_one("PREV_HIGH_REVIEW_RISK", review_risk, ["review_candidate_count", "candidate_count"], [False, False]) pick_one("MINUTE_COVERAGE_START_BOUNDARY", near_start, ["entry_dt"], [True]) pick_one("MINUTE_COVERAGE_END_BOUNDARY", near_end, ["entry_dt"], [False]) pick_one("MID_RANGE_NORMAL_OPEN", open_with_candidates, ["entry_dt"], [True]) case_rows = [] for i, row in enumerate(selected, 1): case_id = f"WUJI-PILOT-CASE-{i:02d}-{row['entry_trade_date'].replace('-', '')}" status = "SELECTED_FOR_PILOT_REPLAY" if not row["market_gate_open_flag"]: status = "SELECTED_NO_TRADE_MARKET_GATE_CLOSED" elif row["strict_candidate_count"] == 0: status = "SELECTED_NO_STRICT_CANDIDATE" case_rows.append( { "case_id": case_id, "entry_trade_date": row["entry_trade_date"], "signal_trade_date": row["signal_trade_date"], "selection_bucket": row["selection_bucket"], "case_status": status, "market_gate_status": row["market_gate_status"], "candidate_count": int(row["candidate_count"]), "strict_candidate_count": int(row["strict_candidate_count"]), "review_candidate_count": int(row["review_candidate_count"]), "up_count": "" if pd.isna(row["up_count"]) else int(row["up_count"]), "down_count": "" if pd.isna(row["down_count"]) else int(row["down_count"]), "selection_reason": ( "覆盖小样本分层:" + row["selection_bucket"] ), } ) case_index = pd.DataFrame(case_rows) case_path = write_csv(case_index, "case_index.csv") result = { "schema_version": "1.0", "run_id": RUN_ID, "generated_at": "2026-06-08T00:40:00+08:00", "stage": "CANDIDATE_POOL_GENERATED", "date_constraints": { "entry_start": minute_start, "entry_end": minute_end, "reason": "minute replay evidence required for buy/sell decision views", }, "input_rows": { "daily_rows": int(len(daily)), "breadth_rows": int(len(breadth)), "entry_dates": int(len(entry_links)), }, "candidate_counts": { "candidate_rows": int(len(candidate_ledger)), "candidate_entry_dates": int(candidate_ledger["entry_trade_date"].nunique()) if len(candidate_ledger) else 0, "strict_candidate_rows": int(candidate_ledger["strict_candidate_flag"].sum()) if len(candidate_ledger) else 0, "review_candidate_rows": int((candidate_ledger["candidate_status"] != "PASS").sum()) if len(candidate_ledger) else 0, "market_gate_open_entry_dates": open_dates, "market_gate_closed_entry_dates": closed_dates, }, "candidate_sort_policy": { "group_key": "signal_trade_date", "sort_keys": ["upper_shadow_pct", "volume_ratio", "amount"], "ascending": [False, False, False], "evidence_field": "amount", "top_n": cfg["candidate_rules"]["upper_shadow_top_n"], }, "previous_high_reference_volume_policy": { "policy_id": "FIRST_PREVIOUS_HIGH_IN_60D_WINDOW", "description": "When the signal day touches the previous 60-trading-day high, prev60_high_volume uses the first occurrence of the maximum high in the prior 60-trading-day window. This run preserves the audited pilot implementation and exposes prev60_high_ref_date / prev60_high_ref_policy for review.", "ledger_fields": ["prev60_high", "prev60_high_volume", "prev60_high_ref_date", "prev60_high_ref_policy"], }, "selected_cases": case_rows, "artifacts": { "candidate_ledger.csv": { "path": str(candidate_path.relative_to(ROOT)).replace("\\", "/"), "size": candidate_path.stat().st_size, "sha256": sha256_file(candidate_path), }, "candidate_date_summary.csv": { "path": str(summary_path.relative_to(ROOT)).replace("\\", "/"), "size": summary_path.stat().st_size, "sha256": sha256_file(summary_path), }, "case_index.csv": { "path": str(case_path.relative_to(ROOT)).replace("\\", "/"), "size": case_path.stat().st_size, "sha256": sha256_file(case_path), }, }, "limitations": [ "limit-up memory uses high/previous close >= 9.5% as a code proxy for first-pass candidate generation; later image/manual review must confirm semantics.", "flat60_flag uses a 60-day high/low range <= 35% only as a stratification signal, not a hard baseline buy rule.", "candidate_rank uses amount as the final same-day tie-breaker; candidate_ledger.csv now retains amount as sorting evidence.", "prev60_high_volume uses FIRST_PREVIOUS_HIGH_IN_60D_WINDOW and records prev60_high_ref_date for audit.", "candidate pool means observation pool only; it is not a buy list and produces no return conclusion.", ], "next_step": "Generate candidate daily K-line image package for selected pilot cases.", } (ROOT / "candidate_generation_summary.json").write_text( json.dumps(result, ensure_ascii=False, indent=2) + "\n", encoding="utf-8" ) summary_md = ROOT / "candidate_generation_summary.md" summary_md.write_text( "\n".join( [ "# candidate_generation_summary", "", f"run_id:`{RUN_ID}`", "阶段:`CANDIDATE_POOL_GENERATED`", "", "## 结果", "", f"- 入场日期范围:`{minute_start}` 至 `{minute_end}`", f"- 可评估 entry dates:{len(entry_links)}", f"- 候选行数:{len(candidate_ledger)}", f"- 有候选 entry dates:{candidate_ledger['entry_trade_date'].nunique() if len(candidate_ledger) else 0}", f"- strict candidate 行数:{int(candidate_ledger['strict_candidate_flag'].sum()) if len(candidate_ledger) else 0}", f"- review/risk candidate 行数:{int((candidate_ledger['candidate_status'] != 'PASS').sum()) if len(candidate_ledger) else 0}", f"- 市场闸门打开 entry dates:{open_dates}", f"- 市场闸门关闭 entry dates:{closed_dates}", f"- 已选择 pilot cases:{len(case_index)}", f"- 候选排名排序:`upper_shadow_pct desc -> volume_ratio desc -> amount desc`", f"- 前高参考量口径:`FIRST_PREVIOUS_HIGH_IN_60D_WINDOW`,账本保留 `prev60_high_ref_date`", "", "## 边界", "", "候选池只表示进入观察,不代表买入,不产生收益率、成功率、胜率或回撤结论。", "", "分钟线可回放边界为 `2023-03-24` 至 `2026-04-20`,后续小样本买卖图和回放不得超出该范围。", "", "## 下一步", "", "生成已选 pilot cases 的日 K 选股图片包和 `image_manifest.csv`。", "", ] ), encoding="utf-8", ) if __name__ == "__main__": main()