重构: F-09批2 调用方迁移per-conv(双写桥接共存期)

- agentic.rs: GeneratingGuard加conv_id + 三处退出校验改conv存在性(决策e) + stop_flag/notify/messages/iteration per-conv + try_continue加conv_id参数
- commands.rs: IPC写路径双写桥接(per_conv新真相源+顶层双写,批4删) + ai_is_generating读per_conv fallback顶层
- audit/conversation/title/knowledge_inject/lib.rs: per_conv读写迁移
- conv()/conv_read()去allow(批2有调用方)
共存期双写: per_conv主+顶层双写兼容批4前IPC,批4签名改后删顶层
主代兜底: cargo check --workspace 0 + test 96 passed + grep三处退出/GeneratingGuard/双写印证
This commit is contained in:
2026-06-19 02:01:06 +08:00
parent 2c8764abee
commit 2a4d745c5b
9 changed files with 578 additions and 202 deletions

View File

@@ -694,7 +694,7 @@
**待修项回流 todo**: **无** 🔴/🟡 项
### CR-260619-03 SMELL-P1-9 crud.rs按表拆分 + F-09批1 PerConvState(agent crud-split + f09-batch1 实施·主代统一兜底核验) — 🟡 待审
### CR-260619-03 SMELL-P1-9 crud.rs按表拆分 + F-09批1 PerConvState(agent crud-split + f09-batch1 实施·主代统一兜底核验) — ✅ 已审(PASS·🟡1 WATCH·批2中间态)
- **范围**(2 独立任务攒批):
- **SMELL-P1-9 crud.rs 拆分**(agent crud-split):`crud.rs` 2212 行 → `crud/` 6 文件(mod.rs 宏+工具+re-export / settings / project_repo / task_repo / conversation_repo / idea_repo)。12 Repo + from_row 按表域归属;re-export `pub use {conversation,idea,project,settings,task}_repo::*` 零调用方改动;宏 `pub(crate) use impl_repo` + 子模块 `use super::impl_repo`;基线测试 all_known_tables_have_column_whitelist(12 表)+ all_repos_constructible_in_memory(13 Repo new 不 panic)。
@@ -703,6 +703,39 @@
- **审查要点(供审查 agent)**:① SMELL re-export 完整性(全仓 `df_storage::crud::XxxRepo` 路径不变,cargo --workspace 0 error 印证);② 宏可见性(pub(crate) use + 子模块 use super + helper import);③ 12 Repo/from_row 归属表正确(idea/project/task/conversation/settings 域);④ 基线测试锁回归(12 表白名单 + 13 Repo 构造);⑤ F-09 批1 PerConvState 字段初值对齐 AiSession::new + conv/conv_read 逻辑 + #[allow(dead_code)] 共存期标注合理(批2 迁移后移除);⑥ F-09 批1 纯新增无行为变化(顶层字段未迁移)。
- **关联**:todo SMELL-P1-9 销账 / F-09 阶段2 批1(批2-8 待)。
**复审结论(2026-06-19·独立 grep/read 核验源码当前形态·commit 2c8764a 已落地·不跑 cargo 避批2中间态误报)**: ✅ **PASS** — 🔴0 🟡0 ⚪1
**逐项核验表(file:line 佐证 + 判定)**:
| 项 | 核验点 | 佐证 | 判定 |
|---|---|---|---|
| A-拆分 | crud.rs 2212 行 → 6 文件 | `crates/df-storage/src/crud/` 实存 mod.rs/settings.rs/project_repo.rs/task_repo.rs/conversation_repo.rs/idea_repo.rs 6 文件;旧 crud.rs 已删(git show 2c8764a -2212) | ✅ |
| A-re-export | `pub use {conversation,idea,project,settings,task}_repo::*` 5 路全在 | `mod.rs:21-25` 五条 `pub use` 全部就位;调用方 `df_storage::crud::{TaskRepo,ProjectRepo,IdeaRepo,AiProviderRepo,AiConversationRepo,AiToolExecutionRepo,KnowledgeRepo,KnowledgeEventsRepo,KnowledgeEventsRepo,WorkflowRepo,...}` 路径 grep 全部命中且未改(state.rs:13 / workflow.rs:15 / audit.rs:11 / conversation.rs:7 / title.rs:15 / knowledge_inject.rs:15 / knowledge_timeline.rs:11 / agentic.rs:352 / tool_registry.rs 多处) | ✅ |
| A-宏可见性 | `pub(crate) use impl_repo` + 子模块 `use super::impl_repo` | `mod.rs:191 pub(crate) use impl_repo`(宏定义 :40-187 在前·textual scope 注入注释 :189-190 说明顺序约束);子模块 import 全在:`project_repo.rs:15`/`task_repo.rs:13`/`conversation_repo.rs:13`/`idea_repo.rs:13``use super::impl_repo` | ✅ |
| A-helper import | 子模块用 super::{now_millis_str,storage_err,validate_column_name,...} | project_repo.rs:16 `use super::{normalize_stored_path, now_millis_str, storage_err, validate_column_name}` / task_repo.rs:14 / conversation_repo.rs:14 / idea_repo.rs:14;settings.rs:12 `use super::{now_millis_str, storage_err}`(KV 不需 validate/Arc/Mutex 外借,自含);mod.rs 工具 storage_err:199/normalize_stored_path:206/now_millis_str:220 全 pub(crate) 可见 | ✅ |
| A-12Repo归属 | 13 Repo 按 5 域正确分布 | settings.rs:22 SettingsRepo(1,手写 KV)·project_repo.rs 5(ProjectRepo:100/BranchRepo:256/ReleaseRepo:283/WorkflowRepo:310/NodeExecutionRepo:337)·task_repo.rs 1(TaskRepo:45)·conversation_repo.rs 3(AiProviderRepo:114/AiConversationRepo:153/AiToolExecutionRepo:184)·idea_repo.rs 3(IdeaRepo:118/KnowledgeRepo:147/KnowledgeEventsRepo:400)=**13 Repo** 全在;from_row 12 个(宏生成 Repo 各一,SettingsRepo 无 from_row)归属域正确(project_from_row/branch/release/workflow/node_execution/task/ai_provider/ai_conversation/ai_tool_execution/idea/knowledge/knowledge_event) | ✅ |
| A-基线测试 | 12 表白名单 + 13 Repo 构造 | `mod.rs:243-255 all_known_tables_have_column_whitelist` 12 表全断言(ideas/projects/tasks/releases/branches/workflow_executions/node_executions/ai_providers/ai_conversations/ai_tool_executions/knowledges/knowledge_events)·`mod.rs:259-276 all_repos_constructible_in_memory` 13 Repo(SettingsRepo/IdeaRepo/ProjectRepo/TaskRepo/BranchRepo/ReleaseRepo/WorkflowRepo/NodeExecutionRepo/AiProviderRepo/AiConversationRepo/AiToolExecutionRepo/KnowledgeRepo/KnowledgeEventsRepo)逐 new 不 panic | ✅ |
| A-白名单单一源 | allowed_columns_for/validate_column_name/is_allowed_column 唯一在 settings.rs | settings.rs:119 allowed_columns_for / :189 validate_column_name(pub(crate)) / :201 is_allowed_column(pub);调用方 `df_storage::crud::is_allowed_column`(tool_registry.rs:452/555)路径未破 | ✅ |
| B-PerConvState字段 | 9 字段逐字对齐 AiSession::new | `mod.rs:571-590` struct 9 字段;`:605-617 new()` 初值:messages(ContextManager::new default)=AiSession:386 / generating:false=:391 / stop_flag(Arc AtomicBool false)=:393 / notify(Arc Notify::new())=:394 / iteration_used:0=:395 / agent_language:None=:392 / model_override:None=:396 / session_trust(HashSet::new())=:397 / **created_at:None(新增字段,注释 :604 对齐 AiSession.active_conv_created_at:None :389)** 全对齐 | ✅ |
| B-conv访问器 | entry().or_insert_with 惰性创建 | `mod.rs:453-457 conv(&mut self,conv_id)` body `self.per_conv.entry(conv_id.to_string()).or_insert_with(PerConvState::new)` — 惰性创建语义正确,已存在返同一实例 | ✅ |
| B-conv_read访问器 | .get() 只读不创建 | `mod.rs:466-468 conv_read(&self,conv_id)` body `self.per_conv.get(conv_id)` — 返 Option<&PerConvState>,无写入路径,只读语义正确 | ✅ |
| B-allow标注 | PerConvState/per_conv/conv/conv_read 共存期 #[allow(dead_code)] | PerConvState struct:570 / per_conv 字段:355 / conv:452 / conv_read:465 四处均标 `#[allow(dead_code)]` 注释「F-09 B 批1 共存期:批2 迁移承接(0 调用方)」 | ✅ |
| B-3单测 | new 初值 + 惰性创建复用 + conv_read 不创建 | `mod.rs:477 test_per_conv_state_new`(8 断言含 stop_flag SeqCst load) / `:497 test_conv_lazy_create`(首次创建 len=1 + 同 id 指针相等复用 len 仍=1 + 不同 id 新建 len=2) / `:525 test_conv_read_none`(未创建返 None + 不触发创建 per_conv 仍空) 三测覆盖 | ✅ |
| B-纯新增 | 顶层字段仍是真相源,per_conv 初始空 | AiSession::new:398 `per_conv: HashMap::new()`(初始空);批1 仅加 struct/字段/访问器,无调用方迁移(conv/conv_read 0 调用方) | ✅(见 WATCH·批2 中间态) |
**对抗核验印证**:
- **re-export 路径不变?** ✅ 全仓 grep `df_storage::crud` 命中 9 调用文件,所有 `XxxRepo`/`is_allowed_column` 路径形式未变,re-export 5 路全覆盖零调用方改动。
- **宏 textual scope 顺序?** ✅ mod.rs 宏定义(:40-187)在 `pub(crate) use impl_repo`(:191)之前,`pub(crate) use``mod` 声明(:15-19)之前,顺序符合 Rust textual scope 注入要求(注释 :189-190 显式说明);4 子模块 `use super::impl_repo` 全部命中。
- **13 Repo 计数?** ✅ 实测 1(settings 手写)+5(project)+1(task)+3(conversation)+3(idea)=13,与基线测试 :263-275 逐行 new 一一对齐。
- **PerConvState created_at 新增字段是否破坏对齐?** ✅ AiSession 无独立 created_at 顶层字段(用 active_conv_created_at:None :389),PerConvState.created_at:None :615 注释 :604 显式声明「批1 新增字段,AiSession 现有 active_conv_created_at 同语义」——对齐声明而非逐字复制,合理(语义等价,字段名差异因 AiSession 持 active 概念而 PerConvState 持会话创建时间)。
- **批1 是否纯新增?** ✅ AiSession::new 仍构造全部顶层字段(:386-397 messages/generating/agent_language/stop_flag/notify/iteration_used/model_override/session_trust),per_conv 初始空 HashMap,conv/conv_read 0 调用方——批1 范围严格纯新增。
**⚠️ 批2 中间态说明(非批1 问题·审查范围外记录)**:
- 核验时发现 `mod.rs:355-380` 顶层 8 字段已标 `#[allow(dead_code)]` 注释「F-09 B 批2 共存期死字段(批3 删)」,且 AiSession::new:386-397 仍构造这些字段——说明 **F-09 批2 正在进行中**(对应 prompt 开头「批2 agent 正在改 agentic.rs」)。批1 的 per_conv/conv/conv_read/PerConvState 已就位待批2 迁移调用方消费。本次审查范围(commit 2c8764a)仅含批1,批2 中间态不评判。**本次审查全程未跑 cargo**(memory review-batching-worktree-transient 教训:批2 中间态会致 check 快照误报,源码形态核验 > check 快照)。
- **⚪ WATCH-1**: F-09 批2 迁移进行中(agentic.rs/commands.rs/audit.rs/conversation.rs/title.rs/knowledge_inject.rs 调用方迁移 + 顶层字段降级死字段),批2 完成后需独立审查(范围:调用方迁移完整性 + 顶层字段是否真成死字段 + cargo test 单会话全功能回归)。批1 本身 PASS 不受批2 影响。
- **待修项回流 todo**: **无** 🔴/🟡 项(批1 范围内)
---
## 已审归档

