修复: unwrap 吞错 + lock 中毒 panic 高危项(数据吞 warn 跳过 / lock 降级返 Err)

This commit is contained in:
lxy
2026-08-01 12:21:52 +08:00
parent 143b859727
commit 484080ac12
6 changed files with 184 additions and 27 deletions
+35 -2
View File
@@ -114,22 +114,55 @@ pub async fn get_provider_secret_async(id: String) -> Option<String> {
}
/// 消费点用:解析 provider 真实密钥 — DB 优先,fallback keyring(兼容未迁移老库)
///
/// **不静默压空**:keyring 无记录/读取故障(None 已合并 keyring Err)时,这里**不**用 `unwrap_or_default()`
/// 静默吞成空串 —— 那会让调用方拿空 api_key 发请求吃 401,错误伪装成「密钥无效」且无线索。改为
/// warn 留痕(区分「真正未配置密钥」与「keyring 后端故障」),仍返空串交由下游 `ensure_resolved_key`
/// 早失败给出用户可读错误 —— 签名不变,调用方零改动。
pub fn resolve_provider_secret(record: &AiProviderRecord) -> String {
if !record.api_key.is_empty() {
return record.api_key.clone();
}
get_provider_secret(&record.id).unwrap_or_default()
match get_provider_secret(&record.id) {
Some(k) => k,
None => {
tracing::warn!(
"[密钥解析] provider {} (id={}) 系统钥匙串无密钥或读取故障 —— \
下游 ensure_resolved_key 将报「密钥缺失」。排查:1) 设置中是否保存过密钥;\
2) OS 钥匙串后端是否可用(Win Credential Manager / macOS Keychain / Linux Secret Service)",
record.name, record.id
);
String::new()
}
}
}
/// [`resolve_provider_secret`] 的 async 版本 — DB 有明文时同步返(不触 keyring),
/// 否则 `spawn_blocking` 调 keyring 防 D-Bus / COM 阻塞 tokio runtime。
///
/// 注:DB 明文路径直接 clone 同步返,只有 fallback keyring 才走 spawn_blocking。
///
/// **不静默压空**:同同步版,keyring 无记录/读取故障时 warn 留痕(不吞成空串致 401 难定位),
/// 仍返空串交由下游 `ensure_resolved_key` 早失败给用户可读错误 —— 签名不变,调用方零改动。
pub async fn resolve_provider_secret_async(record: AiProviderRecord) -> String {
if !record.api_key.is_empty() {
return record.api_key;
}
get_provider_secret_async(record.id).await.unwrap_or_default()
// 先取 id/name 再 await,避免 record 部分移动后无法在 warn 中引用。
let id = record.id.clone();
let name = record.name.clone();
match get_provider_secret_async(id.clone()).await {
Some(k) => k,
None => {
tracing::warn!(
"[密钥解析] provider {} (id={}) 系统钥匙串无密钥或读取故障 —— \
下游 ensure_resolved_key 将报「密钥缺失」。排查:1) 设置中是否保存过密钥;\
2) OS 钥匙串后端是否可用(Win Credential Manager / macOS Keychain / Linux Secret Service)",
name, id
);
String::new()
}
}
}
/// 写入密钥到 keyring(覆盖)
+55 -12
View File
@@ -41,13 +41,21 @@ impl StateMachine {
}
/// 获取节点状态
///
/// 锁中毒时降级返回 `NodeStatus::Pending`(保守默认:视为未启动,
/// 执行器不会误判为已完成/失败),并 `tracing::error!` 记录,不 panic。
pub fn get(&self, node_id: &NodeId) -> NodeStatus {
self.states
.lock()
.expect("状态机锁中毒")
.get(node_id)
.cloned()
.unwrap_or(NodeStatus::Pending)
match self.states.lock() {
Ok(states) => states.get(node_id).cloned().unwrap_or(NodeStatus::Pending),
Err(poisoned) => {
tracing::error!(
"状态机锁中毒,get({}) 降级返回 Pending{}",
node_id,
poisoned
);
NodeStatus::Pending
}
}
}
/// 判断状态转换是否合法
@@ -61,8 +69,23 @@ impl StateMachine {
}
/// 状态转换 — 校验合法性后更新,非法转换返回错误
///
/// 锁中毒时降级返回 `Err`(携带「状态机锁中毒」上下文)并 `tracing::error!` 记录,
/// 由调用方决定如何处理(通常是 `set_running`/`set_completed`/`set_failed`
/// 把 Err 上抛 → 节点执行流捕获后置 Failed),不 panic 拖垮 runtime。
pub fn transition(&self, node_id: NodeId, target: NodeStatus) -> anyhow::Result<()> {
let mut states = self.states.lock().expect("状态机锁中毒");
let mut states = match self.states.lock() {
Ok(guard) => guard,
Err(poisoned) => {
tracing::error!(
"状态机锁中毒,transition({}, {}) 降级返 Err{}",
node_id,
target.as_str(),
poisoned
);
anyhow::bail!("状态机锁中毒,节点 {} 状态转换失败", node_id);
}
};
let current = states
.get(&node_id)
.cloned()
@@ -104,15 +127,35 @@ impl StateMachine {
///
/// 注:原 set_waiting/set_skipped 同为旁路置位但全仓零调用,已删除。
pub fn set_cancelled(&self, node_id: NodeId) {
self.states
.lock()
.expect("状态机锁中毒")
.insert(node_id, NodeStatus::Cancelled);
match self.states.lock() {
Ok(mut states) => {
states.insert(node_id, NodeStatus::Cancelled);
}
Err(poisoned) => {
tracing::error!(
"状态机锁中毒,set_cancelled({}) 降级丢弃取消信号:{}",
node_id,
poisoned
);
}
}
}
/// 获取所有状态快照(clone 返回,调用方持独立副本)
///
/// 锁中毒时降级返回空 HashMap(调用方遍历视为无已完成节点,保守安全),
/// 并 `tracing::error!` 记录,不 panic。
pub fn snapshot(&self) -> HashMap<NodeId, NodeStatus> {
self.states.lock().expect("状态机锁中毒").clone()
match self.states.lock() {
Ok(states) => states.clone(),
Err(poisoned) => {
tracing::error!(
"状态机锁中毒,snapshot() 降级返回空 HashMap{}",
poisoned
);
HashMap::new()
}
}
}
/// 检查节点是否被取消
+14 -1
View File
@@ -45,7 +45,20 @@ pub async fn restore_pending_approvals(state: &AppState) {
// 故对每条 pending 都调一次 conv() 是幂等的(后续命中 entry().or_insert_with 跳过新建)。
let mut convs_restored: HashSet<String> = HashSet::new();
for rec in pending {
let args: serde_json::Value = serde_json::from_str(&rec.arguments).unwrap_or_default();
// 损坏 arguments 不还原(语义保留:数据异常不恢复)——
// unwrap_or_default() 会吞错误置空 {},污染 pending 还原成无参数 tool_call(语义错);
// 改为解析失败 warn! + continue 跳过该条,与下方 risk_level 解析失败的损坏过滤语义对齐。
let args: serde_json::Value = match serde_json::from_str(&rec.arguments) {
Ok(v) => v,
Err(e) => {
tracing::warn!(
"跳过 pending 审批恢复(tool_call_id={}):arguments 解析失败(不还原以免污染无参数语义): {}",
rec.tool_call_id,
e
);
continue;
}
};
// 过滤 risk_level 解析失败的损坏记录(语义保留:数据异常不恢复)·不再绑定 risk(PendingApproval.risk_level 已删)
if risk_from_str(&rec.risk_level).is_none() {
continue;
+26 -2
View File
@@ -585,8 +585,20 @@ async fn extract_knowledge_from_conversation(
.unwrap_or_else(|| "未命名对话".to_string());
// 消息读取:优先 ai_messages 表(消息拆分存储真相源),表空时 fallback 旧 messages JSON 列(老库兼容)
// 注意:DB 查询 Err 不能吞成空 Vec(否则 records.is_empty() 误为真 → 走旧 messages JSON 回退,
// 语义错:本应报 DB 故障)。这里显式 match:Ok 正常流程,Err 记 warn 后跳过本轮知识抽取。
let msg_repo = AiMessageRepo::new(db);
let records = msg_repo.list_by_conversation(conv_id).await.unwrap_or_default();
let records = match msg_repo.list_by_conversation(conv_id).await {
Ok(records) => records,
Err(e) => {
tracing::warn!(
error = %e,
conv_id,
"[KNOWLEDGE-EXTRACT] list_by_conversation 失败,跳过本轮知识抽取"
);
return Ok(0); // DB 故障,不当空数据回退(0 条,不置去重标志,允许下次重试)
}
};
let messages: Vec<ChatMessage> = if !records.is_empty() {
records.iter().map(crate::commands::ai::commands::record_to_message).collect()
} else {
@@ -598,7 +610,19 @@ async fn extract_knowledge_from_conversation(
"[KNOWLEDGE-EXTRACT] ai_messages 表为空,回退旧 messages JSON 列(老库兼容)"
);
}
serde_json::from_str(&conv.messages).unwrap_or_default()
// 损坏 → match Err 分流:勿 unwrap_or_default 吞成空 Vec(空 Vec 会让 knowledge 抽取基于空上下文,
// 误产空知识)。坏数据 warn 后跳过本轮(0 条,不置去重标志,允许下次重试),与 DB 故障同语义。
match serde_json::from_str::<Vec<ChatMessage>>(&conv.messages) {
Ok(v) => v,
Err(e) => {
tracing::warn!(
error = %e,
conv_id,
"[KNOWLEDGE-EXTRACT] 旧 messages JSON 列解析失败,跳过本轮知识抽取(勿基于空上下文抽取)"
);
return Ok(0);
}
}
};
// 过滤 user/assistant,取最后 6 条
let recent: Vec<&ChatMessage> = messages
+11 -1
View File
@@ -446,7 +446,17 @@ async fn route_load_messages(state: &State<'_, AppState>, conversation_id: Strin
conv_id = %conversation_id,
"[remote_bridge] load_messages ai_messages 表空,回退旧 messages JSON 列(老库未迁移)"
);
serde_json::from_str(&rec.messages).unwrap_or_default()
match serde_json::from_str::<Vec<ChatMessage>>(&rec.messages) {
Ok(v) => v,
Err(e) => {
tracing::warn!(
error = %e,
conv_id = %conversation_id,
"[remote_bridge] load_messages 老 messages JSON 解析失败,跳过(推送空列表)"
);
Vec::new()
}
}
}
_ => Vec::new(),
}
+43 -9
View File
@@ -291,12 +291,24 @@ type SkillsGuard = RwLockReadGuard<'static, Option<Vec<SkillInfo>>>;
/// P1-260617-3:`scan_skills` 同步递归 `fs::read_dir` + `read_to_string`(plugins/marketplaces
/// 多层嵌套,Windows 文件多时同步阻塞 tokio runtime)。本函数改 async,慢路径扫盘包
/// `spawn_blocking` 隔离(对齐 commands/project.rs detect_stack 模式)。快路径(读锁命中)仍同步无 fs。
async fn skills_lock_async() -> SkillsGuard {
///
/// 锁中毒(P1-260617-3 加固):读写锁 expect 中毒会 panic,技能加载热路径 panic 不可接受。
/// 中毒 → tracing::error! 记录 + 返 None(调用方 `skills_cached` 得空 Vec / `read_skill_content_stripped`
/// 返 None),不再 panic。中毒通常因持锁 panicking 线程(早期改 *g 时 unwrap)残留,缓存本身可重建,
/// 返空后下次 `invalidate_skills` 或进程重启自愈。
async fn skills_lock_async() -> Option<SkillsGuard> {
// 快路径:读锁命中(无 fs,纯内存)
{
let g = RwLock::read(&SKILLS).expect("SKILLS poisoned");
// 锁中毒不 panic:PoisonError 携 guard 仍可恢复数据,但缓存一致性不保 → 记录后返 None 降级空。
let g = match RwLock::read(&SKILLS) {
Ok(g) => g,
Err(_) => {
tracing::error!("SKILLS 读锁中毒,返空技能列表");
return None;
}
};
if g.is_some() {
return g;
return Some(g);
}
}
// 慢路径:扫盘(spawn_blocking 隔离同步 fs 递归,防阻塞 tokio runtime)
@@ -305,16 +317,28 @@ async fn skills_lock_async() -> SkillsGuard {
.map(|res| res.skills)
.unwrap_or_default();
{
let mut g = RwLock::write(&SKILLS).expect("SKILLS poisoned");
let mut g = match RwLock::write(&SKILLS) {
Ok(g) => g,
Err(_) => {
tracing::error!("SKILLS 写锁中毒,返空技能列表");
return None;
}
};
// 另一线程可能已填,二次检查(双检锁)
if g.is_none() {
*g = Some(scanned);
}
}
// 再取读锁返回(此时必 Some)
let g = RwLock::read(&SKILLS).expect("SKILLS poisoned");
let g = match RwLock::read(&SKILLS) {
Ok(g) => g,
Err(_) => {
tracing::error!("SKILLS 读锁(慢路径后)中毒,返空技能列表");
return None;
}
};
debug_assert!(g.is_some(), "skills_lock 慢路径后必 Some");
g
Some(g)
}
/// 技能扫描结果缓存(进程内;命中即 clone,不重复扫盘)。
@@ -326,7 +350,8 @@ async fn skills_lock_async() -> SkillsGuard {
/// 隔离同步 fs 防阻塞 tokio runtime(Tauri 单线程 runtime)。快路径(读锁命中)无 fs。
pub(crate) async fn skills_cached() -> Vec<SkillInfo> {
let g = skills_lock_async().await;
g.clone().unwrap_or_default()
// 锁中毒 → None → unwrap_or_default() 得空 Vec(对齐 P1-260617-3 中毒降级)。
g.and_then(|g| g.clone()).unwrap_or_default()
}
/// 置缓存为 None,下次 `skills_cached()` 触发重扫。
@@ -334,7 +359,15 @@ pub(crate) async fn skills_cached() -> Vec<SkillInfo> {
/// `ai_reload_skills` IPC 调用:写锁置 None → 紧接 `skills_cached()` 重扫,
/// 实现"改技能不重启即生效"。
pub(crate) fn invalidate_skills() {
let mut g = RwLock::write(&SKILLS).expect("SKILLS poisoned");
// 锁中毒:本就是想置 None 清缓存重建,但中毒时持锁线程已 panic,此处无法恢复一致状态。
// 记录错误并提前返回(下次加载由 skills_lock_async 中毒降级返空,进程重启自愈)。
let mut g = match RwLock::write(&SKILLS) {
Ok(g) => g,
Err(_) => {
tracing::error!("SKILLS 写锁中毒(invalidate_skills),跳过置 None");
return;
}
};
*g = None;
}
@@ -349,8 +382,9 @@ pub(crate) fn invalidate_skills() {
/// 否则非 Send 跨 await 点致 future 不 Send(MentionResolver 要求 Send)。guard 用 { } 限作用域。
pub(crate) async fn read_skill_content_stripped(name: String) -> Option<String> {
// 在作用域内取 path 后立即 drop guard,避免非 Send guard 跨 spawn_blocking await 点
// 锁中毒 → skills_lock_async 返 None → ? 早返 None(对齐 P1-260617-3 中毒降级返空 desc)。
let path: String = {
let g = skills_lock_async().await;
let g = skills_lock_async().await?;
let skills = g.as_ref()?;
let info = skills.iter().find(|s| s.name == name)?;
info.path.clone()