products_restore_missing.py 17 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310
  1. #!/usr/bin/env python3
  2. # -*- coding: utf-8 -*-
  3. r"""补齐"包内没有生成端"的产物 —— 只补缺件, 不动 raw 重算件; 并落一份**来源台账** (2026-09-12)。
  4. ## ★ 本器已退出运行链(用户令 2026-09-17 #1)
  5. **"观澜系统运行不要从含产物的交付包补齐"**。所以:
  6. · `rebuild_all.py` 的第 ⑤ 步**不再调用本器**(改成"反向呼应审计",只报账不搬运);
  7. `products_state` / 门户 / 各组件页里的提示也不再指向"拿交付包补";
  8. · 本器现在**只作研发/离线人工补救**:某台机器上确实缺了一批"包内无生成端"的件、且当场无法重算时,
  9. 由人显式指定来源跑一次。任何自动化的"看旁边有没有 zip"都去掉了 —— 那会把"包"变成隐式依赖。
  10. · 运行期遇到缺件的正确处置是(按优先级):
  11. ① 放原始件到 `data/raw/<场站>/` → `python scripts/rebuild_all.py`(有生成端的那些);
  12. ② 没有生成端的那些 → **由研发补生成端**,逐族可逆性见 `docs/系统设计说明.md` §13.6;
  13. ③ 仍缺 → 页面/自检会如实列出缺哪件(结构化缺件页 + `guanlan.py check` 的「反向呼应」一行)。
  14. ## 为什么它当初存在
  15. `rebuild_from_raw.py` 能从 `data/raw` 重算出一部分产物; 但随包里另有一批产物**全库只有读取方、
  16. 0 处写入方** —— 没有生成端, 从零重算就是拿不到 (工作台主数据 `/api/fleet` 一缺
  17. `pitch/pitch_daily.parquet` / `temp_monthly.parquet` 就整个回 `err=no_products`, 页面看着像"没数据")。
  18. 2026-09-12 的口径是"证明不了的用随包件补齐、但标明它不是重算件";2026-09-17 用户令把这条
  19. **从运行期拿掉**:宁可如实显示缺件,也不让系统依赖交付包。
  20. ## 规则 (保守优先)
  21. 对随包产物来源里的每个文件:
  22. · 目标**已存在** → 跳过 (那是 raw 重算出来的, 不许被随包件覆盖);
  23. · 目标不存在 → 拷过去, 记为 `shipped`(台账 `_provenance.json` 里一眼看得出它不是重算件);
  24. · `.log/.jsonl` → 按用户令 2 落到 `logs/build/<场>/`,不算产物。
  25. 另: 只处理"产物目录"(outputs/<场>/ 下的仓), 不碰 release/ 与 data/。
  26. ## 用法(**人工**,不在重算链里)
  27. python scripts/products_restore_missing.py --dry-run # 只报要补什么
  28. python scripts/products_restore_missing.py --stash <dir> # 从随包产物目录补
  29. python scripts/products_restore_missing.py --stash <包.zip> # 从含产物的 zip 补(离线补救)
  30. python scripts/products_restore_missing.py --find # 只看旁边哪个 zip 里真有产物
  31. """
  32. from __future__ import annotations
  33. import argparse
  34. import json
  35. import pathlib
  36. import shutil
  37. import sys
  38. import time
  39. ROOT = pathlib.Path(__file__).resolve().parents[1]
  40. sys.path.insert(0, str(ROOT))
  41. from src import paths as P # noqa: E402
  42. # 由 data/raw 重算出来的 (重建器 + 验证依据) —— 这些**永不被随包件覆盖**
  43. RAW_DERIVED = {
  44. # 三门台账 (rebuild_from_raw.py; --verify 与随包基线逐值比对过)
  45. 'windscada/alarms.parquet': 'scripts/windscada_alarms_ingest.py',
  46. 'windscada/workorders.parquet': 'scripts/windscada_workorder_ingest.py',
  47. 'windscada/oil_samples_index.parquet': 'scripts/windscada_watch_channels_build.py',
  48. # SCADA 侧 10 个构建器 (rebuild_from_raw.py --scada)
  49. 'windscada/powercurve_dev.parquet': 'src.windscada.perf.powercurve',
  50. 'windscada/powercurve_bins.parquet': 'src.windscada.perf.powercurve',
  51. 'windscada/loss_monthly.parquet': 'src.windscada.perf.availability',
  52. 'windscada/curve_lenses.parquet': 'src.windscada.perf.curves',
  53. 'windscada/curve_liveness.parquet': 'src.windscada.perf.curves',
  54. 'windscada/control_profile.parquet': 'src.windscada.perf.control',
  55. 'windscada/control_schedule.parquet': 'src.windscada.perf.control',
  56. 'windscada/stop_events.parquet': 'src.windscada.perf.faults',
  57. 'windscada/temp_bins.parquet': 'src.windscada.subsys.temp_nbm',
  58. 'windscada/yaw_daily.parquet': 'src.windscada.subsys.yaw',
  59. 'windscada/hydraulic_accum.parquet': 'src.windscada.subsys.hydraulic',
  60. 'windscada/thermal_chain.parquet': 'src.windscada.subsys.thermal_chain',
  61. 'windscada/system_aux.parquet': 'src.windscada.taxonomy',
  62. # 月度派生件 (本轮新增; 逐值对齐随包件后才算数)
  63. 'windscada/temp_monthly.parquet': 'scripts/windscada_monthly_build.py (19494/19494 逐值一致)',
  64. # 本体层 (kb_ingest → populate → chain_ingest → trend_ingest; 全链跑完后 audit 0 问题)
  65. 'ontology/objects.json': 'python -m src.ontology.kb_ingest + populate + chain_ingest + trend_ingest (audit 0 问题)',
  66. 'ontology/retrieval_index.json': 'src.ontology.retrieval.build',
  67. 'ontology/turbine_params.parquet': 'src.ontology.maintenance.refresh_params',
  68. }
  69. def zip_candidates() -> list[pathlib.Path]:
  70. """可能含产物的包(本仓根 + 上一级)—— **仅供离线人工诊断**,运行链不用它。
  71. ★ 用户令 2026-09-17 #1「观澜系统运行不要从含产物的交付包补齐」:本器已**退出运行链**
  72. (`rebuild_all.py` 第 ⑤ 步改成反向呼应审计,只报账不搬运)。自动找包的逻辑一并去掉 ——
  73. 留着它等于让"配一个 zip 在旁边"变成隐式依赖,与"目标机放数据自己重算"这条主线冲突。
  74. 本函数现在只服务于 `--find`(人工看一眼"旁边哪个包里有产物")。
  75. """
  76. out = []
  77. for base in (ROOT, ROOT.parent):
  78. try:
  79. out += [p for p in sorted(base.glob('*.zip')) if p.is_file()]
  80. except OSError:
  81. pass
  82. return list(dict.fromkeys(out))
  83. def find_stash_zip() -> tuple[pathlib.Path, int] | None:
  84. """→ (含产物的包, 该场产物件数);找不到就 None。"""
  85. import zipfile
  86. marker = f'outputs/{P.farm()}/windscada/temp_monthly.parquet'
  87. pref = f'outputs/{P.farm()}/'
  88. best = None
  89. for z in zip_candidates():
  90. try:
  91. with zipfile.ZipFile(z) as zf:
  92. names = set(zf.namelist())
  93. except Exception: # 不是 zip / 读不动 / 加密: 跳过, 不中断
  94. continue
  95. if marker not in names:
  96. continue
  97. n = sum(1 for x in names if x.startswith(pref))
  98. if best is None or n > best[1]:
  99. best = (z, n)
  100. return best
  101. def stash_dir(explicit=None) -> pathlib.Path:
  102. """随包产物原件在哪。
  103. 四个地方都可能是它(取决于用过哪些开关), 所以按"像不像随包整份产物"打分挑:
  104. · `_products_off/<场>/` = products_state --off 挪走的当前产物;
  105. · `_products_off_prev_<时间戳>/<场>/` = --off --archive-old 存档的上一代(= 随包原件);
  106. · `--stash` 显式指定 (目录或交付包 zip);
  107. · **含产物的交付包 zip, 自动找** (2026-09-17 加: 本仓根与上一级 `.zip` 里挑真含随包件的那个)。
  108. 判据用**只在随包里有的件**(windscada/temp_monthly.parquet 这类"包内无生成端"的产物):
  109. 含它的目录才是"标准答案"来源。原先只看第一个候选, 于是 --on 还原后暂存区空了 → 台账被写成 0/0。
  110. """
  111. if explicit:
  112. # 必须 resolve: 传相对路径时下面 `stash.relative_to(ROOT)` 会 ValueError (2026-09-12 实逮)
  113. return pathlib.Path(explicit).resolve()
  114. marker = 'windscada/temp_monthly.parquet'
  115. cands = []
  116. for pat in ('_products_off/*', '_products_off_prev_*/*'):
  117. cands += [p for p in sorted(ROOT.glob(pat)) if p.is_dir() and p.name != '_baseline_kept']
  118. scored = []
  119. for c in cands:
  120. n = sum(1 for _ in c.rglob('*') if _.is_file())
  121. scored.append(((1 if (c / marker).is_file() else 0), n, c))
  122. scored.sort(key=lambda x: (-x[0], -x[1]))
  123. if scored and scored[0][0]: # 只有真含随包件的目录才算数
  124. return scored[0][2]
  125. raise StashMissing('找不到随包产物原件 (暂存区 _products_off/**、或 --stash 显式指定)')
  126. class StashMissing(SystemExit):
  127. """暂存区不存在 —— **不是崩溃, 是"这一步没得做"**。
  128. 2026-09-16 用户令"系统运行中/未运行都能重算"时实逮: `rebuild_all.py` 的第⑤步 (补齐"包内没有生成端"
  129. 的产物) 在**没有暂存区**的检出上必抛 SystemExit → 重建链当场中断, 后面的 ⑦本体链 与 ⑧审计**全不执行**
  130. (脚本的中断语义是"后面的步骤依赖它")。而暂存区只有"清除产物"动作会产生 —— 首次全量重算、或
  131. 清完又还原过的机器上就都没有, 于是"从零重算"在交付包上跑不完。
  132. 处置: 本类带退出码 6, 由 rebuild_all.py 显式容忍并打印"跳过随包补齐" —— 缺的是"随包件"这一路,
  133. 不是重算本身; 如实报出来比假装通过好, 也比把整条链打断好。
  134. """
  135. exit_code = 6
  136. def main() -> int:
  137. ap = argparse.ArgumentParser()
  138. ap.add_argument('--dry-run', action='store_true')
  139. ap.add_argument('--stash', default=None)
  140. ap.add_argument('--find', action='store_true',
  141. help='只看旁边的 .zip 里哪个真含产物 (人工诊断用, 不写盘)')
  142. a = ap.parse_args()
  143. if a.find:
  144. found = find_stash_zip()
  145. if not found:
  146. print('旁边的 .zip 里没有含' + f'outputs/{P.farm()}/ 产物的包')
  147. return 0
  148. print(f'含产物的包: {P.rel(found[0])}({found[1]} 件 outputs/{P.farm()}/ 条目)')
  149. print(' 注意: 本器只作**离线人工补救**; 运行期不要从交付包补齐(用户令 2026-09-17)')
  150. return 0
  151. try:
  152. stash = stash_dir(a.stash)
  153. except StashMissing as e:
  154. # 没有暂存区 → 这一步没得做。**不许静默、也不许把整条重算链打断**(见 StashMissing 的说明)。
  155. print(f'[跳过] {e}')
  156. print(' 处置(运行期): ① 放原始件到 data/raw/<场站>/ 后 python scripts/rebuild_all.py;'
  157. ' ② 没有生成端的族由研发补生成端(逐族可逆性见 docs/系统设计说明.md §13.6)。')
  158. print(' 处置(离线人工补救, 非运行链): --stash <随包产物目录 或 含产物的包.zip>')
  159. return StashMissing.exit_code
  160. # ★2026-09-16: 原写法 `P.STORE if hasattr(P,'STORE') else ROOT/'outputs'/'rudong'` —— P.STORE 根本
  161. # 不存在, 于是**恒**落到硬编码的 rudong, 多场部署下会把 rudong 的随包件补进别的场 (静默串场)。
  162. dest_root = P.out_root()
  163. print(f'随包件: {P.rel(stash)}')
  164. print(f'产物仓: {dest_root.relative_to(ROOT)}\n')
  165. prov = {}
  166. def take(rel: str, src_path=None, src_bytes=None):
  167. """补齐一件 (或在位则只记账)。src_path=目录来源, src_bytes=zip 来源。
  168. ★ 日志不补进产物仓 (2026-09-17): 历史交付包里把构建日志也当产物发过 (如
  169. `m5_cms_tcm/shared_component_bandpower.log`)。用户令 2 定的是"运行日志只在 logs/",
  170. 所以这里把它直接落到 `logs/build/<场>/…` —— 数据没丢, 但不再算产物、也不进台账
  171. (否则每从旧包补一次, `guanlan.py check` 的「日志统一」就红一次: 日志跑到 logs/ 外面了)。
  172. """
  173. low = rel.lower()
  174. if low.endswith(('.log', '.jsonl')):
  175. from src import logfile as _lf
  176. dst_log = _lf.build_log(rel, P.farm())
  177. prov[rel] = dict(source='log-relocated', why='包内日志: 按用户令 2 落到 logs/build/, 不算产物')
  178. if not a.dry_run:
  179. dst_log.parent.mkdir(parents=True, exist_ok=True)
  180. if src_bytes is not None:
  181. dst_log.write_bytes(src_bytes)
  182. elif src_path is not None:
  183. shutil.copy2(src_path, dst_log)
  184. return
  185. dst = dest_root / rel
  186. if dst.exists():
  187. # 已存在的件分两类: 我们自己重算的(raw-derived) 与 上轮已补齐的(shipped)。
  188. # ★不能因为"本轮没拷"就不记账 —— 早先写成"存在即不列", 于是重复跑一次台账就被清成 0/0,
  189. # 页面上的来源统计跟着一起错 (2026-09-12 实逮)。
  190. if rel in RAW_DERIVED:
  191. prov[rel] = dict(source='raw-derived', builder=RAW_DERIVED[rel])
  192. else:
  193. prov[rel] = dict(source='shipped', why='包内无生成端 / 规则未复现 → 随包件补齐 (本轮已在位)')
  194. return
  195. prov[rel] = dict(source='shipped', why='包内无生成端 / 规则未复现 → 用随包件补齐')
  196. if a.dry_run:
  197. return
  198. dst.parent.mkdir(parents=True, exist_ok=True)
  199. if src_bytes is not None:
  200. dst.write_bytes(src_bytes)
  201. else:
  202. shutil.copy2(src_path, dst)
  203. if stash.is_file():
  204. # 交付包 zip 形态 (2026-09-16): 按需从包里取 outputs/<场>/** 补齐 ——
  205. # 这是"清除不留备份"之后**唯一**的随包件来源, 所以把它做成一等输入。
  206. import zipfile
  207. with zipfile.ZipFile(stash) as zf:
  208. pref = f'outputs/{P.farm()}/'
  209. members = [n for n in zf.namelist() if n.startswith(pref) and not n.endswith('/')]
  210. print(f' 从交付包读取 {len(members)} 件 (前缀 {pref})')
  211. for name in sorted(members):
  212. if a.dry_run:
  213. take(name[len(pref):])
  214. else:
  215. take(name[len(pref):], src_bytes=zf.read(name))
  216. else:
  217. for src in sorted(stash.rglob('*')):
  218. if not src.is_file():
  219. continue
  220. take(src.relative_to(stash).as_posix(), src_path=src)
  221. # 重算件里有些**不在随包件里**(如我们新造的 turbine_params.parquet), 也要记进台账
  222. for rel, builder in RAW_DERIVED.items():
  223. if rel not in prov and (dest_root / rel).exists():
  224. prov[rel] = dict(source='raw-derived', builder=builder)
  225. # 构建脚本的**自登记** (src/derived_manifest.py): 振动侧的产物名随窗/分片变, 写不进上面的精确表,
  226. # 按名字通配又会误伤同名旧件 → 由"谁算的谁登记", 这里只认登记。缺失不影响其余记账。
  227. try:
  228. from src.derived_manifest import load as _load_manifest, prune as _prune_manifest
  229. _gone = _prune_manifest(dest_root) # 先清掉"登记了但盘上已删"的条目 (如被删的 _reimport 窗)
  230. if _gone:
  231. print(f' 自登记清理: {_gone} 条指向已删文件的条目')
  232. _man = (_load_manifest(dest_root).get('files') or {})
  233. except Exception:
  234. _man = {}
  235. for rel, info in _man.items():
  236. if (dest_root / rel).exists():
  237. prov[rel] = dict(source='raw-derived',
  238. builder=f"{info.get('builder', '?')} (自登记 by {info.get('by', '?')})")
  239. else:
  240. prov[rel] = dict(source='raw-derived',
  241. builder=f"{info.get('builder', '?')} (自登记, 但该件当前不在盘上)")
  242. by_store = {}
  243. for rel, m in prov.items():
  244. store = rel.split('/')[0]
  245. d = by_store.setdefault(store, {'raw-derived': 0, 'shipped': 0, 'log-relocated': 0})
  246. d.setdefault(m['source'], 0)
  247. d[m['source']] += 1
  248. moved = sum(d['log-relocated'] for d in by_store.values())
  249. print(f'{"仓":26s} {"raw 重算":>9s} {"随包补齐":>9s}' + (f' {"日志归位":>9s}' if moved else ''))
  250. for store, d in sorted(by_store.items()):
  251. print(f' {store:24s} {d["raw-derived"]:9d} {d["shipped"]:9d}'
  252. + (f' {d["log-relocated"]:9d}' if moved else ''))
  253. tot = {k: sum(d[k] for d in by_store.values()) for k in ('raw-derived', 'shipped', 'log-relocated')}
  254. print(f' {"合计":24s} {tot["raw-derived"]:9d} {tot["shipped"]:9d}'
  255. + (f' {tot["log-relocated"]:9d}' if moved else ''))
  256. if moved:
  257. print(f' (日志不落产物仓: {moved} 件按用户令 2 落到 logs/build/, 不计入产物台账)')
  258. if a.dry_run:
  259. print('\n(dry-run, 未写盘)')
  260. return 0
  261. out = dest_root / '_provenance.json'
  262. out.write_text(json.dumps(dict(
  263. at=time.strftime('%Y-%m-%d %H:%M:%S'),
  264. note='逐件来源台账: raw-derived = 由 data/raw 重算(含验证依据); shipped = 包内无生成端, 用随包件补齐',
  265. counts=tot, files=prov), ensure_ascii=False, indent=1), encoding='utf-8')
  266. print(f'\n已写来源台账 {P.rel(out)}')
  267. return 0
  268. if __name__ == '__main__':
  269. # 控制台可能是 GBK: print 里的非 GBK 字形编不出来会抛异常, 干成了却退出码 1
  270. for _s in (sys.stdout, sys.stderr):
  271. try: _s.reconfigure(errors='replace')
  272. except Exception: pass
  273. sys.exit(main())