mirror of
https://gitcode.com/JianFeeeee/HomeAgent.git
synced 2026-10-04 00:03:59 +00:00
fix(webui): SSE writer 批量合并 flush 修复流式 delta 丢包延迟
【根因】SSE writeCh 缓冲仅 64 且 writer 每条 delta 单独 flush。 reasoning/content 增量是高频小包(单轮 200+ 条),socket 写慢时 writeCh 迅速填满,delta 大量 DROPPED——浏览器收不到 逐 token 增量,只能等最终 agent_output 整段到达,体感明显延迟。 实测一轮 16s 纯文本回复:content_delta DROPPED 202 次、 reasoning_delta DROPPED 345 次,前端全程无流式渲染。 【修复】 - writeCh 缓冲 64 → 512 - writer 加 16ms 批量合并窗口:窗口内收集的增量一次性 flush, 或满 64 条立即 flush;done 退出前 flush 残留。 flush 次数从 N 降到约 N/64,socket 写压力骤降。 修复后实测 0 DROPPED,delta 全部实时送达前端。
This commit is contained in:
@ -1693,7 +1693,10 @@ func (h *Handler) handleChatEvents(w http.ResponseWriter, r *http.Request) {
|
|||||||
// writeCh 不 close:Subscribe 回调闭包持有它,handler 退出后回调仍可能被
|
// writeCh 不 close:Subscribe 回调闭包持有它,handler 退出后回调仍可能被
|
||||||
// 总线异步触发,close 后再发送会 panic(send on closed channel,生产日志中
|
// 总线异步触发,close 后再发送会 panic(send on closed channel,生产日志中
|
||||||
// 单日数千次)。writer goroutine 通过 done 退出;发送侧 select on done 防泄漏。
|
// 单日数千次)。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{})
|
writerDone := make(chan struct{})
|
||||||
go func() {
|
go func() {
|
||||||
defer func() {
|
defer func() {
|
||||||
@ -1702,12 +1705,32 @@ func (h *Handler) handleChatEvents(w http.ResponseWriter, r *http.Request) {
|
|||||||
}
|
}
|
||||||
close(writerDone)
|
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 {
|
for {
|
||||||
select {
|
select {
|
||||||
case line := <-writeCh:
|
case line := <-writeCh:
|
||||||
fmt.Fprintf(w, "%s\n", line)
|
pending = append(pending, line)
|
||||||
flusher.Flush()
|
if len(pending) >= 64 {
|
||||||
|
flushPending()
|
||||||
|
}
|
||||||
|
case <-flushTicker.C:
|
||||||
|
flushPending()
|
||||||
case <-done:
|
case <-done:
|
||||||
|
flushPending()
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user