//! run_command 实时流式输出 — task-local sink 机制。 //! //! 治「run_command 执行中黑盒」:execute() 等 exit 才返回整块 stdout/stderr, //! 长命令(cargo/npm 构建)期间前端只看 Started→Completed,中间进度不可见。 //! //! 架构约束:工具 handler 注册为 `Box Future>`(ai_tools.rs:48), //! 签名只收 args 不收 AppHandle/tool_call_id,无法直接 emit 事件。改 schema 把 id //! 塞进 args 会泄漏给 LLM,不可取。 //! //! 本模块用 [`tokio::task_local!`] 解耦:调用方(execute_with_heartbeat / ai_approve) //! 持有 AppHandle + tool_call_id + conv_id,调 `tools.execute` 前用 [`scope`] 把一个 //! [`CommandSink`] 注入当前 task 上下文;run_command handler 在 tools/file.rs 内 //! 读 task-local(同一 task,因 tools.execute 不 spawn 直接 await handler), //! 命中则改走 shell `execute_streaming`,每行回调 [`emit_output`] → AiCommandOutput。 //! //! AC-EFF-S2-5(2026-08-09):AiCommandOutput 合批。run_command 逐行回调不再每行单独 emit + //! publish_event(大输出风暴时 IPC 事件风暴),改为累积进 [`CommandSink`] 共享缓冲(stdout/stderr //! 分桶),由「后台 50ms 定时 flush + 体积阈值(4KB)立即 flush + Drop 兜底 flush」三路输出, //! 每个 stream 字段累积的字符串(多行 '\n' 连接)作为单条 AiCommandOutput 输出(对齐 //! stream_recv DELTA_FLUSH_INTERVAL=50ms 的 delta 合批;消费方按文本渲染,多行输出等价逐行)。 //! //! 未注入 sink(非 run_command / 调用方未配 scope)时 [`emit_output`] 静默 noop, //! 兜底不报错不阻断。 use std::sync::Arc; use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::Mutex as StdMutex; use std::time::Duration; use tauri::{AppHandle, Emitter, Manager}; use super::AiChatEvent; use df_execute::shell::StreamKind; /// AiCommandOutput 合批参数(AC-EFF-S2-5,对齐 stream_recv DELTA_FLUSH_INTERVAL=50ms)。 /// /// 时间维度:后台 flush task 每 50ms 兜底输出一次(与 AR-8 delta 合批同 cadence); /// 体积维度:单次累积超 [`FLUSH_BYTE_THRESHOLD`] 立即 flush(防大输出风暴 50ms 窗口内积压过多)。 /// 语义:stdout/stderr 各累积的字符串(多行以 '\n' 连接)作为**单条** AiCommandOutput 事件输出, /// 不再每行单独 emit + publish_event——砍 IPC 事件风暴(与 delta 合批「合并后前端逐条累加,最终一致」 /// 同思路,消费方按文本渲染,多行输出等价于逐行)。 const FLUSH_INTERVAL: Duration = Duration::from_millis(50); const FLUSH_BYTE_THRESHOLD: usize = 4096; /// 一次 run_command 调用的输出下沉目标(emit 事件所需上下文)。 /// /// AC-EFF-S2-5(2026-08-09):合批。run_command 逐行回调 → 累积进共享缓冲(stdout/stderr 分桶), /// 由「后台 50ms 定时 flush + 体积阈值立即 flush + Drop 兜底 flush」三路输出,不再每行单独 /// emit(治 aichat 效率 S2-5)。缓冲用 `std::sync::Mutex`(短临界区无 await,emit 是同步回调)。 pub struct CommandSink { app: AppHandle, tool_call_id: String, conversation_id: Option, /// 合批共享缓冲(emit 线程写 / flush 定时读)。 buffer: Arc>, /// 后台 flush task 守卫(Drop 时停止 + abort,并兜底 flush 尾部)。 _flusher: FlushGuard, } /// 合批缓冲:stdout/stderr 分开累积(按 stream 字段 emit),bytes 供体积阈值判定。 #[derive(Default)] struct CommandBuffer { stdout: String, stderr: String, bytes: usize, } /// 后台 flush task 守卫:stop 标志 + JoinHandle。Drop 时置停 + abort(任务 50ms 循环读到 stop 退出), /// 尾部输出由 CommandSink::drop 兜底 flush,不丢命令结尾。 struct FlushGuard { stop: Arc, handle: Option>, } impl CommandSink { pub fn new(app: AppHandle, tool_call_id: String, conversation_id: Option) -> Self { let buffer = Arc::new(StdMutex::new(CommandBuffer::default())); let stop = Arc::new(AtomicBool::new(false)); // 后台 50ms flush task。仅当处于 tokio runtime 上下文时 spawn(tokio::spawn 需 runtime)。 // CommandSink::new 由 execute_with_heartbeat / ai_approve(async)调用,正常在 runtime 内; // 非 runtime(如单测直接 new)跳过定时 task,退化为「体积阈值 + Drop 兜底」flush(仍不丢尾部)。 let handle = if tokio::runtime::Handle::try_current().is_ok() { let buf = buffer.clone(); let stop_f = stop.clone(); let app_f = app.clone(); let id_f = tool_call_id.clone(); let conv_f = conversation_id.clone(); Some(tokio::spawn(async move { let mut interval = tokio::time::interval(FLUSH_INTERVAL); interval.tick().await; // 弃首 tick(tokio interval 首 tick 立即返回,对齐心跳弃首模式) loop { interval.tick().await; if stop_f.load(Ordering::SeqCst) { break; } Self::flush_inner(&app_f, &id_f, &conv_f, &buf); } })) } else { None }; Self { app, tool_call_id, conversation_id, buffer, _flusher: FlushGuard { stop, handle }, } } /// 累积一行 + 触发 flush 判定(行加入 stdout/stderr 分桶)。 /// /// 不立即 emit,先累积入缓冲;体积超阈值立即 flush,否则等 50ms 定时 flush /// 或 Drop 兜底 flush。emit 失败静默吞(前端未 listen / 总线无订阅不阻断命令执行)。 fn emit(&self, kind: StreamKind, line: &str) { let over_threshold = { let mut buf = self.buffer.lock().unwrap(); let target = if kind == StreamKind::Stdout { &mut buf.stdout } else { &mut buf.stderr }; if !target.is_empty() { target.push('\n'); } target.push_str(line); buf.bytes += line.len(); buf.bytes >= FLUSH_BYTE_THRESHOLD }; if over_threshold { Self::flush_inner(&self.app, &self.tool_call_id, &self.conversation_id, &self.buffer); } } /// 兜底 flush(把当前缓冲整体取出并 emit)。Drop 时调用,保证尾部输出不丢。 fn flush(&self) { Self::flush_inner(&self.app, &self.tool_call_id, &self.conversation_id, &self.buffer); } /// 内部 flush:单次锁内 take 缓冲(不持锁 emit),stdout/stderr 各输出一条(多行拼接)。 /// 与后台 task / Drop 并发调用安全:take 语义下同一批只被 drain 一次(无重复 emit)。 fn flush_inner( app: &AppHandle, tool_call_id: &str, conversation_id: &Option, buffer: &Arc>, ) { let drained = { let mut buf = buffer.lock().unwrap(); if buf.stdout.is_empty() && buf.stderr.is_empty() { return; } Some(std::mem::take(&mut *buf)) }; if let Some(buf) = drained { if !buf.stdout.is_empty() { Self::emit_line(app, tool_call_id, conversation_id, StreamKind::Stdout, &buf.stdout); } if !buf.stderr.is_empty() { Self::emit_line(app, tool_call_id, conversation_id, StreamKind::Stderr, &buf.stderr); } } } /// emit 一条 AiCommandOutput(stdout/stderr 累积字符串,双写 app.emit + ai_event_bus)。 fn emit_line( app: &AppHandle, tool_call_id: &str, conversation_id: &Option, kind: StreamKind, line: &str, ) { let ev = AiChatEvent::AiCommandOutput { id: tool_call_id.to_string(), stream: kind.as_str().to_string(), line: line.to_string(), conversation_id: conversation_id.clone(), }; let _ = app.emit("ai-chat-event", ev.clone()); let _ = app.state::().ai_event_bus.publish_event(ev); } } impl Drop for CommandSink { fn drop(&mut self) { // 停后台 flush task + abort(任务在 stop 后下次 tick 退出),再兜底 flush 剩余缓冲 // (命令结束尾部输出不丢)。并发安全:flush_inner 的 take 语义保证不重复 emit。 self._flusher.stop.store(true, Ordering::SeqCst); if let Some(h) = self._flusher.handle.take() { h.abort(); } self.flush(); } } // task-local 槽:tools.execute 调用期间(同 task)handler 可读。 // 嵌套 Option:外层是 task-local 是否 set,内层是是否注入 sink(None = 已 scope 但无 sink)。 tokio::task_local! { static SINK: Option; } /// 在 sink 作用域内执行 future。future 完成后自动清理(无残留)。 /// /// 调用方:execute_with_heartbeat / ai_approve 在调 `tools.execute("run_command", args)` /// 前包一层 `command_stream::scope(Some(sink), async { tools.execute(...).await })`。 pub async fn scope(sink: Option, fut: F) -> R where F: std::future::Future, { SINK.scope(sink, fut).await } /// run_command handler 读取:回调每行输出 → emit AiCommandOutput。 /// /// task-local 未 set(非经 scope 调用,如直接单测调 handler)或内层 None → 静默 noop。 /// 返回值忽略(emit 失败不阻断命令)。 pub fn emit_output(kind: StreamKind, line: &str) { // LocalKey::with 在 task-local 未 scope 时返 AccessError,静默吞(noop 兜底)。 let _ = SINK.with(|maybe_sink: &Option| { if let Some(sink) = maybe_sink { sink.emit(kind, line); } }); }