from __future__ import annotations import csv import hashlib import importlib.util import json from datetime import datetime from pathlib import Path from zoneinfo import ZoneInfo import pandas as pd RUN_ID = "RUN-ANA-WUJI-RECENT-NONBJ-ALL-CANDIDATES-20260611-001" ROOT = Path(__file__).resolve().parents[1] PROJECT_ROOT = ROOT.parents[2] SOURCE_SCRIPT = ( PROJECT_ROOT / "ana-data/result/RUN-ANA-WUJI-FULL-2023-2026-20260608-001/tools/generate_candidate_pool.py" ) def load_source_module(): spec = importlib.util.spec_from_file_location("wuji_candidate_source", SOURCE_SCRIPT) if spec is None or spec.loader is None: raise RuntimeError(f"Unable to load source script: {SOURCE_SCRIPT}") module = importlib.util.module_from_spec(spec) spec.loader.exec_module(module) return module 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(rows: pd.DataFrame, name: str) -> Path: path = ROOT / name rows.to_csv(path, index=False, encoding="utf-8-sig") return path def money_gate_from_daily(daily: pd.DataFrame) -> pd.DataFrame: gate = daily[daily["prev_close"].gt(0)].copy() gate["up_flag"] = gate["close_price"].gt(gate["prev_close"]) gate["flat_flag"] = gate["close_price"].eq(gate["prev_close"]) gate["down_flag"] = gate["close_price"].lt(gate["prev_close"]) breadth = ( gate.groupby("trade_date") .agg( stock_count=("symbol", "count"), up_count=("up_flag", "sum"), flat_count=("flat_flag", "sum"), down_count=("down_flag", "sum"), ) .reset_index() ) breadth["market_gate_open_flag"] = breadth["up_count"].ge(3000) breadth["market_gate_status"] = breadth["market_gate_open_flag"].map( { True: "MKT_GATE_OPEN_SIGNAL_DAY_UP_3000_RECALC_FROM_DAILY", False: "NO_TRADE_MARKET_GATE_CLOSED_SIGNAL_DAY_RECALC_FROM_DAILY", } ) breadth = breadth.rename(columns={"trade_date": "signal_trade_date"}) return breadth def add_features(daily: pd.DataFrame) -> pd.DataFrame: 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) prev_high_volume = [] prev_high_ref_date = [] prev_high_ref_policy = [] last_limit_date = [] for _symbol, group in daily.groupby("symbol", sort=False): highs = group["high_price"].to_numpy() vols = group["volume"].to_numpy() dates = group["trade_date"].to_numpy() limit_flags = group["limit_up_event_flag"].to_numpy() for i in range(len(group)): 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"]) return daily def build_candidates( daily: pd.DataFrame, signal_date: pd.Timestamp, entry_date: pd.Timestamp | None, label: str, breadth: pd.DataFrame, ) -> pd.DataFrame: candidates = daily[daily["trade_date"].eq(signal_date)].copy() candidates = candidates[~candidates["symbol"].str.endswith(".BJ")].copy() candidates = candidates[ candidates["prev_close"].gt(0) & candidates["prev5_avg_volume"].gt(0) & candidates["recent_limitup_30_flag"] & candidates["volume_ratio"].ge(1.7) & candidates["upper_shadow_pct"].gt(0) ].copy() candidates["signal_trade_date"] = signal_date candidates["entry_trade_date"] = entry_date if entry_date is not None else pd.NaT candidates = candidates.merge(breadth, on="signal_trade_date", how="left") 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["scan_label"] = label candidates = candidates.sort_values(["candidate_rank", "symbol"]).reset_index(drop=True) candidates["candidate_id"] = [ f"CAND-{RUN_ID}-{label}-{i + 1:04d}" for i in range(len(candidates)) ] return candidates def render_summary_md(summary: dict) -> str: return f"""# Wuji Recent Non-BJ All Candidate Scan run_id: `{summary["run_id"]}` generated_at: `{summary["generated_at"]}` This package rescans the latest available daily data, excludes all Beijing Stock Exchange symbols (`*.BJ`), and removes the old audited top-50 display cap. It is a daily candidate scan only, not a buy recommendation and not a minute-level entry confirmation. ## Views - Entry-ready view: signal date `{summary["latest_entry_ready"]["signal_trade_date"]}`, entry observation date `{summary["latest_entry_ready"]["entry_trade_date"]}`. Market gate: `{summary["market_gate"]["latest_entry_ready"]["market_gate_status"]}`; up_count: `{summary["market_gate"]["latest_entry_ready"]["up_count"]}`. - Latest signal view: signal date `{summary["latest_signal"]["signal_trade_date"]}`; next entry date is not frozen here. Market gate: `{summary["market_gate"]["latest_signal"]["market_gate_status"]}`; up_count: `{summary["market_gate"]["latest_signal"]["up_count"]}`. ## Counts - Entry-ready all non-BJ candidates: `{summary["counts"]["entry_ready_all_non_bj_rows"]}` - Entry-ready PASS non-BJ candidates: `{summary["counts"]["entry_ready_pass_non_bj_rows"]}` - Latest-signal all non-BJ candidates: `{summary["counts"]["latest_signal_all_non_bj_rows"]}` - Latest-signal PASS non-BJ candidates: `{summary["counts"]["latest_signal_pass_non_bj_rows"]}` ## Boundaries - No `*.BJ` symbol is retained in the output files. - `PASS` is the strict stock-level pass status. - `FAKE_BREAKOUT_RISK_REVIEW` rows satisfy the core daily scan but require manual previous-high / fake-breakout risk review. - Minute-level buy point confirmation is not included in this package. """ def write_manifest(files: list[Path]) -> Path: rows = [] for path in files: rows.append( { "path": path.relative_to(ROOT).as_posix(), "size": path.stat().st_size, "sha256": sha256_file(path), } ) manifest = ROOT / "manifest.json" manifest.write_text( json.dumps( { "schema_version": "1.0", "run_id": RUN_ID, "generated_at": datetime.now(ZoneInfo("Asia/Shanghai")).isoformat(timespec="seconds"), "files": rows, }, ensure_ascii=False, indent=2, ), encoding="utf-8", ) return manifest def write_self_check( all_recent: pd.DataFrame, entry_ready: pd.DataFrame, entry_pass: pd.DataFrame, latest_signal: pd.DataFrame, signal_pass: pd.DataFrame, ) -> list[Path]: items = [ { "check_id": "NO_BJ_SYMBOL_RETAINED", "status": "PASS" if int(all_recent["symbol"].str.endswith(".BJ").sum()) == 0 else "FAIL", "detail": f"bj_rows={int(all_recent['symbol'].str.endswith('.BJ').sum())}", }, { "check_id": "ENTRY_READY_PASS_FILE_ONLY_PASS", "status": "PASS" if set(entry_pass["candidate_status"].dropna()) <= {"PASS"} else "FAIL", "detail": f"rows={len(entry_pass)}", }, { "check_id": "LATEST_SIGNAL_PASS_FILE_ONLY_PASS", "status": "PASS" if set(signal_pass["candidate_status"].dropna()) <= {"PASS"} else "FAIL", "detail": f"rows={len(signal_pass)}", }, { "check_id": "TOP_50_CAP_NOT_APPLIED", "status": "PASS" if max( int(entry_ready["candidate_rank"].max() or 0), int(latest_signal["candidate_rank"].max() or 0), ) > 50 else "FAIL", "detail": ( f"entry_max_rank={int(entry_ready['candidate_rank'].max() or 0)}; " f"latest_signal_max_rank={int(latest_signal['candidate_rank'].max() or 0)}" ), }, { "check_id": "ENTRY_READY_MARKET_GATE_OPEN", "status": "PASS" if entry_ready["market_gate_open_flag"].astype(bool).all() else "FAIL", "detail": ( entry_ready["market_gate_status"].iloc[0] if len(entry_ready) else "empty" ), }, { "check_id": "LATEST_SIGNAL_MARKET_GATE_CLOSED_RECORDED", "status": "PASS" if not latest_signal["market_gate_open_flag"].astype(bool).any() else "FAIL", "detail": ( latest_signal["market_gate_status"].iloc[0] if len(latest_signal) else "empty" ), }, ] overall = "PASS" if all(item["status"] == "PASS" for item in items) else "FAIL" csv_path = ROOT / "self_check_items.csv" with csv_path.open("w", encoding="utf-8-sig", newline="") as f: writer = csv.DictWriter(f, fieldnames=["check_id", "status", "detail"]) writer.writeheader() writer.writerows(items) json_path = ROOT / "self_check.json" json_path.write_text( json.dumps( { "schema_version": "1.0", "run_id": RUN_ID, "status": overall, "items": items, }, ensure_ascii=False, indent=2, ), encoding="utf-8", ) return [csv_path, json_path] def main() -> None: ROOT.mkdir(parents=True, exist_ok=True) source_module = load_source_module() generated_at = datetime.now(ZoneInfo("Asia/Shanghai")).isoformat(timespec="seconds") with source_module.get_conn() as conn: latest_date = pd.read_sql( "SELECT MAX(trade_date) AS latest_trade_date FROM a_share_daily_price", conn )["latest_trade_date"].iloc[0] latest_date = pd.Timestamp(latest_date) pull_start = (latest_date - pd.Timedelta(days=280)).strftime("%Y-%m-%d") 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": latest_date.strftime("%Y-%m-%d")}, ) daily["trade_date"] = pd.to_datetime(daily["trade_date"]) for col in ["open_price", "high_price", "low_price", "close_price", "volume", "amount"]: daily[col] = pd.to_numeric(daily[col], errors="coerce") daily = add_features(daily) trade_dates = sorted(daily["trade_date"].drop_duplicates()) if len(trade_dates) < 2: raise RuntimeError("Not enough trade dates to build latest entry-ready view.") latest_signal_date = pd.Timestamp(trade_dates[-1]) entry_ready_signal_date = pd.Timestamp(trade_dates[-2]) latest_entry_date = latest_signal_date breadth = money_gate_from_daily(daily) latest_signal = build_candidates( daily, latest_signal_date, None, "LATEST_SIGNAL_FOR_NEXT_ENTRY", breadth ) entry_ready = build_candidates( daily, entry_ready_signal_date, latest_entry_date, "LATEST_ENTRY_READY", breadth ) all_recent = pd.concat([entry_ready, latest_signal], ignore_index=True) out_cols = [ "candidate_id", "scan_label", "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", "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", ] for frame in [latest_signal, entry_ready, all_recent]: for col in out_cols: if col not in frame.columns: frame[col] = pd.NA entry_pass = entry_ready[entry_ready["candidate_status"].eq("PASS")].copy() signal_pass = latest_signal[latest_signal["candidate_status"].eq("PASS")].copy() files = [ write_csv(all_recent[out_cols], "all_non_bj_recent_candidate_ledger.csv"), write_csv(entry_ready[out_cols], "entry_ready_20260610_all_non_bj_candidates.csv"), write_csv(entry_pass[out_cols], "entry_ready_20260610_pass_non_bj_candidates.csv"), write_csv(latest_signal[out_cols], "latest_signal_20260610_all_non_bj_candidates.csv"), write_csv(signal_pass[out_cols], "latest_signal_20260610_pass_non_bj_candidates.csv"), ] def gate_for(date: pd.Timestamp) -> dict: row = breadth[breadth["signal_trade_date"].eq(date)].iloc[0] return { "signal_trade_date": date.strftime("%Y-%m-%d"), "stock_count": int(row["stock_count"]), "up_count": int(row["up_count"]), "flat_count": int(row["flat_count"]), "down_count": int(row["down_count"]), "market_gate_open_flag": bool(row["market_gate_open_flag"]), "market_gate_status": row["market_gate_status"], } summary = { "schema_version": "1.0", "run_id": RUN_ID, "generated_at": generated_at, "stage": "RECENT_DAILY_ALL_NON_BJ_CANDIDATE_SCAN_NOT_BUY_DECISION", "latest_daily_trade_date": latest_signal_date.strftime("%Y-%m-%d"), "latest_entry_ready": { "entry_trade_date": latest_entry_date.strftime("%Y-%m-%d"), "signal_trade_date": entry_ready_signal_date.strftime("%Y-%m-%d"), }, "latest_signal": {"signal_trade_date": latest_signal_date.strftime("%Y-%m-%d")}, "rules": { "exchange_filter": "exclude symbols ending with .BJ", "volume_ratio_threshold": 1.7, "recent_limitup_window_trading_days": 30, "prev_high_window_trading_days": 60, "upper_shadow_top_n": "not_applied_for_this_all_candidate_scan", "market_gate": "signal_day_up_count_recalculated_from_daily>=3000", }, "market_gate": { "latest_entry_ready": gate_for(entry_ready_signal_date), "latest_signal": gate_for(latest_signal_date), }, "counts": { "entry_ready_all_non_bj_rows": int(len(entry_ready)), "entry_ready_pass_non_bj_rows": int(len(entry_pass)), "entry_ready_fake_breakout_review_rows": int( len(entry_ready) - len(entry_pass) ), "latest_signal_all_non_bj_rows": int(len(latest_signal)), "latest_signal_pass_non_bj_rows": int(len(signal_pass)), "latest_signal_fake_breakout_review_rows": int( len(latest_signal) - len(signal_pass) ), "bj_rows_retained": int(all_recent["symbol"].str.endswith(".BJ").sum()), }, "limitations": [ "Daily candidate scan only; no minute-level buy point confirmation.", "Market breadth is recalculated from a_share_daily_price for the signal day.", "Latest signal date market gate may be closed and must not be read as entry-ready.", ], } summary_json = ROOT / "summary.json" summary_json.write_text(json.dumps(summary, ensure_ascii=False, indent=2), encoding="utf-8") summary_md = ROOT / "summary.md" summary_md.write_text(render_summary_md(summary), encoding="utf-8") files.extend([summary_json, summary_md]) files.extend(write_self_check(all_recent, entry_ready, entry_pass, latest_signal, signal_pass)) files.append(Path(__file__)) manifest = write_manifest(files) print( json.dumps( { "run_id": RUN_ID, "result_dir": ROOT.as_posix(), "entry_ready_all_non_bj_rows": len(entry_ready), "entry_ready_pass_non_bj_rows": len(entry_pass), "latest_signal_all_non_bj_rows": len(latest_signal), "latest_signal_pass_non_bj_rows": len(signal_pass), "bj_rows_retained": int(all_recent["symbol"].str.endswith(".BJ").sum()), "manifest": manifest.as_posix(), }, ensure_ascii=False, indent=2, ) ) if __name__ == "__main__": main()