#!/usr/bin/env python3
|
"""Evaluate current agriculture/new-energy research gaps and rebuild coverage.
|
|
This is a task coordinator, not a second valuation engine. It reuses the
|
registered V2 evidence collectors, the frozen V1 calculation core and the
|
accepted MySQL daily ledger. Agriculture and new-energy memberships are read
|
from their stable current research artifacts at runtime.
|
"""
|
|
from __future__ import annotations
|
|
import argparse
|
import csv
|
import hashlib
|
import importlib.util
|
import json
|
import os
|
import re
|
import subprocess
|
import sys
|
from collections import Counter
|
from concurrent.futures import ThreadPoolExecutor, as_completed
|
from pathlib import Path
|
from typing import Any
|
|
import pymysql
|
|
|
ROOT = Path(__file__).resolve().parents[2]
|
BASE_SCRIPT = ROOT / "ai-valuation-analyst/tools/generate_industry_research_gap_valuations_20260814.py"
|
SPEC = importlib.util.spec_from_file_location("current_research_base", BASE_SCRIPT)
|
if SPEC is None or SPEC.loader is None:
|
raise RuntimeError(f"cannot import {BASE_SCRIPT}")
|
gen = importlib.util.module_from_spec(SPEC)
|
sys.modules[SPEC.name] = gen
|
SPEC.loader.exec_module(gen)
|
|
RESULT_ROOT = ROOT / "ana-data/result/股票估值"
|
SUMMARY_DIR = RESULT_ROOT / "行业调研标的估值覆盖"
|
V2_ROOT = ROOT / "ai-valuation-analyst/tmp/valuation_current_research_20260820"
|
V2_CACHE = ROOT / "ai-valuation-analyst/tmp/v2_cache"
|
BATCH_ID = "BATCH-STOCK-VALUATION-CURRENT-RESEARCH-20260820-001"
|
CASE_ROOT = ROOT / "ana-data/cases/股票估值" / BATCH_ID
|
|
|
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 write_json(path: Path, value: Any) -> None:
|
path.parent.mkdir(parents=True, exist_ok=True)
|
path.write_text(json.dumps(value, ensure_ascii=False, indent=2) + "\n", encoding="utf-8")
|
|
|
def ticker_for(code: str) -> str:
|
if code.startswith("6"):
|
return f"{code}.SH"
|
if code.startswith(("0", "3")):
|
return f"{code}.SZ"
|
return f"{code}.BJ"
|
|
|
def research_targets() -> tuple[dict[str, Any], dict[str, str]]:
|
membership: dict[str, dict[str, Any]] = {}
|
source_notes: dict[str, str] = {}
|
|
def add(ticker: str, name: str, industry: str, layer: str) -> None:
|
item = membership.setdefault(ticker, {"names": [], "industries": set(), "layers": set()})
|
item["names"].append(name.strip())
|
item["industries"].add(industry)
|
item["layers"].add(layer)
|
|
fertilizer_path = ROOT / "ana-data/cases/农业案例/extracted/fertilizer_export_primary_disposition_2025.csv"
|
with fertilizer_path.open(encoding="utf-8-sig", newline="") as handle:
|
for row in csv.DictReader(handle):
|
add(
|
ticker_for(row["security_code"].strip()),
|
row["security_name_at_cutoff"].strip(),
|
"农业",
|
f"FERTILIZER:{row['primary_status'].strip()}",
|
)
|
source_notes["农业-化肥"] = f"{fertilizer_path.relative_to(ROOT).as_posix()}:固定50家公司"
|
|
selection_path = ROOT / "ana-data/cases/农业案例/extracted/candidate_selection.csv"
|
with selection_path.open(encoding="utf-8-sig", newline="") as handle:
|
for row in csv.DictReader(handle):
|
add(
|
ticker_for(row["security_code"].strip()),
|
row["security_name"].strip(),
|
"农业",
|
row["chain_code"].strip(),
|
)
|
source_notes["农业-首批"] = f"{selection_path.relative_to(ROOT).as_posix()}:农药/化肥16家公司"
|
|
company_root = ROOT / "ana-data/cases/新能源案例/核心文档/公司"
|
track_patterns = {
|
"SOLAR": r"光伏|太阳能|硅料|硅片|电池片|组件|逆变器",
|
"BATTERY": r"锂电|动力电池|电池材料",
|
"STORAGE": r"储能",
|
"WIND": r"风电|风机|风塔|叶片",
|
"NUCLEAR": r"核电|核能",
|
"POWER_SYSTEM": r"电力系统|电网|配电|输电",
|
}
|
card_count = 0
|
for path in sorted(company_root.glob("*.md")):
|
text = path.read_text(encoding="utf-8")
|
first_line = text.splitlines()[0] if text else ""
|
match = re.match(r"^#\s+(.+?)((\d{6}))", first_line)
|
if not match:
|
continue
|
name, code = match.groups()
|
tracks = [track for track, pattern in track_patterns.items() if re.search(pattern, text)] or ["NEWENERGY_OTHER"]
|
for track in tracks:
|
add(ticker_for(code), name, "新能源", track)
|
card_count += 1
|
source_notes["新能源"] = f"{company_root.relative_to(ROOT).as_posix()}:{card_count}张稳定公司卡"
|
|
targets = {
|
ticker: gen.Target(
|
ticker=ticker,
|
research_name=next((name for name in item["names"] if name), ticker),
|
industries=tuple(sorted(item["industries"])),
|
research_layers=tuple(sorted(item["layers"])),
|
)
|
for ticker, item in membership.items()
|
}
|
return targets, source_notes
|
|
|
def existing_security() -> dict[str, str]:
|
with mysql_connection("stock_valuation") as connection, connection.cursor() as cursor:
|
cursor.execute("SELECT ticker,company FROM security WHERE active=1")
|
result = {row["ticker"]: row["company"] for row in cursor.fetchall()}
|
manifest_path = CASE_ROOT / "batch_manifest.json"
|
if manifest_path.exists():
|
prior = gen.base.read_json(manifest_path)
|
for row in prior.get("results", []):
|
if row.get("special_valuation_status"):
|
result[row["ticker"]] = row["company"]
|
return result
|
|
|
def package_dir(ticker: str) -> Path:
|
stem = f"v2_{ticker.replace('.', '_')}"
|
candidates = [path for path in V2_ROOT.glob(f"{stem}.failed-*") if (path / "provider_results.json").exists()]
|
exact = V2_ROOT / stem
|
if (exact / "provider_results.json").exists():
|
candidates.append(exact)
|
if not candidates:
|
raise FileNotFoundError(f"V2 package missing for {ticker}")
|
return max(candidates, key=lambda path: path.stat().st_mtime)
|
|
|
def load_provider(ticker: str) -> tuple[dict[str, Any], dict[str, Any], Path]:
|
directory = package_dir(ticker)
|
failure_path = directory / "failure.json"
|
failure = gen.base.read_json(failure_path) if failure_path.exists() else {"status": "COMPLETE"}
|
return gen.base.read_json(directory / "provider_results.json"), failure, directory
|
|
|
def fetch_one(ticker: str, as_of: str) -> dict[str, Any]:
|
try:
|
directory = package_dir(ticker)
|
return {"ticker": ticker, "returncode": 0, "status": "REUSED", "package": directory.name}
|
except FileNotFoundError:
|
pass
|
target = V2_ROOT / 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(V2_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,
|
"status": "FETCHED",
|
"stdout_tail": proc.stdout.strip().splitlines()[-1] if proc.stdout.strip() else "",
|
"stderr_tail": proc.stderr.strip()[-500:],
|
}
|
|
|
def fetch_missing(as_of: str, workers: int) -> dict[str, Any]:
|
targets, _ = research_targets()
|
covered = existing_security()
|
missing = sorted(set(targets) - set(covered))
|
V2_ROOT.mkdir(parents=True, exist_ok=True)
|
rows: list[dict[str, Any]] = []
|
with ThreadPoolExecutor(max_workers=workers) as pool:
|
pending = {pool.submit(fetch_one, ticker, as_of): ticker for ticker in missing}
|
for future in as_completed(pending):
|
ticker = pending[future]
|
try:
|
row = future.result()
|
except Exception as exc:
|
row = {"ticker": ticker, "returncode": 99, "status": "ERROR", "error": repr(exc)}
|
rows.append(row)
|
print(f"[{len(rows):03d}/{len(missing):03d}] {ticker} {row['status']} rc={row['returncode']}", flush=True)
|
rows.sort(key=lambda item: item["ticker"])
|
manifest = {"batch_id": BATCH_ID, "as_of": as_of, "missing_before": len(missing), "items": rows}
|
write_json(V2_ROOT / "fetch_manifest.json", manifest)
|
return manifest
|
|
|
PE_ADDITIONS = {
|
"agri_resource": ((6, 10), (10, 16), (16, 24)),
|
"agri_fertilizer": ((7, 11), (11, 18), (18, 26)),
|
"agri_pesticide": ((10, 15), (15, 24), (24, 34)),
|
"agri_distribution": ((7, 11), (11, 17), (17, 24)),
|
}
|
|
|
def classify(target: Any, company: str) -> tuple[str, str, str, str]:
|
layers = " ".join(target.research_layers)
|
if "农业" in target.industries:
|
if "PESTICIDE" in layers or company in {"海利尔", "美邦股份"}:
|
return "agri_pesticide", "农药原药、制剂或农化服务", "产品价格、销量、登记证、渠道、原料成本和出口需求", "农药周期、环保安监、库存渠道、出口和价格波动"
|
if company in {"辉隆股份", "浙农股份", "天禾股份", "中农立华", "富邦科技", "农发种业"}:
|
return "agri_distribution", "农资流通、助剂或综合农业服务", "销量、周转、渠道效率、费用率和营运资金", "低毛利、应收存货、渠道信用和业务纯度"
|
if company in {"藏格矿业", "盐湖股份", "亚钾国际", "川发龙蟒", "川恒股份", "洛阳钼业", "甘肃能化", "兰花科创", "兖矿能源", "中煤能源"}:
|
return "agri_resource", "资源型化肥原料或跨业务资源品", "资源价格、产销量、现金成本、权益产能和资本开支", "商品周期、业务纯度、资源税费、扩产和峰值利润外推"
|
return "agri_fertilizer", "氮磷钾肥、复合肥或化工联产", "肥价、原料价差、产销量、开工率、出口和产品结构", "化肥周期、煤气磷硫成本、出口政策、库存和现金流"
|
if "新能源" in target.industries:
|
if any(track in layers for track in ("SOLAR", "BATTERY")):
|
return "newenergy_cycle", "光伏或锂电材料、设备、组件与系统", "产品价格、出货、产能利用率、单位成本和库存", "产能过剩、价格战、减值、技术迭代和现金流"
|
return "newenergy_stable", "储能、风电、核电或新型电力系统", "订单、装机、交付、利用小时、项目验收和回款", "项目周期、政策、电价、应收回款和资本开支"
|
return gen._original_classify(target, company)
|
|
|
def bubble_fields(price: float, base_high: float, optimistic_high: float, tier: str, market: dict[str, Any], business: str) -> dict[str, Any]:
|
if not tier.startswith("agri_"):
|
return gen._original_bubble_fields(price, base_high, optimistic_high, tier, market, business)
|
premium = price / base_high - 1 if base_high > 0 else None
|
vs_optimistic = price / optimistic_high - 1 if optimistic_high > 0 else None
|
if premium is None or premium <= 0:
|
return {
|
"status": "未识别估值泡沫",
|
"primary_cause": "当前价未超过基准合理区间上沿",
|
"secondary_cause": "不代表没有周期和经营风险",
|
"nature": "规则判定",
|
"confidence": "中",
|
"reason": f"当前价未超过基准上沿;{business}的价格、成本或现金流变化仍可能令合理区间下修。",
|
}
|
if vs_optimistic is not None and vs_optimistic > 0:
|
status = "泡沫-极端"
|
elif premium <= 0.15:
|
status = "泡沫-轻"
|
elif premium <= 0.50:
|
status = "泡沫-中"
|
else:
|
status = "泡沫-高"
|
causes = {
|
"agri_resource": "资源品稀缺性、价格上行或高景气持续预期",
|
"agri_fertilizer": "肥价修复、出口与成本价差改善预期",
|
"agri_pesticide": "农药周期修复、出口和库存去化预期",
|
"agri_distribution": "渠道整合、周转改善与农业服务成长预期",
|
}
|
cause = causes[tier]
|
secondary = "趋势动量与交易拥挤放大溢价" if (market.get("distance_to_ma60") or 0) > 0.10 else "乐观利润和估值倍数被提前定价"
|
return {
|
"status": status,
|
"primary_cause": cause,
|
"secondary_cause": secondary,
|
"nature": "基于估值反推、业务驱动和前复权量价的推断",
|
"confidence": "中低",
|
"reason": f"当前价高于基准上沿{premium * 100:.1f}%;市场最可能提前交易{cause},{secondary}。周期利润未被按永久高位资本化。",
|
}
|
|
|
def bounded_financials(provider: dict[str, Any]) -> tuple[dict[str, Any], dict[str, Any], dict[str, Any]]:
|
"""Allow one disclosed cash-flow schema omission without inventing profit.
|
|
*ST雅博's registered cash-flow payload omits CONSTRUCT_LONG_ASSET. Capex is
|
therefore set to zero only for the free-cash-flow cross-check; the formal
|
report/evidence manifest records that FCF is an upper bound. Profit,
|
equity, debt and valuation inputs keep their ordinary parser contracts.
|
"""
|
company = (provider.get("announcements", {}).get("identity") or {}).get("company", "")
|
if company not in {"*ST雅博", "ST雅博"}:
|
return gen._original_financials(provider)
|
original_number = gen.base._number
|
|
def optional_capex(row: dict[str, Any], *keys: str) -> float:
|
if keys == ("CONSTRUCT_LONG_ASSET",) and not any(row.get(key) is not None for key in keys):
|
return 0.0
|
return original_number(row, *keys)
|
|
gen.base._number = optional_capex
|
try:
|
return gen._original_financials(provider)
|
finally:
|
gen.base._number = original_number
|
|
|
def disclose_bounded_capex(row: dict[str, Any]) -> None:
|
if row["ticker"] != "002323.SZ":
|
return
|
report_path = ROOT / row["formal_path"]
|
report_path.write_text(
|
report_path.read_text(encoding="utf-8")
|
+ "\n## 有界数据缺口\n\n"
|
+ "- 注册现金流负载未提供 `CONSTRUCT_LONG_ASSET`(购建长期资产现金支出)字段;本报告仅在现金流交叉验证中按0处理,因此自由现金流是上限,不作为抬高合理区间的主锚。\n"
|
+ "- 利润、归母净资产、股本和三情景估值未使用该缺失字段;后续取得法定现金流明细时复核自由现金流质量。\n",
|
encoding="utf-8",
|
)
|
evidence_path = report_path.parent / "source_evidence_manifest.json"
|
evidence = gen.base.read_json(evidence_path)
|
evidence.setdefault("bounded_gaps", []).append(
|
"*ST雅博现金流镜像缺CONSTRUCT_LONG_ASSET;capex=0仅用于FCF上限交叉验证,不进入主估值锚。"
|
)
|
write_json(evidence_path, evidence)
|
|
|
def configure_generator(as_of: str) -> None:
|
stamp = as_of.replace("-", "")
|
gen.AS_OF = as_of
|
gen.STAMP = stamp
|
gen.BATCH_ID = BATCH_ID
|
gen.RESULT_ROOT = RESULT_ROOT
|
gen.SUMMARY_DIR = SUMMARY_DIR
|
gen.CASE_ROOT = CASE_ROOT
|
gen.V2_ROOT = V2_ROOT
|
gen.CACHE = V2_CACHE / "blobs/sha256"
|
gen.base.AS_OF = as_of
|
gen.base.RESULT_ROOT = RESULT_ROOT
|
gen.base.CASE_ROOT = CASE_ROOT
|
gen.base.CACHE = gen.CACHE
|
gen.PE_MULTIPLES.update(PE_ADDITIONS)
|
if not hasattr(gen, "_original_classify"):
|
gen._original_classify = gen.classify
|
if not hasattr(gen, "_original_bubble_fields"):
|
gen._original_bubble_fields = gen.bubble_fields
|
if not hasattr(gen, "_original_financials"):
|
gen._original_financials = gen.financials
|
gen.classify = classify
|
gen.bubble_fields = bubble_fields
|
gen.financials = bounded_financials
|
gen.failed_dir = package_dir
|
gen.load_provider = load_provider
|
|
|
def generate_missing(as_of: str) -> dict[str, Any]:
|
configure_generator(as_of)
|
targets, source_notes = research_targets()
|
before = existing_security()
|
missing = [targets[ticker] for ticker in sorted(set(targets) - set(before))]
|
prior_manifest_path = CASE_ROOT / "batch_manifest.json"
|
prior_manifest = gen.base.read_json(prior_manifest_path) if prior_manifest_path.exists() else {}
|
prior_results = {row["ticker"]: row for row in prior_manifest.get("results", [])}
|
rows: list[dict[str, Any]] = []
|
failures: list[dict[str, str]] = []
|
for index, target in enumerate(missing, 1):
|
try:
|
row = gen.build(target)
|
disclose_bounded_capex(row)
|
rows.append(row)
|
print(f"[{index:03d}/{len(missing):03d}] OK {target.ticker} {row['company']}", flush=True)
|
except Exception as exc:
|
error = f"{type(exc).__name__}: {exc}"
|
failures.append({"ticker": target.ticker, "company": target.research_name, "error": error})
|
print(f"[{index:03d}/{len(missing):03d}] FAIL {target.ticker} {error}", flush=True)
|
|
CASE_ROOT.mkdir(parents=True, exist_ok=True)
|
all_results = dict(prior_results)
|
all_results.update({row["ticker"]: row for row in rows})
|
result_rows = [all_results[ticker] for ticker in sorted(all_results)]
|
manifest = {
|
"batch_id": BATCH_ID,
|
"as_of": as_of,
|
"research_target_count": len(targets),
|
"existing_before": prior_manifest.get("existing_before", len(set(targets) & set(before))),
|
"missing_before": prior_manifest.get("missing_before", len(missing)),
|
"success": len(result_rows),
|
"failures": failures,
|
"sources": source_notes,
|
"results": result_rows,
|
}
|
write_json(CASE_ROOT / "batch_manifest.json", manifest)
|
|
SUMMARY_DIR.mkdir(parents=True, exist_ok=True)
|
detail_path = SUMMARY_DIR / "农业新能源新增估值明细.csv"
|
columns = [
|
"ticker", "company", "research_name", "industries", "research_layers", "price_date", "price",
|
"normalized_profit", "normalized_pe", "pb", "consensus_count", "consensus_2026", "forward_pe",
|
"base_low", "base_high", "optimistic_high", "label", "bubble_status", "bubble_primary_cause",
|
"bubble_reason", "special_valuation_status", "qa", "share_status", "formal_path",
|
]
|
with detail_path.open("w", encoding="utf-8-sig", newline="") as handle:
|
writer = csv.DictWriter(handle, fieldnames=columns, extrasaction="ignore")
|
writer.writeheader()
|
for row in result_rows:
|
item = dict(row)
|
item["industries"] = ";".join(row["industries"])
|
item["research_layers"] = ";".join(row["research_layers"])
|
item["qa"] = row["qa"]["status"]
|
writer.writerow(item)
|
labels = Counter(row["label"] for row in result_rows)
|
lines = [
|
"# 农业与新能源研究标的新增基础估值",
|
"",
|
f"- 估值日:{as_of};初始范围缺口:{manifest['missing_before']}只;已形成结论:{len(result_rows)}只;当前失败:{len(failures)}只。",
|
f"- 标签:偏低{labels['偏低']}、基本合理{labels['基本合理']}、偏贵{labels['偏贵']}、明显偏贵{labels['明显偏贵']}。",
|
"- 估值区间是条件化研究结论,不是目标价、交易指令或收益承诺。",
|
"",
|
"| 公司 | 代码 | 行业 | 收盘价 | 基准合理区间 | 判断 | 泡沫 |",
|
"|---|---|---|---:|---:|---|---|",
|
]
|
for row in result_rows:
|
link = "../../../../" + row["formal_path"]
|
lines.append(
|
f"| [{row['company']}]({link}) | {row['ticker']} | {'、'.join(row['industries'])} | {row['price']:.2f}元 | "
|
f"{row['base_low']:.2f}—{row['base_high']:.2f}元 | {row['label']} | {row['bubble_status']} |"
|
)
|
if failures:
|
lines += ["", "## 明确失败", ""] + [f"- {item['ticker']} {item['company']}:{item['error']}" for item in failures]
|
lines += ["", f"机器明细:`{detail_path.name}`。", ""]
|
(SUMMARY_DIR / "农业与新能源研究标的新增估值汇总.md").write_text("\n".join(lines), encoding="utf-8")
|
return manifest
|
|
|
def merged_research_membership() -> dict[str, dict[str, Any]]:
|
merged: dict[str, dict[str, Any]] = {}
|
|
def add(ticker: str, name: str, industries: list[str], layers: list[str]) -> None:
|
item = merged.setdefault(ticker, {"name": name, "industries": set(), "layers": set()})
|
if name and not item["name"]:
|
item["name"] = name
|
item["industries"].update(value for value in industries if value)
|
item["layers"].update(value for value in layers if value)
|
|
old_path = SUMMARY_DIR / "四行业调研标的估值覆盖清单.csv"
|
with old_path.open(encoding="utf-8-sig", newline="") as handle:
|
for row in csv.DictReader(handle):
|
add(row["ticker"], row.get("research_name") or row.get("company") or "", row.get("industries", "").split(";"), row.get("research_layers", "").split(";"))
|
targets, _ = research_targets()
|
for ticker, target in targets.items():
|
add(ticker, target.research_name, list(target.industries), list(target.research_layers))
|
return merged
|
|
|
def rebuild_coverage(trade_date: str) -> dict[str, Any]:
|
membership = merged_research_membership()
|
with mysql_connection("stock_valuation") as connection, connection.cursor() as cursor:
|
cursor.execute("SELECT ticker,company FROM security WHERE active=1")
|
securities = {row["ticker"]: row["company"] for row in cursor.fetchall()}
|
cursor.execute(
|
"SELECT j.ticker,s.company,j.trade_date,j.close,j.base_low,j.base_high,j.optimistic_high,j.label,"
|
"v.valuation_date,v.report_path FROM daily_judgement j JOIN security s ON s.ticker=j.ticker "
|
"JOIN valuation_version v ON v.valuation_id=j.valuation_id WHERE j.trade_date=%s",
|
(trade_date,),
|
)
|
judgements = {row["ticker"]: row for row in cursor.fetchall()}
|
|
rows: list[dict[str, Any]] = []
|
for ticker, item in sorted(membership.items()):
|
current = judgements.get(ticker, {})
|
if current:
|
status = "已按最新交易日判定"
|
elif ticker in securities:
|
status = "已有基础估值但当日无有效价格"
|
else:
|
status = "尚无数值基础估值或属于特殊结论"
|
rows.append({
|
"ticker": ticker,
|
"company": current.get("company") or securities.get(ticker) or item["name"],
|
"industries": ";".join(sorted(item["industries"])),
|
"research_layers": ";".join(sorted(item["layers"])),
|
"coverage_status": status,
|
"trade_date": str(current.get("trade_date") or ""),
|
"close": current.get("close", ""),
|
"base_low": current.get("base_low", ""),
|
"base_high": current.get("base_high", ""),
|
"optimistic_high": current.get("optimistic_high", ""),
|
"label": current.get("label", ""),
|
"valuation_date": str(current.get("valuation_date") or ""),
|
"report_path": current.get("report_path", ""),
|
})
|
|
SUMMARY_DIR.mkdir(parents=True, exist_ok=True)
|
csv_path = SUMMARY_DIR / "全部行业调研标的估值覆盖清单.csv"
|
with csv_path.open("w", encoding="utf-8-sig", newline="") as handle:
|
writer = csv.DictWriter(handle, fieldnames=list(rows[0]))
|
writer.writeheader()
|
writer.writerows(rows)
|
|
industries = sorted({industry for item in membership.values() for industry in item["industries"]})
|
lines = [
|
"# 全部行业调研标的估值覆盖",
|
"",
|
f"- 当前价格日:{trade_date};研究证券去重后{len(rows)}只。",
|
"- 证券池来自半导体、机器人、军工既有正式覆盖清单,以及农业、新能源当前稳定研究成果;全A股底表和待核验发现候选不纳入。",
|
"- 当前价格判定来自 MySQL `stock_valuation`;合理价值区间为条件化估值,不是目标价或交易指令。",
|
"",
|
"## 覆盖统计",
|
"",
|
"| 行业 | 调研证券 | 最新交易日已判定 | 已有估值但缺当日价 | 尚无数值估值/特殊结论 |",
|
"|---|---:|---:|---:|---:|",
|
]
|
for industry in industries:
|
pool = [row for row in rows if industry in row["industries"].split(";")]
|
counts = Counter(row["coverage_status"] for row in pool)
|
lines.append(
|
f"| {industry} | {len(pool)} | {counts['已按最新交易日判定']} | "
|
f"{counts['已有基础估值但当日无有效价格']} | {counts['尚无数值基础估值或属于特殊结论']} |"
|
)
|
labels = Counter(row["label"] for row in rows if row["label"])
|
lines += [
|
"",
|
"## 最新价格判断分布(研究池去重)",
|
"",
|
f"- 偏低{labels['偏低']}只、基本合理{labels['基本合理']}只、偏贵{labels['偏贵']}只、明显偏贵{labels['明显偏贵']}只。",
|
f"- 逐股完整清单:[{csv_path.name}]({csv_path.name})。",
|
"- 半导体179只按估值由便宜到贵的完整结果:[Markdown](半导体估值排序.md);[CSV](半导体估值排序.csv)。",
|
"- 全部已评估证券(包括用户单独点名、但不属于上述行业研究池的公司)的唯一当前入口仍是 `../估值台账/latest.md` 与 `../估值台账/latest.csv`。",
|
"",
|
]
|
md_path = SUMMARY_DIR / "全部行业调研标的估值覆盖汇总.md"
|
md_path.write_text("\n".join(lines), encoding="utf-8")
|
|
focus_rows = [row for row in rows if set(row["industries"].split(";")) & {"农业", "新能源"}]
|
label_order = {"偏低": 0, "基本合理": 1, "偏贵": 2, "明显偏贵": 3, "": 4}
|
focus_rows.sort(key=lambda row: (label_order.get(row["label"], 5), row["ticker"]))
|
focus_lines = [
|
"# 农业与新能源研究标的当前估值列表",
|
"",
|
f"- 当前完整价格日:{trade_date};农业56只,新能源115只,合计去重{len(focus_rows)}只。",
|
"- 顺序为偏低、基本合理、偏贵、明显偏贵、特殊/缺口;同档按证券代码排列。",
|
"- 合理区间是条件化估值,不是目标价、交易指令或收益承诺。",
|
"",
|
]
|
for industry in ("农业", "新能源"):
|
pool = [row for row in focus_rows if industry in row["industries"].split(";")]
|
counts = Counter(row["label"] or "特殊/缺口" for row in pool)
|
focus_lines.append(
|
f"- {industry}:偏低{counts['偏低']}、基本合理{counts['基本合理']}、偏贵{counts['偏贵']}、"
|
f"明显偏贵{counts['明显偏贵']}、特殊/缺口{counts['特殊/缺口']}。"
|
)
|
focus_lines += [
|
"",
|
"| 公司 | 代码 | 行业 | 价格日 | 收盘价 | 基准合理区间 | 判断 | 状态 | 正式报告 |",
|
"|---|---|---|---|---:|---:|---|---|---|",
|
]
|
for row in focus_rows:
|
close = "—" if row["close"] == "" else f"{float(row['close']):.2f}元"
|
base = "—" if row["base_low"] == "" else f"{float(row['base_low']):.2f}—{float(row['base_high']):.2f}元"
|
report = "—" if not row["report_path"] else f"[报告](../../../../{row['report_path']})"
|
focus_lines.append(
|
f"| {row['company']} | {row['ticker']} | {row['industries'].replace(';', '、')} | "
|
f"{row['trade_date'] or '—'} | {close} | {base} | {row['label'] or '特殊/缺口'} | "
|
f"{row['coverage_status']} | {report} |"
|
)
|
focus_path = SUMMARY_DIR / "农业与新能源研究标的当前估值列表.md"
|
focus_path.write_text("\n".join(focus_lines) + "\n", encoding="utf-8")
|
manifest = {
|
"batch_id": BATCH_ID,
|
"trade_date": trade_date,
|
"target_count": len(rows),
|
"current_judgement_count": sum(bool(row["trade_date"]) for row in rows),
|
"industry_counts": {industry: sum(industry in row["industries"].split(";") for row in rows) for industry in industries},
|
"label_counts": dict(labels),
|
"csv_sha256": hashlib.sha256(csv_path.read_bytes()).hexdigest(),
|
"md_sha256": hashlib.sha256(md_path.read_bytes()).hexdigest(),
|
"agriculture_newenergy_md_sha256": hashlib.sha256(focus_path.read_bytes()).hexdigest(),
|
}
|
ranking_md = SUMMARY_DIR / "半导体估值排序.md"
|
ranking_csv = SUMMARY_DIR / "半导体估值排序.csv"
|
if ranking_md.is_file() and ranking_csv.is_file():
|
manifest["semiconductor_ranking_md_sha256"] = hashlib.sha256(ranking_md.read_bytes()).hexdigest()
|
manifest["semiconductor_ranking_csv_sha256"] = hashlib.sha256(ranking_csv.read_bytes()).hexdigest()
|
write_json(SUMMARY_DIR / "全部行业调研标的估值覆盖manifest.json", manifest)
|
return manifest
|
|
|
def repair_generated_labels(as_of: str) -> dict[str, Any]:
|
"""Align the human conclusion with the frozen V1/ledger interval rule."""
|
configure_generator(as_of)
|
manifest_path = CASE_ROOT / "batch_manifest.json"
|
manifest = gen.base.read_json(manifest_path)
|
changed: list[dict[str, str]] = []
|
for row in manifest.get("results", []):
|
expected = gen.label_for(
|
float(row["price"]),
|
(float(row["base_low"]), float(row["base_high"])),
|
(float(row["base_high"]), float(row["optimistic_high"])),
|
)
|
old = row["label"]
|
if expected == old:
|
continue
|
report_path = ROOT / row["formal_path"]
|
text = report_path.read_text(encoding="utf-8")
|
text = text.replace(f",判断为{old};", f",判断为{expected};", 1)
|
report_path.write_text(text, encoding="utf-8")
|
evidence_path = report_path.parent / "source_evidence_manifest.json"
|
evidence = gen.base.read_json(evidence_path)
|
evidence["checks"]["label"] = expected
|
write_json(evidence_path, evidence)
|
row["label"] = expected
|
changed.append({"ticker": row["ticker"], "old": old, "new": expected})
|
write_json(manifest_path, manifest)
|
generate_missing(as_of)
|
return {"changed_count": len(changed), "changed": changed}
|
|
|
def main() -> None:
|
parser = argparse.ArgumentParser()
|
parser.add_argument("action", choices=("inventory", "fetch", "generate", "coverage", "repair"))
|
parser.add_argument("--as-of", default="2026-08-19")
|
parser.add_argument("--workers", type=int, default=12)
|
args = parser.parse_args()
|
if hasattr(sys.stdout, "reconfigure"):
|
sys.stdout.reconfigure(encoding="utf-8", errors="replace")
|
if args.action == "inventory":
|
targets, notes = research_targets()
|
covered = existing_security()
|
result = {"targets": len(targets), "covered": len(set(targets) & set(covered)), "missing": len(set(targets) - set(covered)), "sources": notes}
|
elif args.action == "fetch":
|
result = fetch_missing(args.as_of, args.workers)
|
elif args.action == "generate":
|
result = generate_missing(args.as_of)
|
elif args.action == "repair":
|
result = repair_generated_labels(args.as_of)
|
else:
|
result = rebuild_coverage(args.as_of)
|
print(json.dumps(result, ensure_ascii=False, default=str))
|
|
|
if __name__ == "__main__":
|
main()
|