"""Scan local A-share K-line/profile data for market reverse-gap candidates. The script uses local MySQL read-only tables: - a_share_profile_snapshot - a_share_daily_price It writes market manifestation audit artifacts into the industry case layout. Use the result as a gap-finding input, not as an investment conclusion. """ from __future__ import annotations import argparse import csv import hashlib import json import os import re import subprocess import sys from datetime import datetime, timezone from pathlib import Path SCRIPT_PATH = Path(__file__).resolve() PROJECT_ROOT = SCRIPT_PATH.parents[2] DEFAULT_MYSQL_EXE = Path(os.environ.get("TIANXIA_MYSQL_EXE", "M:/mysql/server/bin/mysql.exe")) def now_iso() -> str: return datetime.now(timezone.utc).astimezone().isoformat(timespec="seconds") def safe_name(value: str, fallback: str = "artifact", max_len: int = 96) -> str: value = re.sub(r"[^A-Za-z0-9._-]+", "_", value.strip()) value = value.strip("._-") return (value or fallback)[:max_len] def sql_literal(value: object) -> str: if value is None: return "NULL" text = str(value) return "'" + text.replace("\\", "\\\\").replace("'", "''") + "'" def sha256_text(value: str) -> str: return hashlib.sha256(value.encode("utf-8")).hexdigest() def project_relative(path: Path) -> str: try: return path.resolve().relative_to(PROJECT_ROOT.resolve()).as_posix() except ValueError: return path.resolve().as_posix() def keyword_predicate(keywords: list[str]) -> str: clauses = [] for keyword in keywords: pattern = "%" + keyword.replace("\\", "\\\\").replace("%", "\\%").replace("_", "\\_") + "%" lit = sql_literal(pattern) clauses.append( "(" f"p.industry_l1 LIKE {lit} ESCAPE '\\\\' OR " f"p.industry_l2 LIKE {lit} ESCAPE '\\\\' OR " f"p.concepts_raw LIKE {lit} ESCAPE '\\\\' OR " f"p.main_business_raw LIKE {lit} ESCAPE '\\\\' OR " f"p.stock_name LIKE {lit} ESCAPE '\\\\'" ")" ) return " OR ".join(clauses) def build_sql(args: argparse.Namespace) -> str: if not args.keyword: raise SystemExit("pass at least one --keyword") pred = keyword_predicate(args.keyword) start = sql_literal(args.start_date) end = sql_literal(args.end_date) min_return = float(args.min_return_pct) strong_day_pct = float(args.strong_day_pct) amount_ratio = float(args.amount_ratio) return f""" WITH latest_profile AS ( SELECT MAX(snapshot_date) AS snapshot_date FROM a_share_profile_snapshot WHERE snapshot_date <= {end} ), universe AS ( SELECT p.symbol, p.stock_name, p.industry_l1, p.industry_l2, p.concepts_raw, p.main_business_raw FROM a_share_profile_snapshot p JOIN latest_profile lp ON lp.snapshot_date = p.snapshot_date WHERE p.status = 'OK' AND ({pred}) ), price_with_prev AS ( SELECT d.symbol, d.trade_date, d.open_price, d.high_price, d.low_price, d.close_price, d.volume, d.amount, LAG(d.close_price) OVER (PARTITION BY d.symbol ORDER BY d.trade_date) AS prev_close, AVG(d.amount) OVER ( PARTITION BY d.symbol ORDER BY d.trade_date ROWS BETWEEN 20 PRECEDING AND 1 PRECEDING ) AS avg_amount_20 FROM a_share_daily_price d JOIN universe u ON u.symbol = d.symbol WHERE d.trade_date BETWEEN DATE_SUB({start}, INTERVAL 30 DAY) AND {end} ), window_rows AS ( SELECT u.stock_name, u.industry_l1, u.industry_l2, u.concepts_raw, u.main_business_raw, p.*, CASE WHEN p.prev_close IS NULL OR p.prev_close = 0 THEN NULL ELSE (p.close_price / p.prev_close - 1) * 100 END AS daily_return_pct, CASE WHEN p.avg_amount_20 IS NULL OR p.avg_amount_20 = 0 THEN NULL ELSE p.amount / p.avg_amount_20 END AS amount_ratio_20 FROM price_with_prev p JOIN universe u ON u.symbol = p.symbol WHERE p.trade_date BETWEEN {start} AND {end} ), ranked AS ( SELECT w.*, ROW_NUMBER() OVER (PARTITION BY w.symbol ORDER BY w.trade_date) AS rn_first, ROW_NUMBER() OVER (PARTITION BY w.symbol ORDER BY w.trade_date DESC) AS rn_last FROM window_rows w ), agg AS ( SELECT symbol, MAX(stock_name) AS company_name, MAX(industry_l1) AS industry_l1, MAX(industry_l2) AS industry_l2, MAX(concepts_raw) AS concepts_raw, MAX(main_business_raw) AS main_business_raw, MIN(trade_date) AS window_start, MAX(trade_date) AS window_end, MAX(CASE WHEN rn_first = 1 THEN close_price END) AS first_close, MAX(CASE WHEN rn_last = 1 THEN close_price END) AS last_close, COUNT(*) AS trading_days, SUM(CASE WHEN daily_return_pct >= {strong_day_pct} THEN 1 ELSE 0 END) AS strong_up_days, MAX(daily_return_pct) AS max_daily_return_pct, MAX(amount_ratio_20) AS max_amount_ratio_20, MAX(amount) AS max_amount FROM ranked GROUP BY symbol ), scored AS ( SELECT *, CASE WHEN first_close IS NULL OR first_close = 0 THEN NULL ELSE (last_close / first_close - 1) * 100 END AS period_return_pct FROM agg ) SELECT {sql_literal(args.scan_id)} AS scan_id, {sql_literal(args.industry_case)} AS industry_case, symbol, company_name, industry_l1, industry_l2, window_start, window_end, trading_days, ROUND(first_close, 4) AS first_close, ROUND(last_close, 4) AS last_close, ROUND(period_return_pct, 4) AS period_return_pct, strong_up_days, ROUND(max_daily_return_pct, 4) AS max_daily_return_pct, ROUND(max_amount_ratio_20, 4) AS max_amount_ratio_20, CASE WHEN period_return_pct >= {min_return} OR strong_up_days >= 1 OR max_amount_ratio_20 >= {amount_ratio} THEN 'STRONG_MANIFESTATION' ELSE 'WEAK_OR_NORMAL' END AS manifestation_type, 'REVIEW_REQUIRED' AS scope_type, 'REVIEW_REQUIRED' AS suspected_missing_track, 'MARKET_REVERSE_GAP_SCAN' AS gap_type, CASE WHEN period_return_pct >= {min_return} OR strong_up_days >= 1 OR max_amount_ratio_20 >= {amount_ratio} THEN 1 ELSE 0 END AS supplement_required_flag, 'Review core business against industry scheme before adding to evidence queue' AS supplement_question, 'DRAFT_FOR_REVIEW' AS review_status, LEFT(concepts_raw, 500) AS concepts_excerpt, LEFT(main_business_raw, 500) AS main_business_excerpt FROM scored WHERE period_return_pct >= {min_return} OR strong_up_days >= 1 OR max_amount_ratio_20 >= {amount_ratio} ORDER BY supplement_required_flag DESC, period_return_pct DESC, strong_up_days DESC, max_amount_ratio_20 DESC; """.strip() def run_mysql(args: argparse.Namespace, sql: str) -> str: password = os.environ.get("TIANXIA_MYSQL_PASSWORD") if not password: raise SystemExit("TIANXIA_MYSQL_PASSWORD is required in the environment") mysql_exe = Path(args.mysql_exe) if not mysql_exe.exists(): raise SystemExit(f"mysql executable not found: {mysql_exe}") cmd = [ str(mysql_exe), "--protocol=TCP", "--batch", "--raw", "--quick", "--default-character-set=utf8mb4", f"--host={args.host}", f"--port={args.port}", f"--user={args.user}", args.database, "-e", sql, ] env = os.environ.copy() env["MYSQL_PWD"] = password proc = subprocess.run(cmd, text=True, capture_output=True, env=env, check=False) if proc.returncode != 0: raise SystemExit(proc.stderr.strip() or proc.stdout.strip()) return proc.stdout def write_csv_from_tsv(tsv: str, out_path: Path) -> int: out_path.parent.mkdir(parents=True, exist_ok=True) row_count = 0 with out_path.open("w", encoding="utf-8-sig", newline="") as fh: writer = csv.writer(fh) for line_index, line in enumerate(tsv.splitlines()): writer.writerow(line.split("\t")) if line_index > 0: row_count += 1 return row_count def append_manifest(path: Path, row: dict[str, str]) -> None: exists = path.exists() path.parent.mkdir(parents=True, exist_ok=True) with path.open("a", encoding="utf-8-sig", newline="") as fh: writer = csv.DictWriter(fh, fieldnames=list(row.keys())) if not exists: writer.writeheader() writer.writerow(row) def scan(args: argparse.Namespace) -> int: if not args.scan_id: args.scan_id = f"{safe_name(args.run_id, 'run')}_MARKET_GAP_SCAN" industry_root = PROJECT_ROOT / "ana-data" / "cases" / args.industry_case supplement_dir = industry_root / "supplement" manifest_dir = industry_root / "manifest" supplement_dir.mkdir(parents=True, exist_ok=True) manifest_dir.mkdir(parents=True, exist_ok=True) run_safe = safe_name(args.run_id, "run") sql = build_sql(args) sql_hash = sha256_text(sql) audit_path = supplement_dir / f"market_manifestation_gap_audit_{run_safe}_{sql_hash[:12]}.csv" summary_path = supplement_dir / f"market_manifestation_gap_summary_{run_safe}_{sql_hash[:12]}.json" sql_path = manifest_dir / f"market_manifestation_gap_query_{run_safe}_{sql_hash[:12]}.sql" manifest_path = manifest_dir / f"market_manifestation_gap_manifest_{run_safe}.csv" tsv = run_mysql(args, sql) row_count = write_csv_from_tsv(tsv, audit_path) sql_path.write_text(sql + "\n", encoding="utf-8") summary = { "scan_id": args.scan_id, "case_id": args.case_id, "batch_id": args.batch_id, "run_id": args.run_id, "industry_case": args.industry_case, "start_date": args.start_date, "end_date": args.end_date, "keywords": args.keyword, "min_return_pct": args.min_return_pct, "strong_day_pct": args.strong_day_pct, "amount_ratio": args.amount_ratio, "row_count": row_count, "audit_path": project_relative(audit_path), "sql_sha256": sql_hash, "generated_at": now_iso(), "readout": "GAP_SCAN_ONLY_NOT_INVESTMENT_CONCLUSION", } summary_path.write_text(json.dumps(summary, ensure_ascii=False, indent=2), encoding="utf-8") append_manifest( manifest_path, { "scan_id": args.scan_id, "case_id": args.case_id, "batch_id": args.batch_id, "run_id": args.run_id, "industry_case": args.industry_case, "start_date": args.start_date, "end_date": args.end_date, "keywords": "|".join(args.keyword), "sql_sha256": sql_hash, "sql_relative_path": project_relative(sql_path), "audit_relative_path": project_relative(audit_path), "summary_relative_path": project_relative(summary_path), "row_count": str(row_count), "generated_at": summary["generated_at"], "review_status": "DRAFT_FOR_REVIEW", }, ) print(f"audit={project_relative(audit_path)} rows={row_count}") print(f"summary={project_relative(summary_path)}") return 0 def parse_args(argv: list[str]) -> argparse.Namespace: parser = argparse.ArgumentParser(description=__doc__) parser.add_argument("--industry-case", required=True) parser.add_argument("--case-id", required=True) parser.add_argument("--batch-id", required=True) parser.add_argument("--run-id", required=True) parser.add_argument("--scan-id") parser.add_argument("--start-date", required=True) parser.add_argument("--end-date", required=True) parser.add_argument("--keyword", action="append", required=True) parser.add_argument("--min-return-pct", type=float, default=20.0) parser.add_argument("--strong-day-pct", type=float, default=9.5) parser.add_argument("--amount-ratio", type=float, default=2.0) parser.add_argument("--mysql-exe", default=str(DEFAULT_MYSQL_EXE)) parser.add_argument("--host", default=os.environ.get("TIANXIA_MYSQL_HOST", "127.0.0.1")) parser.add_argument("--port", default=os.environ.get("TIANXIA_MYSQL_PORT", "3306")) parser.add_argument("--user", default=os.environ.get("TIANXIA_MYSQL_USER", "root")) parser.add_argument("--database", default=os.environ.get("TIANXIA_MYSQL_DB", "tianxia")) return parser.parse_args(argv) def main(argv: list[str] | None = None) -> int: args = parse_args(argv or sys.argv[1:]) return scan(args) if __name__ == "__main__": raise SystemExit(main())