#!/usr/bin/env python3 # -*- coding: utf-8 -*- r"""补齐"包内没有生成端"的产物 —— 只补缺件, 不动 raw 重算件; 并落一份**来源台账** (2026-09-12)。 ## ★ 本器已退出运行链(用户令 2026-09-17 #1) **"观澜系统运行不要从含产物的交付包补齐"**。所以: · `rebuild_all.py` 的第 ⑤ 步**不再调用本器**(改成"反向呼应审计",只报账不搬运); `products_state` / 门户 / 各组件页里的提示也不再指向"拿交付包补"; · 本器现在**只作研发/离线人工补救**:某台机器上确实缺了一批"包内无生成端"的件、且当场无法重算时, 由人显式指定来源跑一次。任何自动化的"看旁边有没有 zip"都去掉了 —— 那会把"包"变成隐式依赖。 · 运行期遇到缺件的正确处置是(按优先级): ① 放原始件到 `data/raw/<场站>/` → `python scripts/rebuild_all.py`(有生成端的那些); ② 没有生成端的那些 → **由研发补生成端**,逐族可逆性见 `docs/系统设计说明.md` §13.6; ③ 仍缺 → 页面/自检会如实列出缺哪件(结构化缺件页 + `guanlan.py check` 的「反向呼应」一行)。 ## 为什么它当初存在 `rebuild_from_raw.py` 能从 `data/raw` 重算出一部分产物; 但随包里另有一批产物**全库只有读取方、 0 处写入方** —— 没有生成端, 从零重算就是拿不到 (工作台主数据 `/api/fleet` 一缺 `pitch/pitch_daily.parquet` / `temp_monthly.parquet` 就整个回 `err=no_products`, 页面看着像"没数据")。 2026-09-12 的口径是"证明不了的用随包件补齐、但标明它不是重算件";2026-09-17 用户令把这条 **从运行期拿掉**:宁可如实显示缺件,也不让系统依赖交付包。 ## 规则 (保守优先) 对随包产物来源里的每个文件: · 目标**已存在** → 跳过 (那是 raw 重算出来的, 不许被随包件覆盖); · 目标不存在 → 拷过去, 记为 `shipped`(台账 `_provenance.json` 里一眼看得出它不是重算件); · `.log/.jsonl` → 按用户令 2 落到 `logs/build/<场>/`,不算产物。 另: 只处理"产物目录"(outputs/<场>/ 下的仓), 不碰 release/ 与 data/。 ## 用法(**人工**,不在重算链里) python scripts/products_restore_missing.py --dry-run # 只报要补什么 python scripts/products_restore_missing.py --stash # 从随包产物目录补 python scripts/products_restore_missing.py --stash <包.zip> # 从含产物的 zip 补(离线补救) python scripts/products_restore_missing.py --find # 只看旁边哪个 zip 里真有产物 """ from __future__ import annotations import argparse import fnmatch import json import pathlib import shutil import sys import time ROOT = pathlib.Path(__file__).resolve().parents[1] sys.path.insert(0, str(ROOT)) from src import paths as P # noqa: E402 # 由 data/raw 重算出来的 (重建器 + 验证依据) —— 这些**永不被随包件覆盖** RAW_DERIVED = { # 三门台账 (rebuild_from_raw.py; --verify 与随包基线逐值比对过) 'windscada/alarms.parquet': 'scripts/windscada_alarms_ingest.py', 'windscada/workorders.parquet': 'scripts/windscada_workorder_ingest.py', 'windscada/oil_samples_index.parquet': 'scripts/windscada_watch_channels_build.py', # SCADA 侧 10 个构建器 (rebuild_from_raw.py --scada) 'windscada/powercurve_dev.parquet': 'src.windscada.perf.powercurve', 'windscada/powercurve_bins.parquet': 'src.windscada.perf.powercurve', 'windscada/loss_monthly.parquet': 'src.windscada.perf.availability', 'windscada/curve_lenses.parquet': 'src.windscada.perf.curves', 'windscada/curve_liveness.parquet': 'src.windscada.perf.curves', 'windscada/control_profile.parquet': 'src.windscada.perf.control', 'windscada/control_schedule.parquet': 'src.windscada.perf.control', 'windscada/stop_events.parquet': 'src.windscada.perf.faults', 'windscada/temp_bins.parquet': 'src.windscada.subsys.temp_nbm', 'windscada/yaw_daily.parquet': 'src.windscada.subsys.yaw', 'windscada/hydraulic_accum.parquet': 'src.windscada.subsys.hydraulic', 'windscada/thermal_chain.parquet': 'src.windscada.subsys.thermal_chain', 'windscada/system_aux.parquet': 'src.windscada.taxonomy', # 月度派生件 (本轮新增; 逐值对齐随包件后才算数) 'windscada/temp_monthly.parquet': 'scripts/windscada_monthly_build.py (19494/19494 逐值一致)', # 本体层 (kb_ingest → populate → chain_ingest → trend_ingest; 全链跑完后 audit 0 问题) 'ontology/objects.json': 'python -m src.ontology.kb_ingest + populate + chain_ingest + trend_ingest (audit 0 问题)', 'ontology/retrieval_index.json': 'src.ontology.retrieval.build', 'ontology/turbine_params.parquet': 'src.ontology.maintenance.refresh_params', } def zip_candidates() -> list[pathlib.Path]: """可能含产物的包(本仓根 + 上一级)—— **仅供离线人工诊断**,运行链不用它。 ★ 用户令 2026-09-17 #1「观澜系统运行不要从含产物的交付包补齐」:本器已**退出运行链** (`rebuild_all.py` 第 ⑤ 步改成反向呼应审计,只报账不搬运)。自动找包的逻辑一并去掉 —— 留着它等于让"配一个 zip 在旁边"变成隐式依赖,与"目标机放数据自己重算"这条主线冲突。 本函数现在只服务于 `--find`(人工看一眼"旁边哪个包里有产物")。 """ out = [] for base in (ROOT, ROOT.parent): try: out += [p for p in sorted(base.glob('*.zip')) if p.is_file()] except OSError: pass return list(dict.fromkeys(out)) def find_stash_zip() -> tuple[pathlib.Path, int] | None: """→ (含产物的包, 该场产物件数);找不到就 None。""" import zipfile marker = f'outputs/{P.farm()}/windscada/temp_monthly.parquet' pref = f'outputs/{P.farm()}/' best = None for z in zip_candidates(): try: with zipfile.ZipFile(z) as zf: names = set(zf.namelist()) except Exception: # 不是 zip / 读不动 / 加密: 跳过, 不中断 continue if marker not in names: continue n = sum(1 for x in names if x.startswith(pref)) if best is None or n > best[1]: best = (z, n) return best def stash_dir(explicit=None) -> pathlib.Path: """随包产物原件在哪。 四个地方都可能是它(取决于用过哪些开关), 所以按"像不像随包整份产物"打分挑: · `_products_off/<场>/` = products_state --off 挪走的当前产物; · `_products_off_prev_<时间戳>/<场>/` = --off --archive-old 存档的上一代(= 随包原件); · `--stash` 显式指定 (目录或交付包 zip); · **含产物的交付包 zip, 自动找** (2026-09-17 加: 本仓根与上一级 `.zip` 里挑真含随包件的那个)。 判据用**只在随包里有的件**(windscada/temp_monthly.parquet 这类"包内无生成端"的产物): 含它的目录才是"标准答案"来源。原先只看第一个候选, 于是 --on 还原后暂存区空了 → 台账被写成 0/0。 """ if explicit: # 必须 resolve: 传相对路径时下面 `stash.relative_to(ROOT)` 会 ValueError (2026-09-12 实逮) return pathlib.Path(explicit).resolve() marker = 'windscada/temp_monthly.parquet' cands = [] for pat in ('_products_off/*', '_products_off_prev_*/*'): cands += [p for p in sorted(ROOT.glob(pat)) if p.is_dir() and p.name != '_baseline_kept'] scored = [] for c in cands: n = sum(1 for _ in c.rglob('*') if _.is_file()) scored.append(((1 if (c / marker).is_file() else 0), n, c)) scored.sort(key=lambda x: (-x[0], -x[1])) if scored and scored[0][0]: # 只有真含随包件的目录才算数 return scored[0][2] raise StashMissing('找不到随包产物原件 (暂存区 _products_off/**、或 --stash 显式指定)') class StashMissing(SystemExit): """暂存区不存在 —— **不是崩溃, 是"这一步没得做"**。 2026-09-16 用户令"系统运行中/未运行都能重算"时实逮: `rebuild_all.py` 的第⑤步 (补齐"包内没有生成端" 的产物) 在**没有暂存区**的检出上必抛 SystemExit → 重建链当场中断, 后面的 ⑦本体链 与 ⑧审计**全不执行** (脚本的中断语义是"后面的步骤依赖它")。而暂存区只有"清除产物"动作会产生 —— 首次全量重算、或 清完又还原过的机器上就都没有, 于是"从零重算"在交付包上跑不完。 处置: 本类带退出码 6, 由 rebuild_all.py 显式容忍并打印"跳过随包补齐" —— 缺的是"随包件"这一路, 不是重算本身; 如实报出来比假装通过好, 也比把整条链打断好。 """ exit_code = 6 def _family_patterns(glob) -> list[str]: """族 glob(str 或 list)→ 展开花括号后的模式列表 (与 products_reverse_audit._files_of 同规则)。""" out = [] for g in ([glob] if isinstance(glob, str) else list(glob)): if '{' in g: head, tail = g.split('{', 1) opts, rest = tail.split('}', 1) out += [head + o + rest for o in opts.split(',')] else: out.append(g) return out def _attribute(rel: str) -> dict: """盘上有、无人登记的件 → 按**族表**归口 (2026-09-18)。 为什么必须做这一步: 生成端**自登记**(`_derived_manifest.json`) 并非全覆盖 —— 实测 `windscada/*.parquet`(17 件, 由 `rebuild_from_raw.py --scada` 的 10 个构建器写)、 `ontology/*`(3 件, 由本体的 kb_ingest/populate/retrieval 写) 都不自登记。此前本模式 **只**照抄自登记件 ⇒ 台账 1705 条、盘上 1728 件, 差的 23 件既不在账上也不被报出 (账实不符的另一个方向, 而且正好是"页面主取数"那一批)。现在按族表补归口, 归不上 的**如实记 `unregistered`** —— 不猜来源。 深度对齐规则同 `products_reverse_audit._files_of`: 模式无 `**` 时, `/` 个数必须相等 (fnmatch 的 `*` 跨 `/`, 不显式挡会让 `windscada/*.parquet` 吞掉子目录里的同名件)。 """ try: from products_reverse_audit import FAMILIES except Exception as e: # 族表读不到: 全部按未归类记 return dict(source='unregistered', builder='?', why=f'盘上有、生成端未自登记, 且族表读不到({type(e).__name__}) ⇒ 需人工归口') for fam in FAMILIES: for pt in _family_patterns(fam.get('glob')): if '**' not in pt and rel.count('/') != pt.count('/'): continue if rel == pt or fnmatch.fnmatch(rel, pt): kind = fam.get('kind') or 'raw-derived' src = {'raw-derived': 'raw-derived', 'shipped': 'shipped', 'human': 'human', 'source-derived': 'source-derived', 'not-product': 'not-product'}.get(kind, kind) return dict(source=src, builder=fam.get('gen') or '(族表登记: 无生成端)', why=f'盘上在位但生成端未自登记 → 按族表 {fam.get("id")} 归口' f'({fam.get("func") or ""})') return dict(source='unregistered', builder='?', why='盘上有、生成端未自登记、族表也未登记 ⇒ 需人工归口 ' '(本条是软件兜底记账, 不是真来源; 见到了要么补自登记, 要么补族表)') def _reattr_needed(v: dict) -> bool: """这条账是不是"族表归口的产物"(可随族表改进而重写),而不是真来源(自登记/人工补入)。""" return v.get('source') == 'unregistered' or '按族表' in str(v.get('why') or '') def refresh_ledger(farm: str | None = None) -> int: r"""**只维护台账**(不需要任何随包件/交付包):把盘上已经不存在的条目从 `_provenance.json` 里去掉。 为什么需要它(用户令 2026-09-17 #3):把"非产物"的杂物从产物仓移走后,台账里还留着它们的条目 —— 于是反向呼应审计仍把它们当产物统计("帐上 2292 件、盘上 2275 件"这种不一致正是这套台账最该避免的)。 本模式只做"账实相符":**在位且来源已登记的照抄,盘上没了的删掉**,不引入任何新来源。 """ farm = farm or P.farm() root = P.out_root(farm) old = {} f = root / '_provenance.json' if f.is_file(): try: old = (json.loads(f.read_text(encoding='utf-8')) or {}).get('files') or {} except Exception: old = {} kept, dropped = {}, [] for rel, v in old.items(): if (root / rel).is_file(): kept[rel] = v else: dropped.append(rel) # 自登记件 (_derived_manifest.json): 在位且台账里没有的补进来 added = 0 try: from src.derived_manifest import load as _load for rel, info in ((_load(root).get('files') or {})).items(): if (root / rel).is_file() and rel not in kept: kept[rel] = dict(source='raw-derived', builder=f"{info.get('builder', '?')} (自登记 by {info.get('by', '?')})") added += 1 except Exception: pass # ── 另一个方向: 盘上有、账上没有的件 → 按族表归口 (绝不静默漏掉) ────────────── # 台账两本 (_provenance.json / _derived_manifest.json) 都是**这本台账自己**: # 它们不是产物, 不登记 (否则每次写盘都会让"账实相符"永远差两件)。 swept: list[str] = [] for p in sorted(root.rglob('*')): if not p.is_file(): continue rel = p.relative_to(root).as_posix() if rel in ('_provenance.json', '_derived_manifest.json'): continue # 只补"没人登记"的; 但**族表归口过的条目要重跑** —— 否则族表改了(例如给 windcms/cache # 单列一族)这条账会一直停在旧归口上, 而它并不是谁自登记/谁补进来的真来源。 if rel in kept and not _reattr_needed(kept[rel]): continue kept[rel] = _attribute(rel) swept.append(rel) by = {} for v in kept.values(): by[v.get('source', '?')] = by.get(v.get('source', '?'), 0) + 1 print(f'台账: {len(old)} 条 → {len(kept)} 条(盘上已删 {len(dropped)} 条, 自登记补入 {added} 条, ' f'按族表归口 {len(swept)} 条)') for r in dropped[:12]: print(f' - 移出/删除: {r}') if len(dropped) > 12: print(f' … 另有 {len(dropped) - 12} 条') for r in swept[:8]: print(f' + 归口: {r} ← {kept[r].get("builder")}') if len(swept) > 8: print(f' … 另有 {len(swept) - 8} 条归口') unk = [r for r in swept if kept[r].get('source') == 'unregistered'] if unk: print(f' ★ {len(unk)} 件**归不上族**: ' + ', '.join(unk[:6]) + (' …' if len(unk) > 6 else '')) print(' 来源分布: ' + ' · '.join(f'{k}={v}' for k, v in sorted(by.items()))) f.write_text(json.dumps(dict( at=time.strftime('%Y-%m-%d %H:%M:%S'), note='逐件来源台账: raw-derived = 由 data/raw 重算(含验证依据); shipped = 包内无生成端(历史随包件); ' 'source-derived = 由源码重建(如发布清单); human = 人工正本/交证件; ' 'unregistered = 盘上有但无人登记(需人工归口); ' 'log-relocated = 包内日志按用户令 2 落到 logs/build/', counts=dict(by), files=kept), ensure_ascii=False, indent=1), encoding='utf-8') print(f'已写 {P.rel(f)}') return 0 def main() -> int: ap = argparse.ArgumentParser() ap.add_argument('--dry-run', action='store_true') ap.add_argument('--stash', default=None) ap.add_argument('--find', action='store_true', help='只看旁边的 .zip 里哪个真含产物 (人工诊断用, 不写盘)') ap.add_argument('--refresh', action='store_true', help='只维护台账: 把盘上已不存在的条目从 _provenance.json 去掉 (不需要任何随包件)') a = ap.parse_args() if a.refresh: return refresh_ledger() if a.find: found = find_stash_zip() if not found: print('旁边的 .zip 里没有含' + f'outputs/{P.farm()}/ 产物的包') return 0 print(f'含产物的包: {P.rel(found[0])}({found[1]} 件 outputs/{P.farm()}/ 条目)') print(' 注意: 本器只作**离线人工补救**; 运行期不要从交付包补齐(用户令 2026-09-17)') return 0 try: stash = stash_dir(a.stash) except StashMissing as e: # 没有暂存区 → 这一步没得做。**不许静默、也不许把整条重算链打断**(见 StashMissing 的说明)。 print(f'[跳过] {e}') print(' 处置(运行期): ① 放原始件到 data/raw/<场站>/ 后 python scripts/rebuild_all.py;' ' ② 没有生成端的族由研发补生成端(逐族可逆性见 docs/系统设计说明.md §13.6)。') print(' 处置(离线人工补救, 非运行链): --stash <随包产物目录 或 含产物的包.zip>') return StashMissing.exit_code # ★2026-09-16: 原写法 `P.STORE if hasattr(P,'STORE') else ROOT/'outputs'/'rudong'` —— P.STORE 根本 # 不存在, 于是**恒**落到硬编码的 rudong, 多场部署下会把 rudong 的随包件补进别的场 (静默串场)。 dest_root = P.out_root() print(f'随包件: {P.rel(stash)}') print(f'产物仓: {dest_root.relative_to(ROOT)}\n') prov = {} def take(rel: str, src_path=None, src_bytes=None): """补齐一件 (或在位则只记账)。src_path=目录来源, src_bytes=zip 来源。 ★ 日志不补进产物仓 (2026-09-17): 历史交付包里把构建日志也当产物发过 (如 `m5_cms_tcm/shared_component_bandpower.log`)。用户令 2 定的是"运行日志只在 logs/", 所以这里把它直接落到 `logs/build/<场>/…` —— 数据没丢, 但不再算产物、也不进台账 (否则每从旧包补一次, `guanlan.py check` 的「日志统一」就红一次: 日志跑到 logs/ 外面了)。 """ low = rel.lower() if low.endswith(('.log', '.jsonl')): from src import logfile as _lf dst_log = _lf.build_log(rel, P.farm()) prov[rel] = dict(source='log-relocated', why='包内日志: 按用户令 2 落到 logs/build/, 不算产物') if not a.dry_run: dst_log.parent.mkdir(parents=True, exist_ok=True) if src_bytes is not None: dst_log.write_bytes(src_bytes) elif src_path is not None: shutil.copy2(src_path, dst_log) return dst = dest_root / rel if dst.exists(): # 已存在的件分两类: 我们自己重算的(raw-derived) 与 上轮已补齐的(shipped)。 # ★不能因为"本轮没拷"就不记账 —— 早先写成"存在即不列", 于是重复跑一次台账就被清成 0/0, # 页面上的来源统计跟着一起错 (2026-09-12 实逮)。 if rel in RAW_DERIVED: prov[rel] = dict(source='raw-derived', builder=RAW_DERIVED[rel]) else: prov[rel] = dict(source='shipped', why='包内无生成端 / 规则未复现 → 随包件补齐 (本轮已在位)') return prov[rel] = dict(source='shipped', why='包内无生成端 / 规则未复现 → 用随包件补齐') if a.dry_run: return dst.parent.mkdir(parents=True, exist_ok=True) if src_bytes is not None: dst.write_bytes(src_bytes) else: shutil.copy2(src_path, dst) if stash.is_file(): # 交付包 zip 形态 (2026-09-16): 按需从包里取 outputs/<场>/** 补齐 —— # 这是"清除不留备份"之后**唯一**的随包件来源, 所以把它做成一等输入。 import zipfile with zipfile.ZipFile(stash) as zf: pref = f'outputs/{P.farm()}/' members = [n for n in zf.namelist() if n.startswith(pref) and not n.endswith('/')] print(f' 从交付包读取 {len(members)} 件 (前缀 {pref})') for name in sorted(members): if a.dry_run: take(name[len(pref):]) else: take(name[len(pref):], src_bytes=zf.read(name)) else: for src in sorted(stash.rglob('*')): if not src.is_file(): continue take(src.relative_to(stash).as_posix(), src_path=src) # 重算件里有些**不在随包件里**(如我们新造的 turbine_params.parquet), 也要记进台账 for rel, builder in RAW_DERIVED.items(): if rel not in prov and (dest_root / rel).exists(): prov[rel] = dict(source='raw-derived', builder=builder) # 构建脚本的**自登记** (src/derived_manifest.py): 振动侧的产物名随窗/分片变, 写不进上面的精确表, # 按名字通配又会误伤同名旧件 → 由"谁算的谁登记", 这里只认登记。缺失不影响其余记账。 try: from src.derived_manifest import load as _load_manifest, prune as _prune_manifest _gone = _prune_manifest(dest_root) # 先清掉"登记了但盘上已删"的条目 (如被删的 _reimport 窗) if _gone: print(f' 自登记清理: {_gone} 条指向已删文件的条目') _man = (_load_manifest(dest_root).get('files') or {}) except Exception: _man = {} for rel, info in _man.items(): if (dest_root / rel).exists(): prov[rel] = dict(source='raw-derived', builder=f"{info.get('builder', '?')} (自登记 by {info.get('by', '?')})") else: prov[rel] = dict(source='raw-derived', builder=f"{info.get('builder', '?')} (自登记, 但该件当前不在盘上)") by_store = {} for rel, m in prov.items(): store = rel.split('/')[0] d = by_store.setdefault(store, {'raw-derived': 0, 'shipped': 0, 'log-relocated': 0}) d.setdefault(m['source'], 0) d[m['source']] += 1 moved = sum(d['log-relocated'] for d in by_store.values()) print(f'{"仓":26s} {"raw 重算":>9s} {"随包补齐":>9s}' + (f' {"日志归位":>9s}' if moved else '')) for store, d in sorted(by_store.items()): print(f' {store:24s} {d["raw-derived"]:9d} {d["shipped"]:9d}' + (f' {d["log-relocated"]:9d}' if moved else '')) tot = {k: sum(d[k] for d in by_store.values()) for k in ('raw-derived', 'shipped', 'log-relocated')} print(f' {"合计":24s} {tot["raw-derived"]:9d} {tot["shipped"]:9d}' + (f' {tot["log-relocated"]:9d}' if moved else '')) if moved: print(f' (日志不落产物仓: {moved} 件按用户令 2 落到 logs/build/, 不计入产物台账)') if a.dry_run: print('\n(dry-run, 未写盘)') return 0 out = dest_root / '_provenance.json' out.write_text(json.dumps(dict( at=time.strftime('%Y-%m-%d %H:%M:%S'), note='逐件来源台账: raw-derived = 由 data/raw 重算(含验证依据); shipped = 包内无生成端, 用随包件补齐', counts=tot, files=prov), ensure_ascii=False, indent=1), encoding='utf-8') print(f'\n已写来源台账 {P.rel(out)}') return 0 if __name__ == '__main__': # 控制台可能是 GBK: print 里的非 GBK 字形编不出来会抛异常, 干成了却退出码 1 for _s in (sys.stdout, sys.stderr): try: _s.reconfigure(errors='replace') except Exception: pass sys.exit(main())