service_worker.py 8.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201
  1. #!/usr/bin/env python3
  2. # -*- coding: utf-8 -*-
  3. r"""服务化共用工作体 (2026-09-17 用户令 1: 观澜安装为服务)。
  4. Windows 服务 (`scripts/win_service.py`) 与 Linux systemd (`scripts/service_main.py`) 都调这里,
  5. **同一套起停/守护逻辑,只有外壳不同** —— 免得两平台各写一份、行为慢慢分叉。
  6. ## 它在做什么
  7. 起 : `guanlan.py serve`(幂等:已在跑的复用,缺的补起,等 healthz)
  8. 守护 : 每 `INTERVAL` 秒跑一次同样的 `guanlan.py serve` —— 哪个组件掉了就被拉起来;
  9. 同时看网关 `/healthz`,连续失败则记 ERROR(服务管理器据此判活)
  10. 停 : `guanlan.py stop`(连网关一起停;服务模式下整机停服就该全停)
  11. 为什么用**子进程 + 现成 CLI** 而不是把服务代码直接写进各组件:
  12. 服务管理器的语义是"一个前台进程代表整套系统",而观澜本来就是 6 个独立进程;
  13. 复用 `guanlan.py` 的两条命令,起停口径与手工运维**完全一致**,不会出现"服务起的和手工起的不一样"。
  14. """
  15. from __future__ import annotations
  16. import datetime as dt
  17. import os
  18. import pathlib
  19. import socket
  20. import subprocess
  21. import sys
  22. import time
  23. ROOT = pathlib.Path(__file__).resolve().parents[1]
  24. sys.path.insert(0, str(ROOT))
  25. from src import paths as P # noqa: E402
  26. from src import version as V # noqa: E402
  27. INTERVAL = 15 # 守护巡检间隔(秒)
  28. LOG = None # 由外壳设置为 logs/service.log
  29. def log(msg: str, level: str = 'INFO') -> None:
  30. """服务日志 —— 与 §11.2 同一行格式,落 `logs/<组件>.log` 的 service 一份。"""
  31. from src import logfile as _lf
  32. line = _lf.line(level, 'service', msg)
  33. print(line, flush=True)
  34. try:
  35. with open(_lf.component_log('service'), 'a', encoding='utf-8', newline='\n') as f:
  36. f.write(line + '\n')
  37. except Exception:
  38. pass
  39. def _py() -> str:
  40. return str(P.venv_python() or sys.executable)
  41. def run(cmd: list, timeout: int = 300) -> int:
  42. """跑一条 guanlan.py 子命令(捕获输出写进服务日志)。"""
  43. p = subprocess.run([_py(), 'guanlan.py'] + cmd, cwd=str(ROOT), capture_output=True, text=True,
  44. errors='replace', timeout=timeout)
  45. for ln in (p.stdout or '').splitlines():
  46. log(' ' + ln)
  47. for ln in (p.stderr or '').splitlines():
  48. log(' ' + ln, 'WARN')
  49. return p.returncode
  50. def gateway_up(host: str = '127.0.0.1') -> bool:
  51. try:
  52. with socket.create_connection((host, 28084), timeout=1.5):
  53. return True
  54. except OSError:
  55. return False
  56. def start() -> None:
  57. log(f'启动 {V.NAME} v{V.VERSION}({V.EDITION})')
  58. rc = run(['serve'])
  59. log(f'起组件完成 (rc={rc}; 1 = 有模块未就绪, 属正常)')
  60. def ensure() -> bool:
  61. """一次巡检:把缺的组件补起来。→ 网关是否在听。"""
  62. run(['serve'])
  63. return gateway_up()
  64. def stop() -> None:
  65. log('停止全部组件(含网关)')
  66. rc = run(['stop'])
  67. log(f'停止完成 (rc={rc})')
  68. def raw_watch_cfg() -> dict:
  69. """输入数据看门狗配置(`configs/serve.json` 的 `raw_watch`)。**默认关** —— 自动重算会占机器约一小时,
  70. 该由现场决定开不开。读不到/写坏了按"关"处理(不因为配置问题去动数据)。"""
  71. d = dict(enabled=False, minutes=30, auto_rebuild=False)
  72. try:
  73. import json
  74. c = json.loads((ROOT / 'configs' / 'serve.json').read_text(encoding='utf-8-sig'))
  75. w = c.get('raw_watch') or {}
  76. if isinstance(w, dict):
  77. d.update({k: w[k] for k in ('enabled', 'minutes', 'auto_rebuild') if k in w})
  78. except Exception:
  79. pass
  80. return d
  81. _watch = {'last': 0.0}
  82. def raw_watch(stop_evt) -> None:
  83. """输入数据看门狗(用户令 2026-09-19「增加自动扫描识别机制,以能发现新增数据,并纳入重算」)。
  84. 每 `minutes` 分钟扫一次 `data/raw`:**发现新增/变化就在服务日志里写明"该跑哪几步"**;
  85. `auto_rebuild=true` 时顺手发起一次重算(走 `_ops_run.py`,与运维控制台「执行重算」同一条路)。
  86. ★两道保险:① 重算已在跑(`run/ops_job.json` status=running)就不再发起;② 本函数自己吞异常,
  87. 绝不因为看门狗把组件守护循环带崩。
  88. """
  89. c = raw_watch_cfg()
  90. if not c['enabled']:
  91. return
  92. now = time.time()
  93. if now - _watch['last'] < max(1, int(c['minutes'])) * 60:
  94. return
  95. _watch['last'] = now
  96. try:
  97. env = dict(os.environ, PYTHONUTF8='1', PYTHONIOENCODING='utf-8')
  98. env['WINDSCADA_ROOT'] = str(ROOT)
  99. p = subprocess.run([_py(), 'scripts/raw_scan.py', '--json', '--check'], cwd=str(ROOT), env=env,
  100. capture_output=True, text=True, encoding='utf-8', errors='replace', timeout=1800)
  101. if p.returncode == 5:
  102. log('输入数据扫描: 还没有基线快照 —— 记一份(下次起就能比出新增)')
  103. subprocess.run([_py(), 'scripts/raw_scan.py', '--write'], cwd=str(ROOT), env=env,
  104. capture_output=True, timeout=1800)
  105. return
  106. if p.returncode != 4:
  107. return
  108. import json as _json
  109. d = _json.loads(p.stdout or '{}')
  110. ch = d.get('changes') or []
  111. log(f'输入数据有 {len(ch)} 处新增/变化: ' + ';'.join(f'{x["family"]}({x["kind"]})' for x in ch[:6]))
  112. if d.get('steps'):
  113. log(' 该跑: ' + ' / '.join(d['steps']))
  114. job = ROOT / 'run' / 'ops_job.json'
  115. running = False
  116. try:
  117. running = (_json.loads(job.read_text(encoding='utf-8')).get('status') == 'running')
  118. except Exception:
  119. pass
  120. if not c['auto_rebuild']:
  121. log(' (raw_watch.auto_rebuild=false:只报不跑;要自动重算把它设 true)')
  122. return
  123. if running:
  124. log(' 已有重算在跑 —— 本次不重复发起(新数据会在下一轮或那次重算里被扫到)')
  125. return
  126. log(' 按配置发起一次重算(raw_watch.auto_rebuild=true)')
  127. subprocess.Popen([_py(), 'scripts/_ops_run.py', '--tag', 'rawwatch', '--', 'scripts/rebuild_all.py',
  128. '--auto'], cwd=str(ROOT), env=env, stdout=subprocess.DEVNULL,
  129. stderr=subprocess.DEVNULL, stdin=subprocess.DEVNULL)
  130. except Exception as e:
  131. log(f'输入数据看门狗异常: {type(e).__name__}: {e}', 'WARN')
  132. def supervise(stop_evt, once: bool = False) -> None:
  133. """守护循环:直到 stop_evt 被置位(或 once=True 只跑一轮)。"""
  134. fails = 0
  135. while not stop_evt.is_set():
  136. try:
  137. ok = ensure()
  138. if not ok:
  139. fails += 1
  140. log(f'网关未就绪(连续 {fails} 次)', 'ERROR' if fails >= 3 else 'WARN')
  141. else:
  142. if fails:
  143. log('网关已恢复')
  144. fails = 0
  145. raw_watch(stop_evt) # 输入数据看门狗(默认关; 见 raw_watch_cfg)
  146. except Exception as e:
  147. log(f'巡检异常: {type(e).__name__}: {e}', 'ERROR')
  148. if once:
  149. return
  150. stop_evt.wait(INTERVAL)
  151. def selftest() -> int:
  152. """不需要任何权限的自证:依赖、端口、命令行是否齐备(服务注册前先看这个)。"""
  153. ok = True
  154. log(f'{V.NAME} v{V.VERSION} 服务化自检 @ {ROOT}')
  155. py = _py()
  156. log(f' 解释器: {py} (存在: {pathlib.Path(py).is_file()})')
  157. ok &= pathlib.Path(py).is_file()
  158. log(f' guanlan.py: {(ROOT / "guanlan.py").is_file()}')
  159. ok &= (ROOT / 'guanlan.py').is_file()
  160. log(f' 网关端口 28084 当前: ' + ('在听' if gateway_up() else '未监听'))
  161. log(f' 版本记录: {V.info_path(ROOT)}' + ('(已存在)' if V.info_path(ROOT).is_file() else '(无 —— 全新安装)'))
  162. # 组件端口占用检查(服务接管前这些端口应为空闲或已由本系统占用)
  163. for name, port in (('detail', 18033), ('cms', 18020), ('sim', 18791), ('sim_sys', 18792), ('viewer', 64292)):
  164. with socket.socket() as s:
  165. s.settimeout(0.3)
  166. busy = s.connect_ex(('127.0.0.1', port)) == 0
  167. log(f' 端口 {port} ({name}): ' + ('在听(多半是本系统已起)' if busy else '空闲'))
  168. log('自检结论: ' + ('服务注册前置条件齐备' if ok else '有缺项,见上'), 'INFO' if ok else 'ERROR')
  169. return 0 if ok else 1