| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232 |
- # -*- coding: utf-8 -*-
- """wrapup.py — 三柱收尾管线: 每场分析收尾 → 自动出 工单(柱1)+排程(柱2)+问答索引(柱3).
- 一条命令把三柱串起来: findings.json → cases(柱1生产器) → backfill(柱1回灌) →
- schedule(柱2, cases当任务池) → 运维agent收尾报告(含柱3常见问答预算)。
- = 运维Agent转化评估 §3 四卡流水线的自动收尾件; CLAUDE.md 自动化SOP "每场收尾回路"。
- ★纪律沿用: cases 诚实封顶CANDIDATE / backfill RULE-3只提案 / schedule ¥【参考·参数假定】 /
- assistant grounded+INSUFFICIENT。收尾产出=候选非结论 (过SOP§6 RV-1才交付)。
- 纯组合 (compose_wrapup): day_costs 作参传入 (per-farm 功率/价在 CLI 加载); 无时钟。
- """
- try:
- from app_common.app_common_guanlan.api import install_root as _install_root
- except ImportError: # 理论不可达;包结构异常时回退到按位置上跳
- from pathlib import Path as _P
- def _install_root(_f): return _P(_f).resolve().parents[4]
- from src import paths as P
- from src.sop.cases import produce_cases
- from src.sop.backfill import backfill
- from src.sop.schedule import optimize_schedule, tasks_from_cases
- from src.sop.assistant import answer
- # 收尾问答: 每场自动预算的常见问 (柱3 grounded, 出处/INSUFFICIENT 诚实)
- _WRAPUP_QUERIES = ['哪些台异常', '现场要做什么']
- def compose_wrapup(findings_doc, prev_cases, day_costs, batch_id, now,
- per_day_capacity=2, sched_verdicts=('准定论·预警', '定论', '候选'),
- price_provenance=None):
- """三柱组合 (纯函数): findings + 上批cases + 作业日成本 → {cases, backfill, schedule, qa}.
- findings_doc: findings.json dict。prev_cases: 上批 cases list。
- day_costs: {day: 元/台日} (柱2, None→schedule INSUFFICIENT)。batch_id/now: 批标识/日期。
- price_provenance: 柱2 价来源 (price.build_price_by_hour 出) → stamp 交付 verdict (假定/结算)。
- """
- farm = findings_doc.get('_farm', '?')
- findings = findings_doc.get('findings', [])
- # 柱1: 生产器 + 回灌器
- cases = produce_cases(findings_doc, prev_cases or [], batch_id, now)
- bf = backfill(cases)
- # 柱2: cases → 任务池 → 排程 (放宽含候选=预排程) + 价 provenance 联动交付 verdict
- tasks = tasks_from_cases(cases, verdicts=sched_verdicts)
- if day_costs and tasks['turbine_days'] > 0:
- sched = optimize_schedule(day_costs, tasks['turbine_days'], per_day_capacity)
- if price_provenance:
- from src.sop.price import stamp_schedule_verdict
- sched = stamp_schedule_verdict(sched, price_provenance)
- else:
- sched = {'verdict': 'INSUFFICIENT',
- 'reason': ('无候选停机任务' if tasks['turbine_days'] == 0 else '无功率/价格 day_costs (柱2跳过)')}
- sched['tasks'] = tasks
- # 柱3: 收尾常见问答 (grounded)
- qa = [{'q': q, **answer(q, cases, findings, farm=farm)} for q in _WRAPUP_QUERIES]
- return {'farm': farm, 'cases': cases, 'backfill': bf, 'schedule': sched, 'qa': qa}
- _VD_RANK = {'定论': 5, '准定论·预警': 4, '候选': 3, '参考': 2, 'INSUFFICIENT': 1, '撤回': 0}
- def _ev_val(c):
- v = (c.get('evidence') or [{}])[0].get('value')
- return v if isinstance(v, (int, float)) else -1e9
- def render_md(wrap, farm='?'):
- """收尾报告 markdown (三柱一页): 优先关注引导 → 三柱详情。深挖层(case.analysis)不埋没。"""
- from collections import Counter
- cs, bf, sc = wrap['cases'], wrap['backfill'], wrap['schedule']
- L = [f'> ⚠️ **运维agent收尾报告 · {farm}** — `src/sop/wrapup.py` 生成, 勿手改。'
- f'三柱产出=**候选非结论**, 过 SOP§6 RV-1 + data-first 复现才交付。\n',
- f'# 运维 agent 收尾 · {farm}\n']
- # ── 优先关注 (深挖 ≥准定论·预警 领先; 无深挖则 top 证据) ──
- deep = [c for c in cs if isinstance(c.get('analysis'), dict)
- and _VD_RANK.get(str(c['analysis'].get('verdict', '')).strip('【】'), -1) >= _VD_RANK['准定论·预警']]
- L.append('## ⭐ 优先关注')
- if deep:
- for c in sorted(deep, key=lambda c: -_VD_RANK.get(str(c['analysis'].get('verdict', '')).strip('【】'), -1)):
- an = c['analysis']
- L.append(f'- **{c["turbine"]} · {c["system"]}** 【深挖 {an.get("verdict")}】{an.get("mechanism", "")}')
- if an.get('field_needs'):
- L.append(f' - 现场需求: {an["field_needs"]}' + (f' (据 {str(an.get("ref", "")).split("/")[-1]})' if an.get('ref') else ''))
- else:
- L.append(f'- 无深挖级(≥准定论·预警)结论; {len(cs)} 个候选待现场坐实 (见工单层)。')
- # ── ① 柱1 工单层 (深挖 inline) ──
- L.append('\n## ① 工单层 (柱1: cases)')
- L.append(f'- {len(cs)} 个 case; 状态 {dict(Counter(c["state"] for c in cs))}; '
- f'系统 {dict(Counter(c["system"] for c in cs))}')
- for c in sorted(cs, key=lambda c: -_ev_val(c))[:8]:
- ev = (c.get('evidence') or [{}])[0]
- L.append(f' - {c["turbine"]} · {c["system"]}: {c["state"]}/{c["verdict"]} '
- f'(判据{ev.get("discriminator_id")} 值{ev.get("value")}) · {c["case_id"]}')
- an = c.get('analysis')
- if isinstance(an, dict):
- L.append(f' ▸ 深挖[{an.get("verdict")}]: {str(an.get("mechanism", ""))[:70]}')
- L.append('\n## ① 回灌记分卡 (柱1: backfill, RULE-3只提案不改阈值)')
- L.append(f'- {len(bf["by_discriminator"])} 判据; {len(bf["proposals"])} 提案; '
- f'{len(bf["fp_ledger"])} 误报; {len(bf["relapse_ledger"])} 复发')
- for p in bf['proposals']:
- L.append(f' - 【{p["verdict"]}】{p["discriminator_id"]}: {p["reason"]}')
- # ── ② 柱2 排程层 (delivery verdict 据价 provenance) ──
- L.append('\n## ② 排程层 (柱2: schedule)')
- if sc.get('verdict') == 'INSUFFICIENT':
- L.append(f'- INSUFFICIENT: {sc.get("reason")}')
- else:
- dv = sc.get('delivery_verdict', '【参考·参数假定·待业主结算价】')
- L.append(f'- 任务池(柱1桥) {sc["tasks"]["turbine_days"]}台日: {[t["turbine"] for t in sc["tasks"]["turbines"]]}')
- L.append(f'- 机会成本: 最优 ¥{sc["best_cny"]:,} | 均匀 ¥{sc["mean_cny"]:,} | 最差 ¥{sc["worst_cny"]:,}')
- L.append(f'- **纯排程收益(最优vs均匀): ¥{sc["saving_vs_mean"]:,} ({sc["saving_pct"]}%)** — 不修任何东西 {dv}')
- pp = sc.get('price_provenance')
- if pp:
- L.append(f' - 价来源: {pp.get("note", "")}')
- # ── ③ 柱3 助手层 ──
- L.append('\n## ③ 助手层 (柱3: assistant, grounded/INSUFFICIENT诚实)')
- for item in wrap['qa']:
- L.append(f'- Q: {item["q"]} → [{item["verdict"]}] {item["answer"][:160].replace(chr(10), " ")}')
- return '\n'.join(L)
- # ---------------------------------------------------------------------------
- if __name__ == '__main__':
- # 控制台可能是 GBK(中文 Windows 代码页 936): 正文里的 ✔ ✗ ✅ ⚠ 这类字符编不出来会抛
- # UnicodeEncodeError, 脚本干成了事却以退出码 1 结束(同类坑见 src/console.py)。降级为 '?' 而不是崩;
- # 不用 import 是为了兼顾 python -m 与直接当脚本跑两种启动方式。
- import sys as _sys
- for _s in (_sys.stdout, _sys.stderr):
- try: _s.reconfigure(errors='replace')
- except Exception: pass
- # CLI: python -m src.sop.wrapup <farm> [--price-anchor 300]
- # 读 outputs/<farm>/sop/findings.json (+ cases.json 上批 + cleaned parquet) → 三柱收尾 → 写盘
- import argparse, json, sys
- from pathlib import Path
- from src.sop.schedule import day_costs_from_power_price
- from src.sop.schemas import validate_cases
- ROOT = _install_root(__file__)
- ap = argparse.ArgumentParser()
- ap.add_argument('farm')
- ap.add_argument('--price-anchor', type=float, default=300.0, help='电价锚 元/MWh (假定)')
- ap.add_argument('--power-col', default='发电机有功功率')
- ap.add_argument('--dry', action='store_true')
- a = ap.parse_args()
- sopd = P.sop(a.farm)
- fdoc = json.loads((sopd / 'findings.json').read_text(encoding='utf-8'))
- cpath = sopd / 'cases.json'
- prev = json.loads(cpath.read_text(encoding='utf-8')) if cpath.exists() else []
- batch_id = fdoc.get('_generated', 'unknown')
- now = str(batch_id)[:10]
- # 柱2 day_costs: cleaned parquet 场中位功率 × price.py 参数化价 (config 驱动 + per月季节 + provenance)
- day_costs, provenance = None, None
- # ★ 写读要配对 (2026-09-16 修): 写侧 `src/sop/contract_gate.py:131` 落的是 `cleaned/<section>.parquet`
- # (section 来自契约配置, 可能是 turbine / turbine_10min_cnt / …), 而这里原先**硬编码** 'turbine.parquet'
- # ⇒ section 不是 turbine 时 pq.exists() 恒假 → day_costs 静默变 None、柱2 悄悄降级成 INSUFFICIENT,
- # 一句报错都没有。改为: 优先读 clean_gate.json 里记的 section, 否则取该目录下唯一的 parquet;
- # 仍然取不到就**响亮**说明 (而不是让下游以为"这场的经济性算不出来")。
- _cdir = P.out_root(a.farm) / 'cleaned'
- pq, _why = None, ''
- _gate = _cdir / 'clean_gate.json'
- try:
- if _gate.exists():
- import json as _j
- _sec = (_j.loads(_gate.read_text(encoding='utf-8')) or {}).get('section')
- if _sec and (_cdir / f'{_sec}.parquet').exists():
- pq = _cdir / f'{_sec}.parquet'
- except Exception:
- pass
- if pq is None:
- _cands = sorted(_cdir.glob('*.parquet')) if _cdir.is_dir() else []
- if len(_cands) == 1:
- pq = _cands[0]
- elif _cands:
- pq, _why = _cands[0], f'目录里有 {len(_cands)} 个 parquet, 取第一个 {_cands[0].name}'
- else:
- _why = f'{P.rel(_cdir)} 下没有 cleaned parquet (清洗闸未落盘? persist=False?)'
- if pq is None:
- print(f'[wrapup] 柱2 无 cleaned parquet → day_costs 保持 None: {_why}', flush=True)
- elif _why:
- print(f'[wrapup] 柱2 cleaned parquet 选取说明: {_why}', flush=True)
- vcfg = P.config('value_assumptions', f'{a.farm}.yaml') # 配置唯一取用口: configs/value_assumptions/<场>.yaml
- if pq is not None and pq.exists():
- import pandas as pd
- import yaml
- from src.sop.price import build_price_by_hour
- from src.sop.schema_guard import verified_read
- df, _ = verified_read(pq, columns=['时间', a.power_col]) # A1: 读前核+列断言(power_col缺→响亮)
- df['时间'] = pd.to_datetime(df['时间'], errors='coerce')
- df = df.dropna(subset=['时间'])
- df['date'] = df['时间'].dt.strftime('%Y-%m-%d'); df['hour'] = df['时间'].dt.hour
- df['month'] = df['时间'].dt.month
- med = df.groupby(['date', 'hour', 'month'])[a.power_col].median().reset_index()
- if vcfg.exists():
- price_cfg = yaml.safe_load(vcfg.read_text(encoding='utf-8'))['price']
- else: # 无场配置 → 用锚+默认五段占位 (标假定; 建议 cp templates/value_layer_assumptions.example.yaml)
- price_cfg = {'benchmark_cny_per_mwh': a.price_anchor, 'benchmark_assumed': True,
- 'spot': {'tou_assumed': True, 'tou_factor': {
- 'valley': {'hours': '00-08', 'factor': 0.55},
- 'flat': {'hours': '08-17,22-24', 'factor': 1.0},
- 'peak': {'hours': '17-22', 'factor': 1.75}}}}
- day_costs = {}
- for month, sub in med.groupby('month'):
- price_m, provenance = build_price_by_hour(price_cfg, month=int(month))
- pbdh = {(r.date, r.hour): r[a.power_col] for _, r in sub.iterrows()}
- day_costs.update(day_costs_from_power_price(pbdh, price_m))
- if provenance: # 全年多月, note 统一为 config 展开(assumed/结算 一致)
- provenance = dict(provenance)
- provenance['note'] = f"{a.farm} value_assumptions 参数化价 × per月季节 (全年逐月展开)"
- wrap = compose_wrapup(fdoc, prev, day_costs, batch_id, now, price_provenance=provenance)
- issues = validate_cases(wrap['cases'])
- if issues:
- print(f'❌ validate_cases: {issues[:5]}'); sys.exit(2)
- sc = wrap['schedule']
- print(f'✅ {a.farm} 三柱收尾: {len(wrap["cases"])} case | '
- f'排程 {"¥%s收益" % format(sc.get("saving_vs_mean",0), ",") if sc.get("verdict")!="INSUFFICIENT" else "INSUFFICIENT"} | '
- f'{len(wrap["qa"])} 问答')
- if not a.dry:
- cpath.write_text(json.dumps(wrap['cases'], ensure_ascii=False, indent=2), encoding='utf-8')
- (sopd / 'schedule_report.json').write_text(json.dumps(sc, ensure_ascii=False, indent=2), encoding='utf-8')
- (sopd / '运维agent收尾报告.md').write_text(render_md(wrap, a.farm), encoding='utf-8')
- print(f' 写盘: {sopd.relative_to(ROOT)}/{{cases.json, schedule_report.json, 运维agent收尾报告.md}}')
|