优化: aichat效率剩余(压缩后台化防阻塞/审计批量事务/只读缓存轮内去重/流式增量渲染/AiCommandOutput合批/双渲染合并) + 跨端加固(df-project路径保留大小写/tunnel文档更正supervisor重连/relay固定时间比较与帧上限/启动校验) + 销账
This commit is contained in:
@@ -13,38 +13,179 @@
|
||||
//! 读 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 事件所需上下文)。
|
||||
#[derive(Clone)]
|
||||
///
|
||||
/// 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<String>,
|
||||
/// 合批共享缓冲(emit 线程写 / flush 定时读)。
|
||||
buffer: Arc<StdMutex<CommandBuffer>>,
|
||||
/// 后台 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<AtomicBool>,
|
||||
handle: Option<tokio::task::JoinHandle<()>>,
|
||||
}
|
||||
|
||||
impl CommandSink {
|
||||
pub fn new(app: AppHandle, tool_call_id: String, conversation_id: Option<String>) -> Self {
|
||||
Self { app, tool_call_id, conversation_id }
|
||||
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 },
|
||||
}
|
||||
}
|
||||
|
||||
/// emit 一行 stdout/stderr(AiCommandOutput,双写 app.emit + ai_event_bus)。
|
||||
/// emit 失败静默吞(前端未 listen / 总线无订阅不阻断命令执行)。
|
||||
/// 累积一行 + 触发 flush 判定(行加入 stdout/stderr 分桶)。
|
||||
///
|
||||
/// AC-EFF-S2-5:不立即 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<String>,
|
||||
buffer: &Arc<StdMutex<CommandBuffer>>,
|
||||
) {
|
||||
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<String>,
|
||||
kind: StreamKind,
|
||||
line: &str,
|
||||
) {
|
||||
let ev = AiChatEvent::AiCommandOutput {
|
||||
id: self.tool_call_id.clone(),
|
||||
id: tool_call_id.to_string(),
|
||||
stream: kind.as_str().to_string(),
|
||||
line: line.to_string(),
|
||||
conversation_id: self.conversation_id.clone(),
|
||||
conversation_id: conversation_id.clone(),
|
||||
};
|
||||
let _ = self.app.emit("ai-chat-event", ev.clone());
|
||||
let _ = self.app.state::<crate::state::AppState>().ai_event_bus.publish_event(ev);
|
||||
let _ = app.emit("ai-chat-event", ev.clone());
|
||||
let _ = app.state::<crate::state::AppState>().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();
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user