新增: 消息级溯源 P2 切 ai_messages + AI Chat 跑题改进 P0-P2 + 标题诊断 + 苛刻测

消息级溯源 P2(一次性切读 b + 备份保留):
- 批次A 读路径切 ai_messages(state.rs AppState.ai_messages + 映射 ChatMessage↔AiMessageRecord +
  switch/export load 切 + fallback 兜底)
- 批次B 写路径切 ai_messages(replace_conversation 单事务全量重写 + save/clear 切 + 旧 messages 备份)
- 摘要/压缩不改 updated_at(save_conversation touch_updated_at,compress false 防时间分组跳变)

AI Chat 跑题改进 P0-P2(5 改进,治跑题 5 根因,三原则优雅/可靠/易迭代):
- P0 系统提示聚焦(## 聚焦准则独立段)+ 意图接入 loop(filter_tool_defs 收敛工具 29→5-10,三重 fallback)
- P1 压缩增强(compress prompt 主题锚点 + 失败关键词兜底 extract_keyword_summary)+
  工具结果压缩(should_summarize + extract_key_info view-only 不改持久化)
- P2 主题检测(TrackedMessage.topic + 双高置信保守标记 + tokenize 2-gram 修复中文锚点)
- title LLM complete 失败诊断日志(助定位返空根因)
- 跑题 P0 审查🟡修复(Debug 加 Data 域 + threshold 注释)

苛刻测 38 条(边界/对抗:全漂移/全停用词/全错误行/2KB 边界/连续主题切换/中文 2-gram/emoji)
df-ai 227 passed + workspace EXIT 0
This commit is contained in:
2026-06-20 03:54:58 +08:00
parent 4e3a11f925
commit de049704c6
11 changed files with 2434 additions and 44 deletions

View File

@@ -8,6 +8,13 @@ use tokio::sync::Mutex;
use df_ai::ai_tools::AiToolRegistry;
use df_ai::context::TokenEstimator;
// 改进3 B: 压缩失败兜底关键词摘要(纯函数 extract_keyword_summary)。
// 改进4: 工具结果 view-only 摘要(should_summarize_tool_result / extract_key_info)。
use df_ai::context_helpers::{
extract_key_info, extract_keyword_summary, should_summarize_tool_result,
};
// 改进2 B:意图收敛工具(LLM 可见 tool_defs 按 intent 过滤,执行路径仍走完整 registry)
use df_ai::intent::{filter_tool_defs, IntentRecognizer};
use df_ai::provider::{ChatMessage, CompletionRequest, LlmProvider};
// CR-30-1: 复用 retry::backoff_delay(jitter 1s→2s→4s) + is_status_retryable(Fatal 分类)
// 实现流前失败重试退避对齐(决策 F-260616-07 a1),避免重写退避逻辑。
@@ -57,6 +64,32 @@ const PROTECT_COUNT: usize = 6;
/// partial_text不重试——保文入库 + AiCompleted(incomplete=true) + 系统提示网络中断。
pub const DEFAULT_MAX_AGENT_RETRIES: usize = 3;
/// 改进3 B 常量开关:压缩失败时是否启用关键词摘要兜底(默认 true)。
///
/// true(默认):LLM 压缩失败 → compress_old_messages 标 compressed(释放 token)
/// + insert_at(0, system 关键词摘要),保留用户反复提及的主题词作续接锚点。
/// false(回退):回退原裸裁剪行为(不插关键词摘要,仅 build_for_request 裁剪保最近 6 条)。
/// 排障/对比用:置 false 即可观察无兜底时的裁剪效果。
pub const KEYWORD_FALLBACK_ENABLED: bool = true;
/// 改进4 常量开关:tool_result 大输出是否做 view-only 摘要压缩(默认 true)。
///
/// true(默认):build_for_request 后送 stream 前,遍历 history 中 tool_result,
/// 超阈值(content >2KB 或占历史 token >40%)的 content 应用 extract_key_info
/// 替换(保留错误行 + 首尾各 5 行)。**仅 messages clone 视图,不改 ContextManager
/// 持久化**(对齐 sanitize_messages:DB 原始 tool_result 完整保留)。
/// false(关闭):tool_result 原样送 LLM(旧行为)。排障/对比用。
pub const TOOL_RESULT_COMPRESS_ENABLED: bool = true;
/// 改进5 常量开关:是否启用主题切换系统标记(默认 true)。
///
/// true(默认):push 时末两条 user 消息 topic 都非 None 且不同(双高置信)→
/// pending_topic_marker 置位 → agentic loop 顶部读取并 insert 一条 system 软提示
/// `"── 用户已切换话题(从「{old}」到「{new}」),请以新话题为准 ──"`(软提示,不强制 LLM)。
/// false(关闭):loop 顶部跳过读取/insert(置 false 即可观察无主题标记效果)。
/// 保守:双高置信才标(任一 topic None 不标),不强制 LLM(软提示非硬约束)。
pub const TOPIC_MARKER_ENABLED: bool = true;
// ============================================================
// 重构第一批(2026-06-19):GeneratingGuard 抽离到 guard.rs(纯结构搬迁,行为零变更)。
// run_agentic_loop 内仍 `GeneratingGuard::new(...)`,路径从本模块改 super::guard。
@@ -420,7 +453,55 @@ pub(crate) async fn run_agentic_loop(
}
_ => resolved_model,
};
let tool_defs = tools_arc.tool_definitions();
// 改进2 B:意图收敛工具(LLM 可见 tool_defs 按 intent 过滤,执行路径仍走完整 registry)
//
// 机制化收敛跑题:取末条 active user 消息 → IntentRecognizer 识别 → 置信 ≥ 阈值
// 则按 intent domain 过滤工具子集(减少 LLM 在无关工具上分心/误用)。
// 三重 fallback(改进2 B §可靠):
// 1. 置信 < INTENT_CONF_THRESHOLD(0.7)→ 全量
// 2. subset 空(Chat/Unknown/Conversation)→ filter_tool_defs 内部回全量
// 3. 过滤后 < 3 条(疑似漂移/误收敛)→ 回全量
// 关键安全:filter 仅改 LLM 可见 tool_defs,**不改执行**——audit 走 tools_arc.get/execute
// 完整 registry,LLM 即使幻觉一个被滤掉的工具名,audit 也能查到/拒绝。
//
// INTENT_CONF_THRESHOLD(常量开关):阈值,低于此值不收敛(回全量)。
// 注意:conf 截断到 1.0(intent.rs),单关键词 weight=1.0 即达 conf=1.0,故阈值=1.0
// 时 conf>=1.0 仍过滤(非关闭)。真关闭收敛:置 >1.0(如 1.1)。调低 = 更激进收敛。
let user_text: String = {
let session = session_arc.lock().await;
session
.conv_read(&conv_id)
.and_then(|c| {
let msgs = c.messages.all_messages_clone();
msgs.iter()
.rev()
.find(|m| matches!(m.role, df_ai::provider::MessageRole::User))
.map(|m| m.content.clone())
})
.unwrap_or_default()
};
let (intent, conf) = IntentRecognizer::recognize(&user_text);
const INTENT_CONF_THRESHOLD: f32 = 0.7;
let all_defs = tools_arc.tool_definitions();
let total = all_defs.len(); // 提前记录全量数(all_defs 将 move 进 tool_defs)
let tool_defs = if conf >= INTENT_CONF_THRESHOLD {
let filtered = filter_tool_defs(&all_defs, &intent);
if filtered.len() < 3 {
all_defs // 兜底:过滤<3(漂移/误收敛)回全量
} else {
filtered
}
} else {
all_defs // 低置信 fallback 全量
};
tracing::info!(
conv_id = %conv_id,
intent = intent.as_str(),
conf,
filtered = tool_defs.len(),
total,
"[ai] 意图收敛工具"
);
// 停止信号副本stream_llm 与每轮迭代共享读取,避免重复加锁
// notify 同取一份 Arc 引用B-260615-14stream_llm select! 监听 notified() 即时唤醒
// F-260616-09 B 批2:取 per_conv 的 stop_flag/notify(设计 §4.2 :446)。
@@ -467,7 +548,7 @@ pub(crate) async fn run_agentic_loop(
total_tokens: tokens.total(),
};
// 入口 stop:本轮可能尚未 stream(首轮即停),不记 model——避免把未实际生成的 model 写入 models 数组
save_conversation(&session_arc, &db, &conv_id, Some(&usage), None).await;
save_conversation(&session_arc, &db, &conv_id, Some(&usage), None, true).await;
// 标题生成后台化:不阻塞 Completed emit失败有 extract_title 兜底)
spawn_ensure_title(&provider_config, &db, &conv_id, &app_handle, &session_arc, &llm_concurrency);
guard.reset().await;
@@ -476,6 +557,56 @@ pub(crate) async fn run_agentic_loop(
return;
}
// 改进5: 主题切换系统标记(保守,双高置信才标,软提示非强制)。
//
// push 时若末两条 user 消息的 topic 都非 None 且不同(Intent 高置信推断的双 topic),
// ContextManager 已置位 pending_topic_marker("old|new")。loop 顶部读取并消费
// (take 一次性清空,防重复 insert),insert 一条 system 软提示告知 LLM 用户已切换话题。
//
// 保守设计:
// - 双高置信:两条 topic 都非 None(都达 0.7 阈值)才标,任一 None(低置信未标)不标。
// - 软提示:仅 insert 一条 system 消息,不强制 LLM 行为(LLM 仍可按自己理解响应)。
// - TOPIC_MARKER_ENABLED(常量开关)false → 跳过(排障/对比用)。
// topic 不参与裁剪/压缩(只供检测),insert_at(0) 同压缩摘要定位(首位 system)。
if TOPIC_MARKER_ENABLED {
let topic_marker_raw: Option<String> = {
let mut session = session_arc.lock().await;
if !session.per_conv.contains_key(&conv_id) {
tracing::warn!(
stale_conv = %conv_id,
"[ai] conv 已删除,loop 退出(主题标记段入口)"
);
return;
}
let conv = session.conv(&conv_id);
conv.messages.take_topic_marker()
};
if let Some((old_topic, new_topic)) = topic_marker_raw.and_then(|s| {
// 解析 "old|new" 格式;splitn 防 topic 名内含 '|' 误切(仅切首 '|' 一次)。
let mut parts = s.splitn(2, '|');
let old = parts.next()?.to_string();
let new = parts.next()?.to_string();
Some((old, new))
}) {
let marker_text = format!(
"── 用户已切换话题(从「{}」到「{}」),请以新话题为准 ──",
old_topic, new_topic
);
let mut session = session_arc.lock().await;
if session.per_conv.contains_key(&conv_id) {
let conv = session.conv(&conv_id);
conv.messages.insert_at(0, ChatMessage::system(&marker_text));
tracing::info!(
conv_id = %conv_id,
iteration,
old_topic = %old_topic,
new_topic = %new_topic,
"[ai] 主题切换标记已 insert(软提示,保守双高置信)"
);
}
}
}
// B-260615-11: 旧 loop 污染防护——每轮开始校验对话一致性。
// 用户新建/切换对话后 active_conversation_id 变更,本 loop(conv_id 快照)成陈旧,
// 继续跑会往新对话 push 消息/pending 造成污染。检测到即退出(guard Drop 复位 generating)。
@@ -577,6 +708,15 @@ pub(crate) async fn run_agentic_loop(
(active_msgs, lang)
};
// 改进3 B:在 active_msgs move 进 compress_via_llm 前,先算关键词摘要兜底文本。
// LLM 压缩失败时仍想保留用户反复提及的主题词(续接锚点),避免裸裁剪丢主题。
// KEYWORD_FALLBACK_ENABLED=false → 跳过(回退原裸裁剪行为,排障/对比用)。
let keyword_fallback: String = if KEYWORD_FALLBACK_ENABLED {
extract_keyword_summary(&active_msgs)
} else {
String::new()
};
// 压缩调用(复用 loop 顶部已 build 的 provider,api_key 经 secret 闭环)。
// 成功 → Some(summary);失败 → Err;无 active 可压缩(active_msgs 空)→ 视为 noop。
let compress_outcome: Result<Option<String>, String> = if active_msgs.is_empty() {
@@ -621,15 +761,46 @@ pub(crate) async fn run_agentic_loop(
session_arc.lock().await.conv(&conv_id).messages.set_compressing(false);
}
Err(e) => {
// LLM 失败 → 消息状态完全不变(未改 status / 未扣 token)。
// set_compressing(false) 复位 + emit AiError(message 不含 api_key)。
// 不阻塞 loop:继续走下方 build_for_request 裁剪路径(保最近 6 条)。
// LLM 失败 → 改进3 B:仍标 compressed 释放 token + 关键词摘要塞回首条(非裸裁剪)。
//
// 旧行为:消息状态完全不变,降级走 build_for_request 裁剪(丢主题)。
// 新行为(KEYWORD_FALLBACK_ENABLED=true 默认):
// - compress_old_messages 标 [..protect_start] active 为 compressed(释放 token,
// 与成功路径一致,后续 build_for_request 不再把它们进 LLM 上下文);
// - keyword_fallback 非空 → insert_at(0, system 关键词摘要)作续接锚点;
// - keyword_fallback 空(无 user 消息/无可提取词)→ 不插,等价旧裁剪(保底)。
// 持久化语义不变:compressed 仍软删可追溯(DB 全量保留),与成功路径一致。
// KEYWORD_FALLBACK_ENABLED=false → 跳过兜底,等价旧行为(set_compressing(false) +
// 消息状态不变,降级 build_for_request 裁剪)。
session_arc.lock().await.conv(&conv_id).messages.set_compressing(false);
tracing::warn!(
conv_id = %conv_id,
error = %e,
"[ai] 自动压缩失败,降级走原裁剪(build_for_request)"
keyword_fallback_len = keyword_fallback.len(),
"[ai] 自动压缩失败,降级走关键词摘要兜底(KEYWORD_FALLBACK_ENABLED={})",
KEYWORD_FALLBACK_ENABLED,
);
if KEYWORD_FALLBACK_ENABLED {
// 标 compressed 释放 token + 关键词摘要塞首位(若非空)。
let inserted = {
let mut session = session_arc.lock().await;
let conv = session.conv(&conv_id);
let _compressed = conv.messages.compress_old_messages(protect_start);
if !keyword_fallback.is_empty() {
conv.messages.insert_at(0, ChatMessage::system(&keyword_fallback));
true
} else {
false
}
};
if inserted {
tracing::info!(
conv_id = %conv_id,
pre_tokens = pre_compress_tokens,
"[ai] 压缩失败兜底:关键词摘要已插首位(compressed 标记已扣 token)"
);
}
}
let _ = app_handle.emit("ai-chat-event", AiChatEvent::AiError {
error: format!("自动上下文压缩失败,已降级为普通裁剪: {}", e),
error_type: Some(ErrorType::Unknown),
@@ -660,6 +831,68 @@ pub(crate) async fn run_agentic_loop(
msgs
};
// 改进4: tool_result view-only 摘要压缩(build_for_request 后,送 stream 前)。
//
// 遍历 messages(history clone,已含 system prompt + sanitize 后历史)中的 role=Tool 消息,
// 超阈值(content >2KB 或占历史 token >40%)的 content 应用 extract_key_info 替换:
// 保留错误行(error/panic/失败/.rs:N)+ 首/尾各 5 行,中间省略。
//
// **view-only**:messages 是 build_for_request 返回的 clone,改它只影响本轮 LLM 请求视图,
// 不改 ContextManager 持久化(对齐 sanitize_messages:DB 原始 tool_result 完整保留)。
// 故即使摘要有误/过度压缩,下次 build_for_request 仍从 DB 全量重建,可自愈。
//
// TOOL_RESULT_COMPRESS_ENABLED=false(常量开关)→ 跳过(原样送 LLM,排障/对比用)。
let messages: Vec<ChatMessage> = if TOOL_RESULT_COMPRESS_ENABLED {
// 历史总 token 用于占比判定(history_tokens 即 build_for_request 前的快照,
// 此处 messages 已裁剪过,用 estimated_prompt 不准;改用 history_tokens 快照更稳)。
// 注:用 conv.messages.history_tokens() 快照做占比基准(裁剪前的真相),避免循环依赖。
let history_tokens_snapshot: u32 = {
let session = session_arc.lock().await;
session
.conv_read(&conv_id)
.map(|c| c.messages.history_tokens())
.unwrap_or(0)
};
let est = TokenEstimator::default();
let mut compressed_bytes: usize = 0;
let mut original_bytes: usize = 0;
let mut summarized_count: usize = 0;
let msgs: Vec<ChatMessage> = messages
.into_iter()
.map(|mut m| {
if !matches!(m.role, df_ai::provider::MessageRole::Tool) {
return m;
}
let content_len = m.content.len();
let content_tokens = est.estimate_text(&m.content);
if !should_summarize_tool_result(content_len, history_tokens_snapshot, content_tokens) {
return m;
}
original_bytes += content_len;
// tool_name:无 tool_calls 关联(本消息是 tool_result,无 name 字段),用 call_id 或占位。
let tool_name = m.tool_call_id.clone().unwrap_or_else(|| "tool".to_string());
let compressed = extract_key_info(&m.content, &tool_name);
compressed_bytes += compressed.len();
summarized_count += 1;
m.content = compressed;
m
})
.collect();
if summarized_count > 0 {
tracing::info!(
conv_id = %conv_id,
iteration,
summarized_count,
original_bytes,
compressed_bytes,
"[ai] tool_result view-only 摘要压缩(不改持久化)"
);
}
msgs
} else {
messages
};
// 预估输入 token(兜底:部分 provider 如 GLM 流式 usage 不报 prompt_tokens,后段用它补)
// 注:stream_one_provider 内每次重试重建 request(因 provider.stream 消费 body),
// 此处不再预构建 request(旧 request 变量已废弃),仅保留 messages 供 estimated_prompt。
@@ -829,7 +1062,7 @@ pub(crate) async fn run_agentic_loop(
}
}
save_conversation(&session_arc, &db, &conv_id, Some(&usage), Some(&resolved_model)).await;
save_conversation(&session_arc, &db, &conv_id, Some(&usage), Some(&resolved_model), true).await;
// 标题生成后台化(失败有 extract_title 兜底)
spawn_ensure_title(&provider_config, &db, &conv_id, &app_handle, &session_arc, &llm_concurrency);
guard.reset().await;
@@ -895,7 +1128,7 @@ pub(crate) async fn run_agentic_loop(
completion_tokens: tokens.completion(),
total_tokens: tokens.total(),
};
save_conversation(&session_arc, &db, &conv_id, Some(&usage), Some(&resolved_model)).await;
save_conversation(&session_arc, &db, &conv_id, Some(&usage), Some(&resolved_model), true).await;
// 标题生成后台化:不阻塞 Completed emit失败有 extract_title 兜底)
spawn_ensure_title(&provider_config, &db, &conv_id, &app_handle, &session_arc, &llm_concurrency);
guard.reset().await;
@@ -920,7 +1153,7 @@ pub(crate) async fn run_agentic_loop(
completion_tokens: tokens.completion(),
total_tokens: tokens.total(),
};
save_conversation(&session_arc, &db, &conv_id, Some(&usage), Some(&resolved_model)).await;
save_conversation(&session_arc, &db, &conv_id, Some(&usage), Some(&resolved_model), true).await;
// B-260615-26: 审批等待 return 前 disarm guard——保持 generating=true 留 try_continue 续生成,
// 同时 Drop 因 done=true 跳过复位 spawn(避免误复位审批态 generating 致 ai_approve→try_continue 不续)
guard.disarm();
@@ -949,7 +1182,7 @@ pub(crate) async fn run_agentic_loop(
completion_tokens: tokens.completion(),
total_tokens: tokens.total(),
};
save_conversation(&session_arc, &db, &conv_id, Some(&usage), Some(&resolved_model)).await;
save_conversation(&session_arc, &db, &conv_id, Some(&usage), Some(&resolved_model), true).await;
// 暂停态保持 generating=true(防其他 send 抢占,仿审批),disarm guard 跳过 Drop 兜底复位
guard.disarm();
let _ = app_handle.emit("ai-chat-event", AiChatEvent::AiMaxRoundsReached {
@@ -979,7 +1212,7 @@ pub(crate) async fn run_agentic_loop(
let llm_concurrency = llm_concurrency.clone();
let resolved_model = resolved_model.clone();
tauri::async_runtime::spawn(async move {
save_conversation(&session_arc, &db, &conv_id, Some(&usage), Some(&resolved_model)).await;
save_conversation(&session_arc, &db, &conv_id, Some(&usage), Some(&resolved_model), true).await;
// 知识提炼:需读已落库的对话消息,故在 save 之后
if let Err(e) = maybe_spawn_extraction(&session_arc, &db, &conv_id, &provider_config, &knowledge_config, llm_concurrency.clone()).await {
tracing::warn!("知识提炼触发失败(非阻断): {}", e);

View File

@@ -155,7 +155,7 @@ pub async fn ai_regenerate(
}
// 落库:弹出后的历史先持久化(前端立即反映已删旧回复;loop 内再 save 覆盖)
save_conversation(&state.ai_session, &state.db, &conv_id, None, None).await;
save_conversation(&state.ai_session, &state.db, &conv_id, None, None, true).await;
let session_arc = state.ai_session.clone();
let tools_arc = state.ai_tools.clone();
@@ -402,7 +402,7 @@ pub async fn ai_approve(
// 拒绝结果立即落库(含 recovered 积压审批)——switch 时已 restore_from_messages 载完整历史,
// messages 非空,save 不会污染老对话;原 if !recovered 守卫前提不成立已移除。
if let Some(ref cid) = conv_id {
save_conversation(&state.ai_session, &state.db, cid, None, None).await;
save_conversation(&state.ai_session, &state.db, cid, None, None, true).await;
}
// 审计:拒绝(决策者=human
audit_finalize(&state, &tool_call_id, "rejected", None).await;
@@ -504,7 +504,7 @@ pub async fn ai_approve(
// 含 recovered 积压审批——switch 时已 restore_from_messages 载完整历史,messages 非空,
// save 不污染老对话;原 if !recovered 守卫前提不成立已移除。
if let Some(ref cid) = conv_id {
save_conversation(&state.ai_session, &state.db, cid, None, None).await;
save_conversation(&state.ai_session, &state.db, cid, None, None, true).await;
}
// F-260616-11 决策 a: 审批续跑 iteration 累计(不重置)——读 per_conv.iteration_used 透传
@@ -576,7 +576,7 @@ pub async fn ai_authorize_dir(
});
audit_finalize(&state, &tool_call_id, "rejected", Some(err_msg)).await;
if let Some(ref cid) = conv_id {
save_conversation(&state.ai_session, &state.db, cid, None, None).await;
save_conversation(&state.ai_session, &state.db, cid, None, None, true).await;
}
let cont_conv_id = conv_id.clone().unwrap_or_default();
let start_iter = {
@@ -635,7 +635,7 @@ pub async fn ai_authorize_dir(
conversation_id: conv_id.clone(),
});
if let Some(ref cid) = conv_id {
save_conversation(&state.ai_session, &state.db, cid, None, None).await;
save_conversation(&state.ai_session, &state.db, cid, None, None, true).await;
}
let cont_conv_id = conv_id.clone().unwrap_or_default();
let start_iter = {
@@ -688,13 +688,20 @@ pub async fn ai_chat_clear(state: State<'_, AppState>) -> Result<(), String> {
}
session.pending_approvals.retain(|_, a| a.conversation_id.as_deref() != active_id.as_deref());
drop(session);
// 真删 DB:清空该对话 messages JSON + 清零 token(保留对话壳),刷新不再恢复(AR-7)
// 真删 DB:清空该对话 messages(JSON 备份列同步清空 + 清零 token,保留对话壳),刷新不再恢复(AR-7)
// F-260619-03 批次 B:同时清空 ai_messages 表(全删,delete_range min_seq=0 max=None)
// 旧 clear_messages(置 messages='[]')保留调用:备份列同步清空,防 reload fallback 读旧脏数据。
if let Some(id) = active_id {
state
.ai_conversations
.clear_messages(&id)
.await
.map_err(err_str)?;
state
.ai_messages
.delete_range(&id, 0, None)
.await
.map_err(err_str)?;
}
Ok(())
}
@@ -740,7 +747,7 @@ pub async fn ai_chat_clear_context(
conversation_id
};
// 落库持久化新 status(照 save_conversation 模式,DB 持久化新 status)
save_conversation(&state.ai_session, &state.db, &conv_id, None, None).await;
save_conversation(&state.ai_session, &state.db, &conv_id, None, None, true).await;
let _ = app.emit("ai-chat-event", AiChatEvent::AiContextCleared {
conversation_id: Some(conv_id),
});
@@ -867,7 +874,7 @@ pub async fn ai_chat_compress_context(
conv.messages.insert_at(0, ChatMessage::system(&summary));
conv.messages.set_compressing(false);
}
save_conversation(&state.ai_session, &state.db, &conv_id, None, None).await;
save_conversation(&state.ai_session, &state.db, &conv_id, None, None, false).await;
let _ = app.emit("ai-chat-event", AiChatEvent::AiCompressed {
conversation_id: Some(conv_id),
summary,
@@ -955,7 +962,7 @@ pub async fn ai_chat_edit(
}
// 落库:编辑+截断后的历史先持久化(前端立即反映已截断旧回复)
save_conversation(&state.ai_session, &state.db, &conv_id, None, None).await;
save_conversation(&state.ai_session, &state.db, &conv_id, None, None, true).await;
let session_arc = state.ai_session.clone();
let tools_arc = state.ai_tools.clone();

View File

@@ -15,7 +15,8 @@ use std::sync::atomic::Ordering;
use tauri::{AppHandle, State};
use df_ai::provider::ChatMessage;
use df_ai::provider::{ChatMessage, MessageRole};
use df_storage::models::AiMessageRecord;
use df_types::types::new_id;
use crate::state::AppState;
@@ -26,6 +27,124 @@ use super::super::conversation::save_conversation;
use super::super::prompt::get_active_provider;
use super::super::title::spawn_ensure_title;
// ============================================================
// ChatMessage ↔ AiMessageRecord 映射(F-260619-03 消息拆分存储)
// ============================================================
//
// 读路径(批次 A)核心:`record_to_message` 把 AiMessageRecord(单行)还原成 ChatMessage。
// 写路径(批次 B)备好但本批次不调:`message_to_record`(save_conversation 写 ai_messages 用)。
//
// round-trip 严格性:
// - role 枚举 ↔ 小写字符串双向(serialize rename_all=lowercase)
// - parts / tool_calls:JSON 字符串 ↔ Option<Vec> 双向反序列化
// - status:None ↔ "active"(对齐 sanitize is_active 归一化)
// - id:record.id 必非空(主键),还原为 ChatMessage.id=Some;
// 反向 ChatMessage.id=None 时兜底 `msg_{conv}_{seq}`(save 写路径用,本批次不触发)
//
// 映射工具放 src-tauri(非 df-ai-core):df-ai-core 不依赖 df-storage(防循环依赖,
// df-storage 的迁移层注释明确不依赖 df-ai-core),只能在上层 src-tauri 桥接两者类型。
/// role 字符串 → MessageRole 枚举(小写,对齐 serialize rename_all="lowercase")。
///
/// 未知值兜底 User(防御性,实际 DB role 列由 message_to_record 写入,值域受控)。
fn role_from_str(s: &str) -> MessageRole {
match s {
"system" => MessageRole::System,
"user" => MessageRole::User,
"assistant" => MessageRole::Assistant,
"tool" => MessageRole::Tool,
_ => MessageRole::User,
}
}
/// MessageRole 枚举 → role 字符串(小写)。
fn role_to_str(r: &MessageRole) -> &'static str {
match r {
MessageRole::System => "system",
MessageRole::User => "user",
MessageRole::Assistant => "assistant",
MessageRole::Tool => "tool",
}
}
/// AiMessageRecord → ChatMessage(读路径核心:list_by_conversation 还原消息)。
///
/// - role 字符串 → 枚举
/// - parts/tool_calls:JSON 字符串 → Vec 反序列化(空/解析失败 → None)
/// - status:"active" → None(对齐 sanitize is_active 归一化),其他原样保留
/// - id:record.id 必非空(主键)→ Some
pub fn record_to_message(rec: &AiMessageRecord) -> ChatMessage {
let parts = rec
.parts
.as_deref()
.filter(|s| !s.is_empty())
.and_then(|s| serde_json::from_str(s).ok());
let tool_calls = rec
.tool_calls
.as_deref()
.filter(|s| !s.is_empty())
.and_then(|s| serde_json::from_str(s).ok());
// status "active" → None(归一化,对齐 ChatMessage 默认语义 None=正常可见)
let status = if rec.status == "active" {
None
} else {
Some(rec.status.clone())
};
ChatMessage {
id: Some(rec.id.clone()),
role: role_from_str(&rec.role),
content: rec.content.clone(),
parts,
tool_call_id: rec.tool_call_id.clone(),
tool_calls,
model: rec.model.clone(),
status,
reasoning_content: rec.reasoning_content.clone(),
timestamp: rec.timestamp,
}
}
/// ChatMessage → AiMessageRecord(写路径:save_conversation 全量重写 ai_messages 用)。
///
/// 调用方传入 conversation_id + seq(对话内序号,按索引)。
/// - id None 兜底:`msg_{conv}_{seq}`(老数据无 id 时,落库时补主键)
/// - role 枚举 → 小写字符串
/// - parts/tool_calls → JSON 字符串
/// - status None → "active"(归一化,对齐 record 默认列值)
/// - created_at 调用方传入(save_conversation 用对话 created_at 或 now)
pub fn message_to_record(
msg: &ChatMessage,
conversation_id: &str,
seq: i64,
created_at: &str,
) -> AiMessageRecord {
let id = msg.id.clone().unwrap_or_else(|| format!("msg_{}_{}", conversation_id, seq));
let parts = msg
.parts
.as_ref()
.map(|p| serde_json::to_string(p).unwrap_or_default());
let tool_calls = msg
.tool_calls
.as_ref()
.map(|t| serde_json::to_string(t).unwrap_or_default());
let status = msg.status.clone().unwrap_or_else(|| "active".to_string());
AiMessageRecord {
id,
conversation_id: conversation_id.to_string(),
seq,
role: role_to_str(&msg.role).to_string(),
content: msg.content.clone(),
parts,
tool_call_id: msg.tool_call_id.clone(),
tool_calls,
model: msg.model.clone(),
status,
reasoning_content: msg.reasoning_content.clone(),
timestamp: msg.timestamp,
created_at: created_at.to_string(),
}
}
// ============================================================
// 对话管理
// ============================================================
@@ -61,7 +180,7 @@ pub async fn ai_conversation_create(
};
if let Some(ref oc) = old_conv {
if old_has_msgs {
save_conversation(&state.ai_session, &state.db, oc.as_str(), None, None).await;
save_conversation(&state.ai_session, &state.db, oc.as_str(), None, None, true).await;
}
}
@@ -130,10 +249,36 @@ pub async fn ai_conversation_switch(
.map_err(err_str)?
.ok_or_else(|| format!("对话不存在: {}", conversation_id))?;
let messages: Vec<ChatMessage> = serde_json::from_str(&record.messages)
.map_err(|e| format!("解析消息失败: {}", e))?;
// F-260619-03 消息拆分存储(批次 A 读路径):优先从 ai_messages 表加载(每条消息一行),
// 替代旧 ai_conversations.messages 整对话 JSON 反序列化。映射 Vec<AiMessageRecord>
// → Vec<ChatMessage> → restore_from_messages。
//
// fallback 兜底:ai_messages 表为空(返空 Vec)但旧 messages JSON 列非空 `[]`
// (老库未迁移 / 坏数据 / 批次 B 写路径尚未上线时的新对话)→ 回退读 messages JSON + warn。
// 双向兼容:批次 B 上线后写双轨,读永远先走 ai_messages;迁移未跑的老对话走 fallback。
let records = state.ai_messages.list_by_conversation(&conversation_id).await
.map_err(err_str)?;
let messages: Vec<ChatMessage> = if !records.is_empty() {
records.iter().map(record_to_message).collect()
} else {
// 表空 → fallback 旧 messages JSON 列(若也空则空 Vec,空对话合法)
let has_legacy = record.messages != "[]" && !record.messages.is_empty();
if has_legacy {
tracing::warn!(
"ai_messages 表为空但旧 messages JSON 非空,回退读 JSON 列(conv_id={}): \
老库未迁移或批次 B 写路径未上线",
conversation_id
);
serde_json::from_str(&record.messages)
.map_err(|e| format!("解析消息失败: {}", e))?
} else {
Vec::new()
}
};
let messages_json = record.messages.clone();
// 由 records/chat_messages 重序列化返回前端(前端契约不变,仍吃 JSON 字符串)。
let messages_json = serde_json::to_string(&messages)
.map_err(|e| format!("序列化消息失败: {}", e))?;
let title = record.title.clone();
// B-260617-17 续:历史会话 title 空(显"新对话")→ 切入后触发重新生成(用户诉求)。
// 含 "新对话" 占位(Some 但未生成):title.rs ensure :40 同步排除"新对话"占位不跳过,
@@ -308,8 +453,22 @@ pub async fn ai_conversation_export(
.map_err(err_str)?
.ok_or_else(|| format!("对话不存在: {}", conversation_id))?;
let messages: Vec<ChatMessage> = serde_json::from_str(&record.messages)
.map_err(|e| format!("解析消息失败: {}", e))?;
// F-260619-03 批次 A:导出读路径同样切读 ai_messages 表(与 switch 一致),
// fallback 旧 messages JSON 列(老库未迁移/坏数据)。空对话 → 空 messages。
let records = state.ai_messages.list_by_conversation(&conversation_id).await
.map_err(err_str)?;
let messages: Vec<ChatMessage> = if !records.is_empty() {
records.iter().map(record_to_message).collect()
} else if record.messages != "[]" && !record.messages.is_empty() {
tracing::warn!(
"导出:ai_messages 表为空但旧 messages JSON 非空,回退读 JSON 列(conv_id={})",
conversation_id
);
serde_json::from_str(&record.messages)
.map_err(|e| format!("解析消息失败: {}", e))?
} else {
Vec::new()
};
let body = match fmt {
"markdown" => {
@@ -352,3 +511,377 @@ pub async fn ai_conversation_export(
Ok(body)
}
// ============================================================
// 测试:ChatMessage ↔ AiMessageRecord 映射 round-trip
// ============================================================
#[cfg(test)]
mod tests {
use super::*;
use df_ai::provider::{ChatMessage, ContentPart, MessageRole, ToolCall, ToolCallFunction};
fn base_msg() -> ChatMessage {
ChatMessage {
id: Some("msg_test_1".into()),
role: MessageRole::User,
content: "你好".into(),
parts: None,
tool_call_id: None,
tool_calls: None,
model: None,
status: None,
reasoning_content: None,
timestamp: Some(1700000000000),
}
}
/// 纯文本消息 round-trip:status None ↔ "active" 归一化
#[test]
fn roundtrip_plain_text_status_none() {
let msg = base_msg();
let rec = message_to_record(&msg, "conv_x", 0, "2026-01-01T00:00:00Z");
assert_eq!(rec.id, "msg_test_1");
assert_eq!(rec.conversation_id, "conv_x");
assert_eq!(rec.seq, 0);
assert_eq!(rec.role, "user");
assert_eq!(rec.status, "active", "None 应归一化为 active");
assert!(rec.parts.is_none());
assert!(rec.tool_calls.is_none());
let back = record_to_message(&rec);
assert_eq!(back.id.as_deref(), Some("msg_test_1"));
assert_eq!(back.status, None, "active 应还原为 None");
assert_eq!(back.content, "你好");
assert_eq!(back.timestamp, Some(1700000000000));
assert!(matches!(back.role, MessageRole::User));
}
/// 多模态消息(parts 含 Image base64)round-trip
#[test]
fn roundtrip_multimodal_parts() {
let mut msg = base_msg();
msg.parts = Some(vec![
ContentPart::Text { text: "看这张图".into() },
ContentPart::Image {
url: None,
base64: Some("iVBORw0KGgo=".into()),
media_type: Some("image/png".into()),
alt: None,
},
]);
let rec = message_to_record(&msg, "conv_m", 1, "ts");
// parts 应序列化成 JSON 字符串
assert!(rec.parts.is_some());
let back = record_to_message(&rec);
assert_eq!(back.parts, msg.parts, "parts Image base64 双向一致");
}
/// assistant 工具调用消息(tool_calls 多条)round-trip
#[test]
fn roundtrip_tool_calls() {
let mut msg = base_msg();
msg.role = MessageRole::Assistant;
msg.model = Some("deepseek-chat".into());
msg.tool_calls = Some(vec![
ToolCall {
id: "call_1".into(),
call_type: "function".into(),
function: ToolCallFunction {
name: "read_file".into(),
arguments: r#"{"path":"a.rs"}"#.into(),
},
},
ToolCall {
id: "call_2".into(),
call_type: "function".into(),
function: ToolCallFunction {
name: "write_file".into(),
arguments: r#"{"path":"b.rs"}"#.into(),
},
},
]);
let rec = message_to_record(&msg, "conv_t", 2, "ts");
assert_eq!(rec.role, "assistant");
assert_eq!(rec.model.as_deref(), Some("deepseek-chat"));
assert!(rec.tool_calls.is_some());
let back = record_to_message(&rec);
assert_eq!(back.tool_calls.as_ref().map(|v| v.len()), Some(2));
// ToolCall 无 PartialEq,经 JSON 字符串比对(round-trip 一致)
assert_eq!(
serde_json::to_string(&back.tool_calls).unwrap(),
serde_json::to_string(&msg.tool_calls).unwrap(),
"tool_calls 多条双向一致"
);
assert_eq!(back.model.as_deref(), Some("deepseek-chat"));
}
/// tool 消息(tool_call_id)round-trip
#[test]
fn roundtrip_tool_message() {
let mut msg = base_msg();
msg.role = MessageRole::Tool;
msg.content = "文件内容...".into();
msg.tool_call_id = Some("call_abc".into());
let rec = message_to_record(&msg, "conv_tool", 3, "ts");
assert_eq!(rec.role, "tool");
assert_eq!(rec.tool_call_id.as_deref(), Some("call_abc"));
let back = record_to_message(&rec);
assert!(matches!(back.role, MessageRole::Tool));
assert_eq!(back.tool_call_id.as_deref(), Some("call_abc"));
}
/// status 各值("truncated"/"compressed")round-trip(非 active 原样保留)
#[test]
fn roundtrip_status_non_active_preserved() {
let mut msg = base_msg();
msg.status = Some("truncated".into());
let rec = message_to_record(&msg, "c", 0, "ts");
assert_eq!(rec.status, "truncated");
let back = record_to_message(&rec);
assert_eq!(back.status.as_deref(), Some("truncated"));
// compressed
msg.status = Some("compressed".into());
let rec = message_to_record(&msg, "c", 0, "ts");
assert_eq!(rec.status, "compressed");
}
/// id=None 兜底:`msg_{conv}_{seq}`
#[test]
fn roundtrip_id_none_fallback() {
let mut msg = base_msg();
msg.id = None;
let rec = message_to_record(&msg, "conv_y", 5, "ts");
assert_eq!(rec.id, "msg_conv_y_5", "id None 应兜底 msg_conv_seq");
// 读回:id 来自 record.id(非空)→ Some
let back = record_to_message(&rec);
assert_eq!(back.id.as_deref(), Some("msg_conv_y_5"));
}
/// reasoning_content round-trip
#[test]
fn roundtrip_reasoning_content() {
let mut msg = base_msg();
msg.role = MessageRole::Assistant;
msg.reasoning_content = Some("思考过程...".into());
let rec = message_to_record(&msg, "c", 0, "ts");
assert_eq!(rec.reasoning_content.as_deref(), Some("思考过程..."));
let back = record_to_message(&rec);
assert_eq!(back.reasoning_content.as_deref(), Some("思考过程..."));
}
/// 完整多字段混合 round-trip(覆盖全字段一致性)
#[test]
fn roundtrip_full_fields() {
let mut msg = base_msg();
msg.role = MessageRole::Assistant;
msg.content = "结果".into();
msg.parts = Some(vec![ContentPart::Text { text: "t".into() }]);
msg.tool_calls = Some(vec![ToolCall {
id: "c1".into(),
call_type: "function".into(),
function: ToolCallFunction {
name: "f".into(),
arguments: "{}".into(),
},
}]);
msg.model = Some("m".into());
msg.status = Some("compressed".into());
msg.reasoning_content = Some("r".into());
msg.timestamp = Some(123);
let rec = message_to_record(&msg, "conv_full", 7, "ts_full");
let back = record_to_message(&rec);
assert_eq!(back.id.as_deref(), Some("msg_test_1"));
assert!(matches!(back.role, MessageRole::Assistant));
assert_eq!(back.content, "结果");
assert_eq!(back.parts, msg.parts);
assert_eq!(
serde_json::to_string(&back.tool_calls).unwrap(),
serde_json::to_string(&msg.tool_calls).unwrap()
);
assert_eq!(back.model.as_deref(), Some("m"));
assert_eq!(back.status.as_deref(), Some("compressed"));
assert_eq!(back.reasoning_content.as_deref(), Some("r"));
assert_eq!(back.timestamp, Some(123));
}
// ============================================================
// F-260619-03 批次 B 端到端 round-trip(写路径全链:映射 → replace → list → 还原)
// ============================================================
//
// 覆盖 save_conversation 写路径核心链路(不依赖 AiSession 夹具):
// Vec<ChatMessage> --message_to_record--> Vec<AiMessageRecord>
// --replace_conversation--> ai_messages 表
// --list_by_conversation--> Vec<AiMessageRecord>
// --record_to_message--> Vec<ChatMessage>
// 验证:长对话 50+ 轮全字段(parts/tool_calls/status/id/timestamp/reasoning_content)一致。
#[tokio::test]
async fn batch_b_save_load_roundtrip_50_rounds_full_fields() {
use df_storage::crud::AiMessageRepo;
use df_storage::db::Database;
let db = Database::open_in_memory().await.expect("open_in_memory");
let repo = AiMessageRepo::new(&db);
// 构造 50 轮 user/assistant/tool 混合消息(共 150 条),覆盖全字段
let conv_id = "conv_rt";
let created_at = "2026-06-20T00:00:00Z";
let original: Vec<ChatMessage> = (0..150)
.map(|i| {
let role = i % 3;
let mut m = match role {
0 => ChatMessage {
id: Some(format!("u_{i}")),
role: MessageRole::User,
content: format!("用户提问 {i},带中文"),
parts: Some(vec![ContentPart::Text { text: format!("片 {i}") }]),
timestamp: Some(1_700_000_000_000 + i),
..base_msg()
},
1 => ChatMessage {
id: Some(format!("a_{i}")),
role: MessageRole::Assistant,
content: format!("助手回答 {i}"),
model: Some("deepseek-chat".into()),
reasoning_content: Some(format!("思考 {i}")),
tool_calls: Some(vec![ToolCall {
id: format!("call_{i}"),
call_type: "function".into(),
function: ToolCallFunction {
name: "read_file".into(),
arguments: format!(r#"{{"path":"{i}.rs"}}"#),
},
}]),
status: if i % 6 == 0 {
Some("compressed".into())
} else {
None
},
timestamp: Some(1_700_000_000_000 + i),
..base_msg()
},
_ => ChatMessage {
id: Some(format!("t_{i}")),
role: MessageRole::Tool,
content: format!("工具结果 {i}:大段内容"),
tool_call_id: Some(format!("call_{}", i - 1)),
timestamp: Some(1_700_000_000_000 + i),
..base_msg()
},
};
m.id = Some(format!("msg_{conv_id}_{i}"));
m
})
.collect();
// 写路径:映射 + replace_conversation(全量重写)
let records: Vec<_> = original
.iter()
.enumerate()
.map(|(seq, m)| message_to_record(m, conv_id, seq as i64, created_at))
.collect();
repo.replace_conversation(conv_id, records)
.await
.expect("replace");
// 读路径:list + 还原
let got_records = repo.list_by_conversation(conv_id).await.expect("list");
assert_eq!(got_records.len(), 150, "应读回全部 150 条");
let restored: Vec<ChatMessage> =
got_records.iter().map(record_to_message).collect();
// 全字段逐一比对
assert_eq!(restored.len(), original.len());
for (i, (orig, back)) in original.iter().zip(restored.iter()).enumerate() {
assert_eq!(back.id, orig.id, "id 不一致 @ {i}");
// role 枚举比对(MessageRole 非 Copy,用 std::mem::discriminant 判变体相等)
assert_eq!(
std::mem::discriminant(&back.role),
std::mem::discriminant(&orig.role),
"role 不一致 @ {i}"
);
assert_eq!(back.content, orig.content, "content 不一致 @ {i}");
assert_eq!(back.parts, orig.parts, "parts 不一致 @ {i}");
assert_eq!(back.tool_call_id, orig.tool_call_id, "tool_call_id 不一致 @ {i}");
// tool_calls 无 PartialEq → JSON 比对
assert_eq!(
serde_json::to_string(&back.tool_calls).unwrap(),
serde_json::to_string(&orig.tool_calls).unwrap(),
"tool_calls 不一致 @ {i}"
);
assert_eq!(back.model, orig.model, "model 不一致 @ {i}");
assert_eq!(back.status, orig.status, "status 不一致 @ {i}");
assert_eq!(
back.reasoning_content, orig.reasoning_content,
"reasoning_content 不一致 @ {i}"
);
assert_eq!(back.timestamp, orig.timestamp, "timestamp 不一致 @ {i}");
}
// seq 顺序校验
for (i, rec) in got_records.iter().enumerate() {
assert_eq!(rec.seq, i as i64, "seq 应连续递增 @ {i}");
assert_eq!(rec.conversation_id, conv_id, "conversation_id 应一致 @ {i}");
}
}
/// save 全量重写后再次 save 变更内容(模拟 compress/edit 后下轮 save 覆盖)
#[tokio::test]
async fn batch_b_save_overwrite_reflects_inmemory_change() {
use df_storage::crud::AiMessageRepo;
use df_storage::db::Database;
let db = Database::open_in_memory().await.expect("open_in_memory");
let repo = AiMessageRepo::new(&db);
let conv_id = "conv_overwrite";
let created_at = "ts";
// 第一轮:3 条消息
let v1: Vec<ChatMessage> = (0..3)
.map(|i| ChatMessage {
id: Some(format!("m_{i}")),
role: MessageRole::User,
content: format!("v1_{i}"),
..base_msg()
})
.collect();
let recs: Vec<_> = v1
.iter()
.enumerate()
.map(|(s, m)| message_to_record(m, conv_id, s as i64, created_at))
.collect();
repo.replace_conversation(conv_id, recs).await.expect("save v1");
// 第二轮:内存变化——压缩成 1 条(status=compressed)+ 删除 2 条 + 新增 1 条
let v2: Vec<ChatMessage> = vec![
ChatMessage {
id: Some("m_summary".into()),
role: MessageRole::Assistant,
content: "压缩摘要".into(),
status: Some("compressed".into()),
..base_msg()
},
ChatMessage {
id: Some("m_new".into()),
role: MessageRole::User,
content: "压缩后新提问".into(),
..base_msg()
},
];
let recs2: Vec<_> = v2
.iter()
.enumerate()
.map(|(s, m)| message_to_record(m, conv_id, s as i64, created_at))
.collect();
repo.replace_conversation(conv_id, recs2).await.expect("save v2");
// 读回:应完全是 v2,v1 的 3 条已删
let got = repo.list_by_conversation(conv_id).await.expect("list");
assert_eq!(got.len(), 2, "v2 全量重写后应只剩 2 条");
assert_eq!(got[0].id, "m_summary");
assert_eq!(got[0].status, "compressed");
assert_eq!(got[1].id, "m_new");
assert_eq!(got[1].content, "压缩后新提问");
}
}

View File

@@ -9,6 +9,7 @@ use df_storage::db::Database;
use crate::commands::now_millis;
use super::commands::message_to_record;
use super::AiSession;
/// Token 用量累加器(agent loop 生命周期内各轮叠加)
@@ -130,15 +131,22 @@ pub(crate) fn truncate_parts_for_persist(parts: &[df_ai::provider::ContentPart])
/// 保存对话到数据库(按 conv_id 写库,不受 active_conversation_id 切换影响)
///
/// 写 messages + updated_at + 累加 token 用量 + 首次落库的 model标题由 ensure_conversation_title 单独生成。
/// 写 ai_messages(每条消息一行,F-260619-03 拆分存储)+ updated_at + 累加 token 用量 +
/// 首次落库的 model;标题由 ensure_conversation_title 单独生成。
/// token 走累加模式:upsert 读旧值叠加,保证审批暂停→恢复跨 loop 实例不覆盖丢失。
/// model 仅首次落库写入 + 旧记录缺值时补填(不覆盖历史已存值,兼容本次改造前的老对话)。
///
/// F-260619-03 批次 B:写路径切 ai_messages(全量重写,replace_conversation)。
/// 内存 ContextManager 是运行时真相源,save 是同步点——每轮全量回写 ai_messages,
/// compress/replace/edit/clear_context 内存改 status/content 后下轮 save 自动覆盖。
/// 旧 messages JSON 列**保留不写**(作备份,防 reload fallback 读旧脏数据),不赋新值也不置空。
pub(crate) async fn save_conversation(
session_arc: &Arc<Mutex<AiSession>>,
db: &Arc<Database>,
conv_id: &str,
usage: Option<&df_ai::provider::TokenUsage>,
model: Option<&str>,
touch_updated_at: bool,
) {
// 取 messages + 懒创建首次落库所需的 provider_id/created_at
// 工具结果(content)超 50KB 时截断头尾各 ~20KB + 中段标注,防大体量结果(read_file 1MB 洞 /
@@ -149,7 +157,7 @@ pub(crate) async fn save_conversation(
// loop 内 save 由 run_agentic_loop 入参 conv_id 透传;IPC 路径(commands.rs)save 也传 conv_id。
// conv() 惰性建:save 路径 conv 必然已建(send/regenerate/edit/switch 均先 conv());若极端
// 未建(如启动恢复无 live conv),conv() 建空 PerConvState,save 空 messages(幂等不污染)。
let (messages_json, provider_id, created_at) = {
let (persist_msgs, provider_id, created_at) = {
let mut session = session_arc.lock().await;
let mut msgs = session.conv(conv_id).messages.all_messages_clone();
for m in &mut msgs {
@@ -161,18 +169,34 @@ pub(crate) async fn save_conversation(
}
}
(
serde_json::to_string(&msgs).unwrap_or_else(|_| "[]".to_string()),
msgs,
session.active_provider_id.clone(),
session.active_conv_created_at.clone(),
)
};
// F-260619-03 批次 B:映射 Vec<ChatMessage> → Vec<AiMessageRecord>(带 seq 索引 + conv_id)
// 全量重写 ai_messages(单事务 DELETE + INSERT OR IGNORE,原子无中间空窗)。
// created_at 用对话级 created_at(老对话 None 时 now 兜底),保证消息创建时间与对话一致。
let now = now_millis();
let msg_created_at = created_at.clone().unwrap_or_else(|| now.clone());
let records: Vec<df_storage::models::AiMessageRecord> = persist_msgs
.iter()
.enumerate()
.map(|(seq, m)| message_to_record(m, conv_id, seq as i64, &msg_created_at))
.collect();
let conv_repo = AiConversationRepo::new(db);
let msg_repo = df_storage::crud::AiMessageRepo::new(db);
match conv_repo.get_by_id(conv_id).await {
Ok(Some(mut rec)) => {
// 已落库:更新 messages + updated_at;token 累加(读旧值+新值,跨 loop 实例防覆盖)
rec.messages = messages_json;
rec.updated_at = now_millis();
// 已落库:更新对话元数据(token/model/updated_at);messages JSON 列**不赋新值**(保留旧值作备份)。
// updated_at 按 touch_updated_at 条件改(用户活跃=true 反映最后活跃,
// 系统摘要/压缩=false 防会话时间分组"昨天→今天"跳变);
// token 累加(读旧值+新值,跨 loop 实例防覆盖)。
if touch_updated_at {
rec.updated_at = now_millis();
}
if let Some(u) = usage {
rec.prompt_tokens = accumulate_tokens(rec.prompt_tokens, u.prompt_tokens);
rec.completion_tokens = accumulate_tokens(rec.completion_tokens, u.completion_tokens);
@@ -187,13 +211,21 @@ pub(crate) async fn save_conversation(
if !list.iter().any(|x| x == m) { list.push(m.to_string()); }
rec.models = Some(serde_json::to_string(&list).unwrap_or_else(|_| "[]".to_string()));
}
// update_full 仍写 messages 列(保留旧值,本批不改 messages 字段),写元数据 + updated_at
if let Err(e) = conv_repo.update_full(&rec).await {
tracing::warn!("更新对话失败 {conv_id}: {e}");
tracing::warn!("更新对话元数据失败 {conv_id}: {e}");
}
// 消息拆分存储:全量重写 ai_messages
if let Err(e) = msg_repo.replace_conversation(conv_id, records).await {
tracing::warn!("全量重写 ai_messages 失败 {conv_id}: {e}");
}
}
Ok(None) => {
// 懒创建首次落库(此为空对话不落库的落库点:走到这里 messages 必非空)
let now = now_millis();
// messages JSON 列首次落库也写(兼容未跑迁移的老库 fallback 读路径),
// 同时写 ai_messages(新读路径真相源)。
let messages_json = serde_json::to_string(&persist_msgs).unwrap_or_else(|_| "[]".to_string());
let conv_created = created_at.unwrap_or_else(|| now.clone());
let rec = df_storage::models::AiConversationRecord {
id: conv_id.to_string(),
title: None,
@@ -205,12 +237,17 @@ pub(crate) async fn save_conversation(
pinned: false,
prompt_tokens: usage.map(|u| u.prompt_tokens as i64),
completion_tokens: usage.map(|u| u.completion_tokens as i64),
created_at: created_at.unwrap_or_else(|| now.clone()),
created_at: conv_created.clone(),
updated_at: now,
};
if let Err(e) = conv_repo.insert(rec).await {
tracing::warn!("落库对话失败 {conv_id}: {e}");
}
// 消息拆分存储:首次落库同样全量写 ai_messages
// (records 已用 conv_created 作 created_at,与对话记录一致)
if let Err(e) = msg_repo.replace_conversation(conv_id, records).await {
tracing::warn!("首次写 ai_messages 失败 {conv_id}: {e}");
}
}
Err(e) => tracing::warn!("读取对话 {conv_id} 失败: {e}"),
}

View File

@@ -68,7 +68,12 @@ fn system_prompt_parts(lang: &str) -> (&'static str, &'static str, &'static str)
- Briefly explain your intent before executing actions\n\
- Ask for clarification if the user's intent is unclear\n\
- Prefer using tools to complete actions rather than just describing steps\n\
- When a tool call fails, clearly tell the user it failed and why. Never disguise a fallback action as the original intent's success (e.g. don't write to description to fake a directory binding), and never falsely report success\n",
- When a tool call fails, clearly tell the user it failed and why. Never disguise a fallback action as the original intent's success (e.g. don't write to description to fake a directory binding), and never falsely report success\n\
## Focus\n\
- Always center your response on the core goal of the user's current request; the previous round's topic is only background, not the current task.\n\
- When the user switches topics (a clear new intent), follow the latest request; do not drag the old topic into the new answer.\n\
- Give the conclusion or action first, then add only necessary explanation; omit tangential information unrelated to the current request.\n\
- Do not proactively expand context (files/data) that is irrelevant to the current request.\n",
"\n## Current Projects\n",
"\n## Current Tasks\n",
),
@@ -86,7 +91,12 @@ fn system_prompt_parts(lang: &str) -> (&'static str, &'static str, &'static str)
- 执行操作前简要说明你的意图\n\
- 如果不确定用户意图,先提问\n\
- 优先使用工具完成操作,而不是只描述步骤\n\
- 工具调用失败时必须明确告知用户失败原因,严禁用替代操作冒充原意图成功(如绑定目录失败不得改写描述冒充已绑定),也绝不谎报成功\n",
- 工具调用失败时必须明确告知用户失败原因,严禁用替代操作冒充原意图成功(如绑定目录失败不得改写描述冒充已绑定),也绝不谎报成功\n\
## 聚焦准则\n\
- 始终围绕用户当前请求的核心目标回答;上一轮的主题只是背景,不是当前任务。\n\
- 用户切换话题(明显的新意图)时,以最新请求为准,不要把旧话题带进新回答。\n\
- 回答先给结论/动作,再补必要的解释;与当前请求无关的扩展信息省略。\n\
- 需要的上下文(文件/数据)若与当前请求无关,不要主动展开。\n",
"\n## 当前项目\n",
"\n## 当前任务\n",
),
@@ -168,6 +178,9 @@ pub(crate) fn compress_prompt(lang: &str) -> &'static str {
- Be concise; prefer bullet points.\n\
- Drop small talk and transient pleasantries; keep only technically load-bearing facts.\n\
- Preserve file paths, identifiers, and error messages verbatim.\n\
- Must preserve core topic words, entity names, and technical terms the user \
repeatedly mentions; they are anchors for continuing the conversation and losing \
them breaks the context thread.\n\
- Do NOT invent facts not present in the conversation.\n\
- Output the four sections only, no preamble or extra commentary.",
_ => "你是对话总结器。请把以下对话压缩为结构化摘要,保留继续推进工作所必需的关键上下文。\
@@ -189,7 +202,84 @@ pub(crate) fn compress_prompt(lang: &str) -> &'static str {
- 简洁,优先用要点。\n\
- 去掉寒暄、过渡性客套,只保留技术上有价值的事实。\n\
- 文件路径、标识符、错误信息等照原样保留。\n\
- 必须保留用户反复提及的核心主题词、实体名、技术名词(它们是对话续接的锚点,丢失会致上下文断裂)。\n\
- 不要编造对话中没有的事实。\n\
- 只输出上述四段内容,不要前言、解释或额外评论。",
}
}
#[cfg(test)]
mod tests {
use super::*;
// 改进1 系统提示聚焦段:中文版含「聚焦准则」独立段
#[test]
fn system_prompt_zh_has_focus_section() {
let (prefix, _, _) = system_prompt_parts("zh-CN");
assert!(
prefix.contains("聚焦准则"),
"中文 system prompt 应含聚焦准则段,实际: {}",
prefix
);
// 独立段标题,非稀释在行为准则里
assert!(prefix.contains("## 聚焦准则"));
// 核心条款抽样
assert!(prefix.contains("始终围绕用户当前请求的核心目标"));
assert!(prefix.contains("也绝不谎报成功"));
}
// 改进1 系统提示聚焦段:英文版含「Focus」独立段
#[test]
fn system_prompt_en_has_focus_section() {
let (prefix, _, _) = system_prompt_parts("en");
assert!(
prefix.contains("Focus"),
"英文 system prompt 应含 Focus 段,实际: {}",
prefix
);
assert!(prefix.contains("## Focus"));
assert!(prefix.contains("never falsely report success"));
assert!(prefix.contains("core goal of the user's current request"));
}
// 兜底:lang 未匹配回落中文,聚焦段仍存在
#[test]
fn system_prompt_unknown_lang_falls_back_zh_with_focus() {
let (prefix, _, _) = system_prompt_parts("fr");
assert!(prefix.contains("聚焦准则"));
}
#[test]
fn compress_prompt_unaffected_by_focus_addition() {
// 压缩 prompt 是独立函数,聚焦段改动不应波及
let zh = compress_prompt("zh-CN");
assert!(zh.contains("意图"));
let en = compress_prompt("en");
assert!(en.contains("Intent"));
}
// 改进3 A: 压缩 prompt 主题保留段(锚点词防上下文断裂)
#[test]
fn compress_prompt_zh_keeps_topic_anchor_clause() {
let zh = compress_prompt("zh-CN");
assert!(
zh.contains("主题词"),
"中文 compress_prompt 应含主题词/锚点保留要求,实际: {}",
zh
);
assert!(zh.contains("锚点"));
assert!(zh.contains("上下文断裂"));
}
#[test]
fn compress_prompt_en_keeps_topic_anchor_clause() {
let en = compress_prompt("en");
assert!(
en.contains("anchors"),
"英文 compress_prompt 应含 anchors 保留要求,实际: {}",
en
);
assert!(en.contains("topic words"));
assert!(en.contains("breaks the context"));
}
}

View File

@@ -177,7 +177,15 @@ async fn generate_title_via_llm(
// F-09 B 批5: per_conv 改 HashMap<conv_id>,标题针对本对话,用 conv_id 共享限流槽。
let _global_permit = llm_concurrency.acquire_global().await;
let _per_conv_permit = llm_concurrency.acquire_per_conv(conv_id).await;
let resp = provider.complete(request).await.ok()?;
let resp = match provider.complete(request).await {
Ok(r) => r,
Err(e) => {
// 诊断:标题 LLM complete 失败原因(网络/模型/超时)。ensure :121 防御已保留 extract 兜底,
// 此日志助定位"为何标题未 LLM 精炼生成"(complete 失败 vs resp.text 空走 clean 兜底"新对话")。
tracing::warn!("标题 LLM complete 失败(conv_id={}, model={}): {}", conv_id, model, e);
return None;
}
};
Some(clean_title(&resp.text))
}

View File

@@ -11,9 +11,9 @@ use tokio::sync::{Mutex, RwLock, Semaphore};
use df_ai::ai_tools::AiToolRegistry;
use df_storage::crud::{
AiConversationRepo, AiProviderRepo, AiToolExecutionRepo, IdeaRepo, KnowledgeEventsRepo,
KnowledgeRepo, NodeExecutionRepo, ProjectRepo, ReleaseRepo, SettingsRepo, TaskRepo,
WorkflowRepo,
AiConversationRepo, AiMessageRepo, AiProviderRepo, AiToolExecutionRepo, IdeaRepo,
KnowledgeEventsRepo, KnowledgeRepo, NodeExecutionRepo, ProjectRepo, ReleaseRepo,
SettingsRepo, TaskRepo, WorkflowRepo,
};
use df_storage::db::Database;
use df_workflow::eventbus::EventBus;
@@ -261,6 +261,8 @@ pub struct AppState {
pub ai_providers: AiProviderRepo,
/// AI 对话历史 Repo
pub ai_conversations: AiConversationRepo,
/// AI 消息 Repo(F-260619-03 消息拆分存储:ai_messages 表,读路径批次 A 切读用)
pub ai_messages: AiMessageRepo,
/// AI 工具执行审计 Repo
pub ai_tool_executions: AiToolExecutionRepo,
/// AI 工具注册表
@@ -484,6 +486,7 @@ impl AppState {
node_executions: NodeExecutionRepo::new(&db),
ai_providers: AiProviderRepo::new(&db),
ai_conversations: AiConversationRepo::new(&db),
ai_messages: AiMessageRepo::new(&db),
ai_tool_executions: AiToolExecutionRepo::new(&db),
ai_session: Arc::new(Mutex::new(AiSession::new())),
knowledge: KnowledgeRepo::new(&db),