service_worker.py 8.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202
  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. from src import paths as P # 配置路径的唯一取用口(不手拼 configs/)
  75. c = json.loads(P.config('serve.json').read_text(encoding='utf-8-sig'))
  76. w = c.get('raw_watch') or {}
  77. if isinstance(w, dict):
  78. d.update({k: w[k] for k in ('enabled', 'minutes', 'auto_rebuild') if k in w})
  79. except Exception:
  80. pass
  81. return d
  82. _watch = {'last': 0.0}
  83. def raw_watch(stop_evt) -> None:
  84. """输入数据看门狗(用户令 2026-09-19「增加自动扫描识别机制,以能发现新增数据,并纳入重算」)。
  85. 每 `minutes` 分钟扫一次 `data/raw`:**发现新增/变化就在服务日志里写明"该跑哪几步"**;
  86. `auto_rebuild=true` 时顺手发起一次重算(走 `_ops_run.py`,与运维控制台「执行重算」同一条路)。
  87. ★两道保险:① 重算已在跑(`run/ops_job.json` status=running)就不再发起;② 本函数自己吞异常,
  88. 绝不因为看门狗把组件守护循环带崩。
  89. """
  90. c = raw_watch_cfg()
  91. if not c['enabled']:
  92. return
  93. now = time.time()
  94. if now - _watch['last'] < max(1, int(c['minutes'])) * 60:
  95. return
  96. _watch['last'] = now
  97. try:
  98. env = dict(os.environ, PYTHONUTF8='1', PYTHONIOENCODING='utf-8')
  99. env['WINDSCADA_ROOT'] = str(ROOT)
  100. p = subprocess.run([_py(), 'scripts/raw_scan.py', '--json', '--check'], cwd=str(ROOT), env=env,
  101. capture_output=True, text=True, encoding='utf-8', errors='replace', timeout=1800)
  102. if p.returncode == 5:
  103. log('输入数据扫描: 还没有基线快照 —— 记一份(下次起就能比出新增)')
  104. subprocess.run([_py(), 'scripts/raw_scan.py', '--write'], cwd=str(ROOT), env=env,
  105. capture_output=True, timeout=1800)
  106. return
  107. if p.returncode != 4:
  108. return
  109. import json as _json
  110. d = _json.loads(p.stdout or '{}')
  111. ch = d.get('changes') or []
  112. log(f'输入数据有 {len(ch)} 处新增/变化: ' + ';'.join(f'{x["family"]}({x["kind"]})' for x in ch[:6]))
  113. if d.get('steps'):
  114. log(' 该跑: ' + ' / '.join(d['steps']))
  115. job = ROOT / 'run' / 'ops_job.json'
  116. running = False
  117. try:
  118. running = (_json.loads(job.read_text(encoding='utf-8')).get('status') == 'running')
  119. except Exception:
  120. pass
  121. if not c['auto_rebuild']:
  122. log(' (raw_watch.auto_rebuild=false:只报不跑;要自动重算把它设 true)')
  123. return
  124. if running:
  125. log(' 已有重算在跑 —— 本次不重复发起(新数据会在下一轮或那次重算里被扫到)')
  126. return
  127. log(' 按配置发起一次重算(raw_watch.auto_rebuild=true)')
  128. subprocess.Popen([_py(), 'scripts/_ops_run.py', '--tag', 'rawwatch', '--', 'scripts/rebuild_all.py',
  129. '--auto'], cwd=str(ROOT), env=env, stdout=subprocess.DEVNULL,
  130. stderr=subprocess.DEVNULL, stdin=subprocess.DEVNULL)
  131. except Exception as e:
  132. log(f'输入数据看门狗异常: {type(e).__name__}: {e}', 'WARN')
  133. def supervise(stop_evt, once: bool = False) -> None:
  134. """守护循环:直到 stop_evt 被置位(或 once=True 只跑一轮)。"""
  135. fails = 0
  136. while not stop_evt.is_set():
  137. try:
  138. ok = ensure()
  139. if not ok:
  140. fails += 1
  141. log(f'网关未就绪(连续 {fails} 次)', 'ERROR' if fails >= 3 else 'WARN')
  142. else:
  143. if fails:
  144. log('网关已恢复')
  145. fails = 0
  146. raw_watch(stop_evt) # 输入数据看门狗(默认关; 见 raw_watch_cfg)
  147. except Exception as e:
  148. log(f'巡检异常: {type(e).__name__}: {e}', 'ERROR')
  149. if once:
  150. return
  151. stop_evt.wait(INTERVAL)
  152. def selftest() -> int:
  153. """不需要任何权限的自证:依赖、端口、命令行是否齐备(服务注册前先看这个)。"""
  154. ok = True
  155. log(f'{V.NAME} v{V.VERSION} 服务化自检 @ {ROOT}')
  156. py = _py()
  157. log(f' 解释器: {py} (存在: {pathlib.Path(py).is_file()})')
  158. ok &= pathlib.Path(py).is_file()
  159. log(f' guanlan.py: {(ROOT / "guanlan.py").is_file()}')
  160. ok &= (ROOT / 'guanlan.py').is_file()
  161. log(f' 网关端口 28084 当前: ' + ('在听' if gateway_up() else '未监听'))
  162. log(f' 版本记录: {V.info_path(ROOT)}' + ('(已存在)' if V.info_path(ROOT).is_file() else '(无 —— 全新安装)'))
  163. # 组件端口占用检查(服务接管前这些端口应为空闲或已由本系统占用)
  164. for name, port in (('detail', 18033), ('cms', 18020), ('sim', 18791), ('sim_sys', 18792), ('viewer', 64292)):
  165. with socket.socket() as s:
  166. s.settimeout(0.3)
  167. busy = s.connect_ex(('127.0.0.1', port)) == 0
  168. log(f' 端口 {port} ({name}): ' + ('在听(多半是本系统已起)' if busy else '空闲'))
  169. log('自检结论: ' + ('服务注册前置条件齐备' if ok else '有缺项,见上'), 'INFO' if ok else 'ERROR')
  170. return 0 if ok else 1