Files
Thing/src-tauri/src/monitor_kernel.rs
T

1155 lines
47 KiB
Rust
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
// 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<u16> {
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::<TOKEN_ELEVATION>() 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::<MonitorSettings>(&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::<serde_json::Value>(&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<bool>,
pub uptime_ms: f64,
pub group_count: u32,
pub sensor_count: u32,
pub providers: Vec<String>,
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<u32>,
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 管理)。
/// 实现 Clonesub_handle / listener_ids 用 Arc<Mutex> 共享,
/// 这样从 tauri::State clone 出的实例与原实例共享订阅控制状态。
#[derive(Clone)]
pub struct MonitorKernel {
root: PathBuf,
client: Client,
/// 无超时 client,专用于 SSE 长连接(/stream
sse_client: Client,
/// SSE 订阅任务句柄,用于在 stop 时取消订阅
sub_handle: Arc<Mutex<Option<tauri::async_runtime::JoinHandle<()>>>>,
/// 监听 process-status-changed 的句柄,用于在 stop 时取消监听
listener_ids: Arc<Mutex<Vec<tauri::EventId>>>,
/// 是否以管理员权限运行(提权模式下进程不受 ProcessManager 管控,停止走 /shutdown
elevated: Arc<AtomicBool>,
/// SSE 订阅循环活跃标志(重入守卫,防止并发/重复订阅导致双 emit monitor-data
subscribing: Arc<AtomicBool>,
}
/// SSE 订阅循环标志守卫:Drop 时复位 subscribing,覆盖所有 return / abort 路径
struct SubGuard(Arc<AtomicBool>);
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<MonitorKernelInfo, String> {
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<StartProcessParams, String> {
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<KernelStatus, String> {
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<SensorSnapshot, String> {
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::<KernelStatus>().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_sincenow 超过 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::<SensorSnapshot>(&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::<Vec<_>>();
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::<Vec<_>>();
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::<serde_json::Value>(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 在运行(可能是上次提权遗留的进程):
/// - 如果是提权 KernelisAdmin=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<crate::process_manager::ProcessInfo, String> {
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::<ProcessManager>();
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<Mutex> 状态)启动订阅,
// 确保 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::<ProcessManager>();
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,
&params,
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 未提权时才需要 /shutdownProcessManager 无法 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<String> {
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<MonitorKernelInfo, String> {
state.prepare_kernel(&app)
}
#[tauri::command]
pub async fn monitor_status(
state: tauri::State<'_, MonitorKernel>,
pm: tauri::State<'_, ProcessManager>,
) -> Result<MonitorStatus, String> {
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=trueThingHK 在管理员权限运行)
elevated: thing_elevated,
thing_elevated,
})
}
#[tauri::command]
pub async fn monitor_start(
state: tauri::State<'_, MonitorKernel>,
app: AppHandle,
) -> Result<crate::process_manager::ProcessInfo, String> {
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<KernelStatus, String> {
state.get_status().await
}
/// 检测并修复 PawnIO 驱动:确保安装器随内核部署 → 静默安装 → 重启监控内核
/// (LHM 打开一次后不会重新发现驱动,装完必须重启才能恢复 CPU 温度/功耗读取)。
/// 返回 { installed, needReboot }installed=false 说明未提权或安装失败,交由内核启动自装。
#[tauri::command]
pub async fn monitor_repair_pawnio(
state: tauri::State<'_, MonitorKernel>,
pm: tauri::State<'_, ProcessManager>,
app: AppHandle,
) -> Result<serde_json::Value, String> {
// 1. 确保 PawnIO_setup.exe 随内核部署
state.prepare_kernel(&app)?;
let setup = state.cores_dir().join("PawnIO_setup.exe");
if !setup.exists() {
return Err("未找到 PawnIO_setup.exe 安装器(binaries 资源未随包部署),请重新部署监控内核".into());
}
// 2. 静默安装驱动(继承当前进程权限;Thing 已提权则直接成功)
let mut cmd = std::process::Command::new(&setup);
cmd.args(["-install", "-silent"]);
crate::process_manager::setup_creation_flags(&mut cmd);
let (installed, need_reboot) = match cmd.status() {
Ok(status) => {
let code = status.code().unwrap_or(-1);
match code {
3010 => (true, true), // ERROR_SUCCESS_REBOOT_REQUIRED
0 => (true, false),
_ => (false, false),
}
}
Err(e) => return Err(format!("运行 PawnIO 安装器失败: {}", e)),
};
// 3. 重启监控内核,使 LHM 以 PawnIO 重新打开传感器
state.stop_subscription(&app).await;
if state.is_elevated() && !is_thing_elevated() {
state.shutdown_kernel().await?;
state.elevated.store(false, Ordering::SeqCst);
state.start_elevated(&app).await?;
} else {
pm.stop(PROCESS_ID)?;
state.start_with_subscription(&app).await?;
}
Ok(serde_json::json!({ "installed": installed, "needReboot": need_reboot }))
}
#[tauri::command]
pub async fn monitor_get_snapshot(state: tauri::State<'_, MonitorKernel>) -> Result<SensorSnapshot, String> {
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<serde_json::Value, String> {
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<std::collections::HashMap<String, bool>>,
sensor_types: Option<std::collections::HashMap<String, bool>>,
) -> Result<serde_json::Value, String> {
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))
}