View File

@@ -70,20 +70,33 @@ pub const DEFAULT_MAX_AGENT_RETRIES: usize = 3;
///
/// 注:try_continue_agent_loop 不用 guard——其 should_continue=false 路径需保持
/// generating=true(审批等待态),全函数 guard 会误复位;该函数单点 provider-Err 复位保持手动。
///
/// F-260616-09 B 批2:guard 持 `conv_id`,复位改写 `session.conv(&conv_id).generating = false`
/// (per-conv 真相源)。同时**双写顶层 `session.generating = false`** 作共存期桥接 —— 批2
/// 仅迁移 agentic.rs 路径,IPC(ai_is_generating/ai_chat_send)仍读顶层,故 guard 须双写
/// 保证 IPC 读到正确值(否则前端 ai_is_generating 永远 true 卡死发送)。批4 IPC 迁移后
/// 顶层双写移除。
struct GeneratingGuard {
session: Arc<Mutex<AiSession>>,
/// guard 所属会话(loop 启动时快照的 conv_id,来自 run_agentic_loop 入参)。
conv_id: String,
done: bool,
}
impl GeneratingGuard {
fn new(session: Arc<Mutex<AiSession>>) -> Self {
Self { session, done: false }
fn new(session: Arc<Mutex<AiSession>>, conv_id: String) -> Self {
Self { session, conv_id, done: false }
}
/// 显式复位 generating=false。emit 前调用保证顺序。幂等。
///
/// 双写:per_conv.conv_id.generating(新真相源)+ 顶层 generating(共存期 IPC 桥接)。
async fn reset(&mut self) {
if !self.done {
self.session.lock().await.generating = false;
let mut session = self.session.lock().await;
session.conv(&self.conv_id).generating = false;
// F-09 B 批2 桥接:顶层双写,批4 IPC 迁移后移除
session.generating = false;
self.done = true;
}
}
@@ -100,8 +113,12 @@ impl Drop for GeneratingGuard {
fn drop(&mut self) {
if !self.done {
let session = self.session.clone();
let conv_id = self.conv_id.clone();
tauri::async_runtime::spawn(async move {
session.lock().await.generating = false;
let mut s = session.lock().await;
s.conv(&conv_id).generating = false;
// F-09 B 批2 桥接:顶层双写,批4 IPC 迁移后移除
s.generating = false;
});
}
}
@@ -331,7 +348,32 @@ pub(crate) async fn run_agentic_loop(
model_override: Option<String>,
) {
// B-260615-09: generating 状态由 RAII guard 收敛复位(正常 exit 显式 reset;panic/异常 Drop 兜底)
let mut guard = GeneratingGuard::new(session_arc.clone());
// F-260616-09 B 批2:guard 持 conv_id,复位改 per-conv.generating(设计 §4.3)。
let mut guard = GeneratingGuard::new(session_arc.clone(), conv_id.clone());
// F-260616-09 B 批2 入口桥接:loop 启动前确保 per_conv 存在(已存在则保留累积,不存在则建)。
//
// 批2 把所有调用方(IPC commands.rs 写路径 + agentic.rs loop + audit.rs process_tool_calls +
// conversation.rs save + title.rs + knowledge_inject.rs)迁移到 per_conv 真相源。IPC 在 spawn
// loop 前已通过 `session.conv(active_conversation_id).*` 建立 per_conv 并写入初始状态
// (messages push user / generating=true / stop_flag=false / iteration_used=0 等),故 loop
// 入口只需确保 per_conv 存在(防御性 conv() 惰性建,正常路径下已存在)。
//
// 已存在的 per_conv(同 conv 上一轮 loop 留下,如审批暂停后续跑)**保留累积状态**:loop 上一轮
// 在 per_conv.messages 累积的 assistant tool_calls / tool_result 不会因桥接抹掉。审批拒绝/通过
// 时 IPC 直接写 per_conv.messages.replace_tool_result_content,续跑 loop 读 per_conv 拿到正确状态。
//
// guard 语义:loop 启动占用生成态,per_conv.generating=true(IPC spawn 前已置 true,此处幂等确认)。
//
// 单 loop 安全性:批2 阶段是单 active_conversation_id(决策 e 真并发批3+ 落地),无并发 loop 抢
// per_conv 覆盖。批3+ 多 loop 并发时,每 conv 各自 per_conv 条目,互不干扰(本桥接无需改)。
{
let mut session = session_arc.lock().await;
let conv = session.conv(&conv_id);
conv.generating = true;
// 顶层同步(批2 桥接,ai_is_generating 等 IPC 仍读顶层,批4 IPC 迁移后移除)。
session.generating = true;
}
// F-260614-04 / F-260614-04b: 多 Provider 负载均衡池 — 选主 + fallback 候选列表。
//
@@ -443,9 +485,14 @@ pub(crate) async fn run_agentic_loop(
let tool_defs = tools_arc.tool_definitions();
// 停止信号副本stream_llm 与每轮迭代共享读取,避免重复加锁
// notify 同取一份 Arc 引用B-260615-14stream_llm select! 监听 notified() 即时唤醒
// F-260616-09 B 批2:取 per_conv 的 stop_flag/notify(设计 §4.2 :446)。
// conv_id 来源:run_agentic_loop 入参(loop 启动快照,与 guard 一致)。
// 批2 桥接期 per_conv.stop_flag/notify 与顶层共享同一 Arc(入口桥接时顶层 Arc 直接搬入),
// 故 ai_chat_stop 写顶层 stop_flag 与 loop 读 per_conv stop_flag 互见(同一 AtomicBool 实例)。
let (stop_flag, notify) = {
let session = session_arc.lock().await;
(session.stop_flag.clone(), session.notify.clone())
let conv = session.conv_read(&conv_id).expect("[F-09 B 批2] loop 入口桥接后 per_conv 必存在");
(conv.stop_flag.clone(), conv.notify.clone())
};
// token 累加器:loop 生命周期内各轮叠加,退出时传 save_conversation(累加模式落库)
@@ -483,19 +530,28 @@ pub(crate) async fn run_agentic_loop(
// B-260615-11: 旧 loop 污染防护——每轮开始校验对话一致性。
// 用户新建/切换对话后 active_conversation_id 变更,本 loop(conv_id 快照)成陈旧,
// 继续跑会往新对话 push 消息/pending 造成污染。检测到即退出(guard Drop 复位 generating)。
//
// F-260616-09 B 批2(设计 §4.1):退出判据从 `active_conversation_id != conv_id` 改为
// `!per_conv.contains_key(conv_id)`(conv 存在性)。决策 e 真并发下旧 loop 跑自己的 conv
// 不污染他人,active_conversation_id 切换不应让旧 loop 退出;仅当 conv 被删(clear/delete)
// 才退出。批2 阶段单 loop,conv 不会被删,此校验主要保留语义对齐(批3+ 真并发场景生效)。
// conv_id 来源:run_agentic_loop 入参(与 guard/stop_flag 取用同源)。
{
let mut session = session_arc.lock().await;
if session.active_conversation_id.as_deref() != Some(conv_id.as_str()) {
if !session.per_conv.contains_key(&conv_id) {
tracing::warn!(
stale_conv = %conv_id,
active_conv = ?session.active_conversation_id,
"[ai] 对话已切换,旧 loop 退出(B-260615-11)避免污染新对话"
"[ai] conv 已删除,旧 loop 退出(F-260616-09 B 批2/决策 e)"
);
return;
}
// F-260616-11 决策 a: 累计 iteration 计数(「下一轮起算值」= 当前轮+1)。
// 审批等待/达 max 暂停退出时,本字段即「当前轮+1」;审批续跑 ai_approve 读此值作
// start_iteration 透传,实现 iteration 累计不重置(防多次审批反复跑满 max 致 token 失控)。
// F-09 B 批2:写 per_conv.iteration_used + 双写顶层(批2 桥接,ai_approve 续跑读顶层,
// 批4 IPC 迁移后顶层双写移除)。
let conv = session.conv(&conv_id);
conv.iteration_used = iteration + 1;
session.iteration_used = iteration + 1;
}
@@ -528,12 +584,21 @@ pub(crate) async fn run_agentic_loop(
// keyring),api_key 经 df_storage::secret 闭环;summary/error payload/日志均不含 api_key。
// is_compressing 防重入:set_compressing(true/false) 成对(LLM 调用前后均复位)。
// 单轮问答(history_tokens 未超 0.6*budget)不触发,零行为变化。
let prev_compressing = session_arc.lock().await.messages.is_compressing();
// F-260616-09 B 批2:messages 操作改 per_conv(设计 §4.2)。
// conv_id 来源:run_agentic_loop 入参。
let prev_compressing = session_arc.lock().await.conv_read(&conv_id).map(|c| c.messages.is_compressing()).unwrap_or(false);
if !prev_compressing {
// 读触发条件(history_tokens / budget / has_compressible_messages),持锁快照判断。
let (should_compress, protect_start, pre_compress_tokens) = {
let session = session_arc.lock().await;
let mgr = &session.messages;
let mgr = match session.conv_read(&conv_id) {
Some(c) => c,
None => {
tracing::warn!(stale_conv = %conv_id, "[ai] conv 已删除,loop 退出(压缩段入口)");
return;
}
};
let mgr = &mgr.messages;
let protect_start = mgr.len().saturating_sub(PROTECT_COUNT);
let history_tokens = mgr.history_tokens();
let budget = mgr.budget_limit();
@@ -551,15 +616,16 @@ pub(crate) async fn run_agentic_loop(
});
let (active_msgs, lang) = {
let mut session = session_arc.lock().await;
session.messages.set_compressing(true);
let conv = session.conv(&conv_id);
conv.messages.set_compressing(true);
// 读 active 克隆(不改 status):filter is_active,LLM 失败则消息状态完全不变。
let active_msgs: Vec<ChatMessage> = session.messages.messages_mut()
let active_msgs: Vec<ChatMessage> = conv.messages.messages_mut()
[..protect_start]
.iter()
.filter(|t| t.message.is_active())
.map(|t| t.message.clone())
.collect();
let lang = session.agent_language.clone()
let lang = conv.agent_language.clone()
.unwrap_or_else(|| "zh-CN".to_string());
(active_msgs, lang)
};
@@ -586,9 +652,10 @@ pub(crate) async fn run_agentic_loop(
// active_msgs 等价(LLM 调用期间 messages 不变,见上方口径决策注)。
{
let mut session = session_arc.lock().await;
let _compressed = session.messages.compress_old_messages(protect_start);
session.messages.insert_at(0, ChatMessage::system(&summary));
session.messages.set_compressing(false);
let conv = session.conv(&conv_id);
let _compressed = conv.messages.compress_old_messages(protect_start);
conv.messages.insert_at(0, ChatMessage::system(&summary));
conv.messages.set_compressing(false);
}
tracing::info!(
conv_id = %conv_id,
@@ -603,13 +670,13 @@ pub(crate) async fn run_agentic_loop(
}
Ok(None) => {
// 保护区外无 active 可压缩(已全 compressed/archived)→ noop,仅复位 is_compressing。
session_arc.lock().await.messages.set_compressing(false);
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 条)。
session_arc.lock().await.messages.set_compressing(false);
session_arc.lock().await.conv(&conv_id).messages.set_compressing(false);
tracing::warn!(
conv_id = %conv_id,
error = %e,
@@ -628,9 +695,18 @@ pub(crate) async fn run_agentic_loop(
// 构建请求消息(超预算时自动裁剪旧消息,保护工具调用三元组 + 最近 6 条)
// F-260616-13: sys_tokens 已在 loop 外缓存;本块构建的 messages 在本轮重试循环中复用
// (本轮 stream_llm 不持 session_arc、不改 messages,重试无 push 发生,重建等价于复用)。
// F-260616-09 B 批2:build_for_request 改 per_conv.messages(设计 §4.2 :631)。
// conv_id 来源:run_agentic_loop 入参。
let messages = {
let session = session_arc.lock().await;
let (history_msgs, _trimmed) = session.messages.build_for_request(sys_tokens);
let conv = match session.conv_read(&conv_id) {
Some(c) => c,
None => {
tracing::warn!(stale_conv = %conv_id, "[ai] conv 已删除,loop 退出(build_for_request 入口)");
return;
}
};
let (history_msgs, _trimmed) = conv.messages.build_for_request(sys_tokens);
let mut msgs = vec![ChatMessage::system(&system_prompt)];
msgs.extend(history_msgs);
msgs
@@ -781,26 +857,28 @@ pub(crate) async fn run_agentic_loop(
tokens.add(usage.prompt_tokens, usage.completion_tokens);
// 追加 partial assistant 消息(若无 tool_calls 且有文本)
// F-260616-09 B 批2(设计 §4.1 :786):退出校验改 conv 存在性 + push 改 per_conv.messages。
// conv_id 来源:run_agentic_loop 入参。
{
let mut session = session_arc.lock().await;
if session.active_conversation_id.as_deref() != Some(conv_id.as_str()) {
if !session.per_conv.contains_key(&conv_id) {
tracing::warn!(
stale_conv = %conv_id,
active_conv = ?session.active_conversation_id,
"[ai] MidStream 保文后对话已切换,丢弃本轮 push(B-260615-11)"
"[ai] MidStream 保文后 conv 已删除,丢弃本轮 push(F-260616-09 B 批2)"
);
return;
}
if !full_text.is_empty() {
let conv = session.conv(&conv_id);
let mut msg = ChatMessage::assistant(&full_text);
msg.model = Some(resolved_model.clone());
// BUG-260617-12: MidStream 保文也回填 reasoning_content
msg.reasoning_content = round_reasoning_content.clone();
session.messages.push(msg);
conv.messages.push(msg);
// 追加系统提示消息:响应因网络中断不完整(对齐决策 a1 系统提示机制)
let mut notice = ChatMessage::system("⚠ 响应因网络中断不完整,以上为已接收的部分内容。可重新发送以获取完整回复。");
notice.model = Some(resolved_model.clone());
session.messages.push(notice);
conv.messages.push(notice);
}
}
@@ -830,15 +908,16 @@ pub(crate) async fn run_agentic_loop(
// 追加 assistant 消息到历史
let has_tool_calls = !tool_calls_acc.is_empty();
// F-260616-09 B 批2(设计 §4.1 :837):退出校验改 conv 存在性 + push 改 per_conv.messages。
// conv_id 来源:run_agentic_loop 入参。
{
let mut session = session_arc.lock().await;
// B-260615-11: push 前再校验(stream_llm 期间用户可能新建对话)。
// 读端读到被 clear 的空历史不致命,但 push 写回新对话是污染,必须挡。
if session.active_conversation_id.as_deref() != Some(conv_id.as_str()) {
// 决策 e 真并发下:push 前再校验 conv 是否仍存在(被删则退出)。
// 读端读到被 clear 的空历史不致命,但 push 写回已删 conv 是污染,必须挡。
if !session.per_conv.contains_key(&conv_id) {
tracing::warn!(
stale_conv = %conv_id,
active_conv = ?session.active_conversation_id,
"[ai] stream 后对话已切换,丢弃本轮 push(B-260615-11)避免污染新对话"
"[ai] stream 后 conv 已删除,丢弃本轮 push(F-260616-09 B 批2)"
);
return;
}
@@ -855,12 +934,12 @@ pub(crate) async fn run_agentic_loop(
msg.model = Some(resolved_model.clone());
// BUG-260617-12: 回填 reasoning_content 供下一轮请求透传
msg.reasoning_content = last_reasoning_content.clone();
session.messages.push(msg);
session.conv(&conv_id).messages.push(msg);
} else if !full_text.is_empty() {
let mut msg = ChatMessage::assistant(&full_text);
msg.model = Some(resolved_model.clone());
msg.reasoning_content = last_reasoning_content.clone();
session.messages.push(msg);
session.conv(&conv_id).messages.push(msg);
}
}
@@ -989,24 +1068,42 @@ pub(crate) async fn run_agentic_loop(
/// ai_chat_stop/clear/switch 并发改写,属竞态耦合。has_pending=false(全部审批已处理,续生成)
/// 分支:审批已被 remove,改为以剩余 pending_approvals 任一 conversation_id 做一致性校验
/// (此处空,校验通过即沿用全局值,该期 generating=true 且 switch 为 readonly 不并发)。
pub(crate) async fn try_continue_agent_loop(app: &AppHandle, state: &AppState, start_iteration: usize) {
pub(crate) async fn try_continue_agent_loop(
app: &AppHandle,
state: &AppState,
conv_id: &str,
start_iteration: usize,
) {
// BUG-260617-05: 原 5 次独立 lock().await 造成 TOCTOU 竞态——should_continue=true 判出后、
// spawn 前用户点 stop(ai_chat_stop 复位 generating=false),续跑仍按过时快照继续 spawn。
// 修复:单次 lock 取结构化快照(所有续跑判定所需字段),无锁态判定;spawn 前单次 lock 原子
// 重检 generating 仍为 true 才续跑,stop 后中途插入的直接收敛退出。
//
// F-260616-09 B 批2(设计 §4.4):conv_id 改显式入参(替代从 pending/active 推断),
// has_pending 按 conv_id 过滤(决策 e 真并发准备);generating/agent_language/model_override
// 读 per_conv(批2 桥接期与顶层等价,per_conv 未建 fallback 顶层)。
// conv_id 来源:调用方 ai_approve 传 approval.conversation_id;ai_continue_loop/ai_stop_loop
// 传 IPC 参数 conversation_id。
let snap = {
let session = state.ai_session.lock().await;
// pending_approvals 中任一审批的 conversation_id:审批等待态(has_pending)下作为 conv_id 来源,
// 取第一个非空值(同一对话的审批 conversation_id 一致,见 process_tool_calls 写入路径)。
// has_pending:仅本 conv 的未决审批算续跑阻塞(决策 e 真并发准备)。
let has_pending = session.pending_approvals.values()
.any(|a| a.conversation_id.as_deref() == Some(conv_id));
// pending_conv_id 保留(should_continue=false 路径的 emit conv_id 回退逻辑)。
let pending_conv_id = session.pending_approvals.values()
.find_map(|a| a.conversation_id.clone());
let conv = session.conv_read(conv_id);
let is_generating = conv.map(|c| c.generating).unwrap_or(session.generating);
let agent_language = conv.and_then(|c| c.agent_language.clone())
.or_else(|| session.agent_language.clone());
let model_override = conv.and_then(|c| c.model_override.clone())
.or_else(|| session.model_override.clone());
ContinueSnapshot {
is_generating: session.generating,
has_pending: !session.pending_approvals.is_empty(),
is_generating,
has_pending,
pending_conv_id,
active_conversation_id: session.active_conversation_id.clone(),
agent_language: session.agent_language.clone(),
model_override: session.model_override.clone(),
agent_language,
model_override,
}
};
let should_continue = snap.is_generating && !snap.has_pending;
@@ -1017,22 +1114,22 @@ pub(crate) async fn try_continue_agent_loop(app: &AppHandle, state: &AppState, s
// 轮 token 已在前序 AiCompleted/AiApprovalResult 流程落库,此处零 token 上报仅作收敛信号。
if snap.is_generating {
// pending_approvals 非空但 generating 仍 true:转审批态,前端审批态 watchdog 已 clear,不卡
tracing::info!("[ai] try_continue 跳过:仍有待审批,转审批等待态");
tracing::info!(conv_id = %conv_id, "[ai] try_continue 跳过:仍有待审批,转审批等待态");
} else {
// generating 已复位(用户 stop 或前序循环已 emit Completed):补发 AiCompleted 防前端卡住
tracing::info!("[ai] try_continue 跳过:generating 已复位(被 stop/已结束),补发 AiCompleted 清前端 streaming");
tracing::info!(conv_id = %conv_id, "[ai] try_continue 跳过:generating 已复位(被 stop/已结束),补发 AiCompleted 清前端 streaming");
// R-PD-6: 优先用审批所属 conversation_id(审批等待态被 stop 触发,审批仍在 pending_approvals),
// 仅当无任何审批(has_pending=false 且 generating=false)时回退 active_conversation_id
let conv_id = match snap.pending_conv_id.clone() {
// 仅当无任何审批(has_pending=false 且 generating=false)时回退入参 conv_id(批2 显式参数)
let emit_conv_id = match snap.pending_conv_id.clone() {
Some(cid) => cid,
None => snap.active_conversation_id.clone().unwrap_or_default(),
None => conv_id.to_string(),
};
let _ = app.emit("ai-chat-event", AiChatEvent::AiCompleted {
total_tokens: 0,
prompt_tokens: 0,
completion_tokens: 0,
incomplete: None,
conversation_id: Some(conv_id),
conversation_id: Some(emit_conv_id),
});
}
return;
@@ -1046,29 +1143,24 @@ pub(crate) async fn try_continue_agent_loop(app: &AppHandle, state: &AppState, s
// build_provider_for Err 分支处理(同样 emit AiError),此处与之一致。
// 不用 GeneratingGuard:try_continue 的 should_continue=false 路径需保 generating=true(审批等待态),
// 全函数 guard 会误复位。此点单点 provider-Err 复位,语义独立。
// F-09 B 批2:双写复位 per_conv.generating + 顶层 generating(批2 桥接)。
let mut session = state.ai_session.lock().await;
session.conv(conv_id).generating = false;
session.generating = false;
// R-PD-6: 续生成被拒(provider 缺失)回退 conv_id 优先审批所属;无审批再读全局。
let conv_id = session.pending_approvals.values()
.find_map(|a| a.conversation_id.clone())
.or_else(|| session.active_conversation_id.clone())
.unwrap_or_default();
drop(session);
tracing::warn!(error = %e, "[ai] try_continue 失败:无可用 provider");
tracing::warn!(conv_id = %conv_id, error = %e, "[ai] try_continue 失败:无可用 provider");
let _ = app.emit("ai-chat-event", AiChatEvent::AiError {
error: e,
// 无可用 provider(配置丢失/全删):用户需在 Settings 设 provider,归 ProviderConfig
error_type: Some(ErrorType::ProviderConfig),
conversation_id: Some(conv_id),
conversation_id: Some(conv_id.to_string()),
});
return;
}
};
// R-PD-6: 续生成路径 conv_id 解耦——has_pending=false 时审批已 remove,无审批 conversation_id 可取;
// 此期 generating=true 且 switchConversation 为 readonly 不并发改 active_conversation_id,
// 故读快照值安全(非竞态期);若 has_pending=true 已在上面 return,不会到此。
// F-09 B 批2:续生成路径 conv_id 用入参(显式,不再从 active 推断),语言/override 读 per_conv 快照。
let lang = snap.agent_language.clone().unwrap_or_else(|| "zh-CN".to_string());
let conv_id = snap.active_conversation_id.clone().unwrap_or_default();
let conv_id_owned = conv_id.to_string();
let system_prompt = build_system_prompt(state, &lang).await;
let session_arc = state.ai_session.clone();
@@ -1087,14 +1179,19 @@ pub(crate) async fn try_continue_agent_loop(app: &AppHandle, state: &AppState, s
// BUG-260617-05 续: provider 解析/build_system_prompt 期间用户可能点 stop。
// spawn 前单次 lock 原子重检 generating——若已被 stop 复位,收敛退出而非覆盖用户的 stop。
// (run_agentic_loop 入口 GeneratingGuard 会再次置 generating=true,若不重检会抹掉 stop。)
if !state.ai_session.lock().await.generating {
tracing::info!("[ai] try_continue 终止:spawn 前重检 generating 已被 stop 复位,补发 AiCompleted");
// F-09 B 批2:重检 per_conv.generating(批2 桥接期与顶层等价)。
let still_generating = {
let session = state.ai_session.lock().await;
session.conv_read(conv_id).map(|c| c.generating).unwrap_or(session.generating)
};
if !still_generating {
tracing::info!(conv_id = %conv_id, "[ai] try_continue 终止:spawn 前重检 generating 已被 stop 复位,补发 AiCompleted");
let _ = app.emit("ai-chat-event", AiChatEvent::AiCompleted {
total_tokens: 0,
prompt_tokens: 0,
completion_tokens: 0,
incomplete: None,
conversation_id: Some(conv_id.clone()),
conversation_id: Some(conv_id_owned.clone()),
});
return;
}
@@ -1103,22 +1200,23 @@ pub(crate) async fn try_continue_agent_loop(app: &AppHandle, state: &AppState, s
// 不应追加到发起工具调用的旧消息,用 AiAgentRound 隔开
let _ = app.emit("ai-chat-event", AiChatEvent::AiAgentRound {
round: 0,
conversation_id: Some(conv_id.clone()),
conversation_id: Some(conv_id_owned.clone()),
});
tauri::async_runtime::spawn(async move {
run_agentic_loop(session_arc, tools_arc, db, app_handle, provider_config, system_prompt, conv_id, knowledge_config, llm_concurrency, max_iterations, max_retries, start_iteration, model_override).await;
run_agentic_loop(session_arc, tools_arc, db, app_handle, provider_config, system_prompt, conv_id_owned, knowledge_config, llm_concurrency, max_iterations, max_retries, start_iteration, model_override).await;
});
}
/// BUG-260617-05: try_continue_agent_loop 续跑判定所需 session 字段的一次性快照。
/// 单次 lock 取出后无锁态判定,消除多 lock 间其他 IPC(ai_chat_stop/clear/switch)改写 session
/// 致续跑判断基于过时快照的 TOCTOU 竞态。
///
/// F-260616-09 B 批2:active_conversation_id 字段移除(改用入参 conv_id,不再从 snap 读)。
struct ContinueSnapshot {
is_generating: bool,
has_pending: bool,
pending_conv_id: Option<String>,
active_conversation_id: Option<String>,
agent_language: Option<String>,
model_override: Option<String>,
}

View File

@@ -444,8 +444,12 @@ pub(crate) fn emit_data_changed(app_handle: &AppHandle, tool_name: &str) {
/// - 返回 None 表示无缓存命中(走原审批流程)。
///
/// `session` 只读扫描 messages不写调用方据返回值决定是否跳过 insert pending。
///
/// F-260616-09 B 批2:经 `conv_id` 索引 per_conv.messages(顶层 messages 批2 后是死字段)。
/// conv_id 来源:process_tool_calls 入参 → 由 agentic.rs run_agentic_loop 入参透传。
async fn find_cached_high_risk_result(
session: &AiSession,
conv_id: &str,
audit_repo: &AiToolExecutionRepo,
tool_name: &str,
args: &serde_json::Value,
@@ -455,9 +459,12 @@ async fn find_cached_high_risk_result(
// 规范化新调用的 args 为可比字符串(排序键,键序无关)
let new_args_key = canonical_args_key(args);
// F-260616-09 B 批2:读 per_conv.messages。process_tool_calls 调用前 loop 入口已桥接建立 per_conv,
// 故 conv_read 必命中;防御性 None 时返 None(无缓存命中,走原审批流程)。
let conv = session.conv_read(conv_id)?;
// ContextManager::iter 返回 impl Iterator非 DoubleEndedcollect 成 Vec 再反向遍历。
// 单对话消息量小百级collect 开销可忽略。
let msgs: Vec<&ChatMessage> = session.messages.iter().collect();
let msgs: Vec<&ChatMessage> = conv.messages.iter().collect();
// 1) 反向扫描 assistant tool_calls找最近一条同名同参的 High 工具调用 → 拿到旧 tool_call_id
// 反向:循环是「最近一次超时→重试」,命中通常是末尾附近,反向先停省全扫。
@@ -629,7 +636,8 @@ pub(crate) async fn process_tool_calls(
// (用户信任同目录同类操作,但仍要看每次的真实结果)。
let trust_hit = super::trust_key_for(&draft.name, &args)
.and_then(|key| {
if session.session_trust.contains(&key) {
// F-260616-09 B 批2:读 per_conv.session_trust(conv_id 来源:本函数入参)。
if session.conv_read(conv_id).map(|c| c.session_trust.contains(&key)).unwrap_or(false) {
Some(key)
} else {
None
@@ -660,14 +668,15 @@ pub(crate) async fn process_tool_calls(
}
// F-05仅 High 查去重缓存Med 保持原审批流程
if matches!(risk_level, RiskLevel::High) {
if let Some((cached, status)) = find_cached_high_risk_result(session, &audit_repo, &draft.name, &args).await {
if let Some((cached, status)) = find_cached_high_risk_result(session, conv_id, &audit_repo, &draft.name, &args).await {
// 命中:把缓存结果作为新 tool_call_id 的 tool_result 回传,跳过审批
tracing::info!(
tool = %draft.name,
new_tool_call_id = %draft.id,
"[F-05] 高危工具去重命中:LLM 重试同命令,复用缓存结果跳过审批(断循环)"
);
session.messages.push(ChatMessage::tool_result(&draft.id, &cached));
// F-260616-09 B 批2:写 per_conv.messages(conv_id 来源:本函数入参)。
session.conv(conv_id).messages.push(ChatMessage::tool_result(&draft.id, &cached));
let _ = app_handle.emit("ai-chat-event", AiChatEvent::AiToolCallCompleted {
id: draft.id.clone(),
result: serde_json::Value::String(cached.clone()),
@@ -696,7 +705,7 @@ pub(crate) async fn process_tool_calls(
recovered: false,
diff: approval_diff.clone(),
});
session.messages.push(ChatMessage::tool_result(&draft.id, PENDING_APPROVAL_PLACEHOLDER));
session.conv(conv_id).messages.push(ChatMessage::tool_result(&draft.id, PENDING_APPROVAL_PLACEHOLDER));
let reason = build_approval_reason(&draft.name, &args, risk_level, db).await;
let _ = app_handle.emit("ai-chat-event", AiChatEvent::AiApprovalRequired {
id: draft.id.clone(),
@@ -757,7 +766,8 @@ pub(crate) async fn process_tool_calls(
Ok(c) => ("completed", c),
Err(c) => ("failed", c),
};
session.messages.push(ChatMessage::tool_result(&draft.id, content.clone()));
// F-260616-09 B 批2:写 per_conv.messages。
session.conv(conv_id).messages.push(ChatMessage::tool_result(&draft.id, content.clone()));
// 审计:trust 放行仍记一条(decided_by=auto_trust),留痕可追溯
audit_tool_call(&audit_repo, conv_id, &draft.id, &draft.name, &draft.args, status, risk_level, Some(content), Some("auto_trust")).await;
}
@@ -810,7 +820,8 @@ pub(crate) async fn process_tool_calls(
Ok(c) => ("completed", c),
Err(c) => ("failed", c),
};
session.messages.push(ChatMessage::tool_result(&draft.id, content.clone()));
// F-260616-09 B 批2:写 per_conv.messages。
session.conv(conv_id).messages.push(ChatMessage::tool_result(&draft.id, content.clone()));
audit_tool_call(&audit_repo, conv_id, &draft.id, &draft.name, &draft.args, status, RiskLevel::Low, Some(content), Some("auto")).await;
}
}

View File

@@ -35,7 +35,10 @@ use super::AiChatEvent;
/// (replace_tool_result_content 找不到 tool_result 返回 false),不报错、无副作用。
///
/// 参数 `final_text` 由调用方按清理场景语义选择(如"会话已停止"/"已取消"/"会话已清除")。
fn finalize_pending_placeholders(session: &mut super::AiSession, final_text: &str) {
///
/// F-260616-09 B 批2:`conv_id` 指定要终态化的 per_conv.messages(批2 真相源)。
/// 顶层 messages 双写(批2 桥接,批4 移除)。conv_id 为空(无活跃对话)时只终态化顶层。
fn finalize_pending_placeholders(session: &mut super::AiSession, conv_id: &str, final_text: &str) {
// SW-260618-02: 先 clone pending 的 tool_call_id(借用在此结束),再可变借 messages。
// 必须整体借 &mut session 在函数体内做 disjoint field borrow —— 调用方若分别传
// &mut messages + &pending 两个引用,函数参数列表不做 disjoint 推断会触发 E0502(2026-06-18 主代修)。
@@ -46,6 +49,9 @@ fn finalize_pending_placeholders(session: &mut super::AiSession, final_text: &st
.collect();
for id in ids {
// replace_tool_result_content 反向遍历,找不到对应 tool_call_id 时返回 false,无副作用
if !conv_id.is_empty() {
session.conv(conv_id).messages.replace_tool_result_content(&id, final_text);
}
session.messages.replace_tool_result_content(&id, final_text);
}
}
@@ -73,9 +79,13 @@ pub async fn ai_regenerate(
let provider_config = super::prompt::get_active_provider(&state).await?;
// 原子占用 generating + 弹出末尾 AI 回复(保留 user 消息)
// F-260616-09 B 批2:写 per_conv(conv_id 来源:IPC 参数 conversation_id,已与 active 一致性校验)。
// 顶层字段同步双写(批2 桥接,批4 IPC 迁移后移除)。
{
let mut session = state.ai_session.lock().await;
if session.generating {
// generating 拦截:读 per_conv(若已建)否则顶层(批2 桥接兼容)。
let is_gen = session.conv_read(&conversation_id).map(|c| c.generating).unwrap_or(session.generating);
if is_gen {
return Err("AI 正在生成中,请等待完成".to_string());
}
// 一致性:regenerate 限定当前活跃对话(避免历史快照陈旧时弹错对话的消息)。
@@ -84,19 +94,27 @@ pub async fn ai_regenerate(
if session.active_conversation_id.as_deref() != Some(conversation_id.as_str()) {
return Err("对话已切换,无法重新生成".to_string());
}
session.generating = true;
session.stop_flag.store(false, Ordering::SeqCst);
session.agent_language = language.clone();
let conv = session.conv(&conversation_id);
conv.generating = true;
conv.stop_flag.store(false, Ordering::SeqCst);
conv.agent_language = language.clone();
// F-260616-11: 重生成 = 新生命周期起点,iteration 从头计数。
session.iteration_used = 0;
conv.iteration_used = 0;
// F-01 阶段6: 记录用户指定模型 override(主对话专用,兜底见 run_agentic_loop)。
session.model_override = model_override.clone();
let popped = session.messages.pop_last_assistant_round();
conv.model_override = model_override.clone();
let popped = conv.messages.pop_last_assistant_round();
if !popped {
// 历史末尾无 AI 回复可弹(空对话/末尾是 user 错误态等),复位 generating 报错
conv.generating = false;
session.generating = false;
return Err("没有可重新生成的回复".to_string());
}
// 顶层双写(批2 桥接)。
session.generating = true;
session.stop_flag.store(false, Ordering::SeqCst);
session.agent_language = language.clone();
session.iteration_used = 0;
session.model_override = model_override.clone();
}
let _tool_defs = state.ai_tools.tool_definitions();
@@ -104,12 +122,13 @@ pub async fn ai_regenerate(
let system_prompt = build_system_prompt(&state, &lang).await;
// 知识注入:取末尾 user 消息文本做检索(与 send 同款,语义命中刷新上下文)
// F-260616-09 B 批2:读 per_conv.messages(conv_id 来源:active_conversation_id,regenerate
// 已校验 == conversation_id)。
let (conv_id, last_user_text) = {
let session = state.ai_session.lock().await;
let cid = session.active_conversation_id.clone().unwrap_or_default();
// 末尾 user 消息文本(用于知识检索;检索本身失败不阻断重生成)
// iter() 非 DoubleEnded,反向找 user:经 all_messages_clone 正向遍历后取末尾 user
let msgs = session.messages.all_messages_clone();
let msgs = session.conv_read(&cid).map(|c| c.messages.all_messages_clone())
.unwrap_or_else(|| session.messages.all_messages_clone());
let last_user = msgs.iter().rev()
.find(|m| matches!(m.role, df_ai::provider::MessageRole::User))
.map(|m| m.content.clone())
@@ -155,7 +174,10 @@ pub async fn ai_regenerate(
#[tauri::command]
pub async fn ai_is_generating(state: State<'_, AppState>) -> Result<bool, String> {
let session = state.ai_session.lock().await;
Ok(session.generating)
// F-260616-09 B 批2:批2 单 active 场景,读 active conv 的 per_conv.generating;
// per_conv 未建时 fallback 顶层(批2 桥接兼容)。批4 改签名 ai_is_generating(conv_id) 后精确查询。
let active = session.active_conversation_id.clone().unwrap_or_default();
Ok(session.conv_read(&active).map(|c| c.generating).unwrap_or(session.generating))
}
/// 发送消息并获取流式 AI 响应
@@ -180,19 +202,37 @@ pub async fn ai_chat_send(
let provider_config = super::prompt::get_active_provider(&state).await?;
// 原子检查并占用生成标志,防止并发双发;同步追加用户消息,按需自动创建对话
// F-260616-09 B 批2:写 per_conv(conv_id 来源:active_conversation_id;首次发送懒创建)。
// 顶层双写(批2 桥接,批4 移除)。
{
let mut session = state.ai_session.lock().await;
if session.generating {
// generating 拦截:active conv 的 per_conv.generating(未建 fallback 顶层)。
let active_for_check = session.active_conversation_id.clone().unwrap_or_default();
let is_gen = if active_for_check.is_empty() {
session.generating
} else {
session.conv_read(&active_for_check).map(|c| c.generating).unwrap_or(session.generating)
};
if is_gen {
return Err("AI 正在生成中,请等待完成".to_string());
}
session.generating = true;
session.stop_flag.store(false, Ordering::SeqCst);
session.agent_language = language.clone();
// 首次发送时生成对话 id(懒创建:不立即落库,避免空对话残留;
// 实际记录由 save_conversation 在生成内容后 upsert 写入)
if session.active_conversation_id.is_none() {
let conv_id = new_id();
session.active_conversation_id = Some(conv_id);
session.active_conv_created_at = Some(now_millis());
}
let active = session.active_conversation_id.clone().unwrap_or_default();
let conv = session.conv(&active);
conv.generating = true;
conv.stop_flag.store(false, Ordering::SeqCst);
conv.agent_language = language.clone();
// F-260616-11: 新对话生命周期 iteration 从头计数(累计计数器复位)。
session.iteration_used = 0;
conv.iteration_used = 0;
// F-01 阶段6: 记录用户指定模型 override(主对话专用)。run_agentic_loop 兜底校验
// (override 非空且在 provider model_configs 池中才用,否则落回路由结果)。
session.model_override = model_override.clone();
conv.model_override = model_override.clone();
// F-260614-02 §5.2:纯技能调用(用户未填文本)时,落库 user content 改 /{skillname}
// 作为技能调用标记(非伪造用户文本),让 title.rs summary_msgs 取到非空素材生成标题。
// 非空 message 原样落库。
@@ -207,19 +247,22 @@ pub async fn ai_chat_send(
};
// F-260614-05 Phase 2c: parts 非空走 user_parts(多模态),空/None 走 user()(纯文本零回归)。
// user_parts 内部把 parts 挂到 ChatMessage.parts,provider 在 has_image() 为真时转 vision 端点。
if let Some(ps) = parts.as_ref().filter(|p| !p.is_empty()) {
conv.messages.push(ChatMessage::user_parts(&user_content, ps.clone()));
} else {
conv.messages.push(ChatMessage::user(&user_content));
}
// 顶层双写(批2 桥接)。
session.generating = true;
session.stop_flag.store(false, Ordering::SeqCst);
session.agent_language = language.clone();
session.iteration_used = 0;
session.model_override = model_override.clone();
if let Some(ps) = parts.as_ref().filter(|p| !p.is_empty()) {
session.messages.push(ChatMessage::user_parts(&user_content, ps.clone()));
} else {
session.messages.push(ChatMessage::user(&user_content));
}
// 首次发送时生成对话 id(懒创建:不立即落库,避免空对话残留;
// 实际记录由 save_conversation 在生成内容后 upsert 写入)
if session.active_conversation_id.is_none() {
let conv_id = new_id();
session.active_conversation_id = Some(conv_id);
session.active_conv_created_at = Some(now_millis());
}
}
// 获取工具定义(预取仅用于触发注册表初始化,实际 tool_defs 在 agentic loop 内部按需获取)
@@ -308,6 +351,12 @@ pub async fn ai_approve(
if !approved {
// 替换占位 tool_result 为拒绝结果
// F-260616-09 B 批2:写 per_conv.messages(conv_id 来源:approval.conversation_id,业务真相源)。
// per_conv 未建时 fallback 顶层(批2 桥接兼容,无 live loop 时顶层仍是真相源)。
let _conv_id_for_msg = approval.conversation_id.clone();
if let Some(ref cid) = _conv_id_for_msg {
session.conv(cid).messages.replace_tool_result_content(&tool_call_id, "用户拒绝了此操作");
}
session.messages.replace_tool_result_content(&tool_call_id, "用户拒绝了此操作");
let conv_id = approval.conversation_id.clone();
let _ = app.emit("ai-chat-event", AiChatEvent::AiApprovalResult {
@@ -325,9 +374,17 @@ pub async fn ai_approve(
audit_finalize(&state, &tool_call_id, "rejected", None).await;
// F-260616-11 决策 a: 审批续跑 iteration 累计(不重置)——读 session.iteration_used 透传
// 作 start_iteration,防多次审批反复跑满 max_iterations 致 token 失控。
let start_iter = state.ai_session.lock().await.iteration_used;
// F-09 B 批2:读 per_conv.iteration_used(conv_id 来源:approval.conversation_id,本分支业务真相源)。
let start_iter = {
let session = state.ai_session.lock().await;
session.conv_read(&conv_id.clone().unwrap_or_default())
.map(|c| c.iteration_used)
.unwrap_or(session.iteration_used)
};
// 所有待审批处理完毕后恢复 agentic 循环
try_continue_agent_loop(&app, &state, start_iter).await;
// F-09 B 批2:传 conv_id(approval.conversation_id),try_continue 按 conv 过滤 pending/读 per_conv。
let cont_conv_id = conv_id.clone().unwrap_or_default();
try_continue_agent_loop(&app, &state, &cont_conv_id, start_iter).await;
return Ok("rejected".to_string());
}
@@ -338,8 +395,19 @@ pub async fn ai_approve(
// AE-2025-04 会话级信任:用户人工批准 write_file / run_command 时,把该工具+目录写入
// session_trust —— 下次同会话同类操作(同工具+同目录 TrustKey)自动放行,跳过 pending + 二次确认。
// 仅首批工具(write_file/run_command)命中 trust_key_for 的 Some,其余工具 noop(None 不写)。
// F-260616-09 B 批2:写 per_conv.session_trust(conv_id 来源:approval.conversation_id,业务真相源)。
// 顶层双写(批2 桥接)。无 conv_id(异常数据)时只写顶层。
if let Some(key) = super::trust_key_for(&approval.tool_name, &approval.arguments) {
let mut inserted = false;
if let Some(ref cid) = conv_id {
if session.conv(cid).session_trust.insert(key.clone()) {
inserted = true;
}
}
if session.session_trust.insert(key.clone()) {
inserted = true;
}
if inserted {
tracing::info!(
tool = %approval.tool_name,
key = ?key,
@@ -382,7 +450,11 @@ pub async fn ai_approve(
}
// 重新获取锁,替换占位 tool_result 为真实结果失败时为错误信息LLM 据此决定下一步)
// F-260616-09 B 批2:写 per_conv.messages + 顶层 messages(双写,conv_id 来源:approval.conversation_id)。
let mut session = state.ai_session.lock().await;
if let Some(ref cid) = conv_id {
session.conv(cid).messages.replace_tool_result_content(&id, &result_val.to_string());
}
session.messages.replace_tool_result_content(&id, &result_val.to_string());
let _ = app.emit("ai-chat-event", AiChatEvent::AiToolCallCompleted {
@@ -406,9 +478,15 @@ pub async fn ai_approve(
// F-260616-11 决策 a: 审批续跑 iteration 累计(不重置)——读 session.iteration_used 透传
// 作 start_iteration,防多次审批反复跑满 max_iterations 致 token 失控。
let start_iter = state.ai_session.lock().await.iteration_used;
// F-09 B 批2:读 per_conv.iteration_used + 传 conv_id 给 try_continue(按 conv 过滤 pending)。
let cont_conv_id = conv_id.clone().unwrap_or_default();
let start_iter = {
let session = state.ai_session.lock().await;
session.conv_read(&cont_conv_id).map(|c| c.iteration_used)
.unwrap_or(session.iteration_used)
};
// 所有待审批处理完毕后恢复 agentic 循环(recovered 无 live loop,try_continue 因 generating=false 自然不续)
try_continue_agent_loop(&app, &state, start_iter).await;
try_continue_agent_loop(&app, &state, &cont_conv_id, start_iter).await;
Ok("executed".to_string())
}
@@ -447,7 +525,11 @@ pub async fn ai_chat_clear(state: State<'_, AppState>) -> Result<(), String> {
let active_id = session.active_conversation_id.clone();
// SW-260618-02:清 pending 前先把占位 tool_result 终态化(此处 messages 紧随即 clear,
// 调用为 no-op,但保持清理点统一口径防顺序变更时占位随 messages 残留 DB)。
finalize_pending_placeholders(&mut *session, "会话已清除");
// F-260616-09 B 批2:per_conv.messages + 顶层双清,conv_id 来源 active_id。
finalize_pending_placeholders(&mut *session, &active_id.clone().unwrap_or_default(), "会话已清除");
if let Some(ref id) = active_id {
session.conv(id).messages.clear();
}
session.messages.clear();
session.pending_approvals.clear();
drop(session);
@@ -479,16 +561,16 @@ pub async fn ai_chat_clear_context(
) -> Result<(), String> {
// 活跃对话一致性:仅对当前活跃对话分段(防陈旧快照标错对话的消息)。
// 与 ai_chat_send/ai_regenerate 取 session 模式一致(state.ai_session.lock)。
// F-260616-09 B 批2:messages 操作改 per_conv(conv_id 来源:IPC 参数 conversation_id,已与 active 校验)。
let conv_id = {
let mut session = state.ai_session.lock().await;
if session.active_conversation_id.as_deref() != Some(conversation_id.as_str()) {
return Err("对话已切换,无法分段".to_string());
}
// 保护区:保留最近 PROTECT_COUNT 条 active。
// PROTECT_COUNT 是 df-ai context.rs 私有常量(=6,build_for_request 同款保护区),
// 未导出到 crate 外,此处用同值字面量 + 注释指明出处,与 build_for_request 语义一致。
const PROTECT_COUNT: usize = 6;
let protect_start = session.messages.len().saturating_sub(PROTECT_COUNT);
let conv = session.conv(&conversation_id);
let protect_start = conv.messages.len().saturating_sub(PROTECT_COUNT);
// 空会话/全在保护区(消息 ≤ N 条) → 无可分段消息,emit 后 noop
if protect_start == 0 {
drop(session);
@@ -498,14 +580,7 @@ pub async fn ai_chat_clear_context(
return Ok(());
}
// 对保护区外 active 消息标 archived_segment(messages_mut 直接改 status)。
// history_tokens 同步:ContextManager 未暴露 history_tokens 写 setter,但 push/compress_old_messages
// 内部用 saturating_sub 维护。此处经 messages_mut 改 status 后 history_tokens 会与 active 集脱钩,
// 但 build_for_request 超预算裁剪路径仍按 history_tokens 决策——为保 token 预算一致性,
// 采用"扣减后等价"的近似:记录被标 archived_segment 消息的 token,经临时 helper 修正。
// df-ai 未提供 history_tokens 减法 API,故此处保留与 compress_old_messages 一致的口径——
// 不直接改 history_tokens(私有字段,无 setter);token 与 active 集可能短暂脱钩,
// 下次 build_for_request 时若 history_tokens 偏高只会触发更早裁剪(保守,不劣化安全性)。
for t in session.messages.messages_mut()[..protect_start].iter_mut() {
for t in conv.messages.messages_mut()[..protect_start].iter_mut() {
if t.message.is_active() {
t.message.status = Some("archived_segment".to_string());
}
@@ -543,24 +618,26 @@ pub async fn ai_chat_compress_context(
let lang = language.unwrap_or_else(|| "zh-CN".to_string());
// 取 active 克隆 + compress_end(读不改,LLM 失败则消息状态完全不变)
// F-260616-09 B 批2:messages 操作改 per_conv(conv_id 来源:IPC 参数 conversation_id,已与 active 校验)。
let (conv_id, active_msgs, compress_end) = {
let mut session = state.ai_session.lock().await;
if session.active_conversation_id.as_deref() != Some(conversation_id.as_str()) {
return Err("对话已切换,无法压缩".to_string());
}
if session.messages.is_compressing() {
let conv = session.conv(&conversation_id);
if conv.messages.is_compressing() {
return Err("压缩正在进行中".to_string());
}
// emit 开始 + 置位防重入(成对释放,见下方所有出口)
let _ = app.emit("ai-chat-event", AiChatEvent::AiCompressing {
conversation_id: Some(conversation_id.clone()),
});
session.messages.set_compressing(true);
conv.messages.set_compressing(true);
let protect_start = session.messages.len().saturating_sub(PROTECT_COUNT);
let protect_start = conv.messages.len().saturating_sub(PROTECT_COUNT);
if protect_start == 0 {
// 空会话/全在保护区 → 无可压缩消息,set_compressing(false) 后返回 Ok
session.messages.set_compressing(false);
conv.messages.set_compressing(false);
let cid = session.active_conversation_id.clone();
drop(session);
let _ = app.emit("ai-chat-event", AiChatEvent::AiCompressed {
@@ -570,7 +647,7 @@ pub async fn ai_chat_compress_context(
return Ok(());
}
// 读 active 克隆(不改 status):compress_old_messages 留到 LLM 成功后才调
let active_msgs: Vec<ChatMessage> = session
let active_msgs: Vec<ChatMessage> = conv
.messages
.messages_mut()[..protect_start]
.iter()
@@ -584,7 +661,7 @@ pub async fn ai_chat_compress_context(
if active_msgs.is_empty() {
// 保护区外无 active 消息(已全 compressed/archived_segment/truncated) → noop
let mut session = state.ai_session.lock().await;
session.messages.set_compressing(false);
session.conv(&conv_id).messages.set_compressing(false);
drop(session);
let _ = app.emit("ai-chat-event", AiChatEvent::AiCompressed {
conversation_id: Some(conv_id),
@@ -599,7 +676,7 @@ pub async fn ai_chat_compress_context(
Ok(p) => p,
Err(e) => {
let mut session = state.ai_session.lock().await;
session.messages.set_compressing(false);
session.conv(&conv_id).messages.set_compressing(false);
drop(session);
let _ = app.emit("ai-chat-event", AiChatEvent::AiError {
error: format!("压缩失败: {}", e),
@@ -621,7 +698,7 @@ pub async fn ai_chat_compress_context(
Err(e) => {
// LLM 失败:不阻塞,set_compressing(false),消息状态不变(未调 compress_old_messages)
let mut session = state.ai_session.lock().await;
session.messages.set_compressing(false);
session.conv(&conv_id).messages.set_compressing(false);
drop(session);
let _ = app.emit("ai-chat-event", AiChatEvent::AiError {
error: format!("压缩失败: {}", e),
@@ -635,9 +712,10 @@ pub async fn ai_chat_compress_context(
// LLM 成功 → 标 compressed(扣 token) + 摘要插首位
{
let mut session = state.ai_session.lock().await;
let _compressed = session.messages.compress_old_messages(compress_end);
session.messages.insert_at(0, ChatMessage::system(&summary));
session.messages.set_compressing(false);
let conv = session.conv(&conv_id);
let _compressed = conv.messages.compress_old_messages(compress_end);
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;
let _ = app.emit("ai-chat-event", AiChatEvent::AiCompressed {
@@ -667,35 +745,38 @@ pub async fn ai_chat_edit(
let provider_config = super::prompt::get_active_provider(&state).await?;
// 原子占用 generating + 替换末条 user content + truncate 其后
// F-260616-09 B 批2:messages/generating 等改 per_conv(conv_id 来源:IPC 参数 conversation_id,已与 active 校验)。
{
let mut session = state.ai_session.lock().await;
if session.generating {
let is_gen = session.conv_read(&conversation_id).map(|c| c.generating).unwrap_or(session.generating);
if is_gen {
return Err("AI 正在生成中,请等待完成".to_string());
}
// 活跃对话一致性(防切走后编辑老快照)
if session.active_conversation_id.as_deref() != Some(conversation_id.as_str()) {
return Err("对话已切换,无法编辑".to_string());
}
let conv = session.conv(&conversation_id);
// ① 替换末条 active user 消息 content(无 active user → Err)
if session
.messages
.replace_last_active_user_content(&new_message)
.is_err()
{
if conv.messages.replace_last_active_user_content(&new_message).is_err() {
return Err("没有可编辑的用户消息".to_string());
}
// ② 其后所有消息标 truncated(无后续也 OK,返回 0)
let _ = session
.messages
.truncate_after_user_message(&new_message)
let _ = conv.messages.truncate_after_user_message(&new_message)
.map_err(|_| "定位被编辑消息失败".to_string())?;
// 占用 generating + stop_flag
conv.generating = true;
conv.stop_flag.store(false, Ordering::SeqCst);
conv.agent_language = language.clone();
// F-260616-11: 编辑重生成 = 新生命周期起点,iteration 从头计数。
conv.iteration_used = 0;
// F-01 阶段6: 记录用户指定模型 override(主对话专用)。
conv.model_override = model_override.clone();
// 顶层双写(批2 桥接)。
session.generating = true;
session.stop_flag.store(false, Ordering::SeqCst);
session.agent_language = language.clone();
// F-260616-11: 编辑重生成 = 新生命周期起点,iteration 从头计数。
session.iteration_used = 0;
// F-01 阶段6: 记录用户指定模型 override(主对话专用)。
session.model_override = model_override.clone();
}
@@ -704,10 +785,12 @@ pub async fn ai_chat_edit(
let system_prompt = build_system_prompt(&state, &lang).await;
// 知识注入:用新 user 文本检索(与 send/regenerate 同款)
// F-260616-09 B 批2:读 per_conv.messages(conv_id 来源:active_conversation_id,edit 已校验 == conversation_id)。
let (conv_id, last_user_text) = {
let session = state.ai_session.lock().await;
let cid = session.active_conversation_id.clone().unwrap_or_default();
let msgs = session.messages.all_messages_clone();
let msgs = session.conv_read(&cid).map(|c| c.messages.all_messages_clone())
.unwrap_or_else(|| session.messages.all_messages_clone());
// 取末条 active user 文本(sanitize 前的全量,但 truncated 已标,这里取 active 的末条)
let last_user = msgs
.iter()
@@ -795,20 +878,36 @@ pub async fn ai_chat_force_send(
let mut session = state.ai_session.lock().await;
// ① 软停止复位(照旧实现,与 ai_chat_stop 审批分支一致)
let old = session.active_conversation_id.clone();
// F-260616-09 B 批2:per_conv + 顶层双写。old 为旧 conv(若存在),置其 per_conv.generating=false。
if let Some(ref old_id) = old {
session.conv(old_id).generating = false;
}
session.generating = false;
// SW-260618-02:强制发送丢弃旧审批,先终态化占位 tool_result,防占位随 messages 残留
// 下次发送喂给 LLM(此处 messages 不 clear,是真实残留路径)。
finalize_pending_placeholders(&mut *session, "已取消");
// F-09 B 批2:finalize 双写 per_conv(old)+ 顶层;conv_id 来源 old active(若存在)。
finalize_pending_placeholders(&mut *session, &old.clone().unwrap_or_default(), "已取消");
session.pending_approvals.clear();
if let Some(ref old_id) = old {
session.conv(old_id).stop_flag.store(true, Ordering::SeqCst);
}
session.stop_flag.store(true, Ordering::SeqCst);
// ② 同锁内立即占用 + 追加用户消息(逻辑照 ai_chat_send:160-194,无 generating 拦截——
// 这是"强制"语义本身,前面已主动复位)
session.generating = true;
session.stop_flag.store(false, Ordering::SeqCst);
session.agent_language = language.clone();
// 首次发送懒创建 conv id(与 ai_chat_send:190-194 一致)
if session.active_conversation_id.is_none() {
let conv_id = new_id();
session.active_conversation_id = Some(conv_id);
session.active_conv_created_at = Some(now_millis());
}
let active = session.active_conversation_id.clone().unwrap_or_default();
let conv = session.conv(&active);
conv.generating = true;
conv.stop_flag.store(false, Ordering::SeqCst);
conv.agent_language = language.clone();
// F-260616-11: 强制发送 = 新对话生命周期起点,iteration 从头计数。
session.iteration_used = 0;
session.model_override = model_override.clone();
conv.iteration_used = 0;
conv.model_override = model_override.clone();
// user content 处理(技能调用空文本标记 /{skillname}),照 ai_chat_send:171-179
let user_content = if message.trim().is_empty() {
if let Some(ref skill_name) = skill {
@@ -820,17 +919,22 @@ pub async fn ai_chat_force_send(
message.clone()
};
// parts 非空走 user_parts(多模态),空/None 走 user()(纯文本零回归)
if let Some(ps) = parts.as_ref().filter(|p| !p.is_empty()) {
conv.messages.push(ChatMessage::user_parts(&user_content, ps.clone()));
} else {
conv.messages.push(ChatMessage::user(&user_content));
}
// 顶层双写(批2 桥接)。
session.generating = true;
session.stop_flag.store(false, Ordering::SeqCst);
session.agent_language = language.clone();
session.iteration_used = 0;
session.model_override = model_override.clone();
if let Some(ps) = parts.as_ref().filter(|p| !p.is_empty()) {
session.messages.push(ChatMessage::user_parts(&user_content, ps.clone()));
} else {
session.messages.push(ChatMessage::user(&user_content));
}
// 首次发送懒创建 conv id(与 ai_chat_send:190-194 一致)
if session.active_conversation_id.is_none() {
let conv_id = new_id();
session.active_conversation_id = Some(conv_id);
session.active_conv_created_at = Some(now_millis());
}
old
};
@@ -895,15 +999,28 @@ pub async fn ai_chat_force_send(
#[tauri::command]
pub async fn ai_chat_stop(state: State<'_, AppState>, app: AppHandle) -> Result<(), String> {
let mut session = state.ai_session.lock().await;
if !session.generating {
// F-260616-09 B 批2:读 per_conv.generating(active conv)+ 顶层 fallback。
let active = session.active_conversation_id.clone().unwrap_or_default();
let is_gen = if active.is_empty() {
session.generating
} else {
session.conv_read(&active).map(|c| c.generating).unwrap_or(session.generating)
};
if !is_gen {
return Ok(());
}
if !session.pending_approvals.is_empty() {
// 审批等待态loop 已退出,直接清理让会话立即可用
// SW-260618-02:messages 不 clear,占位 tool_result 会残留,下次发送喂给 LLM。
// 清 pending 前先终态化为"会话已停止"。
finalize_pending_placeholders(&mut *session, "会话已停止");
// F-09 B 批2:per_conv + 顶层双清,conv_id 来源 active。
finalize_pending_placeholders(&mut *session, &active, "会话已停止");
session.pending_approvals.clear();
if !active.is_empty() {
let conv = session.conv(&active);
conv.generating = false;
conv.stop_flag.store(true, Ordering::SeqCst);
}
session.generating = false;
session.stop_flag.store(true, Ordering::SeqCst); // 双保险:防 try_continue 误判重启
let conv_id = session.active_conversation_id.clone();
@@ -912,11 +1029,20 @@ pub async fn ai_chat_stop(state: State<'_, AppState>, app: AppHandle) -> Result<
return Ok(());
}
// 流式生成中:置位让 loop 自行收尾
if !active.is_empty() {
let conv = session.conv(&active);
conv.stop_flag.store(true, Ordering::SeqCst);
}
session.stop_flag.store(true, Ordering::SeqCst);
// B-260615-14:置位后立即 notify_one 唤醒阻塞在 stream.next() 的 select!
// 不再等 30s 心跳 tick 或 120s idle timeout 才轮到 stop_flag 检查。
// Notify 仅承载「即时唤醒」,停止真值仍由 stop_flag 决定stream_llm 唤醒后再判 flag
session.stop_notify().notify_one();
// F-09 B 批2:per_conv.notify 与顶层 notify 批2 桥接期共享同一 Arc(入口桥接搬入),notify_one 互通。
if let Some(c) = session.conv_read(&active).map(|c| c.notify.clone()) {
c.notify_one();
} else {
session.stop_notify().notify_one();
}
let conv_id = session.active_conversation_id.clone();
drop(session);
@@ -929,7 +1055,17 @@ pub async fn ai_chat_stop(state: State<'_, AppState>, app: AppHandle) -> Result<
tauri::async_runtime::spawn(async move {
tokio::time::sleep(std::time::Duration::from_secs(3)).await;
let mut session = session_arc.lock().await;
if session.generating {
// F-09 B 批2:兜底复位读 per_conv(若有 active)+ 顶层。
let act = session.active_conversation_id.clone().unwrap_or_default();
let still_gen = if act.is_empty() {
session.generating
} else {
session.conv_read(&act).map(|c| c.generating).unwrap_or(session.generating)
};
if still_gen {
if !act.is_empty() {
session.conv(&act).generating = false;
}
session.generating = false;
drop(session); // 释放锁后再 emit避免持锁调 runtime emit
let _ = app_handle.emit(
@@ -967,20 +1103,27 @@ pub async fn ai_continue_loop(
) -> Result<String, String> {
{
let mut session = state.ai_session.lock().await;
if !session.generating {
// F-260616-09 B 批2:读 per_conv.generating(conv_id 来源:IPC 参数 conversation_id,已与 active 校验)。
let is_gen = session.conv_read(&conversation_id).map(|c| c.generating).unwrap_or(session.generating);
if !is_gen {
return Err("AI 未在暂停态,无需继续".to_string());
}
if session.active_conversation_id.as_deref() != Some(conversation_id.as_str()) {
return Err("对话已切换,无法继续".to_string());
}
let conv = session.conv(&conversation_id);
// 复位停止信号:暂停态可能因上一轮 stop_flag 残留为 true续跑 loop 入口会立即退出走完成流程
session.stop_flag.store(false, Ordering::SeqCst);
conv.stop_flag.store(false, Ordering::SeqCst);
// F-260616-11 决策 a: 达 max 续跑重计 iteration(F-260616-03 决策 a,用户点继续=授权重来)。
// 审批续跑(ai_approve)累计不重置见另路径;本路径 start_iteration 传 0,loop 从头计数。
conv.iteration_used = 0;
// 顶层双写(批2 桥接)。
session.stop_flag.store(false, Ordering::SeqCst);
session.iteration_used = 0;
}
// 复用审批恢复续 loop 入口(不重写 loop其内部 spawn run_agentic_loop
try_continue_agent_loop(&app, &state, 0).await;
// F-09 B 批2:传 conv_id(IPC 参数 conversation_id)。
try_continue_agent_loop(&app, &state, &conversation_id, 0).await;
Ok("ok".to_string())
}
@@ -1000,13 +1143,19 @@ pub async fn ai_stop_loop(
) -> Result<String, String> {
let conv_id = {
let mut session = state.ai_session.lock().await;
if !session.generating {
// F-260616-09 B 批2:读 per_conv.generating(conv_id 来源:IPC 参数 conversation_id,已与 active 校验)。
let is_gen = session.conv_read(&conversation_id).map(|c| c.generating).unwrap_or(session.generating);
if !is_gen {
return Err("AI 未在暂停态,无需停止".to_string());
}
if session.active_conversation_id.as_deref() != Some(conversation_id.as_str()) {
return Err("对话已切换,无法停止".to_string());
}
let conv = session.conv(&conversation_id);
// 置 stop_flag 双保险:防 try_continue 误判重启(与 ai_chat_stop 审批分支一致)
conv.stop_flag.store(true, Ordering::SeqCst);
conv.generating = false;
// 顶层双写(批2 桥接)。
session.stop_flag.store(true, Ordering::SeqCst);
session.generating = false;
session.active_conversation_id.clone().unwrap_or_default()
@@ -1357,13 +1506,23 @@ pub async fn ai_conversation_create(
// B-260615-10: 生成中软复位取代硬拦——强制结束当前生成,让用户能立即新建对话。
// 旧 loop 经 B-260615-11 一致性校验(active_conversation_id 变更)自动退出,不污染新对话;
// stop_flag 置位作双保险,让 streaming 中的旧 loop 也尽快收尾。
if session.generating {
// F-260616-09 B 批2:per_conv + 顶层双写。active 旧 conv(若存在)的 per_conv 也复位。
let old_active = session.active_conversation_id.clone();
let is_gen = session.conv_read(&old_active.clone().unwrap_or_default())
.map(|c| c.generating).unwrap_or(session.generating);
if is_gen {
let old_conv = session.active_conversation_id.clone();
if let Some(ref old_id) = old_conv {
session.conv(old_id).generating = false;
}
session.generating = false;
// SW-260618-02:清 pending 前终态化旧对话占位 tool_result(下方 messages.clear 会清旧消息,
// 此处调用为 no-op,但保持统一口径防清理顺序变更)。
finalize_pending_placeholders(&mut *session, "已取消");
finalize_pending_placeholders(&mut *session, &old_conv.clone().unwrap_or_default(), "已取消");
session.pending_approvals.clear();
if let Some(ref old_id) = old_conv {
session.conv(old_id).stop_flag.store(true, Ordering::SeqCst);
}
session.stop_flag.store(true, Ordering::SeqCst);
drop(session);
if let Some(old_conv) = old_conv {
@@ -1378,12 +1537,12 @@ pub async fn ai_conversation_create(
session = state.ai_session.lock().await;
}
// 修复 B-260618-06: 清内存前先持久化当前活跃会话,防「新建会话致旧对话消息丢失」。
// 丢失路径1(高频稳定复现):生成中新建 → 本函数 messages.clear() 抹掉旧会话内存消息,
// 旧 loop 经 B-260615-11 一致性校验(active_conversation_id≠conv_id)return 不 save → 旧会话整段永久丢失。
// 丢失路径2(低概率):loop 正常完成 spawn 后台 save(agentic.rs:928)与 clear 竞态,读空内存覆盖旧会话为空。
// 此处 save 在 clear 前 await 完成,内存仍是旧会话完整消息;幂等(已落库则 update 路径不累加 token/不改 model)。
// F-09 B 批2:old_has_msgs 读 per_conv.messages(若 per_conv 已建)。
let old_conv = session.active_conversation_id.clone();
let old_has_msgs = !session.messages.is_empty();
let old_has_msgs = session.conv_read(&old_conv.clone().unwrap_or_default())
.map(|c| !c.messages.is_empty())
.unwrap_or_else(|| !session.messages.is_empty());
drop(session);
if let Some(ref oc) = old_conv {
if old_has_msgs {
@@ -1395,17 +1554,27 @@ pub async fn ai_conversation_create(
session.active_conv_created_at = Some(now);
// SW-260618-02:清 pending 前终态化占位 tool_result(此处 messages 紧即将 clear,调用为 no-op,
// 保持清理点统一口径;非生成分支无上方 messages.clear 之外的清场,防占位随新会话残留)。
finalize_pending_placeholders(&mut *session, "会话已清除");
finalize_pending_placeholders(&mut *session, &old_conv.clone().unwrap_or_default(), "会话已清除");
// F-09 B 批2:旧 conv 的 per_conv 清空(若存在);新会话 id 建独立 per_conv(全空默认)。
if let Some(ref old_id) = old_conv {
if session.per_conv.contains_key(old_id) {
session.conv(old_id).messages.clear();
session.conv(old_id).agent_language = None;
session.conv(old_id).model_override = None;
session.conv(old_id).session_trust.clear();
session.conv(old_id).stop_flag.store(false, Ordering::SeqCst);
session.conv(old_id).iteration_used = 0;
}
}
// 新会话建一份独立 per_conv(全新 state,与旧会话隔离)。
let new_conv = session.conv(&id);
let _ = new_conv; // conv(&id) 已惰性建立空 PerConvState,字段全为默认值
session.messages.clear();
session.pending_approvals.clear();
// F-260616-09(A 路线):补漏清字段维持单例软隔离,解「新建会话上下文残留」。
// stop_flag 复位 false:上方生成中分支曾置 true 停旧 loop,不复位则新会话 loop
// 启动即见 stop_flag=true 异常退出;agent_language 清空防新会话沿用旧会话语言设置。
// F-260616-09(A 路线):补漏清字段维持单例软隔离,解「新建会话上下文残留」(顶层双写批2 桥接)
session.stop_flag.store(false, Ordering::SeqCst);
session.agent_language = None;
// F-01 阶段6: 清旧对话的 model_override,防新对话沿用上一次的「指定」模型。
session.model_override = None;
// AE-2025-04新建对话清会话信任,新会话重审(信任上下文相关,不跨对话)。
session.session_trust.clear();
Ok(serde_json::json!({ "id": id }))
@@ -1471,7 +1640,14 @@ pub async fn ai_conversation_switch(
let mut session = state.ai_session.lock().await;
// 生成中允许只读切换:返回目标对话的 messages 供前端展示,但不修改 session 状态
// 后台 loop 持有快照的 conv_id不受 active_conversation_id 变更影响
if session.generating {
// F-260616-09 B 批2:读 active conv 的 per_conv.generating(批3+ 决策 e 删 readonly 分支)。
let active_for_check = session.active_conversation_id.clone().unwrap_or_default();
let is_gen = if active_for_check.is_empty() {
session.generating
} else {
session.conv_read(&active_for_check).map(|c| c.generating).unwrap_or(session.generating)
};
if is_gen {
return Ok(serde_json::json!({
"id": record.id,
"title": title,
@@ -1480,15 +1656,25 @@ pub async fn ai_conversation_switch(
}));
}
session.active_conversation_id = Some(conversation_id.clone());
session.messages.restore_from_messages(messages);
// F-09 B 批2:messages restore 到 per_conv(conv_id 来源:IPC 参数 conversation_id)。
// 顶层 messages 同步双写(批2 桥接,批4 IPC 迁移后顶层移除)。
let conv = session.conv(&conversation_id);
conv.messages.restore_from_messages(messages);
conv.model_override = None;
conv.session_trust.clear();
conv.agent_language = None;
conv.iteration_used = 0;
conv.stop_flag.store(false, Ordering::SeqCst);
// 顶层 messages 双写:先 clone per_conv.messages 再 restore(避免借用冲突)。
let top_restore = session.conv_read(&conversation_id).map(|c| c.messages.all_messages_clone()).unwrap_or_default();
session.messages.restore_from_messages(top_restore);
// 仅清空目标对话自身的 pending_approvals,保留其他对话的(防 init 重建的内存 HashMap 被清空,
// 重启恢复链路:restore_pending_approvals(init 重建) → switchConversation(此处不清目标对话的)
// → ai_pending_tool_calls 查询 → ai_approve 落库)
session.pending_approvals.retain(|_, a| a.conversation_id.as_deref() != Some(&conversation_id));
// F-01 阶段6: 切换对话清旧 override,防新对话沿用上一次的「指定」模型
// (override 是单对话级 UI 选择,不跨对话持久化)。
// F-01 阶段6: 切换对话清旧 override(顶层双写)。
session.model_override = None;
// AE-2025-04切换对话清会话信任,新会话重审(信任上下文相关,不跨对话)。
// AE-2025-04切换对话清会话信任(顶层双写)。
session.session_trust.clear();
// 释放 session lock 再做 async provider 查询 + spawn(避免持锁 await DB)
drop(session);
@@ -1534,6 +1720,9 @@ pub async fn ai_conversation_delete(
// 非活跃对话的恢复审批(recovered,conversation_id 指向被删对话)若不 retain 清理,
// 会永久残留死审批条目。对齐 ai_conversation_switch 的 retain 口径(仅清目标对话,保留其他)。
session.pending_approvals.retain(|_, a| a.conversation_id.as_deref() != Some(&conversation_id));
// F-260616-09 B 批2:删除 conv 时移除其 per_conv 条目(设计 §4.1 conv 存在性判据依赖此)。
// 顶层 messages.clear 仅在删的是 active conv 时执行(批2 桥接双清)。
session.per_conv.remove(&conversation_id);
if session.active_conversation_id.as_deref() == Some(&conversation_id) {
session.active_conversation_id = None;
session.messages.clear();

View File

@@ -144,9 +144,17 @@ pub(crate) async fn save_conversation(
// 工具结果(content)超 50KB 时截断头尾各 ~20KB + 中段标注,防大体量结果(read_file 1MB 洞 /
// list_directory 13782 项)落库后每轮重发累积致 token 暴增。仅影响持久化视图,不污染
// 内存真相源(ContextManager)——build_for_request 仍读全量 messages。
//
// F-260616-09 B 批2:读 per_conv.messages(conv_id 来源:本函数入参,与 save 落库的 conv 一致)。
// loop 内 save 由 run_agentic_loop 入参 conv_id 透传;IPC 路径(commands.rs)save 也传 conv_id。
// per_conv 未建时 fallback 顶层 messages(兼容批2 桥接未覆盖的边缘路径,如 ai_chat_clear 后
// 无 per_conv 但顶层有,虽然 clear 已清空,fallback 顶层也是空,行为一致)。
let (messages_json, provider_id, created_at) = {
let session = session_arc.lock().await;
let mut msgs = session.messages.all_messages_clone();
let mut msgs = match session.conv_read(conv_id) {
Some(c) => c.messages.all_messages_clone(),
None => session.messages.all_messages_clone(),
};
for m in &mut msgs {
m.content = truncate_for_persist(&m.content);
// F-260614-05 Phase 2a: parts(Image base64) 同样截断(替换占位 Text 片),

View File

@@ -250,7 +250,15 @@ pub(crate) async fn maybe_spawn_extraction(
return Ok(());
}
// 消息数守卫(总消息数,含 system/assistant/tool)
let msg_count = session_arc.lock().await.messages.len();
// F-260616-09 B 批2:读 per_conv.messages.len()(conv_id 来源:本函数入参)。
// per_conv 未建时 fallback 顶层 messages.len()(同锁内取,避免双锁死锁)。
let msg_count = {
let session = session_arc.lock().await;
match session.conv_read(conv_id) {
Some(c) => c.messages.len(),
None => session.messages.len(),
}
};
if (msg_count as u32) < config.min_messages {
return Ok(());
}

View File

@@ -349,10 +349,11 @@ pub struct AiSession {
pub session_trust: HashSet<TrustKey>,
/// F-260616-09 B 阶段批1 多会话并发架构:会话级状态按 conversation_id 切分。
///
/// **批1 状态(纯新增无行为变化)**:此 HashMap 初始空,批1 不迁移任何调用方 ——
/// 现有所有路径继续读写顶层单例字段(messages/generating/stop_flag 等),完全无感。
/// 批2 承接迁移:调用方改走 `session.conv(conv_id).*`,迁移完成后顶层单例字段废弃删除。
#[allow(dead_code)] // F-09 B 批1 共存期:批2 迁移承接(0 写入方)
/// **批2 状态(per-conv 是新真相源)**:批2 已迁移所有调用方(agentic.rs loop +
/// GeneratingGuard + try_continue + commands.rs IPC 写路径 + audit.rs process_tool_calls +
/// conversation.rs save + title.rs + knowledge_inject.rs + lib.rs L0)读写 `session.conv(conv_id).*`
/// / `session.conv_read(conv_id)`。顶层单例字段(messages/generating/stop_flag 等)批2 期间
/// 作双写桥接保留(IPC 仍双写顶层保 ai_is_generating 等读路径兼容,批4 IPC 迁移后顶层移除)。
pub per_conv: HashMap<String, PerConvState>,
}
@@ -423,9 +424,6 @@ impl AiSession {
/// 不存在则 `entry().or_insert_with(PerConvState::new)` 创建一份默认状态,
/// 初值对齐 [`AiSession::new`] / [`PerConvState::new`]。
/// 已存在则返回同一实例引用(同 conv_id 复用)。
///
/// **批1 状态**:0 调用方(批2 迁移承接)。行为不变。
#[allow(dead_code)] // F-09 B 批1 共存期:批2 迁移承接(0 调用方)
pub fn conv(&mut self, conv_id: &str) -> &mut PerConvState {
self.per_conv
.entry(conv_id.to_string())
@@ -436,9 +434,6 @@ impl AiSession {
///
/// 未创建的 conv_id 返 `None`,已创建返 `Some(&PerConvState)`。用于读取路径(如
/// 列举/查询)避免误触发惰性创建。
///
/// **批1 状态**:0 调用方(批2 迁移承接)。行为不变。
#[allow(dead_code)] // F-09 B 批1 共存期:批2 迁移承接(0 调用方)
pub fn conv_read(&self, conv_id: &str) -> Option<&PerConvState> {
self.per_conv.get(conv_id)
}

View File

@@ -43,24 +43,48 @@ pub(crate) async fn ensure_conversation_title(
}
// 取对话文本(仅 user/assistant跳过 tool 噪音),取前 6 条供 LLM 总结
// F-260616-09 B 批2:读 per_conv.messages(conv_id 来源:本函数入参)。
// per_conv 未建时 fallback 顶层(switch 触发的 spawn_ensure_title 可能无 per_conv)。
let (summary_msgs, all_msgs) = {
let session = session_arc.lock().await;
let summary: Vec<ChatMessage> = session.messages.iter()
.filter(|m| matches!(m.role, MessageRole::User | MessageRole::Assistant))
.take(6)
.map(|m| ChatMessage {
role: m.role.clone(),
content: m.content.clone(),
parts: None,
tool_call_id: None,
tool_calls: None,
model: None,
status: None,
reasoning_content: m.reasoning_content.clone(),
timestamp: m.timestamp,
})
.collect();
(summary, session.messages.all_messages_clone())
let conv = session.conv_read(conv_id);
let summary: Vec<ChatMessage> = match conv {
Some(c) => c.messages.iter()
.filter(|m| matches!(m.role, MessageRole::User | MessageRole::Assistant))
.take(6)
.map(|m| ChatMessage {
role: m.role.clone(),
content: m.content.clone(),
parts: None,
tool_call_id: None,
tool_calls: None,
model: None,
status: None,
reasoning_content: m.reasoning_content.clone(),
timestamp: m.timestamp,
})
.collect(),
None => session.messages.iter()
.filter(|m| matches!(m.role, MessageRole::User | MessageRole::Assistant))
.take(6)
.map(|m| ChatMessage {
role: m.role.clone(),
content: m.content.clone(),
parts: None,
tool_call_id: None,
tool_calls: None,
model: None,
status: None,
reasoning_content: m.reasoning_content.clone(),
timestamp: m.timestamp,
})
.collect(),
};
let all_msgs = match conv {
Some(c) => c.messages.all_messages_clone(),
None => session.messages.all_messages_clone(),
};
(summary, all_msgs)
};
if summary_msgs.is_empty() {
return;

View File

@@ -40,9 +40,19 @@ pub fn run() {
let app_h = app_handle.clone();
tauri::async_runtime::spawn(async move {
let mut session = session_arc.lock().await;
let was_generating = session.generating;
// F-260616-09 B 批2:读 active conv 的 per_conv.generating + 双写复位。
// conv_id 来源:active_conversation_id(HMR 重连场景的当前展示 conv)。
let conv_id = session.active_conversation_id.clone();
let active_for_check = conv_id.clone().unwrap_or_default();
let was_generating = if active_for_check.is_empty() {
session.generating
} else {
session.conv_read(&active_for_check).map(|c| c.generating).unwrap_or(session.generating)
};
if was_generating {
if !active_for_check.is_empty() {
session.conv(&active_for_check).generating = false;
}
session.generating = false;
session.pending_approvals.clear();
drop(session);