// monitor_kernel.rs — ThingHK 硬件监控内核的 Tauri 侧集成 // // 职责(见 ThingHK_GUIDE.md 第 2.1 节 + 阶段三): // 1. 准备 Kernel 可执行文件(从 binaries/ 资源目录复制到工作目录) // 2. 构造 StartProcessParams 交由 ProcessManager 拉起/重启(不自己管生命周期) // 3. 作为 HTTP 客户端:轮询 /status 判断就绪 → 订阅 /stream SSE → emit "monitor-data" // 4. 写入熔断:SSE 断开后停止转发,监听 process-status-changed 在 Kernel 恢复后重新订阅 // // 数据流:Kernel --SSE--> MonitorKernel(本文件) --emit--> 前端 MonitorModule.vue // // 注意:Kernel 的 stdout/stderr 被 ProcessManager 设为 null,所有数据交互走 HTTP。 use std::fs; use std::path::PathBuf; use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::Arc; use std::time::Duration; use futures_util::StreamExt; use reqwest::Client; use serde::{Deserialize, Serialize}; use tauri::path::BaseDirectory; use tauri::{AppHandle, Emitter, Listener, Manager}; use tokio::sync::Mutex; use crate::process_manager::{ProcessManager, StartProcessParams}; // ===================== 常量 ===================== /// Kernel 进程在 ProcessManager 中的 id(与 monitor 模块 index.ts 的 process.name 对应) const PROCESS_ID: &str = "monitor"; /// Kernel 监听端口(与 ThingHK 默认端口一致,见 ThingHK Program.cs DefaultPort) const KERNEL_PORT: u16 = 8730; /// 冷启动就绪轮询间隔(与 mihomo 经验一致) const READY_POLL_INTERVAL_MS: u64 = 500; /// 冷启动就绪总超时(阶段一实测冷启动约 5s,留 5s 余量) const READY_TIMEOUT_MS: u64 = 10_000; // ===================== Windows 提权启动 FFI ===================== /// ShellExecuteW 的 nShowCmd 取值:隐藏窗口。 /// ThingHK.exe 是控制台程序,提权启动时隐藏控制台窗口避免黑框。 #[cfg(windows)] const SW_HIDE: i32 = 0; // ShellExecuteW 返回值 <= 32 表示错误(HINSTANCE 被强制为 isize)。 // 1223 = ERROR_CANCELLED(用户取消了 UAC 弹窗)。 #[cfg(windows)] #[link(name = "shell32")] extern "system" { fn ShellExecuteW( hwnd: *mut std::ffi::c_void, lp_operation: *const u16, lp_file: *const u16, lp_parameters: *const u16, lp_directory: *const u16, n_show_cmd: i32, ) -> isize; } /// 以管理员权限启动可执行文件(弹 UAC)。 /// 通过 ShellExecuteW "runas" 动词实现,返回的进程句柄不可用(跨权限级别), /// 因此提权后的进程不加入 ProcessManager,停止时通过 Kernel 的 /shutdown 接口。 #[cfg(windows)] fn shell_execute_elevated(file: &str, params: &str, dir: Option<&str>) -> Result<(), String> { use std::os::windows::ffi::OsStrExt; fn to_wide(s: &str) -> Vec { std::ffi::OsStr::new(s).encode_wide().chain(std::iter::once(0)).collect() } let operation = to_wide("runas"); let file_w = to_wide(file); let params_w = to_wide(params); let dir_w = dir.map(to_wide); let dir_ptr = dir_w.as_ref().map_or(std::ptr::null(), |d| d.as_ptr()); let result = unsafe { ShellExecuteW( std::ptr::null_mut(), operation.as_ptr(), file_w.as_ptr(), params_w.as_ptr(), dir_ptr, SW_HIDE, ) }; if result <= 32 { // 1223 = ERROR_CANCELLED:用户点了"否"或关闭了 UAC 弹窗 if result == 1223 { return Err("用户取消了 UAC 提权".into()); } Err(format!("ShellExecuteW 失败,错误码: {}", result)) } else { Ok(()) } } /// 检测当前 Thing 进程是否以管理员权限运行。 /// 用于"永久提权"模式:Thing 本身是管理员时,ThingHK 子进程自动继承权限, /// 无需 ShellExecuteW "runas",ProcessManager 可直接管理(同权限级别可 kill)。 #[cfg(windows)] pub fn is_thing_elevated() -> bool { use windows_sys::Win32::Foundation::{CloseHandle, HANDLE}; use windows_sys::Win32::Security::{ GetTokenInformation, TokenElevation, TOKEN_ELEVATION, TOKEN_QUERY, }; use windows_sys::Win32::System::Threading::{GetCurrentProcess, OpenProcessToken}; unsafe { let mut token: HANDLE = 0; if OpenProcessToken(GetCurrentProcess(), TOKEN_QUERY, &mut token) == 0 { return false; } let mut elevation = TOKEN_ELEVATION { TokenIsElevated: 0 }; let mut ret_len = 0u32; let success = GetTokenInformation( token, TokenElevation, &mut elevation as *mut _ as *mut _, std::mem::size_of::() as u32, &mut ret_len, ); CloseHandle(token); success != 0 && elevation.TokenIsElevated != 0 } } #[cfg(not(windows))] pub fn is_thing_elevated() -> bool { false } // ===================== 监控设置持久化 ===================== /// 监控模块设置(存储于 {app_data_dir}/monitor/settings.json) /// 参照代理模块 ProxySettings 的模式,仅包含需要持久化的运行时开关。 #[derive(Serialize, Deserialize, Clone, Default)] #[serde(rename_all = "camelCase")] pub struct MonitorSettings { /// 应用启动时自动启动监控内核(默认 false) #[serde(default)] pub auto_start: bool, } /// 监控设置文件路径: {app_data_dir}/monitor/settings.json fn settings_path(app_data_dir: &std::path::Path) -> PathBuf { app_data_dir.join("monitor").join("settings.json") } /// 读取监控设置(文件不存在或解析失败时返回默认值) pub fn load_monitor_settings(app_data_dir: &std::path::Path) -> MonitorSettings { fs::read_to_string(settings_path(app_data_dir)) .ok() .and_then(|s| serde_json::from_str::(&s).ok()) .unwrap_or_default() } /// 写入监控设置 pub fn save_monitor_settings(app_data_dir: &std::path::Path, settings: &MonitorSettings) -> Result<(), String> { let path = settings_path(app_data_dir); if let Some(parent) = path.parent() { fs::create_dir_all(parent).map_err(|e| format!("创建目录失败: {}", e))?; } let content = serde_json::to_string_pretty(settings).map_err(|e| format!("序列化设置失败: {}", e))?; fs::write(path, content).map_err(|e| format!("写入设置失败: {}", e)) } // ===================== 永久提权标志持久化 ===================== /// 永久提权标志文件路径: {app_data_dir}/monitor/elevate.json /// 内容: { "enabled": true } fn elevate_flag_path(app_data_dir: &std::path::Path) -> PathBuf { app_data_dir.join("monitor").join("elevate.json") } /// 读取永久提权标志 fn read_elevate_flag(flag_path: &std::path::Path) -> bool { match fs::read_to_string(flag_path) { Ok(content) => serde_json::from_str::(&content) .ok() .and_then(|v| v.get("enabled").and_then(|e| e.as_bool())) .unwrap_or(false), Err(_) => false, } } /// 写入永久提权标志 fn write_elevate_flag(flag_path: &std::path::Path, enabled: bool) -> Result<(), String> { if let Some(parent) = flag_path.parent() { fs::create_dir_all(parent).map_err(|e| format!("创建目录失败: {}", e))?; } let content = serde_json::json!({ "enabled": enabled }).to_string(); fs::write(flag_path, content).map_err(|e| format!("写入提权标志失败: {}", e)) } /// 启动时检查永久提权标志,如果已设置且当前非管理员,则 ShellExecute "runas" 重启自身。 /// 在 `setup` 早期调用(MonitorKernel 初始化之前)。 /// 返回 true 表示已触发提权重启,调用方应立即退出当前进程。 pub fn check_and_relaunch_if_needed(app_data_dir: &std::path::Path) -> bool { let flag_path = elevate_flag_path(app_data_dir); if !read_elevate_flag(&flag_path) { return false; } if is_thing_elevated() { return false; // 已是管理员,无需重启 } let exe = match std::env::current_exe() { Ok(e) => e, Err(_) => return false, }; crate::logger::log_info("monitor", "检测到永久提权标志,以管理员权限重启 Thing"); match shell_execute_elevated(&exe.to_string_lossy(), "", None) { Ok(()) => true, Err(e) => { crate::logger::log_warn("monitor", &format!("永久提权重启失败(用户可能取消了 UAC): {}", e)); false } } } // ===================== 数据结构 ===================== /// Kernel /status 响应(与 ThingHK Contracts.cs KernelStatus 对应) #[derive(Serialize, Deserialize, Clone, Debug)] #[serde(rename_all = "camelCase")] pub struct KernelStatus { pub ready: bool, pub is_admin: bool, /// PawnIO 驱动是否已安装;旧版内核无此字段,Option 兼容 pub pawn_io_installed: Option, pub uptime_ms: f64, pub group_count: u32, pub sensor_count: u32, pub providers: Vec, pub schema_version: u32, } /// Kernel /snapshot 与 /stream 推送的传感器快照(与 ThingHK Contracts.cs SensorSnapshot 对应) /// 这里用 serde_json::Value 透传,避免 Rust 侧重复定义完整 schema: /// Kernel 的 schemaVersion=1 契约由 Kernel 维护,前端按 schemaVersion 解析。 pub type SensorSnapshot = serde_json::Value; /// 返回给前端的 Kernel 信息 #[derive(Serialize, Clone)] #[serde(rename_all = "camelCase")] pub struct MonitorKernelInfo { pub path: String, pub exists: bool, pub port: u16, } /// 返回给前端的运行状态 #[derive(Serialize, Clone)] #[serde(rename_all = "camelCase")] pub struct MonitorStatus { pub running: bool, pub pid: Option, pub ready: bool, pub sensor_count: u32, pub restart_count: u32, /// 是否为提权模式(通过 UAC 以管理员权限启动) pub elevated: bool, /// Thing 自身是否以管理员权限运行(永久提权模式) pub thing_elevated: bool, } // ===================== MonitorKernel ===================== /// Tauri 侧的 Kernel 客户端。 /// 持有 HTTP client 和订阅控制句柄,不持有进程句柄(进程由 ProcessManager 管理)。 /// 实现 Clone:sub_handle / listener_ids 用 Arc 共享, /// 这样从 tauri::State clone 出的实例与原实例共享订阅控制状态。 #[derive(Clone)] pub struct MonitorKernel { root: PathBuf, client: Client, /// 无超时 client,专用于 SSE 长连接(/stream) sse_client: Client, /// SSE 订阅任务句柄,用于在 stop 时取消订阅 sub_handle: Arc>>>, /// 监听 process-status-changed 的句柄,用于在 stop 时取消监听 listener_ids: Arc>>, /// 是否以管理员权限运行(提权模式下进程不受 ProcessManager 管控,停止走 /shutdown) elevated: Arc, /// SSE 订阅循环活跃标志(重入守卫,防止并发/重复订阅导致双 emit monitor-data) subscribing: Arc, } /// SSE 订阅循环标志守卫:Drop 时复位 subscribing,覆盖所有 return / abort 路径 struct SubGuard(Arc); impl Drop for SubGuard { fn drop(&mut self) { self.0.store(false, Ordering::SeqCst); } } impl MonitorKernel { pub fn new(app_data_dir: PathBuf) -> Self { let root = app_data_dir.join("monitor"); fs::create_dir_all(&root).ok(); Self { root, client: Client::builder() .timeout(Duration::from_secs(30)) .build() .unwrap_or_else(|_| Client::new()), sse_client: Client::builder() .build() .unwrap_or_else(|_| Client::new()), sub_handle: Arc::new(Mutex::new(None)), listener_ids: Arc::new(Mutex::new(Vec::new())), elevated: Arc::new(AtomicBool::new(false)), subscribing: Arc::new(AtomicBool::new(false)), } } /// 当前是否处于提权模式 pub fn is_elevated(&self) -> bool { self.elevated.load(Ordering::SeqCst) } /// 永久提权标志是否已启用(后续启动自动触发 UAC) pub fn elevate_on_launch(&self) -> bool { read_elevate_flag(&self.root.join("elevate.json")) } /// 设置/清除永久提权标志 pub fn set_elevate_on_launch(&self, enabled: bool) -> Result<(), String> { write_elevate_flag(&self.root.join("elevate.json"), enabled) } /// 读取监控设置(auto_start 等) pub fn load_settings(&self) -> MonitorSettings { load_monitor_settings(&self.root.parent().unwrap_or(&self.root).to_path_buf()) } /// 写入监控设置 pub fn save_settings(&self, settings: &MonitorSettings) -> Result<(), String> { save_monitor_settings(&self.root.parent().unwrap_or(&self.root).to_path_buf(), settings) } /// 应用启动时检查是否需要自动启动监控内核 pub fn auto_start_on_launch(&self, app: &AppHandle) { let settings = self.load_settings(); if !settings.auto_start { return; } let monitor = self.clone(); let app_handle = app.clone(); tauri::async_runtime::spawn(async move { match monitor.start_with_subscription(&app_handle).await { Ok(info) => crate::logger::log_info("monitor", &format!("自动启动成功, pid={:?}", info.pid)), Err(e) => crate::logger::log_warn("monitor", &format!("自动启动跳过: {}", e)), } }); } fn cores_dir(&self) -> PathBuf { self.root.join("cores") } pub fn kernel_path(&self) -> PathBuf { self.cores_dir().join("ThingHK.exe") } fn kernel_url(&self) -> String { format!("http://127.0.0.1:{}", KERNEL_PORT) } /// 硬件监控配置文件路径:{app_data_dir}/monitor/hardware-config.json pub fn hardware_config_path(&self) -> PathBuf { self.root.join("hardware-config.json") } /// 确保内核就位:若 cores/ 无内核或版本过期(源文件较新),从资源目录复制。 /// 同时把 PawnIO_setup.exe(可选资源)复制过去——内核提权启动时会静默安装它, /// 作为 WinRing0 被系统/杀软拦截时读取温度/频率的替代驱动。 pub fn prepare_kernel(&self, app: &AppHandle) -> Result { let kernel = self.kernel_path(); if let Ok(src) = app.path().resolve("binaries/ThingHK.exe", BaseDirectory::Resource) { if src.exists() { let need_copy = !kernel.exists() || fs::metadata(&src) .and_then(|s| fs::metadata(&kernel).map(|d| s.len() != d.len())) .unwrap_or(true); if need_copy { fs::create_dir_all(self.cores_dir()).ok(); fs::copy(&src, &kernel).map_err(|e| format!("复制 Kernel 失败: {}", e))?; } } } // PawnIO 安装器:可选资源,缺失时仅影响自动安装能力(不影响内核运行) if let Ok(setup_src) = app.path().resolve("binaries/PawnIO_setup.exe", BaseDirectory::Resource) { if setup_src.exists() { let setup_dest = self.cores_dir().join("PawnIO_setup.exe"); let need_copy = !setup_dest.exists() || fs::metadata(&setup_src) .and_then(|s| fs::metadata(&setup_dest).map(|d| s.len() != d.len())) .unwrap_or(true); if need_copy { fs::create_dir_all(self.cores_dir()).ok(); if let Err(e) = fs::copy(&setup_src, &setup_dest) { crate::logger::log_warn("monitor", &format!("复制 PawnIO_setup.exe 失败: {}", e)); } } } } Ok(MonitorKernelInfo { path: kernel.to_string_lossy().to_string(), exists: kernel.exists(), port: KERNEL_PORT, }) } /// 构造启动 Kernel 的进程参数(交由 ProcessManager.start 拉起) pub fn prepare_for_start(&self, app: &AppHandle) -> Result { let info = self.prepare_kernel(app)?; if !info.exists { return Err(format!( "ThingHK Kernel 未安装。请将 ThingHK.exe 放置到 src-tauri/binaries/ 后重新构建,或直接放到:\n{}", self.cores_dir().to_string_lossy() )); } Ok(StartProcessParams { id: PROCESS_ID.into(), executable: self.kernel_path().to_string_lossy().to_string(), args: vec![ "serve".into(), "--port".into(), KERNEL_PORT.to_string(), "--config".into(), self.hardware_config_path().to_string_lossy().to_string(), ], cwd: Some(self.cores_dir().to_string_lossy().to_string()), name: "ThingHK".into(), restart_on_crash: true, max_restarts: 3, }) } /// 查询 Kernel /status(不启动订阅) pub async fn get_status(&self) -> Result { let url = format!("{}/status", self.kernel_url()); let resp = self .client .get(&url) .timeout(Duration::from_secs(3)) .send() .await .map_err(|e| format!("请求 Kernel /status 失败: {}", e))?; if !resp.status().is_success() { return Err(format!("Kernel /status 返回 HTTP {}", resp.status())); } resp.json().await.map_err(|e| format!("解析 Kernel /status 失败: {}", e)) } /// 一次性拉取 /snapshot pub async fn get_snapshot(&self) -> Result { let url = format!("{}/snapshot", self.kernel_url()); let resp = self .client .get(&url) .timeout(Duration::from_secs(5)) .send() .await .map_err(|e| format!("请求 Kernel /snapshot 失败: {}", e))?; if !resp.status().is_success() { return Err(format!("Kernel /snapshot 返回 HTTP {}", resp.status())); } resp.json().await.map_err(|e| format!("解析 Kernel /snapshot 失败: {}", e)) } /// 启动 SSE 订阅循环。 /// 流程:轮询 /status 等 ready → 订阅 /stream → 解析 SSE 事件 → emit "monitor-data" /// 写入熔断:SSE 断开后停止 emit,等待外部调用 reconnect 或 process-status-changed 触发重连 pub async fn start_subscription(self: Self, app: AppHandle) { // 重入守卫:已有订阅循环在运行则跳过,防止并发/重复订阅(双 emit monitor-data) if self.subscribing.swap(true, Ordering::SeqCst) { crate::logger::log_warn("monitor", "已有 SSE 订阅循环运行中,跳过重复订阅"); return; } // Drop 时复位标志,覆盖所有 return 路径(含 stop 时 abort 取消) let _guard = SubGuard(self.subscribing.clone()); // 1. 轮询等待 Kernel ready(冷启动约 5s) if let Err(e) = self.wait_for_ready(&app).await { crate::logger::log_warn("monitor", &format!("等待 Kernel ready 失败,订阅不启动: {}", e)); let _ = app.emit(crate::constants::events::MONITOR_ERROR, serde_json::json!({ "stage": "ready", "message": e })); return; } // 2. 订阅 SSE self.run_sse_loop(app).await; } /// 轮询 /status 直到 ready 或超时 async fn wait_for_ready(&self, app: &AppHandle) -> Result<(), String> { let url = format!("{}/status", self.kernel_url()); let deadline = std::time::Instant::now() + Duration::from_millis(READY_TIMEOUT_MS); let mut last_err = String::new(); while std::time::Instant::now() < deadline { match self .client .get(&url) .timeout(Duration::from_secs(2)) .send() .await { Ok(resp) if resp.status().is_success() => { match resp.json::().await { Ok(s) if s.ready => { // PawnIO 诊断:已提权但驱动缺失时,温度/频率等 ring0 传感器大概率无法读取 if s.is_admin && s.pawn_io_installed == Some(false) { crate::logger::log_warn( "monitor", "Kernel 已提权但 PawnIO 驱动未安装,CPU 温度/频率可能无法读取(检查 cores/PawnIO_setup.exe 是否随包部署)", ); } let _ = app.emit( crate::constants::events::MONITOR_READY, serde_json::json!({ "isAdmin": s.is_admin, "sensorCount": s.sensor_count, "providers": s.providers, "pawnIoInstalled": s.pawn_io_installed, }), ); return Ok(()); } Ok(_) => {} // 还没 ready,继续轮询 Err(e) => last_err = e.to_string(), } } Ok(resp) => last_err = format!("HTTP {}", resp.status()), Err(e) => last_err = e.to_string(), } // 通知前端正在加载(前端可显示 "Kernel 启动中...") // 用 saturating_duration_since:now 超过 deadline 时返回 0,避免 duration_since panic let remaining = deadline.saturating_duration_since(std::time::Instant::now()).as_millis() as u64; let elapsed = READY_TIMEOUT_MS.saturating_sub(remaining); let _ = app.emit(crate::constants::events::MONITOR_LOADING, serde_json::json!({ "elapsedMs": elapsed })); tokio::time::sleep(Duration::from_millis(READY_POLL_INTERVAL_MS)).await; } Err(format!("Kernel 在 {}ms 内未就绪: {}", READY_TIMEOUT_MS, last_err)) } /// SSE 订阅主循环。 /// 断开后自动重试(带退避),实现写入熔断 + 自动重连。 async fn run_sse_loop(self: Self, app: AppHandle) { let url = format!("{}/stream", self.kernel_url()); loop { match self.subscribe_once(&url, &app).await { // 正常结束(客户端取消或服务端关闭) Ok(()) => { crate::logger::log_info("monitor", "SSE 流正常结束"); break; } Err(e) => { crate::logger::log_warn("monitor", &format!("SSE 流异常断开: {},3s 后重试", e)); let _ = app.emit( crate::constants::events::MONITOR_DISCONNECTED, serde_json::json!({ "message": e }), ); tokio::time::sleep(Duration::from_secs(3)).await; // 重连前先确认 Kernel 是否还活着(可能已被 stop) if !self.is_kernel_alive().await { crate::logger::log_info("monitor", "Kernel 已停止,退出 SSE 循环"); // 提权模式下 Kernel 崩溃/退出后重置 elevated 标志, // 否则 monitor_status 会一直认为提权模式但 Kernel 已死,用户无法重启 if self.is_elevated() { self.elevated.store(false, Ordering::SeqCst); crate::logger::log_info("monitor", "提权 Kernel 已退出,重置 elevated 标志"); } break; } } } } } /// 订阅一次 SSE 流,直到断开。 /// 解析 `event: snapshot\ndata: {json}\n\n` 格式,emit "monitor-data"。 async fn subscribe_once(&self, url: &str, app: &AppHandle) -> Result<(), String> { let resp = self .sse_client .get(url) .header("Accept", "text/event-stream") .send() .await .map_err(|e| format!("请求 /stream 失败: {}", e))?; if !resp.status().is_success() { return Err(format!("/stream 返回 HTTP {}", resp.status())); } let mut stream = resp.bytes_stream(); let mut buffer = String::new(); while let Some(chunk) = stream.next().await { let chunk = chunk.map_err(|e| format!("读取 SSE chunk 失败: {}", e))?; // SSE 是文本协议,按 UTF-8 解码追加到缓冲区 buffer.push_str(&String::from_utf8_lossy(&chunk)); // 按双换行分割事件(SSE 事件以空行分隔) while let Some(pos) = buffer.find("\n\n") { let event_str = buffer[..pos].to_string(); buffer.drain(..pos + 2); if let Some(json_str) = parse_sse_data(&event_str) { if let Ok(snap) = serde_json::from_str::(&json_str) { let _ = app.emit(crate::constants::events::MONITOR_DATA, snap); } } } } Ok(()) } /// 检查 Kernel 是否还在响应(用于 SSE 断开后判断是否应重连) async fn is_kernel_alive(&self) -> bool { let url = format!("{}/status", self.kernel_url()); self.client .get(&url) .timeout(Duration::from_secs(2)) .send() .await .map(|r| r.status().is_success()) .unwrap_or(false) } /// 停止 SSE 订阅(进程由 ProcessManager.stop 负责) pub async fn stop_subscription(&self, app: &AppHandle) { // 取消 SSE 任务并等待其退出,确保 subscribing 重入守卫复位。 // 否则 stop 后立即 start 时旧任务仍在 Drop,新订阅会被守卫误跳过。 if let Some(handle) = self.sub_handle.lock().await.take() { handle.abort(); let _ = handle.await; } // 取消 process-status-changed 监听 let ids = self.listener_ids.lock().await.drain(..).collect::>(); for id in ids { app.unlisten(id); } } /// 注册 process-status-changed 监听:当 Kernel 进程被自动重启恢复 Running 时, /// 自动重新启动 SSE 订阅(实现崩溃恢复后的自愈)。 /// 注册前先清理旧 listener,防止重复注册导致多次 spawn SSE 订阅。 pub async fn register_auto_reconnect(self: Self, app: AppHandle) { // 清理旧 listener(防止重复注册) let old_ids = self.listener_ids.lock().await.drain(..).collect::>(); for id in old_ids { app.unlisten(id); } let app_clone = app.clone(); let this = self.clone(); let id = app.listen("process-status-changed", move |event| { // 只关心 monitor 进程的状态变化 // ProcessInfo 只有 Serialize,这里用 Value 解析 if let Ok(v) = serde_json::from_str::(event.payload()) { if v.get("id").and_then(|i| i.as_str()) == Some(PROCESS_ID) { let is_running = v.get("status").and_then(|s| s.as_str()) == Some("running"); if is_running { let this = this.clone(); let app = app_clone.clone(); tauri::async_runtime::spawn(async move { crate::logger::log_info("monitor", "检测到 Kernel 重启恢复,重新订阅 SSE"); // 重启后需要重新等待 ready(冷启动约 5s) this.clone().start_subscription(app).await; }); } } } }); self.listener_ids.lock().await.push(id); } /// 启动 Kernel 进程并开始 SSE 订阅(命令和 setup 自动启动共用)。 /// ProcessManager 通过 app.state 获取,无需外部传入。 /// /// 启动前先检测端口 8730 是否已有 Kernel 在运行(可能是上次提权遗留的进程): /// - 如果是提权 Kernel(isAdmin=true):直接接管,设置 elevated=true,跳过 ProcessManager 启动 /// - 如果是普通权限 Kernel:调用 /shutdown 停止后重新启动 /// /// 永久提权模式(Thing 自身是管理员): /// - ThingHK 通过 ProcessManager 启动,子进程继承管理员权限 /// - 设置 elevated=true 表示 ThingHK 在管理员权限下运行 /// - ProcessManager 可直接 kill(同权限级别),无需 /shutdown 接口 pub async fn start_with_subscription( &self, app: &AppHandle, ) -> Result { let thing_elevated = is_thing_elevated(); // 检测是否已有 Kernel 在运行(端口被占用) if let Ok(status) = self.get_status().await { if status.ready { if thing_elevated { // Thing 是管理员:停止已有 Kernel(无论什么权限),用 ProcessManager 重启以继承权限 crate::logger::log_info("monitor", "Thing 已提权,重启 ThingHK 以继承管理员权限"); let _ = self.shutdown_kernel().await; tokio::time::sleep(Duration::from_millis(500)).await; } else if status.is_admin { // Thing 非管理员,但已有提权 Kernel:直接接管 crate::logger::log_info("monitor", "检测到已有提权 Kernel 运行中,直接接管"); self.elevated.store(true, Ordering::SeqCst); let kernel = self.clone(); kernel.clone().register_auto_reconnect(app.clone()).await; // 已有订阅循环在运行则不覆盖句柄(否则旧任务句柄丢失,stop 无法取消) if !kernel.subscribing.load(Ordering::SeqCst) { let app_clone = app.clone(); let handle = tauri::async_runtime::spawn(async move { kernel.start_subscription(app_clone).await; }); *self.sub_handle.lock().await = Some(handle); } // 提权模式下无 ProcessInfo,返回一个占位的 return Ok(crate::process_manager::ProcessInfo { id: PROCESS_ID.into(), name: "ThingHK".into(), status: crate::process_manager::ProcessStatus::Running, pid: None, restart_count: 0, }); } else { // Thing 非管理员,已有普通 Kernel:先停止 crate::logger::log_info("monitor", "检测到已有普通权限 Kernel 运行中,先停止再重启"); let _ = self.shutdown_kernel().await; tokio::time::sleep(Duration::from_millis(500)).await; } } } let pm = app.state::(); let params = self.prepare_for_start(app)?; let info = pm.start(params)?; // Thing 是管理员时,ThingHK 继承权限,标记 elevated(仍由 ProcessManager 管理) if thing_elevated { self.elevated.store(true, Ordering::SeqCst); crate::logger::log_info("monitor", "ThingHK 已以管理员权限启动(继承自 Thing)"); } // 用 self 的 clone(共享 Arc 状态)启动订阅, // 确保 register_auto_reconnect 注册的 listener 与 sub_handle 共享, // stop_subscription 时才能正确清理 listener。 let kernel = self.clone(); kernel.clone().register_auto_reconnect(app.clone()).await; // 已有订阅循环在运行则不覆盖句柄(否则旧任务句柄丢失,stop 无法取消) if !kernel.subscribing.load(Ordering::SeqCst) { let app_clone = app.clone(); let handle = tauri::async_runtime::spawn(async move { kernel.start_subscription(app_clone).await; }); *self.sub_handle.lock().await = Some(handle); } Ok(info) } /// 以管理员权限启动 Kernel(弹 UAC)。 /// /// 流程: /// 1. 停止当前普通权限 Kernel(如果在运行)+ 取消 SSE 订阅 /// 2. 等待端口释放 /// 3. ShellExecuteW "runas" 启动提权 Kernel(隐藏控制台窗口) /// 4. 标记 elevated=true,注册自动重连 + 启动 SSE 订阅 /// /// 提权模式限制(与普通模式的差异): /// - 进程不加入 ProcessManager(跨权限级别句柄不可用) /// - 停止通过 Kernel 的 POST /shutdown 接口(普通权限无法 TerminateProcess 管理员进程) /// - 崩溃不支持自动重启(ProcessManager 看不到此进程) /// - pid 不可查询(返回 None) pub async fn start_elevated(&self, app: &AppHandle) -> Result<(), String> { // 1. 停止当前 Kernel(可能是普通权限,也可能是已提权的) self.stop_subscription(app).await; let pm = app.state::(); if pm.get_status(PROCESS_ID).is_some() { let _ = pm.stop(PROCESS_ID); } if self.elevated.load(Ordering::SeqCst) { // 已是提权模式,先通过 /shutdown 停止旧进程 let _ = self.shutdown_kernel().await; self.elevated.store(false, Ordering::SeqCst); } // 2. 准备 Kernel 二进制 let info = self.prepare_kernel(app)?; if !info.exists { return Err(format!( "ThingHK Kernel 未安装。请将 ThingHK.exe 放置到 src-tauri/binaries/ 后重新构建。" )); } // 3. 等待端口释放(pm.stop 异步 kill 后 Windows 端口释放有延迟) tokio::time::sleep(Duration::from_millis(800)).await; // 4. ShellExecute "runas" 启动提权 Kernel // 提权模式下控制台窗口隐藏(SW_HIDE),stderr 不可见, // 通过 --log-file 将日志写入文件,便于排查崩溃问题 let log_path = self.root.join("kernel-elevated.log"); let config_path = self.hardware_config_path(); let params = format!( "serve --port {} --config \"{}\" --log-file \"{}\"", KERNEL_PORT, config_path.to_string_lossy(), log_path.to_string_lossy() ); #[cfg(windows)] { shell_execute_elevated( &info.path, ¶ms, Some(&self.cores_dir().to_string_lossy()), )?; } #[cfg(not(windows))] { let _ = params; return Err("提权启动仅支持 Windows".into()); } self.elevated.store(true, Ordering::SeqCst); crate::logger::log_info("monitor", "提权启动已发起,等待 Kernel ready"); // 5. 启动 SSE 订阅(提权模式不注册 register_auto_reconnect: // 提权进程不归 ProcessManager 管,process-status-changed 事件不会触发, // 注册了反而可能在其他进程状态变化时误触发 SSE 重连) let kernel = self.clone(); let app_clone = app.clone(); // 已有订阅循环在运行则不覆盖句柄(否则旧任务句柄丢失,stop 无法取消) if !kernel.subscribing.load(Ordering::SeqCst) { let handle = tauri::async_runtime::spawn(async move { kernel.start_subscription(app_clone).await; }); *self.sub_handle.lock().await = Some(handle); } Ok(()) } /// 调用 Kernel 的 POST /shutdown 接口,让提权 Kernel 自行退出。 /// 普通权限 Tauri 无法 TerminateProcess 管理员进程,必须通过 HTTP 优雅关闭。 async fn shutdown_kernel(&self) -> Result<(), String> { let url = format!("{}/shutdown", self.kernel_url()); let resp = self .client .post(&url) .timeout(Duration::from_secs(3)) .send() .await .map_err(|e| format!("请求 /shutdown 失败: {}", e))?; if !resp.status().is_success() { return Err(format!("/shutdown 返回 HTTP {}", resp.status())); } // 等待 Kernel 进程退出 + 端口释放 tokio::time::sleep(Duration::from_millis(500)).await; Ok(()) } /// 程序退出时清理:取消 SSE 订阅 + 停止提权 Kernel。 /// 普通权限 Kernel 和永久提权模式(Thing 也是管理员)由 ProcessManager.stop_all 统一清理; /// 仅"提权 ThingHK"模式(Thing 普通、ThingHK 管理员)需通过 /shutdown 接口停止, /// 否则程序关闭后提权 Kernel 会残留并占用端口,导致下次启动端口冲突。 pub async fn cleanup_on_exit(&self, app: &AppHandle) { self.stop_subscription(app).await; // 仅当 ThingHK 提权但 Thing 未提权时才需要 /shutdown(ProcessManager 无法 kill 管理员进程) if self.is_elevated() && !is_thing_elevated() { crate::logger::log_info("monitor", "退出清理:停止提权 Kernel(/shutdown)"); let _ = self.shutdown_kernel().await; self.elevated.store(false, Ordering::SeqCst); } } } /// 解析 SSE 事件文本,提取 data: 字段的 JSON 内容。 /// 格式:`event: snapshot\ndata: {...json...}` fn parse_sse_data(event_str: &str) -> Option { let mut data = String::new(); for line in event_str.lines() { if let Some(rest) = line.strip_prefix("data:") { data.push_str(rest.trim()); } } if data.is_empty() { None } else { Some(data) } } // ===================== Tauri 命令 ===================== #[tauri::command] pub fn monitor_kernel_info( state: tauri::State<'_, MonitorKernel>, app: AppHandle, ) -> Result { state.prepare_kernel(&app) } #[tauri::command] pub async fn monitor_status( state: tauri::State<'_, MonitorKernel>, pm: tauri::State<'_, ProcessManager>, ) -> Result { let thing_elevated = is_thing_elevated(); // 提权模式(仅 ThingHK 提权,Thing 普通):进程不归 ProcessManager 管 if state.is_elevated() && !thing_elevated { let kernel_status = state.get_status().await.ok(); let running = kernel_status.is_some(); return Ok(MonitorStatus { running, pid: None, ready: kernel_status.as_ref().map(|s| s.ready).unwrap_or(false), sensor_count: kernel_status.as_ref().map(|s| s.sensor_count).unwrap_or(0), restart_count: 0, elevated: true, thing_elevated: false, }); } // 普通权限模式 或 永久提权模式(Thing 管理员):通过 ProcessManager 获取进程状态 let info = pm.get_status(PROCESS_ID); let running = info .as_ref() .map(|i| matches!(i.status, crate::process_manager::ProcessStatus::Running)) .unwrap_or(false); let kernel_status = if running { state.get_status().await.ok() } else { None }; Ok(MonitorStatus { running, pid: info.as_ref().and_then(|i| i.pid), ready: kernel_status.as_ref().map(|s| s.ready).unwrap_or(false), sensor_count: kernel_status.as_ref().map(|s| s.sensor_count).unwrap_or(0), restart_count: info.as_ref().map(|i| i.restart_count).unwrap_or(0), // 永久提权模式下 elevated=true(ThingHK 在管理员权限运行) elevated: thing_elevated, thing_elevated, }) } #[tauri::command] pub async fn monitor_start( state: tauri::State<'_, MonitorKernel>, app: AppHandle, ) -> Result { state.start_with_subscription(&app).await } /// 以管理员权限启动 Kernel(弹 UAC)。 /// 提权后的进程不归 ProcessManager 管,停止需通过 monitor_stop(内部走 /shutdown)。 #[tauri::command] pub async fn monitor_start_elevated( state: tauri::State<'_, MonitorKernel>, app: AppHandle, ) -> Result<(), String> { state.start_elevated(&app).await } /// 永久提权:以管理员权限重启 Thing 自身,并持久化标志使后续启动自动提权。 /// - 已是管理员:仅设置标志,不重启 /// - 非管理员:设置标志 + ShellExecute "runas" 重启 + 退出当前进程 #[tauri::command] pub async fn monitor_elevate_self( state: tauri::State<'_, MonitorKernel>, app: AppHandle, ) -> Result<(), String> { #[cfg(windows)] { // 持久化标志:后续启动时 setup 检测到此标志会自动 ShellExecute "runas" state.set_elevate_on_launch(true)?; if is_thing_elevated() { // 已是管理员,只需设置标志,无需重启 crate::logger::log_info("monitor", "Thing 已是管理员,仅设置永久提权标志"); return Ok(()); } let exe = std::env::current_exe().map_err(|e| format!("获取当前路径失败: {}", e))?; let exe_str = exe.to_string_lossy().to_string(); crate::logger::log_info("monitor", "永久提权:以管理员权限重启 Thing"); shell_execute_elevated(&exe_str, "", None)?; // 退出当前进程(非提权),新的提权进程会接管 app.exit(0); Ok(()) } #[cfg(not(windows))] { let _ = (state, app); Err("永久提权仅支持 Windows".into()) } } #[tauri::command] pub async fn monitor_stop( state: tauri::State<'_, MonitorKernel>, pm: tauri::State<'_, ProcessManager>, app: AppHandle, ) -> Result<(), String> { state.stop_subscription(&app).await; // 仅"提权 ThingHK"模式(Thing 普通、ThingHK 管理员)需要 /shutdown // 永久提权模式(Thing 管理员)和普通模式都用 ProcessManager kill if state.is_elevated() && !is_thing_elevated() { state.shutdown_kernel().await?; state.elevated.store(false, Ordering::SeqCst); Ok(()) } else { pm.stop(PROCESS_ID) } } #[tauri::command] pub async fn monitor_get_status(state: tauri::State<'_, MonitorKernel>) -> Result { state.get_status().await } #[tauri::command] pub async fn monitor_get_snapshot(state: tauri::State<'_, MonitorKernel>) -> Result { state.get_snapshot().await } /// 查询永久提权标志是否已启用 #[tauri::command] pub fn monitor_get_elevate_on_launch( state: tauri::State<'_, MonitorKernel>, ) -> bool { state.elevate_on_launch() } /// 设置/清除永久提权标志。 /// 清除后,下次启动不再触发 UAC(当前会话权限不变)。 #[tauri::command] pub fn monitor_set_elevate_on_launch( state: tauri::State<'_, MonitorKernel>, enabled: bool, ) -> Result<(), String> { state.set_elevate_on_launch(enabled) } /// 查询"应用启动时自动启动监控内核"是否已启用 #[tauri::command] pub fn monitor_get_auto_start( state: tauri::State<'_, MonitorKernel>, ) -> bool { state.load_settings().auto_start } /// 设置/清除"应用启动时自动启动监控内核"开关 #[tauri::command] pub fn monitor_set_auto_start( state: tauri::State<'_, MonitorKernel>, enabled: bool, ) -> Result<(), String> { let mut settings = state.load_settings(); settings.auto_start = enabled; state.save_settings(&settings) } /// 查询硬件监控配置(透传 Kernel GET /config/hardware)。 /// 返回当前配置 + 可用硬件/传感器类型清单,供前端 Dialog 渲染。 #[tauri::command] pub async fn monitor_get_hardware_config( state: tauri::State<'_, MonitorKernel>, ) -> Result { let url = format!("{}/config/hardware", state.kernel_url()); let resp = state .client .get(&url) .timeout(Duration::from_secs(5)) .send() .await .map_err(|e| format!("请求 /config/hardware 失败: {}", e))?; if !resp.status().is_success() { return Err(format!("/config/hardware 返回 HTTP {}", resp.status())); } resp.json().await.map_err(|e| format!("解析配置失败: {}", e)) } /// 更新硬件监控配置(透传 Kernel POST /config/hardware)。 /// 传感器类型过滤热生效;硬件开关变化返回 restartRequired=true,需重启 Kernel。 #[tauri::command] pub async fn monitor_set_hardware_config( state: tauri::State<'_, MonitorKernel>, hardware: Option>, sensor_types: Option>, ) -> Result { let url = format!("{}/config/hardware", state.kernel_url()); let body = serde_json::json!({ "hardware": hardware, "sensorTypes": sensor_types, }); let resp = state .client .post(&url) .timeout(Duration::from_secs(5)) .json(&body) .send() .await .map_err(|e| format!("请求 POST /config/hardware 失败: {}", e))?; if !resp.status().is_success() { return Err(format!("POST /config/hardware 返回 HTTP {}", resp.status())); } resp.json().await.map_err(|e| format!("解析响应失败: {}", e)) }