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