wrapup.py 13 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232
  1. # -*- coding: utf-8 -*-
  2. """wrapup.py — 三柱收尾管线: 每场分析收尾 → 自动出 工单(柱1)+排程(柱2)+问答索引(柱3).
  3. 一条命令把三柱串起来: findings.json → cases(柱1生产器) → backfill(柱1回灌) →
  4. schedule(柱2, cases当任务池) → 运维agent收尾报告(含柱3常见问答预算)。
  5. = 运维Agent转化评估 §3 四卡流水线的自动收尾件; CLAUDE.md 自动化SOP "每场收尾回路"。
  6. ★纪律沿用: cases 诚实封顶CANDIDATE / backfill RULE-3只提案 / schedule ¥【参考·参数假定】 /
  7. assistant grounded+INSUFFICIENT。收尾产出=候选非结论 (过SOP§6 RV-1才交付)。
  8. 纯组合 (compose_wrapup): day_costs 作参传入 (per-farm 功率/价在 CLI 加载); 无时钟。
  9. """
  10. try:
  11. from app_common.app_common_guanlan.api import install_root as _install_root
  12. except ImportError: # 理论不可达;包结构异常时回退到按位置上跳
  13. from pathlib import Path as _P
  14. def _install_root(_f): return _P(_f).resolve().parents[4]
  15. from src import paths as P
  16. from src.sop.cases import produce_cases
  17. from src.sop.backfill import backfill
  18. from src.sop.schedule import optimize_schedule, tasks_from_cases
  19. from src.sop.assistant import answer
  20. # 收尾问答: 每场自动预算的常见问 (柱3 grounded, 出处/INSUFFICIENT 诚实)
  21. _WRAPUP_QUERIES = ['哪些台异常', '现场要做什么']
  22. def compose_wrapup(findings_doc, prev_cases, day_costs, batch_id, now,
  23. per_day_capacity=2, sched_verdicts=('准定论·预警', '定论', '候选'),
  24. price_provenance=None):
  25. """三柱组合 (纯函数): findings + 上批cases + 作业日成本 → {cases, backfill, schedule, qa}.
  26. findings_doc: findings.json dict。prev_cases: 上批 cases list。
  27. day_costs: {day: 元/台日} (柱2, None→schedule INSUFFICIENT)。batch_id/now: 批标识/日期。
  28. price_provenance: 柱2 价来源 (price.build_price_by_hour 出) → stamp 交付 verdict (假定/结算)。
  29. """
  30. farm = findings_doc.get('_farm', '?')
  31. findings = findings_doc.get('findings', [])
  32. # 柱1: 生产器 + 回灌器
  33. cases = produce_cases(findings_doc, prev_cases or [], batch_id, now)
  34. bf = backfill(cases)
  35. # 柱2: cases → 任务池 → 排程 (放宽含候选=预排程) + 价 provenance 联动交付 verdict
  36. tasks = tasks_from_cases(cases, verdicts=sched_verdicts)
  37. if day_costs and tasks['turbine_days'] > 0:
  38. sched = optimize_schedule(day_costs, tasks['turbine_days'], per_day_capacity)
  39. if price_provenance:
  40. from src.sop.price import stamp_schedule_verdict
  41. sched = stamp_schedule_verdict(sched, price_provenance)
  42. else:
  43. sched = {'verdict': 'INSUFFICIENT',
  44. 'reason': ('无候选停机任务' if tasks['turbine_days'] == 0 else '无功率/价格 day_costs (柱2跳过)')}
  45. sched['tasks'] = tasks
  46. # 柱3: 收尾常见问答 (grounded)
  47. qa = [{'q': q, **answer(q, cases, findings, farm=farm)} for q in _WRAPUP_QUERIES]
  48. return {'farm': farm, 'cases': cases, 'backfill': bf, 'schedule': sched, 'qa': qa}
  49. _VD_RANK = {'定论': 5, '准定论·预警': 4, '候选': 3, '参考': 2, 'INSUFFICIENT': 1, '撤回': 0}
  50. def _ev_val(c):
  51. v = (c.get('evidence') or [{}])[0].get('value')
  52. return v if isinstance(v, (int, float)) else -1e9
  53. def render_md(wrap, farm='?'):
  54. """收尾报告 markdown (三柱一页): 优先关注引导 → 三柱详情。深挖层(case.analysis)不埋没。"""
  55. from collections import Counter
  56. cs, bf, sc = wrap['cases'], wrap['backfill'], wrap['schedule']
  57. L = [f'> ⚠️ **运维agent收尾报告 · {farm}** — `src/sop/wrapup.py` 生成, 勿手改。'
  58. f'三柱产出=**候选非结论**, 过 SOP§6 RV-1 + data-first 复现才交付。\n',
  59. f'# 运维 agent 收尾 · {farm}\n']
  60. # ── 优先关注 (深挖 ≥准定论·预警 领先; 无深挖则 top 证据) ──
  61. deep = [c for c in cs if isinstance(c.get('analysis'), dict)
  62. and _VD_RANK.get(str(c['analysis'].get('verdict', '')).strip('【】'), -1) >= _VD_RANK['准定论·预警']]
  63. L.append('## ⭐ 优先关注')
  64. if deep:
  65. for c in sorted(deep, key=lambda c: -_VD_RANK.get(str(c['analysis'].get('verdict', '')).strip('【】'), -1)):
  66. an = c['analysis']
  67. L.append(f'- **{c["turbine"]} · {c["system"]}** 【深挖 {an.get("verdict")}】{an.get("mechanism", "")}')
  68. if an.get('field_needs'):
  69. L.append(f' - 现场需求: {an["field_needs"]}' + (f' (据 {str(an.get("ref", "")).split("/")[-1]})' if an.get('ref') else ''))
  70. else:
  71. L.append(f'- 无深挖级(≥准定论·预警)结论; {len(cs)} 个候选待现场坐实 (见工单层)。')
  72. # ── ① 柱1 工单层 (深挖 inline) ──
  73. L.append('\n## ① 工单层 (柱1: cases)')
  74. L.append(f'- {len(cs)} 个 case; 状态 {dict(Counter(c["state"] for c in cs))}; '
  75. f'系统 {dict(Counter(c["system"] for c in cs))}')
  76. for c in sorted(cs, key=lambda c: -_ev_val(c))[:8]:
  77. ev = (c.get('evidence') or [{}])[0]
  78. L.append(f' - {c["turbine"]} · {c["system"]}: {c["state"]}/{c["verdict"]} '
  79. f'(判据{ev.get("discriminator_id")} 值{ev.get("value")}) · {c["case_id"]}')
  80. an = c.get('analysis')
  81. if isinstance(an, dict):
  82. L.append(f' ▸ 深挖[{an.get("verdict")}]: {str(an.get("mechanism", ""))[:70]}')
  83. L.append('\n## ① 回灌记分卡 (柱1: backfill, RULE-3只提案不改阈值)')
  84. L.append(f'- {len(bf["by_discriminator"])} 判据; {len(bf["proposals"])} 提案; '
  85. f'{len(bf["fp_ledger"])} 误报; {len(bf["relapse_ledger"])} 复发')
  86. for p in bf['proposals']:
  87. L.append(f' - 【{p["verdict"]}】{p["discriminator_id"]}: {p["reason"]}')
  88. # ── ② 柱2 排程层 (delivery verdict 据价 provenance) ──
  89. L.append('\n## ② 排程层 (柱2: schedule)')
  90. if sc.get('verdict') == 'INSUFFICIENT':
  91. L.append(f'- INSUFFICIENT: {sc.get("reason")}')
  92. else:
  93. dv = sc.get('delivery_verdict', '【参考·参数假定·待业主结算价】')
  94. L.append(f'- 任务池(柱1桥) {sc["tasks"]["turbine_days"]}台日: {[t["turbine"] for t in sc["tasks"]["turbines"]]}')
  95. L.append(f'- 机会成本: 最优 ¥{sc["best_cny"]:,} | 均匀 ¥{sc["mean_cny"]:,} | 最差 ¥{sc["worst_cny"]:,}')
  96. L.append(f'- **纯排程收益(最优vs均匀): ¥{sc["saving_vs_mean"]:,} ({sc["saving_pct"]}%)** — 不修任何东西 {dv}')
  97. pp = sc.get('price_provenance')
  98. if pp:
  99. L.append(f' - 价来源: {pp.get("note", "")}')
  100. # ── ③ 柱3 助手层 ──
  101. L.append('\n## ③ 助手层 (柱3: assistant, grounded/INSUFFICIENT诚实)')
  102. for item in wrap['qa']:
  103. L.append(f'- Q: {item["q"]} → [{item["verdict"]}] {item["answer"][:160].replace(chr(10), " ")}')
  104. return '\n'.join(L)
  105. # ---------------------------------------------------------------------------
  106. if __name__ == '__main__':
  107. # 控制台可能是 GBK(中文 Windows 代码页 936): 正文里的 ✔ ✗ ✅ ⚠ 这类字符编不出来会抛
  108. # UnicodeEncodeError, 脚本干成了事却以退出码 1 结束(同类坑见 src/console.py)。降级为 '?' 而不是崩;
  109. # 不用 import 是为了兼顾 python -m 与直接当脚本跑两种启动方式。
  110. import sys as _sys
  111. for _s in (_sys.stdout, _sys.stderr):
  112. try: _s.reconfigure(errors='replace')
  113. except Exception: pass
  114. # CLI: python -m src.sop.wrapup <farm> [--price-anchor 300]
  115. # 读 outputs/<farm>/sop/findings.json (+ cases.json 上批 + cleaned parquet) → 三柱收尾 → 写盘
  116. import argparse, json, sys
  117. from pathlib import Path
  118. from src.sop.schedule import day_costs_from_power_price
  119. from src.sop.schemas import validate_cases
  120. ROOT = _install_root(__file__)
  121. ap = argparse.ArgumentParser()
  122. ap.add_argument('farm')
  123. ap.add_argument('--price-anchor', type=float, default=300.0, help='电价锚 元/MWh (假定)')
  124. ap.add_argument('--power-col', default='发电机有功功率')
  125. ap.add_argument('--dry', action='store_true')
  126. a = ap.parse_args()
  127. sopd = P.sop(a.farm)
  128. fdoc = json.loads((sopd / 'findings.json').read_text(encoding='utf-8'))
  129. cpath = sopd / 'cases.json'
  130. prev = json.loads(cpath.read_text(encoding='utf-8')) if cpath.exists() else []
  131. batch_id = fdoc.get('_generated', 'unknown')
  132. now = str(batch_id)[:10]
  133. # 柱2 day_costs: cleaned parquet 场中位功率 × price.py 参数化价 (config 驱动 + per月季节 + provenance)
  134. day_costs, provenance = None, None
  135. # ★ 写读要配对 (2026-09-16 修): 写侧 `src/sop/contract_gate.py:131` 落的是 `cleaned/<section>.parquet`
  136. # (section 来自契约配置, 可能是 turbine / turbine_10min_cnt / …), 而这里原先**硬编码** 'turbine.parquet'
  137. # ⇒ section 不是 turbine 时 pq.exists() 恒假 → day_costs 静默变 None、柱2 悄悄降级成 INSUFFICIENT,
  138. # 一句报错都没有。改为: 优先读 clean_gate.json 里记的 section, 否则取该目录下唯一的 parquet;
  139. # 仍然取不到就**响亮**说明 (而不是让下游以为"这场的经济性算不出来")。
  140. _cdir = P.out_root(a.farm) / 'cleaned'
  141. pq, _why = None, ''
  142. _gate = _cdir / 'clean_gate.json'
  143. try:
  144. if _gate.exists():
  145. import json as _j
  146. _sec = (_j.loads(_gate.read_text(encoding='utf-8')) or {}).get('section')
  147. if _sec and (_cdir / f'{_sec}.parquet').exists():
  148. pq = _cdir / f'{_sec}.parquet'
  149. except Exception:
  150. pass
  151. if pq is None:
  152. _cands = sorted(_cdir.glob('*.parquet')) if _cdir.is_dir() else []
  153. if len(_cands) == 1:
  154. pq = _cands[0]
  155. elif _cands:
  156. pq, _why = _cands[0], f'目录里有 {len(_cands)} 个 parquet, 取第一个 {_cands[0].name}'
  157. else:
  158. _why = f'{P.rel(_cdir)} 下没有 cleaned parquet (清洗闸未落盘? persist=False?)'
  159. if pq is None:
  160. print(f'[wrapup] 柱2 无 cleaned parquet → day_costs 保持 None: {_why}', flush=True)
  161. elif _why:
  162. print(f'[wrapup] 柱2 cleaned parquet 选取说明: {_why}', flush=True)
  163. vcfg = P.config('value_assumptions', f'{a.farm}.yaml') # 配置唯一取用口: configs/value_assumptions/<场>.yaml
  164. if pq is not None and pq.exists():
  165. import pandas as pd
  166. import yaml
  167. from src.sop.price import build_price_by_hour
  168. from src.sop.schema_guard import verified_read
  169. df, _ = verified_read(pq, columns=['时间', a.power_col]) # A1: 读前核+列断言(power_col缺→响亮)
  170. df['时间'] = pd.to_datetime(df['时间'], errors='coerce')
  171. df = df.dropna(subset=['时间'])
  172. df['date'] = df['时间'].dt.strftime('%Y-%m-%d'); df['hour'] = df['时间'].dt.hour
  173. df['month'] = df['时间'].dt.month
  174. med = df.groupby(['date', 'hour', 'month'])[a.power_col].median().reset_index()
  175. if vcfg.exists():
  176. price_cfg = yaml.safe_load(vcfg.read_text(encoding='utf-8'))['price']
  177. else: # 无场配置 → 用锚+默认五段占位 (标假定; 建议 cp templates/value_layer_assumptions.example.yaml)
  178. price_cfg = {'benchmark_cny_per_mwh': a.price_anchor, 'benchmark_assumed': True,
  179. 'spot': {'tou_assumed': True, 'tou_factor': {
  180. 'valley': {'hours': '00-08', 'factor': 0.55},
  181. 'flat': {'hours': '08-17,22-24', 'factor': 1.0},
  182. 'peak': {'hours': '17-22', 'factor': 1.75}}}}
  183. day_costs = {}
  184. for month, sub in med.groupby('month'):
  185. price_m, provenance = build_price_by_hour(price_cfg, month=int(month))
  186. pbdh = {(r.date, r.hour): r[a.power_col] for _, r in sub.iterrows()}
  187. day_costs.update(day_costs_from_power_price(pbdh, price_m))
  188. if provenance: # 全年多月, note 统一为 config 展开(assumed/结算 一致)
  189. provenance = dict(provenance)
  190. provenance['note'] = f"{a.farm} value_assumptions 参数化价 × per月季节 (全年逐月展开)"
  191. wrap = compose_wrapup(fdoc, prev, day_costs, batch_id, now, price_provenance=provenance)
  192. issues = validate_cases(wrap['cases'])
  193. if issues:
  194. print(f'❌ validate_cases: {issues[:5]}'); sys.exit(2)
  195. sc = wrap['schedule']
  196. print(f'✅ {a.farm} 三柱收尾: {len(wrap["cases"])} case | '
  197. f'排程 {"¥%s收益" % format(sc.get("saving_vs_mean",0), ",") if sc.get("verdict")!="INSUFFICIENT" else "INSUFFICIENT"} | '
  198. f'{len(wrap["qa"])} 问答')
  199. if not a.dry:
  200. cpath.write_text(json.dumps(wrap['cases'], ensure_ascii=False, indent=2), encoding='utf-8')
  201. (sopd / 'schedule_report.json').write_text(json.dumps(sc, ensure_ascii=False, indent=2), encoding='utf-8')
  202. (sopd / '运维agent收尾报告.md').write_text(render_md(wrap, a.farm), encoding='utf-8')
  203. print(f' 写盘: {sopd.relative_to(ROOT)}/{{cases.json, schedule_report.json, 运维agent收尾报告.md}}')