raw_scan.py 29 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463
  1. #!/usr/bin/env python3
  2. # -*- coding: utf-8 -*-
  3. r"""输入数据自动扫描识别 —— 发现 `<安装目录>/data/raw` 下的**新增/变化**,并给出"该重跑哪几步"。
  4. ## 为什么要有它(用户令 2026-09-19)
  5. 用户令:「对 `<安装目录>/data/raw` 目录下的接入数据处理,增加自动扫描识别机制,以能发现新增数据,并纳入重算。」
  6. 在此之前,"放了新数据要不要重算"全靠人记:`place_raw_data.py` 只管把现场包落位,`raw_data_check.py`
  7. 只管结构与命名对不对,`rebuild_all.py` 无脑全跑(或按人给的 `--skip-*` 跳过)——
  8. **没有任何一处回答"这回新来了什么、因此必须重跑哪几步"**。本脚本补这一环:
  9. 逐族指纹 (件数/体积/最新时间/清单摘要) → 与上次快照比 → 报【新增/变化】 → 映射到链上步骤
  10. ## 三档用法
  11. python scripts/raw_scan.py # 只报当前状态(不写快照)
  12. python scripts/raw_scan.py --write # 写/更新快照 outputs/<场>/_raw_scan.json
  13. python scripts/raw_scan.py --check # 与快照比: rc=0 无变化 · rc=4 有新增/变化 · rc=5 没有快照
  14. python scripts/raw_scan.py --check --write # 比完就把新快照记下(重算链第 ①b 步就是这条)
  15. python scripts/raw_scan.py --plan # 只打印"有变化时该跑哪些步"
  16. python scripts/raw_scan.py --json # 机器可读(rebuild_all 用它做"跳过了该跑的步"提醒)
  17. python scripts/raw_scan.py --deep # 变化判定加上**内容**哈希(默认只比 名字/大小/时间)
  18. python scripts/raw_scan.py --root <别的 raw> # 换个根扫(自测/多场)
  19. python scripts/raw_scan.py --selftest # 族表覆盖自检(与族表/链上步骤对账,防两处漂移)
  20. ## 口径与纪律
  21. · **指纹**默认 = 每族逐件 `相对路径|字节数|mtime(ns)` 排序后串起来取 sha1 —— 快、且"加了一件/改了一件/
  22. 删了一件"都会变;`--deep` 再叠一层内容 sha1(小文件才做,避免拿 153 GB 的 CMS 导出拼哈希)。
  23. · **未归类**:`data/raw/<场>/` 下出现族表里没有的子目录(新来的数据总会有第一次)→ 如实列进
  24. `unclassified` 并给出"没有消费者/需要认领"的话,**不猜**它是哪一族。
  25. · 这里只**发现与建议**,不改数据、不搬文件、不自动重算(要不要跑由 `rebuild_all.py` 或运维决定;
  26. 自动重算的开关见 `configs/serve.json` 的 `raw_watch`)。rc=4 是"有变化",不是失败。
  27. """
  28. from __future__ import annotations
  29. try:
  30. from app_common.app_common_guanlan.api import install_root as _install_root
  31. except ImportError: # 理论不可达;包结构异常时回退到按位置上跳
  32. from pathlib import Path as _P
  33. def _install_root(_f): return _P(_f).resolve().parents[3]
  34. import argparse
  35. import hashlib
  36. import json
  37. import os
  38. import pathlib
  39. import re
  40. import sys
  41. import time
  42. ROOT = _install_root(__file__) # 模块化后按标记找安装根(原 parents[1] 已不成立)
  43. sys.path.insert(0, str(ROOT))
  44. try: # GBK 控制台/日志重定向不再因 '²' 这类字符抛 UnicodeEncodeError
  45. sys.stdout.reconfigure(errors='replace')
  46. except Exception:
  47. pass
  48. SNAP_NAME = '_raw_scan.json'
  49. # ── 族表:raw 子目录 → (计件 glob, 链上步骤, 消费者说明) ──────────────────────────────
  50. # ★步骤标签与 `scripts/rebuild_all.py` 的计划名一一对应(人看到"③ SCADA 侧"就知道去跑哪一步)。
  51. FAMILIES: tuple[tuple[str, tuple[str, ...], tuple[str, ...], str], ...] = (
  52. ('故障报警', ('*.xls', '*.xlsx', '*.csv'), ('② 三门台账',),
  53. '报警事件导出 → windscada/alarms.parquet(报警集中/停机账)。缺 TimeOn/Alarmcode 的统计报表件会被摄入器跳过'),
  54. ('风机故障记录', ('*.xls', '*.xlsx', '*.csv'), ('② 三门台账',),
  55. '检修工单台账 → windscada/workorders.parquet(重复检修/停机损失)'),
  56. ('油样报告', ('*.pdf',), ('② 三门台账',),
  57. '油液化验报告 → windscada/oil_samples_index.parquet(油液面)'),
  58. ('scada_10min', ('*.csv',), ('③ SCADA 侧 10 个构建器', '③b SCADA 10min 窄仓 (按窗重算底座)',
  59. '④ 月度派生件', '④c 变桨面', '⑤b 事实契约 claim'),
  60. '10min 导出 → 功率曲线/损失/温度/停机/偏航/液压等 + 变桨面日粒度(温度/性能/变桨面)+ 窄仓(按窗重算底座)'),
  61. ('scada_1min', ('*.csv',), ('④c 变桨面',),
  62. '1min 导出(内部台号 01E…38B)→ 变桨面零位三口径与越线计数'),
  63. ('scada_mdb', ('*.mdb', '*.zip'), ('③ SCADA 侧 10 个构建器', '④ 月度派生件', '④c 变桨面'),
  64. '现场年度归档 Access 库(按月的 10min/1min)—— CSV 缺失时由 scada_source 回落到它'),
  65. ('windcms', ('*_decode.json',), ('④b 振动侧摄入 + CMS 报告 + 标量 z', '④d 振动在升/换件闭环', '④e 三层基线'),
  66. 'Brande TCM 导出 → 窗索引/谱库/CMS 报告/逐台页 + 在升闭环 + 三层基线(振动面)'),
  67. ('m5_cms_tcm', ('*.json', '*.docx', '*.pdf', '*.md'), ('④b 振动侧摄入 + CMS 报告 + 标量 z',),
  68. '现场给的正本(handoff_vibration_v2.json / component_history.json)与厂家报告——正本在位时优先于观澜自算'),
  69. )
  70. # 机理层资料不在场站目录下(A2 约定): data/raw/西门子4.0技术资料
  71. TECH_DIR, TECH_PATTERNS = '西门子4.0技术资料', ('*',) # 该族的消费者(kb_ingest)读整棵树 ⇒ 白名单就是"全部"
  72. TECH_STEPS = ('⑦ 本体: 码表/手册/文档',)
  73. TECH_NOTE = '厂商技术资料(故障处理手册/维护 WI/图纸/对译表)→ 本体对象库/检索索引/实机参数表'
  74. # raw 根下**随包发的说明件**(不是数据): 别把它们报成"未归类"。口径: 只有名字在白名单里且不参与摄入。
  75. IGNORE_FILES = {'README_把原始数据放这里.txt'}
  76. # ★2026-09-21 用户令「按月 10min 提取件就放在 <场站>/ 下」: 它是**给人核对/交接的导出件**, 不是摄入源
  77. # (摄入走 `<台号>.csv` + 同台补充件)。不加这条, 每次扫描都会把它报成"未归类 + 一处变化" ⇒
  78. # `--auto` / raw_watch 会永远以为"来了新数据", 每轮都触发一次全量重算。
  79. IGNORE_PATTERNS: tuple[tuple[str, str], ...] = (
  80. (r'^scada_\d+min_\d{6}\.csv$', '按月提取/核对件(非摄入源;摄入走 <台号>.csv + 同台补充件)'),
  81. )
  82. # ── 消费者取数口径(用户令 2026-09-20 步骤 C:让"放了没人读"当场可见)────────────────────────
  83. # ★为什么必须有它:族表原来只按**扩展名白名单**计件 —— 往 `scada_10min/` 里放 `WTG01-B2.csv`(856 列、
  84. # 2026-08 数据、与 `<台号>.csv` 命名不同)时,它命中 `*.csv` ⇒ 既算"进了摄入",又不会出现在任何
  85. # "没人读"的提示里;而取数层只认 `<台号>.csv` ⇒ **这批数据从来没被读过**,重算后时间窗自然不变。
  86. # 本表按**消费者真实取数口径**再判一次,并把结果分成两级:
  87. # · 缺口级 (gate):看着是数据(扩展名在 data_exts 里) 但没有任何消费者会读 ⇒ **计入变化(rc=4) + 大声报**
  88. # · 信息级 (info):附件/压缩包/图片/厂商软件这类,本来就不该被摄入 ⇒ 只列信息,不计变化
  89. CONSUME: dict[str, dict] = {
  90. '故障报警': dict(consume=('*.xls',), data_exts=('.xls', '.xlsx', '.csv'),
  91. how='摄入器读 *.xls 的 XML 事件导出(TimeOn/Alarmcode)'),
  92. '风机故障记录': dict(consume=('*.xls', '*.xlsx'), data_exts=('.xls', '.xlsx', '.csv'),
  93. how='台账读取 *.xls*(rar/jpg 是附件,不摄入)'),
  94. '油样报告': dict(consume=('*.pdf',), data_exts=('.pdf',), how='油样摄取读 *.pdf'),
  95. 'scada_10min': dict(consume=('WTG??.csv', 'WTG??-*.csv', 'WTG??_*.csv'), data_exts=('.csv', '.mdb', '.zip'),
  96. how='取数层读 `<台号>.csv` **+ 同台补充件** `<台号>-*.csv`(2026-09-20 起: 按列名对齐、'
  97. '按时间戳去重、主件优先)'),
  98. 'scada_1min': dict(consume=('???.csv', '???-*.csv', '???_*.csv'), data_exts=('.csv', '.mdb', '.zip'),
  99. how='取数层按内部台号取 `<台号>.csv`(01E.csv)+ 同台补充件'),
  100. 'scada_mdb': dict(consume=('*-10min.mdb', '*-1min.mdb', '*.zip'), data_exts=('.mdb', '.zip'),
  101. how='取数层读取按月的库(2026-08-10min.mdb);CSV 缺失时才回落',
  102. # ★2026-09-21: 现场按**类目**交付的月度通道组库(7月-8月交付 zip 里的
  103. # `2026-0X-<类>.mdb`)。它是 `scada_10min` 同台补充件的**上游**:取数层不直接读它,
  104. # 要先转成 `scada_10min/<台号>-<YYYYMM>.csv`。故列为"上游归档"(信息级),
  105. # 既不当缺口(不算变化/不触发重算),也不假装被摄入。
  106. upstream=tuple(f'*-{c}.mdb' for c in
  107. ('cnt', 'din', 'dot', 'flg', 'grd', 'int', 'prs', 'scd', 'std', 'sum', 'tmp', 'tur')),
  108. upstream_how='现场月度**通道组**归档(上游件):先转成同台 10min 补充件 '
  109. '`scada_10min/<台号>-<YYYYMM>.csv` 才被消费(口径见 docs/输入数据放置指导_v0.1.md §4.1)'),
  110. 'windcms': dict(consume=('*_decode.json', '*.pdf', '*.docx', '*.md', '*.json'),
  111. data_exts=('.json', '.zip', '.csv', '.mdb'),
  112. how='振动摄入读 *_decode.json;厂家报告 PDF 由报告转录侧读'),
  113. 'm5_cms_tcm': dict(consume=('*.json', '*.docx', '*.pdf', '*.md'), data_exts=('.json', '.zip'),
  114. how='现场正本 handoff/component_history 与厂家报告'),
  115. TECH_DIR: dict(consume=('*',), data_exts=(), how='本体读整棵树(白名单=全部)'),
  116. }
  117. def _digest(entries: list[tuple[str, int, int]], deep: bool, base: pathlib.Path, files: list[pathlib.Path],
  118. deep_max: int, dirs: list[str] | None = None) -> str:
  119. h = hashlib.sha1()
  120. for rel, size, mt in entries:
  121. h.update(f'{rel}|{size}|{mt}\n'.encode('utf-8', 'replace'))
  122. for d in dirs or []: # 子目录清单也进指纹: "新建了空目录(数据还没落)"也能发现
  123. h.update(f'D:{d}\n'.encode('utf-8', 'replace'))
  124. if deep:
  125. for p, (rel, size, _mt) in zip(files, entries):
  126. if size > deep_max:
  127. continue
  128. try:
  129. fh = hashlib.sha1()
  130. with open(p, 'rb') as f:
  131. for blk in iter(lambda: f.read(1 << 20), b''):
  132. fh.update(blk)
  133. h.update(rel.encode('utf-8', 'replace') + b'=' + fh.digest())
  134. except OSError:
  135. pass
  136. return h.hexdigest()
  137. def fingerprint(d: pathlib.Path, patterns: tuple[str, ...], deep: bool = False, deep_max: int = 8 << 20) -> dict:
  138. """一族目录的指纹。
  139. ★口径 (用户令 2026-09-19/20「对 data/raw 目录(含嵌套子目录)下,文件的增减做到监听」):
  140. · **递归**统计该族目录下**全部文件**(不限扩展名) —— 只按摄入白名单(patterns)统计的话,
  141. 往族目录里丢一个 `.zip`/`.txt` 就"看不见"了, 那不算监听; 白名单命中的另记 `matched`
  142. (才是真进摄入的件数), 差额用 `unmatched` 列出来给人看;
  143. · 指纹里**还含子目录清单**(只放路径名, 不放 mtime —— 目录 mtime 会随子文件变动而变, 那是噪声),
  144. 于是"新建了一个空目录(数据还没落进去)"也动指纹。
  145. """
  146. all_files: list[pathlib.Path] = []
  147. dirs: list[str] = []
  148. if d.is_dir():
  149. for p in d.rglob('*'):
  150. if p.is_file():
  151. all_files.append(p)
  152. elif p.is_dir():
  153. dirs.append(p.relative_to(d).as_posix())
  154. all_files = sorted(set(all_files))
  155. dirs.sort()
  156. matched: list[pathlib.Path] = []
  157. for pat in patterns:
  158. matched += [p for p in d.rglob(pat) if p.is_file()]
  159. mset = {p.relative_to(d).as_posix() for p in set(matched)}
  160. entries: list[tuple[str, int, int]] = []
  161. for p in all_files:
  162. try:
  163. st = p.stat()
  164. except OSError:
  165. continue
  166. entries.append((p.relative_to(d).as_posix(), st.st_size, st.st_mtime_ns))
  167. entries.sort()
  168. newest = ''
  169. if entries:
  170. newest = time.strftime('%Y-%m-%d %H:%M:%S', time.localtime(max(e[2] for e in entries) / 1e9))
  171. _unm = [e[0] for e in entries if e[0] not in mset]
  172. # ── 消费者口径再判一次(步骤 C): 扩展名"像数据"但没人读 ⇒ 缺口; 其余 ⇒ 信息 ──
  173. # ★判据必须落在"**全部文件里没被消费者读的**"上, 不能只看 unmatched(白名单没命中的那些) ——
  174. # 2026-09-20 实逮: `WTG01-B2.csv` 是 `*.csv`,白名单命中了它(所以 unmatched 里没有它),
  175. # 但取数层只认 `WTG01.csv` ⇒ 只看 unmatched 就永远发现不了这类"同名不同款"的静默漏读。
  176. spec = CONSUME.get(d.name) or {}
  177. cg = tuple(spec.get('consume') or ())
  178. dex = tuple(spec.get('data_exts') or ())
  179. consumed: set[str] = set()
  180. for g in cg:
  181. consumed |= {p.relative_to(d).as_posix() for p in d.rglob(g) if p.is_file()}
  182. # ── 上游归档件(客户口径里"已知但不由取数层直接读"的): 归信息级, 不算缺口、不计变化 ──
  183. upat = tuple(spec.get('upstream') or ())
  184. upstream: set[str] = set()
  185. for g in upat:
  186. upstream |= {p.relative_to(d).as_posix() for p in d.rglob(g) if p.is_file()}
  187. not_consumed = [r for r, _s, _m in entries if r not in consumed]
  188. gaps = [r for r in not_consumed
  189. if r not in upstream and pathlib.PurePosixPath(r).suffix.lower() in dex]
  190. others = [r for r in not_consumed if r not in gaps]
  191. return dict(files=len(entries), bytes=sum(e[1] for e in entries), newest=newest,
  192. matched=len(mset), dirs=len(dirs), unmatched=_unm[:10], unmatched_n=len(_unm),
  193. consumed=len(consumed), consume_how=spec.get('how', ''),
  194. upstream=len(upstream), upstream_how=spec.get('upstream_how', ''),
  195. gaps=gaps[:10], gaps_n=len(gaps), others=others[:10], others_n=len(others),
  196. digest=_digest(entries, deep, d, all_files, deep_max, dirs),
  197. top=[e[0] for e in sorted(entries, key=lambda x: -x[2])[:3]])
  198. def scan(raw_root: pathlib.Path, deep: bool = False, station_dir: pathlib.Path | None = None) -> dict:
  199. """扫一个 raw 根:逐族指纹 + 别处不认识的子目录(未归类)。
  200. `station_dir` 不给时按**场配置的辨识结果**定场站目录(`raw_station_dir()`,与摄入侧同一口径);
  201. 显式 `--root` 自测时可传入 `<root>/<场站名>` 之外的位置。
  202. """
  203. from src.windscada.config import raw_station_dir
  204. station = pathlib.Path(station_dir) if station_dir else raw_station_dir()
  205. fam: dict[str, dict] = {}
  206. for name, pats, steps, note in FAMILIES:
  207. fp = fingerprint(station / name, pats, deep)
  208. fp.update(steps=list(steps), note=note, where=str((station / name)))
  209. fam[name] = fp
  210. fp = fingerprint(raw_root / TECH_DIR, TECH_PATTERNS, deep)
  211. fp.update(steps=list(TECH_STEPS), note=TECH_NOTE, where=str(raw_root / TECH_DIR))
  212. fam[TECH_DIR] = fp
  213. known = {n for n, *_ in FAMILIES} | {TECH_DIR}
  214. unclassified = []
  215. # ★2026-09-22 可移植性消缺(check_portability --runtime 报 ERROR): 原先用
  216. # `str(p).startswith(k + os.sep)` 拼系统分隔符判"在不在已知族目录下" —— 一是 os.sep 让代码
  217. # 带上平台特征(换机部署核查报 ERROR),二是前缀匹配会把 `scada_10min_extra/` 误认成
  218. # `scada_10min/` 的子件。改成 pathlib 的祖先判定,语义更准且与平台无关。
  219. known_dirs = {station / n for n, *_ in FAMILIES} | {raw_root / TECH_DIR}
  220. # ★未归类 = 任何**不在已知族目录里**的文件 (递归找, 不只看一层): raw 根下、场站目录下、以及
  221. # 场站目录里嵌套的新子目录(例如新来了 `CMS_202701` 这种整包目录) 都跑不掉。
  222. for p in (sorted(raw_root.rglob('*')) if raw_root.is_dir() else []):
  223. if not p.is_file():
  224. continue
  225. if any(p == k or k in p.parents for k in known_dirs):
  226. continue
  227. if p.name in IGNORE_FILES:
  228. continue # 随包的放置说明件: 不是数据, 不报"未归类"
  229. if any(re.match(rx, p.name) for rx, _why in IGNORE_PATTERNS):
  230. continue # 交付核对用的导出件(位置由用户指定): 非摄入源, 不报"未归类"
  231. rel = p.relative_to(raw_root).as_posix()
  232. unclassified.append(dict(dir=rel, files=1,
  233. note='不在任何已知族的目录里 —— 没有消费者; 是新的输入族就得给它配摄入器'))
  234. # 已知族目录内**消费者不会读**的件 (步骤 C): 两级 —— 缺口级(gaps, 扩展名像数据) 与 信息级(others)
  235. stray, gaps, upstream = [], [], []
  236. for name, fp in fam.items():
  237. if fp.get('others_n'):
  238. stray.append(dict(family=name, n=fp['others_n'], example=fp.get('others')[:5]))
  239. if fp.get('gaps_n'):
  240. gaps.append(dict(family=name, n=fp['gaps_n'], example=fp.get('gaps')[:8],
  241. how=fp.get('consume_how', '')))
  242. if fp.get('upstream'):
  243. upstream.append(dict(family=name, n=fp['upstream'], how=fp.get('upstream_how', '')))
  244. # 场站目录下**新增的未知子目录**(整目录级; 其文件已在上面逐件登记, 这里给个汇总便于人读)
  245. extra_dirs = []
  246. if station.is_dir():
  247. for sub in sorted(station.iterdir()):
  248. if sub.is_dir() and sub.name not in known:
  249. extra_dirs.append(dict(dir=f'{station.name}/{sub.name}',
  250. files=sum(1 for p in sub.rglob('*') if p.is_file()),
  251. note='场站目录下的新子目录(族表里没有)'))
  252. return dict(at=time.strftime('%Y-%m-%d %H:%M:%S'), root=str(raw_root), station=str(station),
  253. deep=deep, families=fam, unclassified=unclassified + extra_dirs,
  254. stray_in_families=stray, gaps=gaps, upstream=upstream)
  255. def diff(old: dict, new: dict) -> list[dict]:
  256. """逐族比:新增/减少/内容变化(`files` 与 `digest` 分开说,便于人看懂)。"""
  257. out = []
  258. of = (old or {}).get('families') or {}
  259. for name, cur in (new.get('families') or {}).items():
  260. prev = of.get(name)
  261. if not prev:
  262. if cur['files']:
  263. out.append(dict(family=name, kind='首次记录', files=cur['files'], steps=cur['steps'],
  264. note=cur['note'], newest=cur['newest']))
  265. continue
  266. d_files = cur['files'] - prev['files']
  267. same_list = cur['digest'] == prev['digest']
  268. if d_files == 0 and same_list:
  269. continue
  270. kind = ('新增 %+d 件' % d_files) if d_files else '内容变化(件数不变)'
  271. _dd = cur.get('dirs', 0) - prev.get('dirs', 0)
  272. if d_files == 0 and _dd:
  273. kind = f'目录结构变化({_dd:+d} 个子目录)'
  274. out.append(dict(family=name, kind=kind, files=cur['files'], delta=d_files,
  275. bytes_delta=cur['bytes'] - prev['bytes'], newest=cur['newest'],
  276. top=cur.get('top'), steps=cur['steps'], note=cur['note'],
  277. matched=cur.get('matched'), unmatched_n=cur.get('unmatched_n', 0)))
  278. for u in new.get('unclassified') or []:
  279. out.append(dict(family=u['dir'], kind='未归类(不在任何已知族的目录里)', files=u['files'],
  280. steps=[], note=u.get('note') or '没有已知消费者 ⇒ 先认领'))
  281. # ★缺口级 (步骤 C, **计入变化**): 扩展名看着是数据、但按消费者口径没人会读 ⇒ 重算也不会纳入。
  282. # 典型: scada_10min/ 里的 `WTG01-B2.csv`(856 列 / 2026-08 数据) —— 取数层只认 `WTG01.csv`。
  283. for g in new.get('gaps') or []:
  284. out.append(dict(family=f'{g["family"]}(数据没人读·缺口)', kind=f'{g["n"]} 件数据不会被摄入',
  285. files=g['n'], steps=[],
  286. note=f'消费者口径: {g["how"]} ⇒ 例: ' + '、'.join(g['example'])
  287. + ' —— 这些件**不会进任何产物**(重算也不会变); 要么按消费者口径改名/转换, '
  288. '要么给取数层加"同台多件"的合并口径'))
  289. # ★附件类(信息级,**不计入变化**): 实测族目录里本来就有大量正常但不该被摄入的东西
  290. # (故障记录里的 .rar/.jpg、windcms 厂家报告 PDF、技术资料里的厂商软件 .exe/.lp …)。
  291. # 它们增删本身已经反映在 files/dirs/digest 上 ⇒ 该报的变化照报, 这里不再重复判一次。
  292. return out
  293. def _steps_of(changes: list[dict]) -> list[str]:
  294. seen, out = set(), []
  295. for c in changes:
  296. for s in c.get('steps') or []:
  297. if s not in seen:
  298. seen.add(s)
  299. out.append(s)
  300. return out
  301. def snap_path(raw_root: pathlib.Path) -> pathlib.Path:
  302. from src import paths as P
  303. return P.out_root() / SNAP_NAME
  304. def main() -> int:
  305. ap = argparse.ArgumentParser(description='输入数据自动扫描识别(发现新增 → 指明该重跑哪几步)')
  306. ap.add_argument('--farm', default=None)
  307. ap.add_argument('--root', default=None, help='覆盖 data/raw 根(自测/多场)')
  308. ap.add_argument('--write', action='store_true', help='把本次指纹写进快照')
  309. ap.add_argument('--check', action='store_true', help='与快照比:rc=0 无变化 · 4 有新增/变化 · 5 无快照')
  310. ap.add_argument('--plan', action='store_true', help='只打印"该跑哪些步"')
  311. ap.add_argument('--json', action='store_true', help='输出机器可读 JSON')
  312. ap.add_argument('--deep', action='store_true', help='变化判定叠内容哈希(小文件)')
  313. ap.add_argument('--selftest', action='store_true', help='族表/步骤覆盖自检(防与族表、链上步骤漂移)')
  314. a = ap.parse_args()
  315. from src import paths as P
  316. from src.windscada import config as C
  317. if a.farm:
  318. os.environ['WINDSCADA_FARM'] = a.farm # 与其它脚本同一口径(场由环境变量选)
  319. raw = pathlib.Path(a.root) if a.root else pathlib.Path(C.RAW_ROOT)
  320. # --root 换成别的 raw 根时,场站目录按"根下唯一/同名"识别,避免拿本场辨识结果去别处找
  321. st_dir = None
  322. if a.root:
  323. from src.windscada.config import station_scan
  324. sc = station_scan(root=raw)
  325. st_dir = (raw / sc['matched']) if sc.get('matched') else None
  326. now = scan(raw, deep=a.deep, station_dir=st_dir)
  327. sp = snap_path(raw)
  328. old = None
  329. if sp.is_file():
  330. try:
  331. old = json.loads(sp.read_text(encoding='utf-8'))
  332. except Exception:
  333. old = None
  334. changes = diff(old, now)
  335. if a.selftest:
  336. return selftest()
  337. if a.json:
  338. print(json.dumps(dict(scan=now, changes=changes, snapshot=str(sp),
  339. steps=_steps_of(changes), has_snapshot=bool(old)), ensure_ascii=False, indent=1))
  340. if a.write:
  341. sp.parent.mkdir(parents=True, exist_ok=True)
  342. sp.write_text(json.dumps(now, ensure_ascii=False, indent=1), encoding='utf-8')
  343. if a.check:
  344. return 0 if (old and not changes) else (4 if old else 5)
  345. return 0
  346. print(f'输入数据扫描 · {now["root"]} ({now["at"]}{" · deep" if a.deep else ""})')
  347. tot_f = sum(f['files'] for f in now['families'].values())
  348. tot_b = sum(f['bytes'] for f in now['families'].values())
  349. print(f' 共 {len(now["families"])} 族 · {tot_f} 件 · {tot_b / 1024 / 1024 / 1024:.1f} GB')
  350. for name, f in now['families'].items():
  351. flag = ''
  352. for c in changes:
  353. if c['family'] == name:
  354. flag = f" ← {c['kind']}"
  355. print(f' {name:16s} {f["files"]:6d} 件 {f["bytes"] / 1024 / 1024:9.1f} MB 最新 {f["newest"] or "—"}{flag}')
  356. if now['unclassified']:
  357. print(' 未归类(族表里没有,先认领):')
  358. for u in now['unclassified']:
  359. print(f' {u["dir"]:16s} {u["files"]:6d} 件 {u.get("note", "没有已知消费者")}')
  360. if now.get('gaps'):
  361. print(' ★[缺口] 族目录里有**看着是数据、但没有任何消费者会读**的件(重算不会纳入它们):')
  362. for g in now['gaps']:
  363. print(f' {g["family"]:16s} {g["n"]:5d} 件 例: ' + '、'.join(g['example'][:4])
  364. + (' …' if g['n'] > 4 else ''))
  365. print(f' 消费者口径: {g["how"]}')
  366. if now.get('upstream'):
  367. print(' (信息,不计入"变化")族目录里的**上游归档件**(要先转换才被消费):')
  368. for u in now['upstream']:
  369. print(f' {u["family"]:16s} {u["n"]:5d} 件')
  370. print(f' {u["how"]}')
  371. if now.get('stray_in_families'):
  372. print(' (信息,不计入"变化")族目录里**消费者不读的附件类**件:')
  373. for s in now['stray_in_families']:
  374. print(f' {s["family"]:16s} {s["n"]:5d} 件 例: ' + '、'.join(s['example'][:3]))
  375. if changes:
  376. print(f'\n== 本次发现 {len(changes)} 处新增/变化 ==')
  377. for c in changes:
  378. print(f' · {c["family"]}: {c["kind"]}'
  379. + (f'(体积 {c["bytes_delta"] / 1024 / 1024:+.1f} MB)' if c.get('bytes_delta') else '')
  380. + f' 最新 {c.get("newest") or "—"}')
  381. if c.get('top'):
  382. print(f' 最近落盘: ' + '、'.join(c['top']))
  383. if c.get('steps'):
  384. print(f' 该跑: ' + ' / '.join(c['steps']))
  385. _real = [c for c in changes if '缺口' not in c['family']]
  386. if _real:
  387. print('\n 建议: 跑一次重算把新数据纳入产物 —— python scripts/rebuild_all.py'
  388. '(默认全跑;若你用了 --skip-*,请确认没跳过上面列出的步)')
  389. if len(_real) < len(changes):
  390. print(' ★缺口类变化**重算也纳入不了**(消费者根本不读这些件): 先按上面写的口径改名/转换, '
  391. '或给取数层加"同台多件"的合并口径(见 docs/输入数据放置指导_v0.1.md §4)')
  392. else:
  393. print('\n== 与上次快照一致,没有新数据 ==' if old else '\n== 还没有快照(跑 --write 记一次基线)==')
  394. if a.write:
  395. sp.parent.mkdir(parents=True, exist_ok=True)
  396. sp.write_text(json.dumps(now, ensure_ascii=False, indent=1), encoding='utf-8')
  397. print(f'\n快照已更新 → {P.rel(sp)}')
  398. if a.check:
  399. if not old:
  400. print(f'[?] 没有快照可比({P.rel(sp)})—— 先跑 --write 记基线; rc=5')
  401. return 5
  402. return 4 if changes else 0
  403. return 0
  404. def selftest() -> int:
  405. """自检:① 族表覆盖 `products_reverse_audit.FAMILIES` 里所有 raw 输入目录 ② 步骤名都能在链上找到。"""
  406. import importlib.util
  407. import re
  408. spec = importlib.util.spec_from_file_location('_pra', ROOT / 'scripts' / 'products_reverse_audit.py')
  409. mod = importlib.util.module_from_spec(spec)
  410. spec.loader.exec_module(mod)
  411. inputs = {f['input'] for f in mod.FAMILIES if f.get('input')}
  412. mine = {n for n, *_ in FAMILIES} | {TECH_DIR}
  413. missing = sorted(inputs - mine)
  414. chain = (ROOT / 'scripts' / 'rebuild_all.py').read_text(encoding='utf-8')
  415. steps = {s for _n, _p, ss, _note in FAMILIES for s in ss} | set(TECH_STEPS)
  416. bad_steps = sorted(s for s in steps if s.split()[0] not in chain)
  417. print(f'族表覆盖: 反向呼应族表的 raw 输入 {len(inputs)} 个 · 本器覆盖 {len(mine)} 个 · 缺 {missing or "无"}')
  418. print(f'步骤标签: {len(steps)} 个 · 在 rebuild_all 里找不到的 {bad_steps or "无"}')
  419. ok = not missing and not bad_steps
  420. print('[OK] 自检通过' if ok else '[X] 自检不过 —— 族表或步骤标签与链上漂移了')
  421. return 0 if ok else 5
  422. if __name__ == '__main__':
  423. sys.exit(main())