diff --git a/internal/plugins/webui/handler.go b/internal/plugins/webui/handler.go index b81270e..f2e5f0e 100644 --- a/internal/plugins/webui/handler.go +++ b/internal/plugins/webui/handler.go @@ -1693,7 +1693,10 @@ func (h *Handler) handleChatEvents(w http.ResponseWriter, r *http.Request) { // writeCh 不 close:Subscribe 回调闭包持有它,handler 退出后回调仍可能被 // 总线异步触发,close 后再发送会 panic(send on closed channel,生产日志中 // 单日数千次)。writer goroutine 通过 done 退出;发送侧 select on done 防泄漏。 - writeCh := make(chan string, 64) + // 缓冲加大到 512 且 writer 做批量合并:reasoning/content 增量是高频小包, + // 每条单独 flush 会因 socket 写慢而填满小缓冲导致 delta 被丢弃(表现为 + // 前端只能等最终的 agent_output 整段,体感延迟)。 + writeCh := make(chan string, 512) writerDone := make(chan struct{}) go func() { defer func() { @@ -1702,12 +1705,32 @@ func (h *Handler) handleChatEvents(w http.ResponseWriter, r *http.Request) { } close(writerDone) }() + // 批量合并窗口:16ms 内收集的增量一次性 flush,降 flush 次数、 + // 避免高频小包拖慢 socket 写导致 writeCh 积压丢 delta。 + pending := make([]string, 0, 64) + flushPending := func() { + if len(pending) == 0 { + return + } + for _, line := range pending { + fmt.Fprintf(w, "%s\n", line) + } + flusher.Flush() + pending = pending[:0] + } + flushTicker := time.NewTicker(16 * time.Millisecond) + defer flushTicker.Stop() for { select { case line := <-writeCh: - fmt.Fprintf(w, "%s\n", line) - flusher.Flush() + pending = append(pending, line) + if len(pending) >= 64 { + flushPending() + } + case <-flushTicker.C: + flushPending() case <-done: + flushPending() return } }