"""桌面 GUI 使用的设备与闭环运行层。 所有 Modbus 访问都在 :class:`ControlWorker` 的后台线程中执行。GUI 主线程只 发送命令并消费事件,避免网络超时阻塞界面。 """ from __future__ import annotations from dataclasses import dataclass import json import logging from logging.handlers import RotatingFileHandler import math from pathlib import Path from queue import Empty, Queue import sys import threading import time import config from controllers import IncrementalPID from flow_control import FlowControlFault, FlowControlLoop from valve_model import ValveModel APP_NAME = "FlowControl" BUILTIN_MODEL_NAME = "valve_model_004.json" MONITOR_PERIOD_S = 0.25 @dataclass(frozen=True) class DeviceSettings: """PLC 连接参数。""" host: str port: int slave_id: int @dataclass(frozen=True) class ModelInfo: """已经校验并加载的阀模型。""" model: ValveModel path: Path display_name: str point_count: int is_builtin: bool is_monotonic: bool def application_directory() -> Path: """返回日志等可写文件的程序目录。""" if getattr(sys, "frozen", False): return Path(sys.executable).resolve().parent return Path(__file__).resolve().parent def resource_path(name: str) -> Path: """同时支持源码运行和 PyInstaller 单文件资源。""" bundle_root = getattr(sys, "_MEIPASS", None) base = Path(bundle_root) if bundle_root else Path(__file__).resolve().parent direct = base / name if direct.exists(): return direct return base / "valve_model" / name def build_runtime_logger() -> logging.Logger: """创建 GUI 运行日志,默认写到 EXE 旁的 logs 目录。""" log_dir = application_directory() / config.LOG_DIRECTORY log_dir.mkdir(parents=True, exist_ok=True) log_path = log_dir / config.LOG_FILE_NAME logger = logging.getLogger("flow_control_gui") logger.setLevel(logging.INFO) logger.propagate = False logger.handlers.clear() handler = RotatingFileHandler( log_path, maxBytes=2 * 1024 * 1024, backupCount=3, encoding="utf-8", ) handler.setFormatter( logging.Formatter( "%(asctime)s.%(msecs)03d %(levelname)s %(message)s", datefmt="%Y-%m-%d %H:%M:%S", ) ) logger.addHandler(handler) return logger def validate_device_settings(host, port, slave_id) -> DeviceSettings: """解析并校验 UI 输入的 PLC 连接参数。""" host_text = str(host).strip() if not host_text or any(char.isspace() for char in host_text): raise ValueError("PLC 地址不能为空或包含空格") try: port_value = int(str(port).strip()) except (TypeError, ValueError) as exc: raise ValueError("端口必须是整数") from exc try: slave_value = int(str(slave_id).strip()) except (TypeError, ValueError) as exc: raise ValueError("站号必须是整数") from exc if not 1 <= port_value <= 65535: raise ValueError("端口必须位于 1~65535") if not 0 <= slave_value <= 247: raise ValueError("站号必须位于 0~247") return DeviceSettings(host_text, port_value, slave_value) def validate_target(target) -> float: """解析并校验目标流量。""" try: value = float(str(target).strip()) except (TypeError, ValueError) as exc: raise ValueError("目标流量必须是数字") from exc if not math.isfinite(value): raise ValueError("目标流量必须是有限数值") if not config.TARGET_FLOW_MIN_SLM <= value <= config.TARGET_FLOW_MAX_SLM: raise ValueError( f"目标流量必须位于 {config.TARGET_FLOW_MIN_SLM:g}~" f"{config.TARGET_FLOW_MAX_SLM:g} SLM" ) return value def load_model(model_path: str | Path | None = None) -> ModelInfo: """加载内置或外部阀模型,并检查其是否匹配当前行程配置。""" is_builtin = model_path is None path = resource_path(BUILTIN_MODEL_NAME) if is_builtin else Path(model_path) path = path.expanduser().resolve() if not path.is_file(): raise ValueError(f"阀模型文件不存在:{path}") try: with path.open("r", encoding="utf-8") as file: raw = json.load(file) except (OSError, json.JSONDecodeError) as exc: raise ValueError(f"无法读取阀模型:{exc}") from exc try: declared_open = float(raw["motor_open"]) declared_closed = float(raw["motor_closed"]) area_table = raw["area_table"] except (KeyError, TypeError, ValueError) as exc: raise ValueError("阀模型缺少有效的 motor_open、motor_closed 或 area_table") from exc if not math.isclose(declared_open, config.MOTOR_OPEN_POSITION, abs_tol=1e-9): raise ValueError( f"模型打开端行程为 {declared_open:g},当前配置要求 " f"{config.MOTOR_OPEN_POSITION:g}" ) if not math.isclose(declared_closed, config.MOTOR_CLOSED_POSITION, abs_tol=1e-9): raise ValueError( f"模型关闭端行程为 {declared_closed:g},当前配置要求 " f"{config.MOTOR_CLOSED_POSITION:g}" ) if not isinstance(area_table, list) or len(area_table) < 2: raise ValueError("阀模型 area_table 至少需要两个标定点") try: model = ValveModel.load( str(path), config.MOTOR_OPEN_POSITION, config.MOTOR_CLOSED_POSITION, ) except (OSError, KeyError, TypeError, ValueError) as exc: raise ValueError(f"阀模型内容无效:{exc}") from exc for stroke, _area in model.area_table: if not config.MOTOR_OPEN_POSITION <= stroke <= config.MOTOR_CLOSED_POSITION: raise ValueError( f"模型行程 {stroke:g} 超出当前有效范围 " f"{config.MOTOR_OPEN_POSITION:g}~{config.MOTOR_CLOSED_POSITION:g}" ) prefix = "内置" if is_builtin else "外部" return ModelInfo( model=model, path=path, display_name=f"{prefix} {path.name}({len(model.area_table)} 点)", point_count=len(model.area_table), is_builtin=is_builtin, is_monotonic=model.is_monotonic, ) def build_control_loop(hardware, valve_model, logger) -> FlowControlLoop: """按 config 构造闭环控制器,不引入命令行记录和绘图依赖。""" pid = IncrementalPID( kp=config.PID_FAR_KP, ki=config.PID_FAR_KI, kd=config.PID_FAR_KD, dt=config.CONTROL_PERIOD_S, out_min=config.OPENING_MIN_PCT, out_max=config.OPENING_MAX_PCT, output_rate_limit=config.MAX_OPENING_RATE_FAR_PCT_S, ) return FlowControlLoop( hardware=hardware, pid=pid, period_s=config.CONTROL_PERIOD_S, flow_channel=config.FLOW_AFTER_ADDR, pressure_channel=config.PRESSURE_AFTER_ADDR, pressure_before_channel=config.PRESSURE_BEFORE_ADDR, motor_channel=config.MOTOR_OUTPUT_ADDR, motor_open_position=config.MOTOR_OPEN_POSITION, motor_closed_position=config.MOTOR_CLOSED_POSITION, target_min_slm=config.TARGET_FLOW_MIN_SLM, target_max_slm=config.TARGET_FLOW_MAX_SLM, zero_flow_threshold_slm=config.ZERO_FLOW_THRESHOLD_SLM, flow_valid_min_slm=config.FLOW_VALID_MIN_SLM, flow_valid_max_slm=config.FLOW_VALID_MAX_SLM, max_pressure_kpa=config.MAX_PRESSURE_KPA, max_control_dt_s=config.MAX_CONTROL_DT_S, max_consecutive_flow_failures=config.MAX_CONSECUTIVE_FLOW_FAILURES, max_consecutive_pressure_failures=config.MAX_CONSECUTIVE_PRESSURE_FAILURES, pid_far_gains=(config.PID_FAR_KP, config.PID_FAR_KI, config.PID_FAR_KD), pid_near_gains=( config.PID_NEAR_KP, config.PID_NEAR_KI, config.PID_NEAR_KD, ), opening_rate_far_pct_s=config.MAX_OPENING_RATE_FAR_PCT_S, opening_rate_near_pct_s=config.MAX_OPENING_RATE_NEAR_PCT_S, pid_near_error_min_slm=config.PID_NEAR_ERROR_MIN_SLM, pid_near_error_max_slm=config.PID_NEAR_ERROR_MAX_SLM, pid_near_error_target_ratio=config.PID_NEAR_ERROR_TARGET_RATIO, pid_far_error_min_slm=config.PID_FAR_ERROR_MIN_SLM, pid_far_error_max_slm=config.PID_FAR_ERROR_MAX_SLM, pid_far_error_target_ratio=config.PID_FAR_ERROR_TARGET_RATIO, pid_mode_switch_confirm_cycles=config.PID_MODE_SWITCH_CONFIRM_CYCLES, valve_model=valve_model, feedforward_enabled=config.FEEDFORWARD_ENABLED, feedforward_correction_band_pct=config.FEEDFORWARD_CORRECTION_BAND_PCT, feedforward_update_error_threshold_slm=( config.FEEDFORWARD_UPDATE_ERROR_THRESHOLD_SLM ), feedforward_update_window_s=config.FEEDFORWARD_UPDATE_WINDOW_S, gain_schedule_enabled=config.GAIN_SCHEDULE_ENABLED, gain_schedule_ref_abs_kpa=config.GAIN_SCHEDULE_REF_ABS_KPA, gain_schedule_floor_abs_kpa=config.GAIN_SCHEDULE_FLOOR_ABS_KPA, gain_schedule_scale_change_threshold=( config.GAIN_SCHEDULE_SCALE_CHANGE_THRESHOLD ), logger=logger, ) def create_hardware(settings: DeviceSettings): """延迟导入硬件模块,便于离线测试 GUI 运行层。""" from PcControl import Easy521ModbusClient return Easy521ModbusClient( host=settings.host, port=settings.port, slave_id=settings.slave_id, flow_addr=config.FLOW_AFTER_ADDR, flow_range=config.FLOW_AFTER_RANGE_SLM, pressure_before_addr=config.PRESSURE_BEFORE_ADDR, pressure_before_range=config.PRESSURE_BEFORE_RANGE_KPA, pressure_after_addr=config.PRESSURE_AFTER_ADDR, pressure_after_range=config.PRESSURE_AFTER_RANGE_KPA, motor_output_addr=config.MOTOR_OUTPUT_ADDR, ) class ControlWorker: """拥有设备连接的后台线程,通过命令/事件队列与 GUI 通信。""" def __init__( self, logger, *, hardware_factory=create_hardware, model_loader=load_model, controller_factory=build_control_loop, ): self.logger = logger self.hardware_factory = hardware_factory self.model_loader = model_loader self.controller_factory = controller_factory self.commands = Queue() self.events = Queue() self.thread = threading.Thread( target=self._run, name="flow-control-worker", daemon=False, ) self._shutdown_requested = False self._state = "DISCONNECTED" self._hardware = None self._controller = None self._next_cycle = 0.0 self._monitor_failures = { "FLOW": 0, "PRESSURE_BEFORE": 0, "PRESSURE_AFTER": 0, } @property def state(self): return self._state def start(self): self.thread.start() def is_alive(self): return self.thread.is_alive() def request_connect(self, settings: DeviceSettings, model_path=None): self.commands.put(("CONNECT", {"settings": settings, "model_path": model_path})) def request_target(self, target): self.commands.put(("SET_TARGET", {"target": target})) def request_disconnect(self): self.commands.put(("DISCONNECT", {})) def request_shutdown(self): self.commands.put(("SHUTDOWN", {})) def _emit(self, kind, **payload): payload["kind"] = kind self.events.put(payload) def _set_state(self, state, message): self._state = state self._emit("state", state=state, message=message) def _run(self): self._set_state("DISCONNECTED", "请选择模型并连接设备") try: while not self._shutdown_requested: timeout = self._command_timeout() command = None try: command = self.commands.get(timeout=timeout) except Empty: pass if command is not None: self._handle_command(*command) if self._shutdown_requested: break now = time.perf_counter() if self._state == "MONITORING" and now >= self._next_cycle: self._monitor_once() self._advance_deadline(now, MONITOR_PERIOD_S) elif self._state == "CONTROLLING" and now >= self._next_cycle: self._control_once(now) self._advance_deadline(now, config.CONTROL_PERIOD_S) finally: errors = self._safe_disconnect() if errors: self.logger.error("程序退出时安全断开存在问题:%s", "; ".join(errors)) self._emit("shutdown_complete", errors=errors) def _command_timeout(self): if self._state in {"MONITORING", "CONTROLLING"}: return max(0.0, min(0.05, self._next_cycle - time.perf_counter())) return 0.1 def _advance_deadline(self, now, period): self._next_cycle += period if self._next_cycle <= now: self._next_cycle = time.perf_counter() + period def _handle_command(self, command, payload): if command == "CONNECT": self._connect(payload["settings"], payload.get("model_path")) elif command == "SET_TARGET": self._set_target(payload["target"]) elif command == "DISCONNECT": self._disconnect() elif command == "SHUTDOWN": self._set_state("STOPPING", "正在安全退出…") self._shutdown_requested = True def _connect(self, settings, model_path): if self._state not in {"DISCONNECTED", "FAULT"}: self._emit("warning", message="当前状态不能重复连接") return self._set_state("CONNECTING", f"正在连接 {settings.host}:{settings.port}…") try: config.validate_config() model_info = self.model_loader(model_path) hardware = self.hardware_factory(settings) self._hardware = hardware if not bool(hardware.connect()): raise RuntimeError("无法连接 PLC") if not bool( hardware.set_motor_position( config.MOTOR_OPEN_POSITION, channel=config.MOTOR_OUTPUT_ADDR, ) ): raise RuntimeError("连接后设置阀门 100% 开度失败") self._controller = self.controller_factory( hardware, model_info.model, self.logger, ) for sensor in self._monitor_failures: self._monitor_failures[sensor] = 0 self._next_cycle = time.perf_counter() self.logger.info( "PLC 已连接 %s:%d,站号=%d;模型=%s", settings.host, settings.port, settings.slave_id, model_info.display_name, ) self._set_state("MONITORING", "设备已连接,正在监测;请设置目标流量") self._emit( "model_active", display_name=model_info.display_name, point_count=model_info.point_count, is_monotonic=model_info.is_monotonic, ) except Exception as exc: errors = self._safe_disconnect() detail = str(exc) if errors: detail += ";安全断开:" + ";".join(errors) self.logger.exception("连接设备失败:%s", detail) self._state = "FAULT" self._emit("fault", message=detail, state="FAULT") def _set_target(self, target): if self._state not in {"MONITORING", "CONTROLLING"} or self._controller is None: self._emit("warning", message="设备尚未连接,无法设置目标流量") return try: value = validate_target(target) self._controller.set_target_flow(value) if not self._controller.running: self._controller.start(initial_opening=config.INITIAL_OPENING_PCT) self._next_cycle = time.perf_counter() self.logger.info("目标流量设置为 %.3f SLM", value) self._set_state("CONTROLLING", f"闭环运行中,目标 {value:g} SLM") self._emit("target_applied", target=value) except Exception as exc: self._emit("warning", message=f"目标流量设置失败:{exc}") def _disconnect(self): if self._state == "DISCONNECTED": return self._set_state("DISCONNECTING", "正在恢复 100% 开度并断开设备…") errors = self._safe_disconnect() if errors: message = ";".join(errors) self.logger.error("设备已断开,但安全收尾存在问题:%s", message) self._emit("warning", message=message) self._set_state("DISCONNECTED", "设备已断开(收尾存在警告)") else: self.logger.info("设备已安全断开") self._set_state("DISCONNECTED", "设备已断开") def _monitor_once(self): try: flow = self._read_monitor_value( "FLOW", "get_flow", config.FLOW_AFTER_ADDR ) pressure_before = self._read_monitor_value( "PRESSURE_BEFORE", "get_pressure", config.PRESSURE_BEFORE_ADDR ) pressure_after = self._read_monitor_value( "PRESSURE_AFTER", "get_pressure", config.PRESSURE_AFTER_ADDR ) if flow is not None and not ( config.FLOW_VALID_MIN_SLM <= flow <= config.FLOW_VALID_MAX_SLM ): raise FlowControlFault( "FLOW_OUT_OF_RANGE", "MONITORING", f"流量读数 {flow:.3f} SLM 超出允许范围", {"measured_flow_slm": flow}, ) for name, value in ( ("阀前压力", pressure_before), ("阀后压力", pressure_after), ): if ( value is not None and config.MAX_PRESSURE_KPA is not None and value > config.MAX_PRESSURE_KPA ): raise FlowControlFault( "PRESSURE_OVER_LIMIT", "MONITORING", f"{name} {value:.3f} kPa 超过上限 " f"{config.MAX_PRESSURE_KPA:.3f} kPa", {"pressure_kpa": value, "sensor": name}, ) status = ( "MONITORING" if all( value is not None for value in (flow, pressure_before, pressure_after) ) else "SENSOR_READ_FAILED" ) self._emit( "telemetry", measured_flow_slm=flow, pressure_before_kpa=pressure_before, pressure_kpa=pressure_after, target_flow_slm=None, opening_pct=100.0, motor_position=config.MOTOR_OPEN_POSITION, feedforward_pct=None, correction_pct=None, pid_mode="--", status=status, ) except Exception as exc: self._handle_runtime_fault(exc) def _read_monitor_value(self, sensor, method_name, address): try: method = getattr(self._hardware, method_name) value = method(address) if value is not None: value = float(value) if value is None or not math.isfinite(value): raise ValueError(f"无效读数:{value!r}") except Exception as exc: self._monitor_failures[sensor] += 1 count = self._monitor_failures[sensor] limit = ( config.MAX_CONSECUTIVE_FLOW_FAILURES if sensor == "FLOW" else config.MAX_CONSECUTIVE_PRESSURE_FAILURES ) self.logger.error( "监测阶段 %s 读取失败 %d/%d:%r", sensor, count, limit, exc, ) if count >= limit: raise FlowControlFault( f"{sensor}_READ_FAILED", "MONITORING", f"{sensor} 连续读取失败 {count} 次", {"failure_detail": repr(exc)}, ) from exc return None self._monitor_failures[sensor] = 0 return value def _control_once(self, now): try: result = self._controller.step(now=now) self._emit( "telemetry", **result.to_dict(), pid_mode=self._controller.pid_mode, ) except Exception as exc: self._handle_runtime_fault(exc) def _handle_runtime_fault(self, exc): if isinstance(exc, FlowControlFault): message = str(exc) else: message = f"未处理的控制异常:{exc}" self.logger.exception("控制故障:%s", message) errors = self._safe_disconnect() if errors: message += ";安全断开:" + ";".join(errors) self._state = "FAULT" self._emit("fault", message=message, state="FAULT") def _safe_disconnect(self): errors = [] controller = self._controller hardware = self._hardware self._controller = None self._hardware = None if controller is not None: try: controller.stop() except Exception as exc: errors.append(f"停止控制失败:{exc}") if hardware is not None: if bool(getattr(hardware, "connected", False)): try: success = hardware.set_motor_position( config.MOTOR_OPEN_POSITION, channel=config.MOTOR_OUTPUT_ADDR, ) if not success: errors.append("断开前设置阀门 100% 开度失败") except Exception as exc: errors.append(f"断开前设置阀门 100% 开度异常:{exc}") try: hardware.disconnect() except Exception as exc: errors.append(f"断开 PLC 失败:{exc}") return errors