对真实 nixmsg 跑 DEVELOPMENT 第 9 节接入清单(跳过仅 JS 跨域),补 README/示例,并修 Paho/HiveMQ 真机联调死锁与鉴权分类。
186 lines
6.5 KiB
Python
186 lines
6.5 KiB
Python
"""真实 nixmsg 进程启动器:临时目录、127.0.0.1:0、admin init、开注册。"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import json
|
||
import os
|
||
import shutil
|
||
import subprocess
|
||
import tempfile
|
||
import time
|
||
import urllib.error
|
||
import urllib.request
|
||
from dataclasses import dataclass
|
||
from http.cookiejar import CookieJar
|
||
from pathlib import Path
|
||
from typing import Any, Optional
|
||
from urllib.parse import urljoin
|
||
|
||
|
||
def _repo_root() -> Path:
|
||
here = Path(__file__).resolve()
|
||
for p in [here] + list(here.parents):
|
||
if (p / "go.mod").is_file() and (p / "cmd" / "nixmsg").is_dir():
|
||
return p
|
||
raise RuntimeError("找不到仓库根(go.mod)")
|
||
|
||
|
||
def ensure_binary() -> Path:
|
||
env = os.environ.get("NIXMSG_BIN")
|
||
if env:
|
||
p = Path(env)
|
||
if p.is_file():
|
||
return p
|
||
root = _repo_root()
|
||
cache = Path(tempfile.gettempdir()) / "nixmsg-s2-python-bin"
|
||
cache.mkdir(parents=True, exist_ok=True)
|
||
name = "nixmsg.exe" if os.name == "nt" else "nixmsg"
|
||
out = cache / name
|
||
# 若已有且较新则复用;否则编译
|
||
need = True
|
||
if out.is_file():
|
||
need = False
|
||
if need or os.environ.get("NIXMSG_REBUILD") == "1":
|
||
cmd = ["go", "build", "-o", str(out), "./cmd/nixmsg"]
|
||
envp = os.environ.copy()
|
||
envp["CGO_ENABLED"] = "0"
|
||
r = subprocess.run(cmd, cwd=str(root), env=envp, capture_output=True, text=True)
|
||
if r.returncode != 0:
|
||
raise RuntimeError(f"go build 失败:\n{r.stdout}\n{r.stderr}")
|
||
return out
|
||
|
||
|
||
@dataclass
|
||
class AdminHTTP:
|
||
base: str
|
||
opener: urllib.request.OpenerDirector
|
||
|
||
def request(self, method: str, path: str, body: Optional[dict] = None) -> tuple[int, dict[str, Any]]:
|
||
data = None
|
||
headers = {"Accept": "application/json"}
|
||
if body is not None:
|
||
data = json.dumps(body, ensure_ascii=False).encode("utf-8")
|
||
headers["Content-Type"] = "application/json"
|
||
if method.upper() in ("POST", "PUT", "PATCH", "DELETE"):
|
||
headers["X-Nixmsg-Request"] = "1"
|
||
req = urllib.request.Request(urljoin(self.base + "/", path.lstrip("/")), data=data, headers=headers, method=method)
|
||
try:
|
||
with self.opener.open(req, timeout=30) as resp:
|
||
raw = resp.read()
|
||
code = resp.getcode()
|
||
except urllib.error.HTTPError as e:
|
||
raw = e.read()
|
||
code = e.code
|
||
if not raw:
|
||
return code, {}
|
||
return code, json.loads(raw.decode("utf-8"))
|
||
|
||
|
||
class NixMsgServer:
|
||
def __init__(self) -> None:
|
||
self.bin = ensure_binary()
|
||
self.data_dir = Path(tempfile.mkdtemp(prefix="nixmsg-s2-py-"))
|
||
self.config_path = self.data_dir / "config.yaml"
|
||
data_slash = self.data_dir.as_posix()
|
||
self.config_path.write_text(
|
||
f'listen: "127.0.0.1:0"\ndata_dir: "{data_slash}"\n',
|
||
encoding="utf-8",
|
||
)
|
||
self.admin_password = self._admin_init()
|
||
self.proc: Optional[subprocess.Popen] = None
|
||
self.addr = ""
|
||
self.http_base = ""
|
||
self.ws_url = ""
|
||
self.reg_code = "s2py-reg-code"
|
||
|
||
def _admin_init(self) -> str:
|
||
env = os.environ.copy()
|
||
env["NIXMSG_CONFIG"] = str(self.config_path)
|
||
r = subprocess.run(
|
||
[str(self.bin), "admin", "init"],
|
||
env=env,
|
||
capture_output=True,
|
||
text=True,
|
||
)
|
||
if r.returncode != 0:
|
||
raise RuntimeError(f"admin init 失败: {r.stdout}\n{r.stderr}")
|
||
text = (r.stdout or "") + "\n" + (r.stderr or "")
|
||
for line in text.splitlines():
|
||
line = line.strip()
|
||
lower = line.lower()
|
||
if lower.startswith("admin password:"):
|
||
return line.split(":", 1)[1].strip()
|
||
if lower.startswith("password:"):
|
||
return line.split(":", 1)[1].strip()
|
||
raise RuntimeError(f"admin init 未解析到密码:\n{text}")
|
||
|
||
def start(self) -> None:
|
||
env = os.environ.copy()
|
||
env["NIXMSG_CONFIG"] = str(self.config_path)
|
||
self.proc = subprocess.Popen(
|
||
[str(self.bin), "serve"],
|
||
env=env,
|
||
stdout=subprocess.DEVNULL,
|
||
stderr=subprocess.DEVNULL,
|
||
)
|
||
addr_file = self.data_dir / "listen.addr"
|
||
deadline = time.time() + 20
|
||
while time.time() < deadline:
|
||
if addr_file.is_file():
|
||
addr = addr_file.read_text(encoding="utf-8").strip()
|
||
if addr:
|
||
self.addr = addr
|
||
self.http_base = f"http://{addr}"
|
||
self.ws_url = f"ws://{addr}/mqtt"
|
||
break
|
||
if self.proc.poll() is not None:
|
||
raise RuntimeError(f"serve 提前退出 code={self.proc.returncode}")
|
||
time.sleep(0.05)
|
||
else:
|
||
self.stop()
|
||
raise RuntimeError("等待 listen.addr 超时")
|
||
self._enable_registration(self.reg_code)
|
||
|
||
def admin(self) -> AdminHTTP:
|
||
jar = CookieJar()
|
||
opener = urllib.request.build_opener(urllib.request.HTTPCookieProcessor(jar))
|
||
admin = AdminHTTP(self.http_base, opener)
|
||
code, body = admin.request(
|
||
"POST",
|
||
"/api/admin/login",
|
||
{"username": "admin", "password": self.admin_password},
|
||
)
|
||
if code != 200 or not body.get("ok"):
|
||
raise RuntimeError(f"admin login 失败: {code} {body}")
|
||
return admin
|
||
|
||
def _enable_registration(self, code: str, enabled: bool = True) -> None:
|
||
admin = self.admin()
|
||
status, body = admin.request(
|
||
"PUT",
|
||
"/api/admin/registration",
|
||
{"enabled": enabled, "code": code},
|
||
)
|
||
if status != 200 or not body.get("ok"):
|
||
raise RuntimeError(f"开启注册失败: {status} {body}")
|
||
|
||
def set_registration(self, *, enabled: bool, code: Optional[str] = None) -> None:
|
||
admin = self.admin()
|
||
payload: dict[str, Any] = {"enabled": enabled}
|
||
if code is not None:
|
||
payload["code"] = code
|
||
status, body = admin.request("PUT", "/api/admin/registration", payload)
|
||
if status != 200 or not body.get("ok"):
|
||
raise RuntimeError(f"改注册设置失败: {status} {body}")
|
||
|
||
def stop(self) -> None:
|
||
if self.proc and self.proc.poll() is None:
|
||
self.proc.kill()
|
||
try:
|
||
self.proc.wait(timeout=5)
|
||
except Exception:
|
||
pass
|
||
self.proc = None
|
||
if self.data_dir.exists():
|
||
shutil.rmtree(self.data_dir, ignore_errors=True)
|