#!/usr/bin/env python3 # -*- coding: utf-8 -*- r"""观澜运维控制台的后端 (2026-09-12) —— 把"停服务 / 起服务 / 重算 / 清产物"做成一个可视化页面。 ## 它解决什么 原来这四件事都得在黑窗口里敲命令 (顺序还不能错), 且没有任何"当前能不能做"的约束。 本模块给网关加三个东西: · `GET /ops` 一个自适应控制台页面(单文件, 无外部依赖); · `GET /ops/api/state` 真实状态: 各服务端口通不通 · 产物在不在 · 有没有任务在跑 · 上次结果; · `POST /ops/api/<动作>` 停服务 · 起服务 · 重算 · 清产物 · 恢复产物。 ## 三条设计纪律 1. **页面活着的服务不能被自己停掉。** 控制台由网关(28084)提供, 所以"停服务"默认**保留网关**, 否则按钮刚点完页面就没了、再也没法"启动服务"。要连网关一起停, 用 `include_gateway=True`, 这时用**分离进程**先停再起(页面会断开十几秒, 之后自动重连)。 2. **一次只允许一个动作。** 四个动作互相冲突(重算要重启服务、清产物要挪走产物), 所以有一把锁: `run/ops_job.json` 里记着当前任务; 任务在跑时, 所有会冲突的按钮在后端**与前端都会被禁掉** (后端拒绝 = 真生效, 前端禁用只是提示)。判断以后端为准。 3. **日志与状态落盘。** 动作全部 `Popen` 到 `logs/ops_<动作>_<时间>.log`, 页面轮询 job 状态并把日志尾巴显示出来 —— 用户看得见"它在干什么", 而不是按钮转圈。 ## 按钮什么时候该灰 (后端 state 直接给出, 前端只管照用) 停服务 : 至少有一个组件服务在监听 启动服务 : 至少有一个组件服务没在监听 执行重算 : 没有任务在跑 清除产物 : 没有任务在跑 且 产物在位 且 暂存区没有同名目录(或允许 --archive-old) 恢复产物 : 没有任务在跑 且 产物被清掉过(暂存区清单存在) """ from __future__ import annotations import datetime as dt import json import os import pathlib import socket import subprocess 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 RUN = ROOT / 'run' LOGS = ROOT / 'logs' JOB = RUN / 'ops_job.json' OFF = ROOT / '_products_off' OFF_MANIFEST = OFF / 'manifest.json' PY = P.venv_python() if hasattr(P, 'venv_python') else sys.executable # 组件服务 (不含网关): 名字 → 端口。端口真源在 configs/serve.json, 这里只是网关不可用时的兜底。 DEFAULT_PORTS = dict(detail=18033, cms=18020, sim=18791, sim_sys=18792, viewer=64292) GATEWAY_PORT = 28084 KEEP_WHEN_STOPPING = 'gateway' # 停服务默认保留网关(控制台自己在上面) def ports() -> dict: p = ROOT / 'configs' / 'serve.json' cfg = json.loads(p.read_text(encoding='utf-8-sig')) if p.exists() else {} out = dict(DEFAULT_PORTS) out['gateway'] = int(cfg.get('gateway') or GATEWAY_PORT) host = cfg.get('host') or '127.0.0.1' for k in list(out): if k in cfg: out[k] = int(cfg[k]) return dict(host=host, **out) def listening(host: str, port: int, timeout=0.35) -> bool: """端口有没有人在听。注意: 本机实测"连已关闭端口会一直等到超时"(不回 RST), 所以超时必须短。""" try: with socket.create_connection((host, port), timeout=timeout): return True except OSError: return False def _job() -> dict: try: return json.loads(JOB.read_text(encoding='utf-8')) except Exception: return {} def job_running() -> tuple[bool, dict]: """(是否在跑, 任务记录)。在跑 = 状态是 running 且执行器进程还活着。 注意: 退出码由执行器 scripts/_ops_run.py 写回 job (不在跑时才可信), 所以这里只判"活没活"; 进程没了就把 running 收尾成 done(执行器正常收尾时会自己写, 这里兜的是被强杀的情况)。 """ j = _job() if not j or j.get('status') != 'running': return False, j pid = j.get('pid') # 心跳优先: 执行器每 15 s 更新 job 文件; 超过 STALE_S 没动静 = 被强杀(pid 可能被复用, 只看 pid 会误判) STALE_S = 150 try: age = time.time() - JOB.stat().st_mtime except OSError: age = 0 if age > STALE_S: j.update(status='done', rc=None, finished=time.strftime('%Y-%m-%d %H:%M:%S'), note=(j.get('note') or '') + f' (心跳停了 {int(age)} s —— 任务多半被强杀, 未拿到退出码)') JOB.write_text(json.dumps(j, ensure_ascii=False, indent=1), encoding='utf-8') return False, j alive = False if pid: try: if os.name == 'nt': # errors='replace': tasklist 输出是控制台代码页(GBK), 而本进程可能是 PYTHONUTF8=1 # → 严格解码会失败并使 stdout 为 None (guanlan.py 的 alive() 踩过同一个坑) out = subprocess.run(['tasklist', '/FI', f'PID eq {pid}'], capture_output=True, text=True, errors='replace').stdout alive = bool(out) and str(pid) in out else: os.kill(pid, 0); alive = True except Exception: alive = False if not alive: j.update(status=j.get('status') if j.get('rc') is not None else 'done', finished=j.get('finished') or time.strftime('%Y-%m-%d %H:%M:%S'), note=(j.get('note') or '') + (' (进程已结束但没写回退出码 —— 多半是被强杀)' if j.get('rc') is None else '')) JOB.write_text(json.dumps(j, ensure_ascii=False, indent=1), encoding='utf-8') return False, j return True, j def _log_tail(path, n=40) -> list: if not path: return [] f = pathlib.Path(path) if not f.is_file(): return [] lines = f.read_text(encoding='utf-8', errors='replace').splitlines() return lines[-n:] def spawn(args: list, tag: str, detached=False) -> dict: """起一个动作: 经 `_ops_launch.py` 二次启动 `_ops_run.py`。 两层是**必须的**(见 _ops_launch.py 注释): 直接 Popen 的话, 任务是网关的子进程, 而 `guanlan.py stop` 用 `taskkill /T` —— 停网关时会连带杀掉正在执行的任务(实测把自己杀了)。 执行器负责写回真实退出码; 日志落 logs/ops__<时间>.log。 """ RUN.mkdir(exist_ok=True); LOGS.mkdir(exist_ok=True) log = LOGS / f'ops_{tag}_{time.strftime("%Y%m%d_%H%M%S")}.log' env = dict(os.environ, PYTHONUTF8='1', PYTHONIOENCODING='utf-8', PYTHONUNBUFFERED='1') r = subprocess.run([PY, 'scripts/_ops_launch.py', '--tag', tag, '--log', str(log), '--'] + list(args), cwd=str(ROOT), env=env, capture_output=True, text=True, errors='replace', timeout=60) pid = None for line in reversed((r.stdout or '').strip().splitlines()): if line.strip().isdigit(): pid = int(line.strip()); break if pid is None: print(' 启动器输出异常:', r.stdout, r.stderr) j = dict(kind=tag, cmd=' '.join(args), pid=pid, log=str(log), status='running', rc=None, started=time.strftime('%Y-%m-%d %H:%M:%S'), note='') JOB.write_text(json.dumps(j, ensure_ascii=False, indent=1), encoding='utf-8') return j # ─────────────────────────────────────────────────────────── 状态 def products_state() -> dict: store = P.store() # outputs/<场>/windscada fam = store.parent # outputs/<场> —— 产物仓的根 n = sum(1 for x in fam.rglob('*') if x.is_file()) if fam.is_dir() else 0 stores = {} if fam.is_dir(): for d in sorted(x for x in fam.iterdir() if x.is_dir()): stores[d.name] = sum(1 for _ in d.rglob('*') if _.is_file()) prov = fam / '_provenance.json' prov_d = json.loads(prov.read_text(encoding='utf-8')) if prov.is_file() else None stash_dirs = {} if OFF.is_dir(): for d in sorted(x for x in OFF.iterdir() if x.is_dir() and x.name != '_baseline_kept'): stash_dirs[d.name] = sum(1 for _ in (d / store.name).rglob('*') if _.is_file()) if (d / store.name).is_dir() else 0 kept = list((OFF / '_baseline_kept').rglob('*')) if (OFF / '_baseline_kept').is_dir() else [] return dict( fam=str(fam.relative_to(ROOT)).replace('\\', '/'), store=str(store.relative_to(ROOT)).replace('\\', '/'), files=n, in_place=n > 0, stores=stores, cleared=n == 0, stash_has_same=any(v > 0 for v in stash_dirs.values()), stash=stash_dirs, can_restore=OFF_MANIFEST.is_file() and n == 0, baseline_kept=sum(1 for x in kept if x.is_file()), provenance=(dict(counts=prov_d.get('counts'), at=prov_d.get('at')) if prov_d else None), ) def state() -> dict: pt = ports(); host = pt['host'] svc = {} for name, port in pt.items(): if name == 'host': continue svc[name] = dict(port=port, up=listening(host, port), is_gateway=(name == 'gateway')) comps = {k: v for k, v in svc.items() if not v['is_gateway']} gw_up = svc.get('gateway', {}).get('up', False) running, j = job_running() pr = products_state() any_comp_up = any(v['up'] for v in comps.values()) any_comp_down = any(not v['up'] for v in comps.values()) return dict( now=dt.datetime.now().isoformat(timespec='seconds'), host=host, services=svc, gateway_up=gw_up, products=pr, job=(dict(kind=j.get('kind'), status=j.get('status'), started=j.get('started'), note=j.get('note'), cmd=j.get('cmd'), rc=j.get('rc'), seconds=j.get('seconds'), finished=j.get('finished'), log_tail=_log_tail(j.get('log'))) if j else None), # 按钮可用性 —— 前端只照这个渲染, 后端还会再拒一次(见 check()) buttons=dict( stop_services=any_comp_up and not running, start_services=any_comp_down and not running, rebuild=not running, products_off=pr['in_place'] and not running, products_on=(not pr['in_place']) and not running, ), anchors=dict(alarms=_rows(f'{pr["fam"]}/windscada/alarms.parquet'), workorders=_rows(f'{pr["fam"]}/windscada/workorders.parquet'), temp_monthly=_rows(f'{pr["fam"]}/windscada/temp_monthly.parquet'), objects=_objects_count(f'{pr["fam"]}/ontology/objects.json')), ) def _rows(rel: str): """某件 parquet 的行数 (路径是相对安装根的), 用于页面上的"验收锚点"自检。""" f = ROOT / rel if not f.is_file(): return None try: import pandas as pd return int(len(pd.read_parquet(f))) except Exception: return None def _objects_count(rel: str): f = ROOT / rel try: return len(json.loads(f.read_text(encoding='utf-8'))) if f.is_file() else None except Exception: return None # ─────────────────────────────────────────────────────────── 动作 def _guard(what: str) -> str | None: running, j = job_running() if running: return f'已有任务在跑 ({j.get("kind")}, 起于 {j.get("started")}) —— 等它结束再操作。' return None def _svc_flags() -> tuple[bool, bool]: """(有组件在跑, 有组件没跑) —— 供动作前置检查, 与 state() 里 buttons 用的是同一判据。""" pt = ports(); host = pt['host'] ups = [listening(host, port) for name, port in pt.items() if name != 'host' and name != 'gateway'] return any(ups), any(not u for u in ups) def act_stop(body: dict) -> tuple[int, dict]: pt = ports(); include_gw = bool(body.get('include_gateway')) if (e := _guard('stop')): return 409, dict(err=e) any_up, _ = _svc_flags() if not any_up: # 按钮在前端就该是灰的, 但后端也必须拒 —— 否则直接打 API/重复提交就能做出无意义动作 return 409, dict(err='组件服务都已停止, 没有可停的 (按钮在页面上是禁用的)。') if include_gw: # 连网关一起停: 必须分离进程"先停 → 睡一会 → 再起", 否则控制台自己没了就再也起不回来。 # 现在经 _ops_launch 二次启动 → 这个任务不是网关的子孙, stop 里的 taskkill /T 杀不到它, # 所以"起回来"那半一定执行得到(此前会让整个栈停在半路)。 j = spawn(['-c', 'import subprocess,sys,time;' f'subprocess.run([r"{PY}", "guanlan.py", "stop"]);' 'time.sleep(2);' f'subprocess.run([r"{PY}", "guanlan.py", "serve"])'], 'restart_all') j['note'] = '含网关的完整重启: 页面会断开约 15 秒, 之后自动重连' JOB.write_text(json.dumps(j, ensure_ascii=False, indent=1), encoding='utf-8') return 200, dict(ok=True, job=j, note=j['note']) j = spawn(['scripts/_ops_stop_keep_gateway.py'], 'stop_keep_gw') return 200, dict(ok=True, job=j, note='已停组件服务, 保留网关(控制台) —— 页面继续可用') def act_start(body: dict) -> tuple[int, dict]: if (e := _guard('start')): return 409, dict(err=e) _, any_down = _svc_flags() if not any_down: return 409, dict(err='组件服务都在运行, 无需启动 (按钮在页面上是禁用的)。') open_browser = body.get('open_browser', True) j = spawn(['scripts/_ops_start_and_open.py'] + ([] if open_browser else ['--no-open']), 'start') return 200, dict(ok=True, job=j, note='正在启动全部服务' + ('; 起来后自动打开 http://127.0.0.1:%d/' % ports()['gateway'] if open_browser else '')) def act_rebuild(body: dict) -> tuple[int, dict]: if (e := _guard('rebuild')): return 409, dict(err=e) args = ['scripts/rebuild_all.py'] if body.get('skip_scada', True): args.append('--skip-scada') if body.get('with_verify'): args.append('--with-verify') src = (body.get('src') or '').strip() if src: args += ['--src', src] j = spawn(args, 'rebuild') return 200, dict(ok=True, job=j, note='重算已开始(含重启服务步骤); 进度看日志尾巴') def act_products_off(body: dict) -> tuple[int, dict]: if (e := _guard('products_off')): return 409, dict(err=e) pr = products_state() if not pr['in_place']: return 409, dict(err='产物已经是清空状态, 无需再清') args = ['scripts/products_state.py', '--off'] if pr['stash_has_same'] or body.get('archive_old'): args.append('--archive-old') # 暂存区已有同名产物时必须加, 否则会嵌套(见该脚本注释) j = spawn(args, 'products_off') return 200, dict(ok=True, job=j, note='正在把产物挪到 _products_off/ (挪完请点"启动服务"让页面呈现空状态)') def act_products_on(body: dict) -> tuple[int, dict]: if (e := _guard('products_on')): return 409, dict(err=e) pr = products_state() if pr['in_place']: return 409, dict(err='产物已在位, 无需恢复') if not OFF_MANIFEST.is_file(): return 409, dict(err='没有 _products_off/manifest.json (没清过或清单已删), 无法恢复') j = spawn(['scripts/products_state.py', '--on'], 'products_on') return 200, dict(ok=True, job=j, note='正在恢复产物; 恢复后请点"启动服务"') ACTIONS = dict(stop_services=act_stop, start_services=act_start, rebuild=act_rebuild, products_off=act_products_off, products_on=act_products_on) def handle(method: str, path: str, body: bytes) -> tuple[int, str, bytes]: """网关调用入口: → (http code, content-type, body bytes)。""" if path in ('/ops', '/ops/'): return 200, 'text/html; charset=utf-8', PAGE.encode('utf-8') if path == '/ops/api/state': return 200, 'application/json; charset=utf-8', json.dumps(state(), ensure_ascii=False).encode('utf-8') if path.startswith('/ops/api/'): name = path[len('/ops/api/'):].strip('/') if method != 'POST': return 405, 'application/json; charset=utf-8', json.dumps( dict(err='这些动作只接受 POST (避免链接被预取/刷新时误触发)'), ensure_ascii=False).encode('utf-8') fn = ACTIONS.get(name) if not fn: return 404, 'application/json; charset=utf-8', json.dumps(dict(err=f'未知动作: {name}'), ensure_ascii=False).encode('utf-8') try: payload = json.loads(body.decode('utf-8')) if body else {} except Exception: payload = {} try: code, obj = fn(payload) except Exception as e: code, obj = 500, dict(err=f'{type(e).__name__}: {e}') return code, 'application/json; charset=utf-8', json.dumps(obj, ensure_ascii=False).encode('utf-8') return 404, 'application/json; charset=utf-8', json.dumps(dict(err='not found'), ensure_ascii=False).encode('utf-8') # ─────────────────────────────────────────────────────────── 页面 PAGE = r""" 观澜 · 运维控制台

观澜 · 运维控制台

停/启服务 · 执行重算 · 清除产物 —— 按钮按真实状态启用; 不可用的动作后端也会拒绝。本页由网关(端口 )提供。

服务

重算

默认跳过 SCADA(只换了台账类数据时的常用档)。重算会自己重启服务。

产物

清除后门户首页仍可打开(它是静态交付页), 但工作台会显示"无产物"; 恢复后请点"启动服务"。

最近一次动作

(无)

命令行等价物: guanlan.py stop/serve · scripts/rebuild_all.py · scripts/products_state.py --off/--on(手册 §0/§5b)。

"""