| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416 |
- # -*- coding: utf-8 -*-
- """cases.py — findings.json(每批快照) → cases.json(跨批次持续档) 生产器 (运维 agent 柱1「智能工单」).
- 为什么: findings.json = 每批全量快照, 无记性 → 同一前兆每月被"重新发现"(B16)。
- case = (farm,turbine,system) 的持续健康档: 跨批次状态延续 + 反向闭环的左半 (evidence 追加)。
- 本模块 = 附录A 状态机的**产线**(schemas.validate_case 是**闸**); 产出必过 validate_cases。
- 设计红线 (与 schemas 不变量同源):
- ① 诚实封顶 CANDIDATE: findings 多为 相对排名(fleet-z), 禁 fleet-z 直升 CONFIRMED(不变量①)。
- 升 CONFIRMED 必 own-data (finding['own_data_confirmed']=True), 生产器不伪造。
- ② 一 key 一档永续: 每 (farm,turbine,system) 全程只一个 case, 跨批改状态, 不重复开案(不变量⑥)。
- ③ TICKETED 及以后 = 人工裁决点①, 生产器**不碰** (只在 NORMAL↔WATCH↔CANDIDATE↔CONFIRMED↔REFUTED
- + CLOSED→RELAPSED 之间自动流转)。工单签发/销案是人的动作, 走独立入口。
- 纯核心 (produce_cases): now/batch_id 作参传入, 不碰文件系统/时钟 (副作用在 CLI __main__)。
- """
- 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
- import json
- import re
- import sys
- from pathlib import Path
- # 状态序 (advance 用): NORMAL<WATCH<CANDIDATE<CONFIRMED
- _STATE_ORDER = {'NORMAL': 0, 'WATCH': 1, 'CANDIDATE': 2, 'CONFIRMED': 3}
- # verdict → 目标态 (诚实封顶: 无 own-data 时 定论/准定论·预警 也只到 CANDIDATE)
- _VERDICT_STATE = {
- '定论': 'CANDIDATE', '准定论·预警': 'CANDIDATE', '候选': 'CANDIDATE', '参考': 'WATCH',
- }
- # 标题 → discriminator_registry 正式 id (有序命中即停; 无命中退 slug 待补)。
- # configs/discriminator_registry.yaml 单源; test_cases_producer_dp 断言映射目标 ∈ 注册表 (外键真)。
- _DISCRIMINATOR_RULES = [
- ('对风', 'static_yaw_assessment'), ('静态偏航', 'static_yaw_assessment'),
- ('三相', 'phase_imbalance'), ('相偏', 'phase_imbalance'), ('相不平衡', 'phase_imbalance'),
- ('1P', 'blade_imbalance_1p'),
- ('变桨', 'pitch_stick'), ('桨距', 'pitch_stick'),
- ('残差', 'nbm_residual'), ('NBM', 'nbm_residual'), ('温度', 'nbm_residual'), ('气动', 'nbm_residual'),
- ('偏航', 'yaw_aging_operators'),
- ]
- def discriminator_id_of(finding):
- """finding.title → 判据注册表正式 id (evidence.discriminator_id 外键真)。无映射→ slug(title) 待补。"""
- title = str(finding.get('title', ''))
- for kw, did in _DISCRIMINATOR_RULES:
- if kw in title:
- return did
- return _slug(title)
- def system_of(finding):
- """finding.title → 规范 system key (跨批次稳定, 跨 OEM 命名前缀感知)。附录A.3 system 枚举。
- ★前缀感知 (2026-07-07, 岳林柏第二样板逮到 farm-tuned gap): 跨 OEM 通道命名歧义 —
- '驱动端' 华电=发电机DE / 盂县'齿轮箱中间轴驱动端'=齿箱IMS-DE。故: **子系统前缀优先** (齿轮箱/发电机)
- 定归属, 再定部件; **裸命名** (无子系统前缀) 按行业约定 (低速/高速轴承=齿箱, 驱动端/自由端=发电机)。"""
- t = str(finding.get('title', ''))
- # 非温度系统 (关键词明确, 先拦)
- if '偏航' in t or '对风' in t:
- return 'yaw'
- if '变桨' in t or '桨距' in t:
- return 'pitch'
- if '欠发' in t:
- return 'underperf'
- if '变流' in t or '变频' in t:
- return 'converter'
- if '安全链' in t:
- return 'safety_chain'
- if '主轴承' in t:
- return 'main_brg'
- # DE/NDE 判别 (非驱动端/自由端 = NDE; 优先于 DE 判)
- is_nde = ('非驱动端' in t) or ('自由端' in t)
- is_de = ('驱动端' in t) and not is_nde
- # 齿轮箱子系统 (前缀优先: '齿轮箱中间轴驱动端'=齿箱 非发电机)
- if '齿轮箱' in t or '齿轮' in t:
- if '油' in t or '润滑' in t:
- return 'gearbox_oil'
- if '高速' in t:
- return 'gearbox_hss_brg'
- if '中速' in t or '中间轴' in t:
- return 'gearbox_ims_de' if is_de else ('gearbox_ims_nde' if is_nde else 'gearbox_ims_brg')
- if '低速' in t:
- return 'gearbox_lss_brg'
- return 'gearbox'
- # 发电机子系统 (显式前缀)
- if '发电机' in t:
- if '定子' in t:
- return 'gen_stator'
- if is_de:
- return 'gen_de_brg'
- if is_nde:
- return 'gen_nde_brg'
- return 'gen'
- # 裸命名 (无子系统前缀) — 行业约定
- if '定子' in t:
- return 'gen_stator'
- if '高速轴' in t:
- return 'gearbox_hss_brg'
- if '中速轴' in t or '中间轴' in t:
- return 'gearbox_ims_brg'
- if '低速轴' in t:
- return 'gearbox_lss_brg'
- if is_de:
- return 'gen_de_brg'
- if is_nde:
- return 'gen_nde_brg'
- if '主轴' in t:
- return 'main_shaft'
- if '传动链' in t:
- return 'driveline'
- if '联轴器' in t:
- return 'coupling'
- slug = re.sub(r'[^a-z0-9一-鿿]+', '_', t.lower()).strip('_')
- return slug[:32] or 'unknown'
- def _slug(s):
- return re.sub(r'[^a-z0-9一-鿿]+', '_', str(s).lower()).strip('_')[:48] or 'x'
- def _disc(finding):
- """finding['discriminator'] 归一为 dict。
- 该字段实测有三种写法: dict (结构化判别器, 53 场) / None (无) /
- **散文字符串** (rudong 13 条把它当"判别方法说明"写)。原写法 `... or {}` 只兜住
- None 与空串, 非空字符串是 truthy 会漏进 .get() -> AttributeError。
- 散文串没有 per_turbine/flagged, 语义上等同"无" -> 与 None 同样处理。
- """
- d = finding.get('discriminator')
- return d if isinstance(d, dict) else {}
- def flagged_turbines(finding):
- """finding 明确 flag 的台号 list (affected_machines 优先, 退 discriminator.flagged)。
- 空 = 该 finding 不制造 case (参考/INSUFFICIENT/健康场级 → affected null → 天然静默)。"""
- am = finding.get('affected_machines')
- if isinstance(am, str):
- return [am]
- if isinstance(am, list) and am:
- return [str(t) for t in am]
- disc = _disc(finding)
- fl = disc.get('flagged')
- if isinstance(fl, list) and fl:
- return [str(t) for t in fl]
- return []
- def _turbine_value(finding, tid):
- """从 discriminator.per_turbine 抽该台的代表值 (z 优先, 退 resid_mean)。无→None."""
- disc = _disc(finding)
- pt = disc.get('per_turbine')
- if isinstance(pt, list):
- for r in pt:
- if isinstance(r, dict) and str(r.get('tid')) == str(tid):
- return r.get('z', r.get('resid_mean'))
- return None
- def _target_state(finding):
- """verdict → 目标态 (own-data 才允许 CONFIRMED; 否则封顶 CANDIDATE)。"""
- from src.sop.schemas import norm_verdict
- v = norm_verdict(finding.get('verdict', ''))
- base = _VERDICT_STATE.get(v)
- if base is None:
- return None # INSUFFICIENT/撤回/未知 → 不开案
- if v in ('定论', '准定论·预警') and finding.get('own_data_confirmed'):
- return 'CONFIRMED'
- return base
- def _open_history(target, batch_id, now):
- """新档: 合成 NORMAL→…→target 的合法迁移链 (附录A.2 边)。"""
- chain = [('NORMAL', 'WATCH', 'signal_raised')]
- if _STATE_ORDER[target] >= _STATE_ORDER['CANDIDATE']:
- chain.append(('WATCH', 'CANDIDATE', 'evidence_added'))
- if target == 'CONFIRMED':
- chain.append(('CANDIDATE', 'CONFIRMED', 'own_data_confirmed'))
- return [{'ts': now, 'from_state': f, 'to_state': t, 'event': e, 'by': 'engine', 'batch_id': batch_id}
- for (f, t, e) in chain]
- def _advance_chain(cur, target):
- """cur→target 的合法逐步边 (只升不降; ≤cur 返回空)。CONFIRMED 不自动降。"""
- if _STATE_ORDER.get(target, 0) <= _STATE_ORDER.get(cur, 0):
- return []
- seq = ['NORMAL', 'WATCH', 'CANDIDATE', 'CONFIRMED']
- i, j = seq.index(cur), seq.index(target)
- events = {('NORMAL', 'WATCH'): 'signal_raised', ('WATCH', 'CANDIDATE'): 'evidence_added',
- ('CANDIDATE', 'CONFIRMED'): 'own_data_confirmed'}
- steps = []
- for k in range(i, j):
- f, t = seq[k], seq[k + 1]
- steps.append((f, t, events[(f, t)]))
- return steps
- def produce_cases(findings_doc, prev_cases, batch_id, now, decay_after=2):
- """核心 (纯函数): 一批 findings + 上批 cases → 新 cases list。产出过 validate_cases。
- findings_doc: findings.json 解析 dict (含 _farm/_generated/findings)。
- prev_cases: 上批 cases list (无则 [])。 batch_id: 本批标识 (findings._generated)。
- now: 本批日期 str (YYYY-MM-DD, 供 opened_at/updated_at/history.ts)。
- decay_after: 连续 N 批未 flag 后 WATCH/CANDIDATE 衰减回 NORMAL (CONFIRMED 不自动衰减)。
- """
- import copy
- farm = findings_doc.get('_farm', '?')
- idx = {} # (farm,turbine,system) → case
- for c in copy.deepcopy(prev_cases or []): # 不就地改入参 (别名安全)
- if isinstance(c, dict):
- idx[(c.get('farm'), c.get('turbine'), c.get('system'))] = c
- flagged_now = set()
- for f in findings_doc.get('findings', []):
- target = _target_state(f)
- if target is None:
- continue
- sysk = system_of(f)
- did = discriminator_id_of(f)
- for tid in flagged_turbines(f):
- key = (farm, tid, sysk)
- flagged_now.add(key)
- ev = {'date': now, 'discriminator_id': did, 'value': _turbine_value(f, tid),
- 'verdict': f.get('verdict'), 'batch_id': batch_id,
- 'own_data': bool(f.get('own_data_confirmed'))}
- prev = idx.get(key)
- if prev is None:
- idx[key] = {
- 'case_id': f'{farm}-{tid}-{sysk}-001',
- 'farm': farm, 'turbine': tid, 'system': sysk, 'state': target,
- 'tier': None, 'verdict': f.get('verdict'),
- 'opened_at': now, 'updated_at': now, 'clean_streak': 0,
- 'evidence': [ev], 'history': _open_history(target, batch_id, now),
- }
- else:
- _advance_existing(prev, target, ev, batch_id, now)
- # 未 flag 的活跃档 → 衰减 (记账, 达阈回 NORMAL; CONFIRMED/终态不动)
- for key, c in idx.items():
- if key in flagged_now:
- continue
- _decay(c, batch_id, now, decay_after)
- return list(idx.values())
- def _advance_existing(c, target, ev, batch_id, now):
- """已存档: 追加 evidence + 合法升态 (CLOSED→复发路; REFUTED 终态不重开=dedup)。"""
- c.setdefault('evidence', []).append(ev)
- c['updated_at'] = now
- c['clean_streak'] = 0
- c['verdict'] = ev.get('verdict')
- st = c.get('state')
- hist = c.setdefault('history', [])
- if st == 'REFUTED':
- return # 终态: 同信号不重登顶 (refuted_basis dedup)
- if st == 'CLOSED': # 复发: CLOSED→RELAPSED→(CONFIRMED 若 own-data)
- hist.append({'ts': now, 'from_state': 'CLOSED', 'to_state': 'RELAPSED',
- 'event': 'relapse_detected', 'by': 'engine', 'batch_id': batch_id})
- c['state'] = 'RELAPSED'
- c['relapse_of'] = c.get('case_id')
- if target == 'CONFIRMED':
- hist.append({'ts': now, 'from_state': 'RELAPSED', 'to_state': 'CONFIRMED',
- 'event': 'own_data_confirmed', 'by': 'engine', 'batch_id': batch_id})
- c['state'] = 'CONFIRMED'
- return
- for (f, t, e) in _advance_chain(st, target): # NORMAL/WATCH/CANDIDATE 升态
- hist.append({'ts': now, 'from_state': f, 'to_state': t, 'event': e,
- 'by': 'engine', 'batch_id': batch_id})
- c['state'] = t
- def _decay(c, batch_id, now, decay_after):
- """未 flag: clean_streak++; 达阈且 state∈{WATCH,CANDIDATE} → 直接回 NORMAL (附录A.2 decay 边)。"""
- st = c.get('state')
- if st not in ('WATCH', 'CANDIDATE'):
- return # CONFIRMED/终态/已 NORMAL 不自动衰减
- c['clean_streak'] = int(c.get('clean_streak', 0)) + 1
- c['updated_at'] = now
- if c['clean_streak'] >= decay_after:
- c.setdefault('evidence', []).append(
- {'date': now, 'discriminator_id': 'decay', 'value': None,
- 'verdict': '参考', 'batch_id': batch_id, 'own_data': False})
- c.setdefault('history', []).append(
- {'ts': now, 'from_state': st, 'to_state': 'NORMAL', 'event': 'decay_to_normal',
- 'by': 'engine', 'batch_id': batch_id})
- c['state'] = 'NORMAL'
- c['clean_streak'] = 0
- def apply_human_action(case, action, now, batch_id='human', **kw):
- """人工侧合法迁移入口 (生产器不碰的右半: 工单/修/回验/证伪/own-data 确诊)。
- 与 produce_cases(engine 侧)对称: 本函数是人的动作 → by=human 审计。返回改后 copy (不改入参)。
- 每步只走附录A.2 合法边; 产出仍须过 schemas.validate_case。action:
- confirm_own_data: CANDIDATE→CONFIRMED (kw: value/discriminator_id → 追加 own_data 证据)
- issue_ticket: CONFIRMED/RELAPSED→TICKETED (kw: ticket={ticket_id,issued_by,actions,...})
- start_repair: TICKETED→IN_REPAIR
- report_verify: TICKETED/IN_REPAIR→VERIFY (kw: baseline,target,window_days)
- close: VERIFY→CLOSED (kw: readings 非空; 自动 verify.result=pass)
- verify_fail: VERIFY→CONFIRMED (未根治)
- refute: CANDIDATE/CONFIRMED→REFUTED (kw: refuted_basis 必填)
- """
- import copy
- from src.sop.schemas import CASE_EDGES
- c = copy.deepcopy(case)
- st = c.get('state')
- hist = c.setdefault('history', [])
- def _step(to, event):
- if to not in CASE_EDGES.get(st, set()):
- raise ValueError(f'apply_human_action: {st}→{to} 非法边 (action={action})')
- hist.append({'ts': now, 'from_state': st, 'to_state': to, 'event': event,
- 'by': 'human', 'batch_id': batch_id})
- c['state'] = to
- c['updated_at'] = now
- if action == 'confirm_own_data':
- c.setdefault('evidence', []).append(
- {'date': now, 'discriminator_id': kw.get('discriminator_id', 'field_measurement'),
- 'value': kw.get('value'), 'verdict': '准定论·预警', 'batch_id': batch_id, 'own_data': True})
- _step('CONFIRMED', 'own_data_confirmed')
- elif action == 'issue_ticket':
- c['ticket'] = kw.get('ticket') or {}
- _step('TICKETED', 'ticket_issued')
- elif action == 'start_repair':
- _step('IN_REPAIR', 'repair_started')
- elif action == 'report_verify':
- c['verify'] = {'baseline': kw.get('baseline'), 'target': kw.get('target'),
- 'window_days': kw.get('window_days', 60), 'readings': kw.get('readings', []),
- 'result': 'pending'}
- _step('VERIFY', 'repair_reported')
- elif action == 'close':
- vf = c.setdefault('verify', {})
- vf['readings'] = kw.get('readings') or vf.get('readings') or []
- vf['result'] = 'pass'
- _step('CLOSED', 'verify_pass')
- elif action == 'verify_fail':
- c.setdefault('verify', {})['result'] = 'fail'
- _step('CONFIRMED', 'verify_fail')
- elif action == 'refute':
- c['refuted_basis'] = kw.get('refuted_basis') or ''
- _step('REFUTED', 'guard_refuted')
- else:
- raise ValueError(f'apply_human_action: 未知 action {action!r}')
- return c
- def apply_analysis(case, verdict, mechanism, now, review=None, field_needs=None,
- ref=None, by='deep-dive'):
- """规范回灌深挖结果 (Phase-2 机制归因 + 独立审) → case.analysis。返回改后 copy (不改入参)。
- 与 apply_human_action(状态迁移) 对称: 本函数只挂**分析层**, 不碰 case.state。
- ★纪律 (不变量①): analysis.verdict 是更深表征(可 > case.verdict, 如 D29 候选→深挖准定论·预警),
- 但升 CONFIRMED 须现场 own-data → 走 apply_human_action('confirm_own_data'), 非本函数。
- 产出过 schemas.validate_analysis (verdict 六枚举 + mechanism 必填)。"""
- import copy
- from src.sop.schemas import norm_verdict, VERDICTS
- if norm_verdict(verdict) not in VERDICTS:
- raise ValueError(f'apply_analysis: verdict "{verdict}" ∉ 六枚举 {sorted(VERDICTS)}')
- if not mechanism:
- raise ValueError('apply_analysis: mechanism 必填 (深挖必给机制归因)')
- c = copy.deepcopy(case)
- an = {'verdict': norm_verdict(verdict), 'mechanism': mechanism, 'review': review,
- 'field_needs': field_needs, 'ref': ref, 'date': now, 'by': by}
- c['analysis'] = {k: v for k, v in an.items() if v is not None} # 去 None 清爽
- c['updated_at'] = now
- return c
- # ---------------------------------------------------------------------------
- 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.cases <farm> [--date YYYY-MM-DD]
- # 读 outputs/<farm>/sop/{findings.json, cases.json(若有)} → 生产 → 校验 → 写 cases.json
- import argparse
- from src.sop.schemas import validate_cases
- ROOT = _install_root(__file__)
- ap = argparse.ArgumentParser()
- ap.add_argument('farm')
- ap.add_argument('--date', default=None, help='本批日期 (默认取 findings._generated 日期)')
- ap.add_argument('--dry', action='store_true', help='只校验不写盘')
- a = ap.parse_args()
- fdir = P.sop(a.farm)
- fdoc = json.loads((fdir / 'findings.json').read_text(encoding='utf-8'))
- cpath = fdir / 'cases.json'
- prev = json.loads(cpath.read_text(encoding='utf-8')) if cpath.exists() else []
- batch_id = fdoc.get('_generated', 'unknown')
- now = a.date or str(batch_id)[:10]
- cases = produce_cases(fdoc, prev, batch_id, now)
- issues = validate_cases(cases)
- if issues:
- print(f'❌ validate_cases {len(issues)} issue:')
- for i in issues[:20]:
- print(' -', i)
- sys.exit(2)
- print(f'✅ {a.farm}: {len(cases)} cases (prev {len(prev)}), validate_cases PASS')
- from collections import Counter
- print(' 状态分布:', dict(Counter(c['state'] for c in cases)))
- if not a.dry:
- cpath.write_text(json.dumps(cases, ensure_ascii=False, indent=2), encoding='utf-8')
- print(f' 写盘: {cpath.relative_to(ROOT)}')
|