From 79b7766ed443a30ffe3d049bc9c5c2f64078bb23 Mon Sep 17 00:00:00 2001 From: JianFeeeee Date: Tue, 25 Aug 2026 12:24:11 +0800 Subject: [PATCH] =?UTF-8?q?fix:=20=E6=B5=81=E5=BC=8F=E6=B8=B2=E6=9F=93?= =?UTF-8?q?=E5=9B=9E=E5=90=88=E7=94=9F=E5=91=BD=E5=91=A8=E6=9C=9F=20+=20LL?= =?UTF-8?q?M=20=E7=9E=AC=E6=96=AD=E9=87=8D=E8=AF=95=E4=B8=8E=20SSE=20body?= =?UTF-8?q?=20=E5=85=9C=E5=BA=95?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 问题一(webui 不是真流式): - sendChat 的 finally 在 POST 结束(15s ackTimer abort)时就复位 chatLoading,但 agent 生成窗口 15~190s,后续 SSE delta 全部走 全量重建路径、停止按钮提前消失、用户误发重复消息。 - GUI app.js 完全没有 content_delta/reasoning_delta 监听器, 只能等聚合帧一次性显示。 修复:三端统一回合生命周期——POST 只是触发,收尾由 SSE 驱动: - dashboard/GUI 新增 endChatTurn/armTurnWatchdog;拿到同步兜底 响应立即收尾,否则保持回合打开等 agent_output final / reset 帧 / 120s watchdog 兜底 - GUI 补齐 delta 监听器;agent_output 聚合分支 += 改覆盖; reasoning 聚合帧改覆盖(多轮工具调用时旧逻辑会重复累加) - agent_output 误杀分支(final 无 source 即 return 丢弃新输出) 改为内容比较去重,多轮连发时新一轮回复不再被吞 - waiter reasoning_delta reset 从清空全部消息改为 sealLastAgent 问题二(三条只成功一条): - handleChat 60s ctx 含排队时间,agent 串行处理下第 N 条必超时 (实测第 3 条 62s 超时 504);放宽到 300s(客户端 abort 时立即取消) - LLM 单 provider 瞬断无重试:process.go provider 循环内加同源 重试(2 次、退避 2s),401/403 凭证错误与用户中断不重试 - llmsproxy auto 链在非流式请求下可能返回 SSE body(上游恢复后 吐已生成的 chunk 流),非流式解析报 invalid character 'd' 丢掉 整段回复;新增 parseOpenAICompatibleSSEBody 拼接为完整响应 - 顺带修 normalizeStreamToolCalls 分片续传 bug:name 不重发时 argsRaw 被顶层 Arguments(nil) 覆盖丢失 function.arguments 验证: - 连发 3 条 + 单条共 4 条全部成功(首条 190s 重试扛住瞬断) - sse_body_test.go 锁定 SSE body 解析契约(content/usage/tool call 分片) --- cmd/gui/renderer/app.js | 142 ++++++++++++++++++++++++-- cmd/waiter/tui.go | 2 +- internal/agent/api/provider.go | 94 ++++++++++++++++- internal/agent/api/sse_body_test.go | 121 ++++++++++++++++++++++ internal/agent/core/process.go | 72 +++++++++---- internal/plugins/webui/dashboard.html | 70 +++++++++++-- internal/plugins/webui/handler.go | 13 ++- 7 files changed, 469 insertions(+), 45 deletions(-) create mode 100644 internal/agent/api/sse_body_test.go diff --git a/cmd/gui/renderer/app.js b/cmd/gui/renderer/app.js index c4dc67d..21d5e5b 100644 --- a/cmd/gui/renderer/app.js +++ b/cmd/gui/renderer/app.js @@ -31,6 +31,7 @@ const state = { messages: [], chatLoading: false, chatStage: "", + _turnWatchdog: null, healthResult: null, starmapInit: false, starmapLoading: false, @@ -2036,6 +2037,46 @@ function buildChatStarmapGraph() { }); } +// 回合收尾:由 SSE 事件(agent_output final / reset 帧)或 watchdog 驱动。 +// POST 结束 ≠ 回合结束:agent 可能还在生成(排队+长生成),提前复位 +// chatLoading 会让后续 delta 走全量重建、停止按钮消失、用户误发重复消息。 +function endChatTurn() { + if (!state.chatLoading) return; + state.chatLoading = false; + state.chatStage = ""; + if (state._turnWatchdog) { + clearTimeout(state._turnWatchdog); + state._turnWatchdog = null; + } + var btn = document.getElementById("chat-send-btn"); + if (btn) { + btn.disabled = false; + btn.textContent = __("发送", "Send"); + } + var sb = document.getElementById("chat-stop-btn"); + if (sb) sb.style.display = "none"; + rerenderChatIfActive(); +} + +// 回合看门狗:POST 已超时且 SSE 迟迟无终帧时兕底收尾(连接不稳/事件丢失), +// 提示用户回复可能已生成、可刷新查看历史。避免回合永久卡在 loading。 +function armTurnWatchdog() { + if (state._turnWatchdog) clearTimeout(state._turnWatchdog); + state._turnWatchdog = setTimeout(() => { + state._turnWatchdog = null; + if (state.chatLoading) { + endChatTurn(); + toast( + __( + "长时间未收到回复,连接可能不稳定;回复可能已生成,可刷新连接后查看", + "No reply for a long time; the reply may have been generated, reconnect to check", + ), + true, + ); + } + }, 120000); +} + async function sendChat() { var inp = document.getElementById("chat-input"); var btn = document.getElementById("chat-send-btn"); @@ -2125,13 +2166,14 @@ async function sendChat() { rerenderChat(); toast(__("请求失败: ", "Request failed: ") + e.message, true); } finally { - state.chatLoading = false; - state.chatStage = ""; - btn.disabled = false; - btn.textContent = __("发送", "Send"); - var sb2 = document.getElementById("chat-stop-btn"); - if (sb2) sb2.style.display = "none"; // 回复完成/失败,隐藏停止按钮 - rerenderChat(); + if (r && r.response) { + // 同步兜底已拿到完整回复:回合结束 + endChatTurn(); + } else { + // 触发式受理(POST 超时/失败):回合仍打开,等 SSE 流式渲染; + // 由 agent_output final / reset 帧 / watchdog 收尾 + armTurnWatchdog(); + } } } @@ -4881,9 +4923,23 @@ async function connectFetchSSE(url) { ? state.messages[state.messages.length - 1] : null; if (last && last.role === "assistant" && !last._final) { + // 聚合最终响应:覆盖 delta 累积的中间内容(以聚合为准,含 stage 插件改写后的文本),置 final 结束本轮流式。 last._grow = true; - last.content += p.content || ""; + last.content = p.content || ""; + last._final = true; rerenderChatIfActive(); + endChatTurn(); + return; + } + if ( + last && + last.role === "assistant" && + last._final && + !last.source && + last.content === (p.content || "") + ) { + // 去重:同一轮的重复帧(如 SSE 重连回放)内容相同则忽略,仅收尾回合 + endChatTurn(); return; } state.messages.push({ @@ -4891,8 +4947,10 @@ async function connectFetchSSE(url) { content: p.content || "", _streaming: true, _grow: true, + _final: true, }); rerenderChatIfActive(); + endChatTurn(); } else if (type === "reasoning") { if (p.content) { state.chatStage = __("AI 思考中...", "AI thinking..."); @@ -4910,10 +4968,74 @@ async function connectFetchSSE(url) { }); last = state.messages[state.messages.length - 1]; } - last.reasoning_content = - (last.reasoning_content || "") + (p.content || ""); + last.reasoning_content = p.content; rerenderChatIfActive(); } + } else if (type === "reasoning_delta") { + // token 级思考流式增量:逐块追加到当前思考内容;reset 帧表示轮次作废 + if (p.channel === "_consolidation_") return; + if (p.reset) { + var lm = state.messages.length + ? state.messages[state.messages.length - 1] + : null; + if (lm && lm.role === "assistant" && !lm._final) { + lm._final = true; + rerenderChatIfActive(); + } + armTurnWatchdog(); + return; + } + if (!p.content) return; + state.chatStage = __("AI 思考中...", "AI thinking..."); + var last = + state.messages.length > 0 + ? state.messages[state.messages.length - 1] + : null; + if (!last || last.role !== "assistant" || last._final) { + state.messages.push({ + role: "assistant", + content: "", + reasoning_content: "", + tool_calls: [], + _streaming: true, + }); + last = state.messages[state.messages.length - 1]; + } + last.reasoning_content = + (last.reasoning_content || "") + p.content; + rerenderChatIfActive(); + } else if (type === "content_delta") { + // token 级回复流式增量:逐块追加到当前回复内容;reset 帧表示轮次作废(中断) + if (p.channel === "_consolidation_") return; + if (p.reset) { + var lm = state.messages.length + ? state.messages[state.messages.length - 1] + : null; + if (lm && lm.role === "assistant" && !lm._final) { + lm._final = true; + rerenderChatIfActive(); + } + armTurnWatchdog(); + return; + } + if (!p.content) return; + state.chatStage = __("AI 回复中...", "AI replying..."); + var last = + state.messages.length > 0 + ? state.messages[state.messages.length - 1] + : null; + if (!last || last.role !== "assistant" || last._final) { + state.messages.push({ + role: "assistant", + content: "", + tool_calls: [], + _streaming: true, + _grow: true, + }); + last = state.messages[state.messages.length - 1]; + } + last.content += p.content; + rerenderChatIfActive(); } else if (type === "tool_call") { if (!p.tool) return; var last = diff --git a/cmd/waiter/tui.go b/cmd/waiter/tui.go index 051b2ba..06d9519 100644 --- a/cmd/waiter/tui.go +++ b/cmd/waiter/tui.go @@ -379,7 +379,7 @@ func (m *tuiModel) handleServerLine(line string) { case "reasoning_delta": // token 级增量:与 reasoning 同样合并到最后一条 reasoning 消息 if rl.Reset { - m.messages = []chatMsg{} + m.sealLastAgent() break } if n := len(m.messages); n > 0 && m.messages[n-1].kind == msgReasoning { diff --git a/internal/agent/api/provider.go b/internal/agent/api/provider.go index dccbd46..368f331 100644 --- a/internal/agent/api/provider.go +++ b/internal/agent/api/provider.go @@ -384,6 +384,11 @@ func (p *LuaAdaptedProvider) Chat(ctx context.Context, req *CompletionRequest) ( if parsed, perr := parseOpenAICompatibleResponse(rawResp); perr == nil { return parsed, nil } + // 网关在非流式请求下返回了 SSE 流 body(上游恢复后吐 chunk 流), + // 拼接为完整响应,避免丢掉已生成的整段回复 + if parsed, ok := parseOpenAICompatibleSSEBody(rawResp); ok { + return parsed, nil + } return nil, fmt.Errorf("lua transform_response: %w", err) } @@ -392,6 +397,9 @@ func (p *LuaAdaptedProvider) Chat(ctx context.Context, req *CompletionRequest) ( if parsed, perr := parseOpenAICompatibleResponse(rawResp); perr == nil { return parsed, nil } + if parsed, ok := parseOpenAICompatibleSSEBody(rawResp); ok { + return parsed, nil + } return nil, fmt.Errorf("unmarshal unified response: %w (body: %s)", err, unifiedJSON) } @@ -460,6 +468,86 @@ func parseOpenAICompatibleResponse(raw []byte) (*CompletionResponse, error) { return out, nil } +// parseOpenAICompatibleSSEBody 将 SSE 格式的响应体("data: {...}" 多行) +// 拼接为完整 CompletionResponse。场景:网关(llmsproxy auto 链等)在非流式 +// 请求下也可能返回流式 body——上游恢复后吐出的是已生成的 chunk 流,若按 +// 普通 JSON 解析会报 "invalid character 'd'" 而丢掉整段完整回复。 +// 返回 false 表示 body 不是 SSE 格式,调用方继续走原有解析路径。 +func parseOpenAICompatibleSSEBody(raw []byte) (*CompletionResponse, bool) { + body := strings.TrimSpace(string(raw)) + if !strings.HasPrefix(body, "data:") && !strings.Contains(body, "\ndata:") { + return nil, false + } + type sseAcc struct { + id string + name string + argsRaw strings.Builder + } + var out CompletionResponse + var contentBuf, reasoningBuf strings.Builder + accs := map[int]*sseAcc{} + toolOrder := []int{} + finish := "" + found := false + + for _, line := range strings.Split(body, "\n") { + line = strings.TrimSpace(line) + if !strings.HasPrefix(line, "data:") { + continue + } + payload := strings.TrimSpace(strings.TrimPrefix(line, "data:")) + if payload == "" || payload == "[DONE]" { + continue + } + ck, ok := parseOpenAICompatibleStreamChunkFull(payload) + if !ok { + continue + } + found = true + contentBuf.WriteString(ck.Content) + reasoningBuf.WriteString(ck.ReasoningContent) + for i, tc := range ck.ToolCalls { + acc := accs[i] + if acc == nil { + acc = &sseAcc{} + accs[i] = acc + toolOrder = append(toolOrder, i) + } + if tc.ID != "" { + acc.id = tc.ID + } + if tc.Name != "" { + acc.name = tc.Name + } + acc.argsRaw.WriteString(tc.RawArguments) + } + if ck.Done && ck.FinishReason != "" { + finish = ck.FinishReason + } + if ck.Usage != nil { + out.TokenUsage = *ck.Usage + } + } + if !found { + return nil, false + } + out.Content = contentBuf.String() + out.ReasoningContent = reasoningBuf.String() + out.FinishReason = finish + for _, i := range toolOrder { + acc := accs[i] + name := strings.TrimSpace(acc.name) + argsStr := strings.TrimSpace(acc.argsRaw.String()) + if name == "" && argsStr == "" && acc.id == "" { + continue + } + tc := ToolCall{ID: acc.id, Type: "function", Name: name, RawArguments: argsStr} + tc.Arguments = parseToolArguments(argsStr) + out.ToolCalls = append(out.ToolCalls, tc) + } + return &out, true +} + type openAIToolCall struct { ID string `json:"id"` Type string `json:"type"` @@ -526,7 +614,11 @@ func normalizeStreamToolCalls(raw []openAIToolCall) []ToolCall { argsRaw := tc.Function.Arguments if name == "" { name = tc.Name - argsRaw = tc.Arguments + // 仅当顶层 Arguments 存在才用扁平格式;否则保留 function.arguments 嵌套值 + // (OpenAI 流式续传 chunk:name 不重发但 function.arguments 继续) + if tc.Arguments != nil { + argsRaw = tc.Arguments + } } typ := tc.Type if typ == "" && (tc.ID != "" || name != "" || argsRaw != nil) { diff --git a/internal/agent/api/sse_body_test.go b/internal/agent/api/sse_body_test.go new file mode 100644 index 0000000..a592f3d --- /dev/null +++ b/internal/agent/api/sse_body_test.go @@ -0,0 +1,121 @@ +package api + +import ( + "encoding/json" + "strings" + "testing" +) + +// 锁定契约:网关(llmsproxy auto 链等)在非流式请求下返回 SSE 流 body 时, +// 必须拼接为完整响应而不是报 "invalid character 'd'" 丢掉已生成的回复。 +// 事故样本取自 2026-08-25 生产日志:上游恢复后吐出完整 chunk 流被非流式解析器丢弃。 +func TestParseOpenAICompatibleSSEBody(t *testing.T) { + body := "data: {\"id\":\"chatcmpl-572\",\"object\":\"chat.completion.chunk\",\"created\":1787630289,\"model\":\"x-preview-f-free\",\"choices\":[{\"index\":0,\"delta\":{\"role\":\"assistant\",\"content\":\"你好\"},\"finish_reason\":null}]}\n" + + "\n" + + "data: {\"id\":\"chatcmpl-572\",\"choices\":[{\"index\":0,\"delta\":{\"role\":\"assistant\",\"content\":\",世界\"},\"finish_reason\":null}]}\n" + + "\n" + + "data: {\"id\":\"chatcmpl-572\",\"choices\":[{\"index\":0,\"delta\":{\"content\":\"\"},\"finish_reason\":\"stop\"}]}\n" + + "data: {\"id\":\"chatcmpl-572\",\"choices\":[],\"usage\":{\"prompt_tokens\":100,\"completion_tokens\":7,\"total_tokens\":107}}\n" + + "data: [DONE]\n" + + resp, ok := parseOpenAICompatibleSSEBody([]byte(body)) + if !ok { + t.Fatal("expected SSE body to be recognized") + } + if resp.Content != "你好,世界" { + t.Errorf("content = %q, want %q", resp.Content, "你好,世界") + } + if resp.FinishReason != "stop" { + t.Errorf("finish_reason = %q, want stop", resp.FinishReason) + } + if resp.TokenUsage.Total != 107 || resp.TokenUsage.Prompt != 100 || resp.TokenUsage.Completion != 7 { + t.Errorf("usage = %+v, want prompt=100 completion=7 total=107", resp.TokenUsage) + } +} + +func TestParseOpenAICompatibleSSEBodyRejectsPlainJSON(t *testing.T) { + plain := `{"choices":[{"message":{"content":"hi"}}]}` + if _, ok := parseOpenAICompatibleSSEBody([]byte(plain)); ok { + t.Fatal("plain JSON body must not be treated as SSE") + } +} + +func TestParseOpenAICompatibleSSEBodyToolCallShards(t *testing.T) { + // 用 json.Marshal 构建测试数据,避免 Go 字面量转义错误 + chunk1 := map[string]interface{}{ + "choices": []map[string]interface{}{{ + "index": 0, + "delta": map[string]interface{}{ + "tool_calls": []map[string]interface{}{{ + "index": 0, + "id": "call_1", + "type": "function", + "function": map[string]interface{}{ + "name": "exec", + "arguments": `{"command":`, + }, + }}, + }, + }}, + } + chunk2 := map[string]interface{}{ + "choices": []map[string]interface{}{{ + "index": 0, + "delta": map[string]interface{}{ + "tool_calls": []map[string]interface{}{{ + "index": 0, + "function": map[string]interface{}{ + "arguments": `"date"}`, + }, + }}, + }, + }}, + } + chunk3 := map[string]interface{}{ + "choices": []map[string]interface{}{{ + "index": 0, + "delta": map[string]interface{}{}, + "finish_reason": "tool_calls", + }}, + } + + var sb strings.Builder + for _, c := range []map[string]interface{}{chunk1, chunk2, chunk3} { + b, _ := json.Marshal(c) + sb.WriteString("data: ") + sb.Write(b) + sb.WriteString("\n") + } + sb.WriteString("data: [DONE]\n") + + t.Logf("SSE body:\n%s", sb.String()) + + resp, ok := parseOpenAICompatibleSSEBody([]byte(sb.String())) + if !ok { + t.Fatal("expected SSE body to be recognized") + } + if len(resp.ToolCalls) != 1 { + t.Fatalf("got %d tool calls, want 1", len(resp.ToolCalls)) + } + tc := resp.ToolCalls[0] + if tc.Name != "exec" || tc.ID != "call_1" { + t.Errorf("tool call name/id = %q/%q, want exec/call_1", tc.Name, tc.ID) + } + args := tc.RawArguments + if args != `{"command":"date"}` { + t.Errorf("raw args = %q", args) + } + if tc.Arguments["command"] != "date" { + t.Errorf("parsed args = %v, want command=date", tc.Arguments) + } + if resp.FinishReason != "tool_calls" { + t.Errorf("finish_reason = %q, want tool_calls", resp.FinishReason) + } +} + +func TestParseOpenAICompatibleSSEBodyEmptyStream(t *testing.T) { + body := "data: \ndata: \n" + if resp, ok := parseOpenAICompatibleSSEBody([]byte(body)); ok && strings.TrimSpace(resp.Content) != "" { + t.Fatalf("empty stream should not parse into non-empty response") + } +} diff --git a/internal/agent/core/process.go b/internal/agent/core/process.go index e746681..5c2eca7 100644 --- a/internal/agent/core/process.go +++ b/internal/agent/core/process.go @@ -7,6 +7,7 @@ import ( "fmt" "log" "strings" + "time" agentAPI "gitcode.com/JianFeeeee/HomeAgent/internal/agent/api" "gitcode.com/JianFeeeee/HomeAgent/internal/events" @@ -116,28 +117,61 @@ func (a *Agent) process(input string, stageCtx *sdk.StageContext) (response stri fbProvider.Name(), pi, len(providers)-1) } - fCtx, fCancel := context.WithCancel(a.ctx) - a.llmMu.Lock() - a.cancelLLM = fCancel - a.llmMu.Unlock() - - resp, llmErr = chatStreamWithFallback(fCtx, fbProvider, req, a) - - a.llmMu.Lock() - a.cancelLLM = nil - a.llmMu.Unlock() - fCancel() - - if llmErr == nil { - a.providerManager.ResetAvailability(fbProvider.Name()) - if fbProvider != a.provider { - a.provider = fbProvider - log.Printf("[agent] switched active provider to %q after fallback", - fbProvider.Name()) + // 同源瞬时错误重试:网关瞬断(502/503/504/429/网络抖动)通常秒级恢复, + // 直接跳下一个 provider(或直接报错)会丢掉本可成功的请求。 + // 凭证错误(401/403)与用户中断不重试。 + const maxAttempts = 2 + for attempt := 1; attempt <= maxAttempts; attempt++ { + if attempt > 1 { + log.Printf("[agent] provider %q transient failure, retry %d/%d in 2s: %v", + fbProvider.Name(), attempt, maxAttempts, llmErr) + select { + case <-time.After(2 * time.Second): + case <-a.ctx.Done(): + llmErr = a.ctx.Err() + } + if llmErr == nil || errors.Is(llmErr, context.Canceled) || errors.Is(llmErr, context.DeadlineExceeded) { + break + } } - break + + fCtx, fCancel := context.WithCancel(a.ctx) + a.llmMu.Lock() + a.cancelLLM = fCancel + a.llmMu.Unlock() + + resp, llmErr = chatStreamWithFallback(fCtx, fbProvider, req, a) + + a.llmMu.Lock() + a.cancelLLM = nil + a.llmMu.Unlock() + fCancel() + + if llmErr == nil { + a.providerManager.ResetAvailability(fbProvider.Name()) + if fbProvider != a.provider { + a.provider = fbProvider + log.Printf("[agent] switched active provider to %q after fallback", + fbProvider.Name()) + } + break + } + + // 用户中断:立即终止,不重试也不换 provider + if errors.Is(llmErr, context.Canceled) { + break + } + // 凭证错误:重试无意义,跳出重试循环进入 provider 标记/切换 + var pe *agentAPI.ProviderError + if errors.As(llmErr, &pe) && (pe.StatusCode == 401 || pe.StatusCode == 403) { + break + } + // 其余错误(含 5xx/429/网络):还有重试机会则继续,否则跳出 } + if llmErr == nil { + break + } if errors.Is(llmErr, context.Canceled) { break } diff --git a/internal/plugins/webui/dashboard.html b/internal/plugins/webui/dashboard.html index 632eea9..17208fd 100644 --- a/internal/plugins/webui/dashboard.html +++ b/internal/plugins/webui/dashboard.html @@ -3633,6 +3633,46 @@ background: }); } + // 回合收尾:由 SSE 事件(agent_output final / reset 帧)或 watchdog 驱动。 + // POST 结束 ≠ 回合结束:agent 可能还在生成(排队+长生成),提前复位 + // chatLoading 会让后续 delta 走全量重建、停止按钮消失、用户误发重复消息。 + function endChatTurn() { + if (!state.chatLoading) return; + state.chatLoading = false; + state.chatStage = ""; + if (state._turnWatchdog) { + clearTimeout(state._turnWatchdog); + state._turnWatchdog = null; + } + var btn = document.getElementById("chat-send-btn"); + if (btn) { + btn.disabled = false; + btn.textContent = __("发送", "Send"); + } + var sb = document.getElementById("chat-stop-btn"); + if (sb) sb.style.display = "none"; + rerenderChat(true); + } + + // 回合看门狗:POST 已 abort 且 SSE 迟迟无终帧时兕底收尾(连接不稳/事件丢失), + // 提示用户回复可能已生成、可刷新查看历史。避免回合永久卡在 loading。 + function armTurnWatchdog() { + if (state._turnWatchdog) clearTimeout(state._turnWatchdog); + state._turnWatchdog = setTimeout(function () { + state._turnWatchdog = null; + if (state.chatLoading) { + endChatTurn(); + toast( + __( + "长时间未收到回复,连接可能不稳定;回复可能已生成,可刷新页面查看", + "No reply received for a long time; the reply may have been generated, refresh to check", + ), + true, + ); + } + }, 120000); + } + async function sendChat() { var inp = document.getElementById("chat-input"); var btn = document.getElementById("chat-send-btn"); @@ -3677,7 +3717,6 @@ background: } finally { clearTimeout(ackTimer); } - state.chatStage = __("AI 回复中...", "AI replying..."); var last = state.messages[state.messages.length - 1]; if (r && r.response) { if (last && last.role === "assistant" && last._streaming) { @@ -3716,13 +3755,14 @@ background: toast(__("请求失败: ", "Request failed: ") + e.message, true); } } finally { - state.chatLoading = false; - state.chatStage = ""; - btn.disabled = false; - btn.textContent = __("发送", "Send"); - var sb2 = document.getElementById("chat-stop-btn"); - if (sb2) sb2.style.display = "none"; // 回复完成/失败,隐藏停止按钮 - rerenderChat(true); + if (r && r.response) { + // 同步兜底已拿到完整回复:回合结束 + endChatTurn(); + } else { + // 触发式受理(POST 已 abort/失败):回合仍打开,等 SSE 流式渲染; + // 由 agent_output final / reset 帧 / watchdog 收尾 + armTurnWatchdog(); + } } } @@ -4146,6 +4186,7 @@ background: lastM2.content = p.content; lastM2._final = true; rerenderChat(); + endChatTurn(); return; } if ( @@ -4154,7 +4195,13 @@ background: lastM2._final && !lastM2.source ) { - return; + // 去重:同一轮的重复帧(如 SSE 重连回放)内容相同则忽略; + // 内容不同视为新一轮输出(上一轮已 final 且无 source),开新消息。 + // 旧逻辑无条件 return 会丢弃多轮连发时新一轮的最终回复。 + if (lastM2.content === p.content) { + endChatTurn(); + return; + } } state.messages.push({ role: "assistant", @@ -4164,6 +4211,7 @@ background: _final: true, }); rerenderChat(); + endChatTurn(); } catch (ex) { console.error("[SSE] agent_output error", ex); } @@ -4183,6 +4231,10 @@ background: lm._final = true; rerenderChat(); } + // 轮次作废(用户中断):定格已显示内容;核心会以 [中断消息] + // 重启轮次,保持回合打开让确认回复继续流式渲染, + // 由其 agent_output final / watchdog 收尾。 + armTurnWatchdog(); return; } if (!p.content) return; diff --git a/internal/plugins/webui/handler.go b/internal/plugins/webui/handler.go index c7ab3ed..a877085 100644 --- a/internal/plugins/webui/handler.go +++ b/internal/plugins/webui/handler.go @@ -1342,8 +1342,11 @@ func (h *Handler) handleChat(w http.ResponseWriter, r *http.Request) { if body.ClientMsgID != "" { payload["client_msg_id"] = body.ClientMsgID } - // 带超时的上下文,防止 InjectInputSync 长时间阻塞 HTTP 请求 - ctx, cancel := context.WithTimeout(r.Context(), 60*time.Second) + // 带超时的上下文,防止 InjectInputSync 长时间阻塞 HTTP 请求。 + // 注意:ctx 派生自 r.Context(),客户端提前断开(前端 15s ackTimer abort)时 + // 立即取消,不会真等满 300s;300s 只约束"连接保持 + agent 排队/长生成"场景 + // (agent 串行处理,后发消息的排队时间也计入,60s 曾导致连发第 3 条必超时)。 + ctx, cancel := context.WithTimeout(r.Context(), 300*time.Second) defer cancel() respCh := make(chan *agentIO.OutputEvent, 1) @@ -1706,8 +1709,8 @@ func (h *Handler) handleOpenAICompletions(w http.ResponseWriter, r *http.Request return } - // 带超时的上下文 - ctx, cancel := context.WithTimeout(r.Context(), 60*time.Second) + // 带超时的上下文(同 handleChat:客户端断开立即取消;300s 约束长生成) + ctx, cancel := context.WithTimeout(r.Context(), 300*time.Second) defer cancel() respCh := make(chan *agentIO.OutputEvent, 1) @@ -1719,7 +1722,7 @@ func (h *Handler) handleOpenAICompletions(w http.ResponseWriter, r *http.Request select { case response = <-respCh: case <-ctx.Done(): - writeJSON(w, http.StatusGatewayTimeout, map[string]string{"error": "agent timeout (60s)"}) + writeJSON(w, http.StatusGatewayTimeout, map[string]string{"error": "agent timeout (300s)"}) return }