MB-X Bilibili Pipeline
6 days ago 873dca205129f123d5ea9dd768bfa1c2e5452830
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
#!/usr/bin/env python3
"""Fetch V2 evidence packages for uncovered industry-research securities.
 
The target list is derived at runtime from the four governed research pools and
the MySQL valuation ledger.  Fetches run concurrently, while each V2 process
retains its own 90-second monotonic deadline and writes only to the analyst's
temporary workspace.
"""
 
from __future__ import annotations
 
import csv
import json
import os
import re
import subprocess
import sys
from concurrent.futures import ThreadPoolExecutor, as_completed
from pathlib import Path
 
import pymysql
 
 
ROOT = Path(__file__).resolve().parents[2]
AS_OF = "2026-08-13"
OUT = ROOT / "ai-valuation-analyst" / "tmp" / "valuation_industry_research_gaps_20260814"
CACHE = ROOT / "ai-valuation-analyst" / "tmp" / "v2_cache"
 
 
def mysql_connection(database: str) -> pymysql.Connection:
    return pymysql.connect(
        host=os.environ.get("MYSQL_HOST", "127.0.0.1"),
        port=int(os.environ.get("MYSQL_PORT", "3306")),
        user=os.environ.get("MYSQL_USER", "root"),
        password=os.environ["MYSQL_PASSWORD"],
        database=database,
        charset="utf8mb4",
        autocommit=True,
        cursorclass=pymysql.cursors.DictCursor,
    )
 
 
def semiconductor_tickers() -> set[str]:
    path = ROOT / "ana-data/cases/半导体案例/ANA-SEMI-20260722-001/outputs/核心文档/国内企业信息表.csv"
    market_map = {"SSE": "SH", "SSE STAR": "SH", "SZSE": "SZ", "SZSE CHINEXT": "SZ", "BSE": "BJ"}
    rows = csv.DictReader(path.open(encoding="utf-8-sig", newline=""))
    result = set()
    for row in rows:
        code = row["ticker"].strip()
        suffix = market_map.get(row["listed_market"].strip())
        if suffix and re.fullmatch(r"\d{6}", code):
            result.add(f"{code}.{suffix}")
    # The North Exchange completed its 920-series code migration after the
    # research snapshot.  The local formal market contract uses the new code.
    result.discard("835179.BJ")
    result.add("920179.BJ")
    return result
 
 
def robot_tickers() -> set[str]:
    path = ROOT / "ana-data/cases/机器人案例/ANA-ROBOT-INDUSTRY-001/evidence/next_robot_036_market_database_company_universe_authority_20260726.csv"
    return {row["symbol"].strip() for row in csv.DictReader(path.open(encoding="utf-8-sig", newline=""))}
 
 
def newenergy_tickers() -> set[str]:
    path = ROOT / "ana-data/cases/新能源案例/核心文档/全量公司横向总表.md"
    result = set()
    for code in re.findall(r"\| (\d{6}) \| \[", path.read_text(encoding="utf-8")):
        suffix = "SH" if code.startswith("6") else ("BJ" if code.startswith(("8", "9")) else "SZ")
        result.add(f"{code}.{suffix}")
    return result
 
 
def missing_tickers() -> list[str]:
    universe = semiconductor_tickers() | robot_tickers() | newenergy_tickers()
    with mysql_connection("stock_valuation") as connection, connection.cursor() as cursor:
        cursor.execute("SELECT ticker FROM security")
        covered = {row["ticker"] for row in cursor.fetchall()}
    return sorted(universe - covered)
 
 
def fetch(ticker: str) -> dict[str, object]:
    target = OUT / f"v2_{ticker.replace('.', '_')}"
    env = os.environ.copy()
    env["PYTHONPATH"] = str(ROOT / "dev/project-dev")
    proc = subprocess.run(
        [
            sys.executable,
            "-m",
            "stock_valuation_pipeline_v2",
            "--ticker",
            ticker,
            "--as-of",
            AS_OF,
            "--cache-dir",
            str(CACHE),
            "--output-dir",
            str(target),
        ],
        cwd=ROOT,
        env=env,
        capture_output=True,
        text=True,
        encoding="utf-8",
        errors="replace",
        timeout=120,
    )
    return {
        "ticker": ticker,
        "returncode": proc.returncode,
        "stdout": proc.stdout.strip().splitlines()[-1] if proc.stdout.strip() else "",
        "stderr": proc.stderr.strip()[-500:],
    }
 
 
def latest_package_exists(ticker: str) -> bool:
    return any(OUT.glob(f"v2_{ticker.replace('.', '_')}.failed-*")) or (OUT / f"v2_{ticker.replace('.', '_')}").exists()
 
 
def main() -> None:
    if hasattr(sys.stdout, "reconfigure"):
        sys.stdout.reconfigure(encoding="utf-8", errors="replace")
    tickers = missing_tickers()
    if "--missing-packages-only" in sys.argv:
        tickers = [ticker for ticker in tickers if not latest_package_exists(ticker)]
    OUT.mkdir(parents=True, exist_ok=True)
    rows: list[dict[str, object]] = []
    with ThreadPoolExecutor(max_workers=12) as pool:
        pending = {pool.submit(fetch, ticker): ticker for ticker in tickers}
        for future in as_completed(pending):
            ticker = pending[future]
            try:
                row = future.result()
            except Exception as exc:  # retain every batch failure
                row = {"ticker": ticker, "returncode": 99, "error": repr(exc)}
            rows.append(row)
            print(f"[{len(rows):03d}/{len(tickers):03d}] {ticker} rc={row.get('returncode')}", flush=True)
    rows.sort(key=lambda item: str(item["ticker"]))
    manifest = {"as_of": AS_OF, "count": len(rows), "items": rows}
    (OUT / "fetch_manifest.json").write_text(json.dumps(manifest, ensure_ascii=False, indent=2) + "\n", encoding="utf-8")
    print(json.dumps({"count": len(rows), "returncodes": {str(code): sum(row.get("returncode") == code for row in rows) for code in sorted({int(row.get("returncode", 99)) for row in rows})}}, ensure_ascii=False))
 
 
if __name__ == "__main__":
    main()