feat: native WireView Pro II monitoring (detection, serial protocol, live tab)
Standalone support for the ThermalGrizzly WireView Pro II without depending on the external wireview_reporter exporter: - wireview.py: serial protocol (STX/ETX + 16-bit CRC16-CCITT) with vendor-data product identification, config version, UID, build string, screen layout, and temperature/power/current sensor reads; hwmon fallback with per-channel index resolution; udev-based detection (vendor 0x2560 / product 0x0101) with fallback port probing; serial read timeout and watchdog reconnect - server.py: startup detection (root only), 1 Hz poller gated on subscribed clients, WS 'wireview' channel, GET /api/wireview, rejected-product memoization, stale-read race guard - frontend: WireView tab (visible when a device is detected) with live temperature/power/current cards, sparkline history, and device info; nullable temp channels - tests: 76 tests covering parser, CRC, fault classification, hwmon resolution, JSON safety, and the serial transport against a pty-based fake device Verified live: tab appears with the device connected, disappears when unplugged, reconnects on re-plug, no serial traffic when idle.
This commit is contained in:
1 parent
cc102f26c1
commit
526b71b45f
10 files changed
+1792
-1
No files matched your search
@@ -72,6 +72,7 @@ from .profiles.native import (
|
||||
rename_profile,
|
||||
save_profile,
|
||||
)
|
||||
from .wireview import create_device, find_wireview_ports
|
||||
from .safety import check_negative_freq_warnings, validate_write
|
||||
|
||||
log = logging.getLogger("nvcurve.server")
|
||||
@@ -103,6 +104,15 @@ def _open_browser_as_user(url: str) -> None:
|
||||
_state: dict[str, Any] = {
|
||||
"gpus": {}, # dict[int, dict] mapping gpu_index -> gpu state
|
||||
"config": default_config,
|
||||
"wireview": {
|
||||
"clients": set(), # connected /ws/wireview clients
|
||||
"device": None, # WireViewSerialDevice | WireViewHwmonDevice | None
|
||||
"info": None, # device identity (static per connection)
|
||||
"last_sample": None, # most recent sensor sample
|
||||
"connected": False, # last read succeeded
|
||||
"failures": 0, # consecutive failed reads
|
||||
"rejected_ports": set(), # ports reporting an unsupported product
|
||||
},
|
||||
}
|
||||
|
||||
|
||||
@@ -258,6 +268,126 @@ async def _fan_poller(gpu_index: int) -> None:
|
||||
await asyncio.sleep(2.0)
|
||||
|
||||
|
||||
# ── WireView Pro II (Thermal Grizzly) ─────────────────────────────────────────
|
||||
# Consecutive failed reads after which a connected device is given up. Same
|
||||
# rule as the exporter: one corrupt frame or a slow read never trips it, and
|
||||
# an unplug is caught at once by the port-node check.
|
||||
WIREVIEW_MAX_FAILED_READS = 6
|
||||
WIREVIEW_WATCHDOG_INTERVAL_S = 5.0
|
||||
|
||||
|
||||
async def _wireview_disconnect() -> None:
|
||||
"""Drop the connected WireView device and tell clients it is gone."""
|
||||
wv = _state["wireview"]
|
||||
device = wv["device"]
|
||||
if device is None:
|
||||
return
|
||||
wv["device"] = None
|
||||
wv["connected"] = False
|
||||
wv["info"] = None
|
||||
wv["last_sample"] = None
|
||||
wv["failures"] = 0
|
||||
await _run(device.close)
|
||||
log.info("WireView disconnected")
|
||||
await _broadcast(wv["clients"], {"type": "unavailable"})
|
||||
|
||||
|
||||
async def _wireview_connect() -> None:
|
||||
"""Detect and connect a WireView Pro II. No-op when none is attached or
|
||||
one is already connected."""
|
||||
wv = _state["wireview"]
|
||||
if wv["device"] is not None:
|
||||
return
|
||||
# Forget rejected ports that disappeared, so a re-plug gets a fresh probe.
|
||||
if wv["rejected_ports"]:
|
||||
present = set(await _run(find_wireview_ports))
|
||||
wv["rejected_ports"] &= present
|
||||
device = await _run(create_device)
|
||||
if device is None:
|
||||
return
|
||||
# A port that reported an unsupported product is not re-probed (and
|
||||
# re-logged) on every watchdog tick.
|
||||
if device.port in wv["rejected_ports"]:
|
||||
return
|
||||
ok = await _run(device.connect)
|
||||
if not ok:
|
||||
await _run(device.close)
|
||||
if device.rejected:
|
||||
wv["rejected_ports"].add(device.port)
|
||||
return
|
||||
wv["device"] = device
|
||||
wv["failures"] = 0
|
||||
wv["connected"] = True
|
||||
wv["info"] = await _run(device.info)
|
||||
# Push the first sample right away so clients don't wait for the next tick.
|
||||
sample = await _run(device.read_sample)
|
||||
if sample is not None:
|
||||
wv["last_sample"] = sample
|
||||
await _broadcast(
|
||||
wv["clients"], {"type": "sample", "info": wv["info"], "sample": sample}
|
||||
)
|
||||
else:
|
||||
wv["connected"] = False
|
||||
|
||||
|
||||
async def _wireview_poller() -> None:
|
||||
"""Read the WireView at poll_interval_s and push to connected WS clients.
|
||||
|
||||
Skips reads while no client is subscribed (like the monitor poller), so
|
||||
idle polling never contends for the port with other tools (official GUI,
|
||||
wireviewd)."""
|
||||
cfg: Config = _state["config"]
|
||||
while True:
|
||||
try:
|
||||
wv = _state["wireview"]
|
||||
device = wv["device"]
|
||||
if device is not None and wv["clients"]:
|
||||
sample = await _run(device.read_sample)
|
||||
if wv["device"] is not device:
|
||||
# The device was disconnected (unplug) while the read was
|
||||
# in flight — drop the stale result instead of
|
||||
# resurrecting state.
|
||||
pass
|
||||
elif sample is not None:
|
||||
wv["failures"] = 0
|
||||
wv["connected"] = True
|
||||
wv["last_sample"] = sample
|
||||
await _broadcast(
|
||||
wv["clients"],
|
||||
{"type": "sample", "info": wv["info"], "sample": sample},
|
||||
)
|
||||
else:
|
||||
wv["failures"] += 1
|
||||
if (
|
||||
not await _run(device.node_exists)
|
||||
or wv["failures"] >= WIREVIEW_MAX_FAILED_READS
|
||||
):
|
||||
await _wireview_disconnect()
|
||||
except asyncio.CancelledError:
|
||||
return
|
||||
except Exception as exc:
|
||||
log.warning("WireView poller error: %s", exc)
|
||||
await asyncio.sleep(cfg.poll_interval_s)
|
||||
|
||||
|
||||
async def _wireview_watchdog() -> None:
|
||||
"""Hot-plug detection: connect when a WireView appears, disconnect when
|
||||
its device node disappears."""
|
||||
while True:
|
||||
await asyncio.sleep(WIREVIEW_WATCHDOG_INTERVAL_S)
|
||||
try:
|
||||
wv = _state["wireview"]
|
||||
device = wv["device"]
|
||||
if device is None:
|
||||
await _wireview_connect()
|
||||
elif not await _run(device.node_exists):
|
||||
await _wireview_disconnect()
|
||||
except asyncio.CancelledError:
|
||||
return
|
||||
except Exception as exc:
|
||||
log.warning("WireView watchdog error: %s", exc)
|
||||
|
||||
|
||||
async def _activate_fan_curve(
|
||||
gpu_index: int, curve: list, fans: list[int] | None = None
|
||||
) -> None:
|
||||
@@ -380,6 +510,16 @@ async def lifespan(app: FastAPI):
|
||||
except Exception as exc:
|
||||
log.error("Failed to initialize GPU %d: %s", idx, exc)
|
||||
|
||||
# ── WireView Pro II detection ─────────────────────────────────────────────
|
||||
# Independent of GPU discovery: the tab appears whenever the connector
|
||||
# monitor is attached, and the watchdog handles hot-plug afterwards.
|
||||
try:
|
||||
await _wireview_connect()
|
||||
except Exception as exc:
|
||||
log.warning("WireView initial detection failed: %s", exc)
|
||||
poller_tasks.append(asyncio.create_task(_wireview_poller()))
|
||||
poller_tasks.append(asyncio.create_task(_wireview_watchdog()))
|
||||
|
||||
# ── Backward Compatibility Bridge ──────────────────────────────────────────
|
||||
# NOTE: This auto-load path is for users running the server directly (e.g.
|
||||
# via an old systemd unit file that lacks the new daemon mode).
|
||||
@@ -465,6 +605,14 @@ async def lifespan(app: FastAPI):
|
||||
with suppress(asyncio.CancelledError):
|
||||
await task
|
||||
|
||||
# Release the WireView port (if connected) so other tools can use it.
|
||||
wv = _state["wireview"]
|
||||
if wv["device"] is not None:
|
||||
with suppress(Exception):
|
||||
await _run(wv["device"].close)
|
||||
wv["device"] = None
|
||||
wv["connected"] = False
|
||||
|
||||
for gpu_index, g_state in _state["gpus"].items():
|
||||
if g_state.get("fan_poller_task"):
|
||||
g_state["fan_poller_task"].cancel()
|
||||
@@ -1442,6 +1590,25 @@ async def api_fans_speed(req: FanSpeedRequest, gpu_index: int = 0):
|
||||
return {"ok": True}
|
||||
|
||||
|
||||
# ── WireView Pro II (Thermal Grizzly) ─────────────────────────────────────────
|
||||
|
||||
|
||||
@app.get("/api/wireview")
|
||||
async def api_wireview():
|
||||
"""WireView Pro II availability and the most recent sensor sample.
|
||||
|
||||
available: a device is connected (the UI shows the WireView tab).
|
||||
sample: None until the first successful read.
|
||||
"""
|
||||
wv = _state["wireview"]
|
||||
return {
|
||||
"available": wv["device"] is not None,
|
||||
"connected": wv["connected"],
|
||||
"info": wv["info"],
|
||||
"sample": wv["last_sample"],
|
||||
}
|
||||
|
||||
|
||||
# ── Write endpoints ────────────────────────────────────────────────────────────
|
||||
|
||||
|
||||
@@ -1832,6 +1999,51 @@ async def ws_curve(ws: WebSocket):
|
||||
g_state["curve_clients"].discard(ws)
|
||||
|
||||
|
||||
@app.websocket("/ws/wireview")
|
||||
async def ws_wireview(ws: WebSocket):
|
||||
"""Stream WireView Pro II sensor samples at poll_interval_s.
|
||||
|
||||
Messages:
|
||||
{"type": "unavailable"} — no device connected
|
||||
{"type": "sample", "info": {...}, "sample": {...}} — new reading
|
||||
"""
|
||||
if not _ws_authenticated(ws):
|
||||
await ws.close(code=1008)
|
||||
return
|
||||
await ws.accept()
|
||||
try:
|
||||
data = await ws.receive_json()
|
||||
if data.get("action") != "subscribe":
|
||||
await ws.close()
|
||||
return
|
||||
except WebSocketDisconnect:
|
||||
return
|
||||
except Exception:
|
||||
await ws.close()
|
||||
return
|
||||
|
||||
wv = _state["wireview"]
|
||||
wv["clients"].add(ws)
|
||||
try:
|
||||
# Send the current state immediately so the client does not have to
|
||||
# wait for the next poll tick.
|
||||
if wv["device"] is not None and wv["last_sample"] is not None:
|
||||
await ws.send_json(
|
||||
{"type": "sample", "info": wv["info"], "sample": wv["last_sample"]}
|
||||
)
|
||||
else:
|
||||
await ws.send_json({"type": "unavailable"})
|
||||
|
||||
while True:
|
||||
await ws.receive_text()
|
||||
except WebSocketDisconnect:
|
||||
pass
|
||||
except Exception:
|
||||
log.debug("wireview ws client error", exc_info=True)
|
||||
finally:
|
||||
wv["clients"].discard(ws)
|
||||
|
||||
|
||||
# ── Frontend SPA ──────────────────────────────────────────────────────────────
|
||||
|
||||
|
||||
|
||||
@@ -0,0 +1,659 @@
|
||||
"""Native reader for the Thermal Grizzly WireView Pro II.
|
||||
|
||||
Talks to the 12 VHPWR connector monitor directly over its USB CDC/ACM
|
||||
serial port (STM32, VID 0483 / PID 5740) — no exporter, no GUI, no kernel
|
||||
module required. When the wireview-hwmon kernel module is loaded, the
|
||||
sysfs node is used instead (the wireviewd daemon owns the port in that
|
||||
case, so direct serial would corrupt frames).
|
||||
|
||||
Protocol notes (matches the firmware's DEVICE_STR_LEN=32 layout):
|
||||
* No framing/CRC — a desynced read corrupts arbitrary fields for one
|
||||
poll. Real frames always carry zero padding bytes and a fan duty
|
||||
<= 100; anything else is discarded and the next poll realigns.
|
||||
* The firmware occasionally stops answering the RTS welcome handshake
|
||||
(observed after USB state changes) while still answering every data
|
||||
command, so identification falls back to the vendor-data reply.
|
||||
"""
|
||||
|
||||
import logging
|
||||
import os
|
||||
import struct
|
||||
import time
|
||||
|
||||
import serial
|
||||
|
||||
log = logging.getLogger("nvcurve.wireview")
|
||||
|
||||
# ── Device identification ─────────────────────────────────────────────────────
|
||||
|
||||
USB_VENDOR_ID = "0483" # STMicroelectronics (CDC/ACM)
|
||||
USB_PRODUCT_ID = "5740" # WireView Pro II normal mode
|
||||
WELCOME_MESSAGE = "Thermal Grizzly WireView Pro II"
|
||||
MAX_WELCOME_LENGTH = 64
|
||||
|
||||
# Vendor/product ids reported by CMD_READ_VENDOR_DATA (not the USB ids).
|
||||
VENDOR_ID_THERMAL_GRIZZLY = 0xEF
|
||||
PRODUCT_ID_PRO2 = 0x05
|
||||
PRODUCT_ID_PRO2_NOCTUA = 0x06
|
||||
|
||||
DEVICE_NAMES = {
|
||||
PRODUCT_ID_PRO2: "WireView Pro II",
|
||||
PRODUCT_ID_PRO2_NOCTUA: "WireView Pro II Noctua Edition",
|
||||
}
|
||||
|
||||
BAUD_RATE = 115200
|
||||
READ_TIMEOUT_S = 1.0
|
||||
|
||||
# ── Serial protocol commands ──────────────────────────────────────────────────
|
||||
|
||||
CMD_READ_VENDOR_DATA = 0x01
|
||||
CMD_READ_UID = 0x02
|
||||
CMD_READ_SENSOR_VALUES = 0x04
|
||||
CMD_READ_CONFIG = 0x05
|
||||
CMD_SCREEN_CHANGE = 0x0C
|
||||
CMD_READ_BUILD_INFO = 0x0D
|
||||
|
||||
SCREEN_RESUME_UPDATES = 0xF1
|
||||
|
||||
# ── Wire layout ───────────────────────────────────────────────────────────────
|
||||
# SensorStruct (100 bytes, little-endian, pack=4):
|
||||
# 4x int16 temperatures (0.1 °C): in, out, ext1, ext2
|
||||
# uint16 Vdd (mV)
|
||||
# uint8 fan duty (%)
|
||||
# pad
|
||||
# 6x { int16 voltage (mV), pad, uint32 current (mA), uint32 power (mW) }
|
||||
# uint32 total power (mW)
|
||||
# uint32 total current (mA)
|
||||
# uint16 avg voltage (mV)
|
||||
# uint8 PSU capability (0=600W, 1=450W, 2=300W, 3=150W)
|
||||
# pad
|
||||
# uint16 fault status mask
|
||||
# uint16 fault log mask
|
||||
SENSOR_STRUCT = struct.Struct(
|
||||
"<4hHBx" + "".join("hxxII" for _ in range(6)) + "IIHBxHH"
|
||||
)
|
||||
SENSOR_STRUCT_SIZE = SENSOR_STRUCT.size # 100
|
||||
|
||||
# BuildStruct: VendorData(3) + ProductName(32) + BuildInfo(32) + NameLength(1)
|
||||
BUILD_STRUCT_SIZE = 3 + 32 + 32 + 1
|
||||
BUILD_INFO_OFFSET = 3 + 32 # 35
|
||||
|
||||
PSU_CAPABILITY_W = {0: 600, 1: 450, 2: 300, 3: 150}
|
||||
|
||||
# Fault bitmask (both the active status and the latched log use these bits).
|
||||
FAULT_BITS = {
|
||||
0: "Chip over-temperature",
|
||||
1: "Sensor over-temperature",
|
||||
2: "Over-current (OCP)",
|
||||
3: "Wire over-current",
|
||||
4: "Over-power (OPP)",
|
||||
5: "Current imbalance",
|
||||
}
|
||||
|
||||
|
||||
def decode_faults(mask: int) -> list[str]:
|
||||
"""Human-readable names of the active fault bits in a status/log mask."""
|
||||
return [
|
||||
FAULT_BITS[bit]
|
||||
for bit in sorted(FAULT_BITS)
|
||||
if mask & (1 << bit)
|
||||
]
|
||||
|
||||
|
||||
def is_supported_product(vendor_id: int, product_id: int) -> bool:
|
||||
"""True for the products the Pro II protocol serves (5 and 6)."""
|
||||
return (
|
||||
vendor_id == VENDOR_ID_THERMAL_GRIZZLY
|
||||
and product_id in (PRODUCT_ID_PRO2, PRODUCT_ID_PRO2_NOCTUA)
|
||||
)
|
||||
|
||||
|
||||
def parse_sensor_struct(buf: bytes) -> dict:
|
||||
"""Decode a 100-byte sensor frame into a JSON-serializable sample.
|
||||
|
||||
Totals are computed from the per-pin readings (voltage * current),
|
||||
matching the exporter's output; the device's own total fields are
|
||||
not used.
|
||||
"""
|
||||
fields = SENSOR_STRUCT.unpack(buf)
|
||||
ts_in, ts_out, ts_ext1, ts_ext2, _vdd, fan_duty = fields[:6]
|
||||
pin_fields = fields[6:24]
|
||||
(
|
||||
_total_power,
|
||||
_total_current,
|
||||
_avg_voltage,
|
||||
psu_cap,
|
||||
fault_status,
|
||||
fault_log,
|
||||
) = fields[24:]
|
||||
|
||||
pins = []
|
||||
power_total = 0.0
|
||||
current_total = 0.0
|
||||
for i in range(6):
|
||||
voltage_v = pin_fields[i * 3] / 1000.0
|
||||
current_a = pin_fields[i * 3 + 1] / 1000.0
|
||||
power_w = pin_fields[i * 3 + 2] / 1000.0
|
||||
pins.append(
|
||||
{
|
||||
"voltage_v": round(voltage_v, 3),
|
||||
"current_a": round(current_a, 3),
|
||||
"power_w": round(power_w, 3),
|
||||
}
|
||||
)
|
||||
power_total += voltage_v * current_a
|
||||
current_total += current_a
|
||||
|
||||
return {
|
||||
"timestamp": time.time(),
|
||||
"power_total_w": round(power_total, 3),
|
||||
"current_total_a": round(current_total, 3),
|
||||
"voltage_avg_v": (
|
||||
round(power_total / current_total, 3) if current_total > 0 else 0.0
|
||||
),
|
||||
"pins": pins,
|
||||
"temp_in_c": ts_in / 10.0,
|
||||
"temp_out_c": ts_out / 10.0,
|
||||
"temp_ext1_c": ts_ext1 / 10.0,
|
||||
"temp_ext2_c": ts_ext2 / 10.0,
|
||||
"fan_duty_pct": fan_duty,
|
||||
"fault_status": fault_status,
|
||||
"fault_log": fault_log,
|
||||
"psu_capability_w": PSU_CAPABILITY_W.get(psu_cap, 0),
|
||||
}
|
||||
|
||||
|
||||
def sensor_frame_is_corrupt(buf: bytes) -> bool:
|
||||
"""Corruption check for a sensor frame: real frames always carry zero
|
||||
padding bytes and a fan duty <= 100. Wire layout: fan duty at offset
|
||||
10, pad1 at 11, pad2 five bytes from the end (before the two 16-bit
|
||||
fault masks)."""
|
||||
if len(buf) < SENSOR_STRUCT_SIZE:
|
||||
return True
|
||||
return buf[10] > 100 or buf[11] != 0 or buf[SENSOR_STRUCT_SIZE - 5] != 0
|
||||
|
||||
|
||||
# ── Device discovery ──────────────────────────────────────────────────────────
|
||||
|
||||
|
||||
def _sysfs_matches_wireview(tty_sysfs_dir: str) -> bool:
|
||||
"""Walk up from a resolved tty sysfs path to the USB device and check
|
||||
its idVendor/idProduct."""
|
||||
path = tty_sysfs_dir
|
||||
while path and path != "/":
|
||||
vid_file = os.path.join(path, "idVendor")
|
||||
pid_file = os.path.join(path, "idProduct")
|
||||
if os.path.isfile(vid_file) and os.path.isfile(pid_file):
|
||||
try:
|
||||
with open(vid_file) as f:
|
||||
vid = f.read().strip().lower()
|
||||
with open(pid_file) as f:
|
||||
pid = f.read().strip().lower()
|
||||
except OSError:
|
||||
return False
|
||||
return vid == USB_VENDOR_ID and pid == USB_PRODUCT_ID
|
||||
path = os.path.dirname(path)
|
||||
return False
|
||||
|
||||
|
||||
def find_wireview_ports() -> list[str]:
|
||||
"""Find /dev nodes of connected WireView Pro II devices.
|
||||
|
||||
Checks the stable /dev/wireview-pro2 symlink (created by the
|
||||
99-wireview.rules udev rule) and falls back to a sysfs scan of all
|
||||
ttyACM* ports matched by USB VID/PID.
|
||||
"""
|
||||
ports: set[str] = set()
|
||||
|
||||
link = "/dev/wireview-pro2"
|
||||
if os.path.islink(link) or os.path.exists(link):
|
||||
try:
|
||||
target = os.path.realpath(link)
|
||||
if os.path.exists(target):
|
||||
ports.add(target)
|
||||
except OSError:
|
||||
pass
|
||||
|
||||
sys_class = "/sys/class/tty"
|
||||
if os.path.isdir(sys_class):
|
||||
try:
|
||||
entries = os.listdir(sys_class)
|
||||
except OSError:
|
||||
entries = []
|
||||
for entry in entries:
|
||||
if not entry.startswith("ttyACM"):
|
||||
continue
|
||||
tty_dir = os.path.join(sys_class, entry)
|
||||
try:
|
||||
resolved = os.path.realpath(tty_dir)
|
||||
except OSError:
|
||||
continue
|
||||
if _sysfs_matches_wireview(resolved):
|
||||
ports.add(f"/dev/{entry}")
|
||||
|
||||
return sorted(ports)
|
||||
|
||||
|
||||
def find_hwmon_path() -> str | None:
|
||||
"""Find the wireview-hwmon sysfs node, if the kernel module is loaded."""
|
||||
base = "/sys/class/hwmon"
|
||||
if not os.path.isdir(base):
|
||||
return None
|
||||
try:
|
||||
entries = os.listdir(base)
|
||||
except OSError:
|
||||
return None
|
||||
for entry in entries:
|
||||
name_path = os.path.join(base, entry, "name")
|
||||
try:
|
||||
with open(name_path) as f:
|
||||
if f.read().strip().lower() == "wireview":
|
||||
return os.path.join(base, entry)
|
||||
except OSError:
|
||||
continue
|
||||
return None
|
||||
|
||||
|
||||
# ── Serial transport ──────────────────────────────────────────────────────────
|
||||
|
||||
|
||||
class WireViewSerialDevice:
|
||||
"""Direct serial access to a WireView Pro II.
|
||||
|
||||
The port is opened and closed per transaction (open → flush → write →
|
||||
read → close), matching the proven behavior of the exporter: it keeps
|
||||
the port unheld between polls so other tools (official GUI, wireviewd)
|
||||
can share the device, and a fresh open realigns a desynced stream.
|
||||
"""
|
||||
|
||||
def __init__(self, port: str, baud: int = BAUD_RATE) -> None:
|
||||
self._port = port
|
||||
self._baud = baud
|
||||
self._connected = False
|
||||
self._rejected = False
|
||||
self._vendor_id = 0
|
||||
self._product_id = 0
|
||||
self._firmware_version = ""
|
||||
self._uid = ""
|
||||
self._build = ""
|
||||
self._config_version = -1
|
||||
|
||||
# ── Identity ──
|
||||
|
||||
@property
|
||||
def connected(self) -> bool:
|
||||
return self._connected
|
||||
|
||||
@property
|
||||
def rejected(self) -> bool:
|
||||
"""True when the device reported an unsupported product id. Callers
|
||||
can memoize this so the port is not re-probed (and re-logged) on
|
||||
every watchdog tick."""
|
||||
return self._rejected
|
||||
|
||||
@property
|
||||
def transport(self) -> str:
|
||||
return "serial"
|
||||
|
||||
@property
|
||||
def port(self) -> str:
|
||||
return self._port
|
||||
|
||||
def info(self) -> dict:
|
||||
return {
|
||||
"device_name": DEVICE_NAMES.get(
|
||||
self._product_id, "WireView Pro II"
|
||||
),
|
||||
"hw_rev": f"{self._vendor_id:02X}{self._product_id:02X}",
|
||||
"firmware_version": self._firmware_version,
|
||||
"uid": self._uid,
|
||||
"build": self._build,
|
||||
"transport": self.transport,
|
||||
"port": self._port,
|
||||
}
|
||||
|
||||
def node_exists(self) -> bool:
|
||||
"""Whether the serial device node still exists (unplug check)."""
|
||||
return os.path.exists(self._port)
|
||||
|
||||
# ── Connection ──
|
||||
|
||||
def connect(self) -> bool:
|
||||
"""Identify the device and prepare it for sensor reads.
|
||||
|
||||
The welcome handshake (RTS edge) is the primary identification,
|
||||
but the firmware occasionally stops answering it while still
|
||||
answering every command — a supported vendor-data reply is
|
||||
equally conclusive, so accept either.
|
||||
"""
|
||||
if self._connected:
|
||||
return True
|
||||
|
||||
# The welcome handshake (RTS edge) is the primary identification, but
|
||||
# the firmware occasionally stops answering it while still answering
|
||||
# every command — the supported vendor-data reply below is equally
|
||||
# conclusive, so the welcome is read for logging only.
|
||||
welcome = self._read_welcome()
|
||||
if welcome and welcome != WELCOME_MESSAGE:
|
||||
log.debug(
|
||||
"WireView: unexpected welcome string %r on %s",
|
||||
welcome,
|
||||
self._port,
|
||||
)
|
||||
|
||||
vd = self._transaction(bytes([CMD_READ_VENDOR_DATA]), 3)
|
||||
if vd is None or len(vd) < 3:
|
||||
return False
|
||||
vendor, product, fw = vd[0], vd[1], vd[2]
|
||||
if not is_supported_product(vendor, product):
|
||||
self._rejected = True
|
||||
log.info(
|
||||
"WireView: unsupported product %02X%02X on %s, skipped",
|
||||
vendor,
|
||||
product,
|
||||
self._port,
|
||||
)
|
||||
return False
|
||||
|
||||
self._vendor_id = vendor
|
||||
self._product_id = product
|
||||
self._firmware_version = str(fw)
|
||||
|
||||
cfg = self._transaction(bytes([CMD_READ_CONFIG]), 4)
|
||||
if cfg is None or len(cfg) < 3:
|
||||
return False
|
||||
self._config_version = cfg[2]
|
||||
|
||||
uid = self._transaction(bytes([CMD_READ_UID]), 12)
|
||||
if uid is not None and len(uid) == 12:
|
||||
self._uid = uid.hex().upper()
|
||||
|
||||
# Enable display updates just in case.
|
||||
self._transaction(
|
||||
bytes([CMD_SCREEN_CHANGE, SCREEN_RESUME_UPDATES]), 0
|
||||
)
|
||||
|
||||
build = self._transaction(bytes([CMD_READ_BUILD_INFO]), BUILD_STRUCT_SIZE)
|
||||
if build is not None and len(build) >= BUILD_INFO_OFFSET + 1:
|
||||
self._build = (
|
||||
build[BUILD_INFO_OFFSET : BUILD_INFO_OFFSET + 32]
|
||||
.split(b"\x00")[0]
|
||||
.decode("ascii", errors="replace")
|
||||
)
|
||||
|
||||
self._connected = True
|
||||
log.info(
|
||||
"WireView connected on %s (%s, fw %s)",
|
||||
self._port,
|
||||
self.info()["hw_rev"],
|
||||
self._firmware_version,
|
||||
)
|
||||
return True
|
||||
|
||||
def close(self) -> None:
|
||||
self._connected = False
|
||||
|
||||
# ── Sensor reads ──
|
||||
|
||||
def read_sample(self) -> dict | None:
|
||||
"""Read one sensor sample, or None when the device is unresponsive
|
||||
or the frame is corrupt."""
|
||||
if not self._connected:
|
||||
return None
|
||||
buf = self._transaction(bytes([CMD_READ_SENSOR_VALUES]), SENSOR_STRUCT_SIZE)
|
||||
if buf is None or sensor_frame_is_corrupt(buf):
|
||||
return None
|
||||
return parse_sensor_struct(buf)
|
||||
|
||||
# ── Transport ──
|
||||
|
||||
def _read_welcome(self) -> str | None:
|
||||
"""Assert RTS and read the NUL-terminated welcome string the device
|
||||
answers with. Null when nothing (or no terminator) arrives in time."""
|
||||
try:
|
||||
ser = serial.Serial(self._port, self._baud, timeout=READ_TIMEOUT_S)
|
||||
except (serial.SerialException, OSError):
|
||||
return None
|
||||
try:
|
||||
ser.reset_input_buffer()
|
||||
ser.rts = False
|
||||
time.sleep(0.01)
|
||||
ser.rts = True
|
||||
time.sleep(0.01)
|
||||
buf = bytearray()
|
||||
deadline = time.monotonic() + READ_TIMEOUT_S
|
||||
while len(buf) < MAX_WELCOME_LENGTH:
|
||||
remaining = deadline - time.monotonic()
|
||||
if remaining <= 0:
|
||||
break
|
||||
ser.timeout = min(READ_TIMEOUT_S, remaining)
|
||||
chunk = ser.read(MAX_WELCOME_LENGTH - len(buf))
|
||||
if not chunk:
|
||||
break
|
||||
buf.extend(chunk)
|
||||
if b"\x00" in chunk:
|
||||
break
|
||||
time.sleep(0.01)
|
||||
ser.rts = False
|
||||
nul = buf.find(b"\x00")
|
||||
if nul >= 0:
|
||||
return bytes(buf[:nul]).decode("ascii", errors="replace")
|
||||
return None
|
||||
except (serial.SerialException, OSError):
|
||||
return None
|
||||
finally:
|
||||
ser.close()
|
||||
|
||||
def _transaction(self, cmd: bytes, response_size: int) -> bytes | None:
|
||||
"""Open the port, send cmd, read exactly response_size bytes (one
|
||||
second budget), close the port. None when the port is unavailable
|
||||
or the reply is incomplete."""
|
||||
try:
|
||||
ser = serial.Serial(self._port, self._baud, timeout=READ_TIMEOUT_S)
|
||||
except (serial.SerialException, OSError):
|
||||
return None
|
||||
try:
|
||||
ser.reset_input_buffer()
|
||||
if cmd:
|
||||
ser.write(cmd)
|
||||
if response_size == 0:
|
||||
return b""
|
||||
return self._read_exact(ser, response_size)
|
||||
except (serial.SerialException, OSError):
|
||||
return None
|
||||
finally:
|
||||
ser.close()
|
||||
|
||||
@staticmethod
|
||||
def _read_exact(ser: serial.Serial, size: int) -> bytes | None:
|
||||
"""Read exactly size bytes within one second, or None."""
|
||||
buf = bytearray()
|
||||
deadline = time.monotonic() + READ_TIMEOUT_S
|
||||
while len(buf) < size:
|
||||
remaining = deadline - time.monotonic()
|
||||
if remaining <= 0:
|
||||
return None
|
||||
ser.timeout = min(READ_TIMEOUT_S, remaining)
|
||||
chunk = ser.read(size - len(buf))
|
||||
if not chunk:
|
||||
return None
|
||||
buf.extend(chunk)
|
||||
return bytes(buf)
|
||||
|
||||
|
||||
# ── hwmon (sysfs) transport ───────────────────────────────────────────────────
|
||||
|
||||
|
||||
class WireViewHwmonDevice:
|
||||
"""Reads the wireview-hwmon sysfs node (kernel module + wireviewd).
|
||||
|
||||
Used when the module is loaded: the daemon owns the serial port in
|
||||
that case, so direct serial would corrupt frames.
|
||||
"""
|
||||
|
||||
def __init__(self, hwmon_path: str) -> None:
|
||||
self._path = hwmon_path
|
||||
self._connected = False
|
||||
|
||||
@property
|
||||
def connected(self) -> bool:
|
||||
return self._connected
|
||||
|
||||
@property
|
||||
def rejected(self) -> bool:
|
||||
return False
|
||||
|
||||
@property
|
||||
def transport(self) -> str:
|
||||
return "hwmon"
|
||||
|
||||
@property
|
||||
def port(self) -> str:
|
||||
return self._path
|
||||
|
||||
def info(self) -> dict:
|
||||
return {
|
||||
"device_name": "WireView Pro II",
|
||||
"hw_rev": "",
|
||||
"firmware_version": "",
|
||||
"uid": "",
|
||||
"build": "",
|
||||
"transport": self.transport,
|
||||
"port": self._path,
|
||||
}
|
||||
|
||||
def node_exists(self) -> bool:
|
||||
return os.path.isdir(self._path)
|
||||
|
||||
def connect(self) -> bool:
|
||||
if self._connected:
|
||||
return True
|
||||
name_path = os.path.join(self._path, "name")
|
||||
try:
|
||||
with open(name_path) as f:
|
||||
if f.read().strip().lower() != "wireview":
|
||||
return False
|
||||
# Probe that the node actually serves data.
|
||||
with open(os.path.join(self._path, "in0_input")) as f:
|
||||
f.read().strip()
|
||||
except OSError:
|
||||
return False
|
||||
self._connected = True
|
||||
log.info("WireView connected via hwmon (%s)", self._path)
|
||||
return True
|
||||
|
||||
def close(self) -> None:
|
||||
self._connected = False
|
||||
|
||||
def read_sample(self) -> dict | None:
|
||||
if not self._connected:
|
||||
return None
|
||||
try:
|
||||
pin_voltage = [
|
||||
self._read_int(f"in{i}_input") / 1000.0 for i in range(6)
|
||||
]
|
||||
pin_current = [
|
||||
self._read_int(f"curr{i + 1}_input") / 1000.0 for i in range(6)
|
||||
]
|
||||
temp_in = self._read_temp("temp1_input")
|
||||
temp_out = self._read_temp("temp2_input")
|
||||
temp_ext1 = self._read_temp("temp3_input")
|
||||
temp_ext2 = self._read_temp("temp4_input")
|
||||
|
||||
fault_status = self._read_int_or("fault_status_raw")
|
||||
if fault_status is None:
|
||||
fault_status = 0xFFFF if self._read_int("intrusion0_alarm") else 0
|
||||
fault_log = self._read_int_or("fault_log_raw")
|
||||
if fault_log is None:
|
||||
fault_log = 0xFFFF if self._read_int("intrusion1_alarm") else 0
|
||||
|
||||
psu_cap_uw = self._read_int_or("power1_cap")
|
||||
if psu_cap_uw is not None:
|
||||
psu_capability = int(round(psu_cap_uw / 1_000_000.0))
|
||||
else:
|
||||
psu_cap = self._read_int_or("psu_cap")
|
||||
psu_capability = PSU_CAPABILITY_W.get(psu_cap or 0, 0)
|
||||
|
||||
pwm = self._read_int_or("pwm1")
|
||||
if pwm is not None:
|
||||
fan_duty = int(round(min(255, max(0, pwm)) * 100 / 255.0))
|
||||
else:
|
||||
fan_duty = self._read_int("fan1_input")
|
||||
|
||||
power_total = sum(
|
||||
v * i
|
||||
for v, i in zip(pin_voltage, pin_current, strict=True)
|
||||
)
|
||||
current_total = sum(pin_current)
|
||||
|
||||
return {
|
||||
"timestamp": time.time(),
|
||||
"power_total_w": round(power_total, 3),
|
||||
"current_total_a": round(current_total, 3),
|
||||
"voltage_avg_v": (
|
||||
round(power_total / current_total, 3)
|
||||
if current_total > 0
|
||||
else 0.0
|
||||
),
|
||||
"pins": [
|
||||
{
|
||||
"voltage_v": round(v, 3),
|
||||
"current_a": round(i, 3),
|
||||
"power_w": round(v * i, 3),
|
||||
}
|
||||
for v, i in zip(pin_voltage, pin_current, strict=True)
|
||||
],
|
||||
"temp_in_c": temp_in,
|
||||
"temp_out_c": temp_out,
|
||||
"temp_ext1_c": temp_ext1,
|
||||
"temp_ext2_c": temp_ext2,
|
||||
"fan_duty_pct": fan_duty,
|
||||
"fault_status": fault_status,
|
||||
"fault_log": fault_log,
|
||||
"psu_capability_w": psu_capability,
|
||||
}
|
||||
except OSError:
|
||||
return None
|
||||
|
||||
def _read_int(self, filename: str) -> int:
|
||||
"""Read an integer sysfs attribute; 0 when missing or unreadable
|
||||
(matches the exporter's ReadIntFile)."""
|
||||
try:
|
||||
with open(os.path.join(self._path, filename)) as f:
|
||||
return int(f.read().strip())
|
||||
except (OSError, ValueError):
|
||||
return 0
|
||||
|
||||
def _read_int_or(self, filename: str) -> int | None:
|
||||
try:
|
||||
with open(os.path.join(self._path, filename)) as f:
|
||||
return int(f.read().strip())
|
||||
except (OSError, ValueError):
|
||||
return None
|
||||
|
||||
def _read_temp(self, filename: str) -> float | None:
|
||||
"""Read a temperature sysfs attribute (m°C); None when missing or
|
||||
unreadable. NaN would poison the whole sample: the WebSocket
|
||||
serializer emits a bare NaN token (invalid JSON) and the REST
|
||||
JSONResponse rejects it with a 500."""
|
||||
try:
|
||||
with open(os.path.join(self._path, filename)) as f:
|
||||
return int(f.read().strip()) / 1000.0
|
||||
except (OSError, ValueError):
|
||||
return None
|
||||
|
||||
|
||||
def create_device() -> WireViewSerialDevice | WireViewHwmonDevice | None:
|
||||
"""Create a device for the first available WireView, or None.
|
||||
|
||||
Preference: the wireview-hwmon sysfs node (the wireviewd daemon owns
|
||||
the serial port in that case), otherwise direct serial on the first
|
||||
matching /dev/ttyACM*.
|
||||
"""
|
||||
hwmon_path = find_hwmon_path()
|
||||
if hwmon_path:
|
||||
return WireViewHwmonDevice(hwmon_path)
|
||||
ports = find_wireview_ports()
|
||||
if ports:
|
||||
return WireViewSerialDevice(ports[0])
|
||||
return None
|
||||
Reference in new issue
Block a user