mirror of
https://gitcode.com/JianFeeeee/HomeAgent.git
synced 2026-10-04 00:03:59 +00:00
fix: 流式渲染回合生命周期 + LLM 瞬断重试与 SSE body 兜底
问题一(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 分片)
This commit is contained in:
@ -31,6 +31,7 @@ const state = {
|
|||||||
messages: [],
|
messages: [],
|
||||||
chatLoading: false,
|
chatLoading: false,
|
||||||
chatStage: "",
|
chatStage: "",
|
||||||
|
_turnWatchdog: null,
|
||||||
healthResult: null,
|
healthResult: null,
|
||||||
starmapInit: false,
|
starmapInit: false,
|
||||||
starmapLoading: 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() {
|
async function sendChat() {
|
||||||
var inp = document.getElementById("chat-input");
|
var inp = document.getElementById("chat-input");
|
||||||
var btn = document.getElementById("chat-send-btn");
|
var btn = document.getElementById("chat-send-btn");
|
||||||
@ -2125,13 +2166,14 @@ async function sendChat() {
|
|||||||
rerenderChat();
|
rerenderChat();
|
||||||
toast(__("请求失败: ", "Request failed: ") + e.message, true);
|
toast(__("请求失败: ", "Request failed: ") + e.message, true);
|
||||||
} finally {
|
} finally {
|
||||||
state.chatLoading = false;
|
if (r && r.response) {
|
||||||
state.chatStage = "";
|
// 同步兜底已拿到完整回复:回合结束
|
||||||
btn.disabled = false;
|
endChatTurn();
|
||||||
btn.textContent = __("发送", "Send");
|
} else {
|
||||||
var sb2 = document.getElementById("chat-stop-btn");
|
// 触发式受理(POST 超时/失败):回合仍打开,等 SSE 流式渲染;
|
||||||
if (sb2) sb2.style.display = "none"; // 回复完成/失败,隐藏停止按钮
|
// 由 agent_output final / reset 帧 / watchdog 收尾
|
||||||
rerenderChat();
|
armTurnWatchdog();
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@ -4881,9 +4923,23 @@ async function connectFetchSSE(url) {
|
|||||||
? state.messages[state.messages.length - 1]
|
? state.messages[state.messages.length - 1]
|
||||||
: null;
|
: null;
|
||||||
if (last && last.role === "assistant" && !last._final) {
|
if (last && last.role === "assistant" && !last._final) {
|
||||||
|
// 聚合最终响应:覆盖 delta 累积的中间内容(以聚合为准,含 stage 插件改写后的文本),置 final 结束本轮流式。
|
||||||
last._grow = true;
|
last._grow = true;
|
||||||
last.content += p.content || "";
|
last.content = p.content || "";
|
||||||
|
last._final = true;
|
||||||
rerenderChatIfActive();
|
rerenderChatIfActive();
|
||||||
|
endChatTurn();
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
if (
|
||||||
|
last &&
|
||||||
|
last.role === "assistant" &&
|
||||||
|
last._final &&
|
||||||
|
!last.source &&
|
||||||
|
last.content === (p.content || "")
|
||||||
|
) {
|
||||||
|
// 去重:同一轮的重复帧(如 SSE 重连回放)内容相同则忽略,仅收尾回合
|
||||||
|
endChatTurn();
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
state.messages.push({
|
state.messages.push({
|
||||||
@ -4891,8 +4947,10 @@ async function connectFetchSSE(url) {
|
|||||||
content: p.content || "",
|
content: p.content || "",
|
||||||
_streaming: true,
|
_streaming: true,
|
||||||
_grow: true,
|
_grow: true,
|
||||||
|
_final: true,
|
||||||
});
|
});
|
||||||
rerenderChatIfActive();
|
rerenderChatIfActive();
|
||||||
|
endChatTurn();
|
||||||
} else if (type === "reasoning") {
|
} else if (type === "reasoning") {
|
||||||
if (p.content) {
|
if (p.content) {
|
||||||
state.chatStage = __("AI 思考中...", "AI thinking...");
|
state.chatStage = __("AI 思考中...", "AI thinking...");
|
||||||
@ -4910,10 +4968,74 @@ async function connectFetchSSE(url) {
|
|||||||
});
|
});
|
||||||
last = state.messages[state.messages.length - 1];
|
last = state.messages[state.messages.length - 1];
|
||||||
}
|
}
|
||||||
last.reasoning_content =
|
last.reasoning_content = p.content;
|
||||||
(last.reasoning_content || "") + (p.content || "");
|
|
||||||
rerenderChatIfActive();
|
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") {
|
} else if (type === "tool_call") {
|
||||||
if (!p.tool) return;
|
if (!p.tool) return;
|
||||||
var last =
|
var last =
|
||||||
|
|||||||
@ -379,7 +379,7 @@ func (m *tuiModel) handleServerLine(line string) {
|
|||||||
case "reasoning_delta":
|
case "reasoning_delta":
|
||||||
// token 级增量:与 reasoning 同样合并到最后一条 reasoning 消息
|
// token 级增量:与 reasoning 同样合并到最后一条 reasoning 消息
|
||||||
if rl.Reset {
|
if rl.Reset {
|
||||||
m.messages = []chatMsg{}
|
m.sealLastAgent()
|
||||||
break
|
break
|
||||||
}
|
}
|
||||||
if n := len(m.messages); n > 0 && m.messages[n-1].kind == msgReasoning {
|
if n := len(m.messages); n > 0 && m.messages[n-1].kind == msgReasoning {
|
||||||
|
|||||||
@ -384,6 +384,11 @@ func (p *LuaAdaptedProvider) Chat(ctx context.Context, req *CompletionRequest) (
|
|||||||
if parsed, perr := parseOpenAICompatibleResponse(rawResp); perr == nil {
|
if parsed, perr := parseOpenAICompatibleResponse(rawResp); perr == nil {
|
||||||
return parsed, 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)
|
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 {
|
if parsed, perr := parseOpenAICompatibleResponse(rawResp); perr == nil {
|
||||||
return parsed, 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)
|
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
|
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 {
|
type openAIToolCall struct {
|
||||||
ID string `json:"id"`
|
ID string `json:"id"`
|
||||||
Type string `json:"type"`
|
Type string `json:"type"`
|
||||||
@ -526,7 +614,11 @@ func normalizeStreamToolCalls(raw []openAIToolCall) []ToolCall {
|
|||||||
argsRaw := tc.Function.Arguments
|
argsRaw := tc.Function.Arguments
|
||||||
if name == "" {
|
if name == "" {
|
||||||
name = tc.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
|
typ := tc.Type
|
||||||
if typ == "" && (tc.ID != "" || name != "" || argsRaw != nil) {
|
if typ == "" && (tc.ID != "" || name != "" || argsRaw != nil) {
|
||||||
|
|||||||
121
internal/agent/api/sse_body_test.go
Normal file
121
internal/agent/api/sse_body_test.go
Normal file
@ -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")
|
||||||
|
}
|
||||||
|
}
|
||||||
@ -7,6 +7,7 @@ import (
|
|||||||
"fmt"
|
"fmt"
|
||||||
"log"
|
"log"
|
||||||
"strings"
|
"strings"
|
||||||
|
"time"
|
||||||
|
|
||||||
agentAPI "gitcode.com/JianFeeeee/HomeAgent/internal/agent/api"
|
agentAPI "gitcode.com/JianFeeeee/HomeAgent/internal/agent/api"
|
||||||
"gitcode.com/JianFeeeee/HomeAgent/internal/events"
|
"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)
|
fbProvider.Name(), pi, len(providers)-1)
|
||||||
}
|
}
|
||||||
|
|
||||||
fCtx, fCancel := context.WithCancel(a.ctx)
|
// 同源瞬时错误重试:网关瞬断(502/503/504/429/网络抖动)通常秒级恢复,
|
||||||
a.llmMu.Lock()
|
// 直接跳下一个 provider(或直接报错)会丢掉本可成功的请求。
|
||||||
a.cancelLLM = fCancel
|
// 凭证错误(401/403)与用户中断不重试。
|
||||||
a.llmMu.Unlock()
|
const maxAttempts = 2
|
||||||
|
for attempt := 1; attempt <= maxAttempts; attempt++ {
|
||||||
resp, llmErr = chatStreamWithFallback(fCtx, fbProvider, req, a)
|
if attempt > 1 {
|
||||||
|
log.Printf("[agent] provider %q transient failure, retry %d/%d in 2s: %v",
|
||||||
a.llmMu.Lock()
|
fbProvider.Name(), attempt, maxAttempts, llmErr)
|
||||||
a.cancelLLM = nil
|
select {
|
||||||
a.llmMu.Unlock()
|
case <-time.After(2 * time.Second):
|
||||||
fCancel()
|
case <-a.ctx.Done():
|
||||||
|
llmErr = a.ctx.Err()
|
||||||
if llmErr == nil {
|
}
|
||||||
a.providerManager.ResetAvailability(fbProvider.Name())
|
if llmErr == nil || errors.Is(llmErr, context.Canceled) || errors.Is(llmErr, context.DeadlineExceeded) {
|
||||||
if fbProvider != a.provider {
|
break
|
||||||
a.provider = fbProvider
|
}
|
||||||
log.Printf("[agent] switched active provider to %q after fallback",
|
|
||||||
fbProvider.Name())
|
|
||||||
}
|
}
|
||||||
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) {
|
if errors.Is(llmErr, context.Canceled) {
|
||||||
break
|
break
|
||||||
}
|
}
|
||||||
|
|||||||
@ -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() {
|
async function sendChat() {
|
||||||
var inp = document.getElementById("chat-input");
|
var inp = document.getElementById("chat-input");
|
||||||
var btn = document.getElementById("chat-send-btn");
|
var btn = document.getElementById("chat-send-btn");
|
||||||
@ -3677,7 +3717,6 @@ background:
|
|||||||
} finally {
|
} finally {
|
||||||
clearTimeout(ackTimer);
|
clearTimeout(ackTimer);
|
||||||
}
|
}
|
||||||
state.chatStage = __("AI 回复中...", "AI replying...");
|
|
||||||
var last = state.messages[state.messages.length - 1];
|
var last = state.messages[state.messages.length - 1];
|
||||||
if (r && r.response) {
|
if (r && r.response) {
|
||||||
if (last && last.role === "assistant" && last._streaming) {
|
if (last && last.role === "assistant" && last._streaming) {
|
||||||
@ -3716,13 +3755,14 @@ background:
|
|||||||
toast(__("请求失败: ", "Request failed: ") + e.message, true);
|
toast(__("请求失败: ", "Request failed: ") + e.message, true);
|
||||||
}
|
}
|
||||||
} finally {
|
} finally {
|
||||||
state.chatLoading = false;
|
if (r && r.response) {
|
||||||
state.chatStage = "";
|
// 同步兜底已拿到完整回复:回合结束
|
||||||
btn.disabled = false;
|
endChatTurn();
|
||||||
btn.textContent = __("发送", "Send");
|
} else {
|
||||||
var sb2 = document.getElementById("chat-stop-btn");
|
// 触发式受理(POST 已 abort/失败):回合仍打开,等 SSE 流式渲染;
|
||||||
if (sb2) sb2.style.display = "none"; // 回复完成/失败,隐藏停止按钮
|
// 由 agent_output final / reset 帧 / watchdog 收尾
|
||||||
rerenderChat(true);
|
armTurnWatchdog();
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@ -4146,6 +4186,7 @@ background:
|
|||||||
lastM2.content = p.content;
|
lastM2.content = p.content;
|
||||||
lastM2._final = true;
|
lastM2._final = true;
|
||||||
rerenderChat();
|
rerenderChat();
|
||||||
|
endChatTurn();
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
if (
|
if (
|
||||||
@ -4154,7 +4195,13 @@ background:
|
|||||||
lastM2._final &&
|
lastM2._final &&
|
||||||
!lastM2.source
|
!lastM2.source
|
||||||
) {
|
) {
|
||||||
return;
|
// 去重:同一轮的重复帧(如 SSE 重连回放)内容相同则忽略;
|
||||||
|
// 内容不同视为新一轮输出(上一轮已 final 且无 source),开新消息。
|
||||||
|
// 旧逻辑无条件 return 会丢弃多轮连发时新一轮的最终回复。
|
||||||
|
if (lastM2.content === p.content) {
|
||||||
|
endChatTurn();
|
||||||
|
return;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
state.messages.push({
|
state.messages.push({
|
||||||
role: "assistant",
|
role: "assistant",
|
||||||
@ -4164,6 +4211,7 @@ background:
|
|||||||
_final: true,
|
_final: true,
|
||||||
});
|
});
|
||||||
rerenderChat();
|
rerenderChat();
|
||||||
|
endChatTurn();
|
||||||
} catch (ex) {
|
} catch (ex) {
|
||||||
console.error("[SSE] agent_output error", ex);
|
console.error("[SSE] agent_output error", ex);
|
||||||
}
|
}
|
||||||
@ -4183,6 +4231,10 @@ background:
|
|||||||
lm._final = true;
|
lm._final = true;
|
||||||
rerenderChat();
|
rerenderChat();
|
||||||
}
|
}
|
||||||
|
// 轮次作废(用户中断):定格已显示内容;核心会以 [中断消息]
|
||||||
|
// 重启轮次,保持回合打开让确认回复继续流式渲染,
|
||||||
|
// 由其 agent_output final / watchdog 收尾。
|
||||||
|
armTurnWatchdog();
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
if (!p.content) return;
|
if (!p.content) return;
|
||||||
|
|||||||
@ -1342,8 +1342,11 @@ func (h *Handler) handleChat(w http.ResponseWriter, r *http.Request) {
|
|||||||
if body.ClientMsgID != "" {
|
if body.ClientMsgID != "" {
|
||||||
payload["client_msg_id"] = body.ClientMsgID
|
payload["client_msg_id"] = body.ClientMsgID
|
||||||
}
|
}
|
||||||
// 带超时的上下文,防止 InjectInputSync 长时间阻塞 HTTP 请求
|
// 带超时的上下文,防止 InjectInputSync 长时间阻塞 HTTP 请求。
|
||||||
ctx, cancel := context.WithTimeout(r.Context(), 60*time.Second)
|
// 注意:ctx 派生自 r.Context(),客户端提前断开(前端 15s ackTimer abort)时
|
||||||
|
// 立即取消,不会真等满 300s;300s 只约束"连接保持 + agent 排队/长生成"场景
|
||||||
|
// (agent 串行处理,后发消息的排队时间也计入,60s 曾导致连发第 3 条必超时)。
|
||||||
|
ctx, cancel := context.WithTimeout(r.Context(), 300*time.Second)
|
||||||
defer cancel()
|
defer cancel()
|
||||||
|
|
||||||
respCh := make(chan *agentIO.OutputEvent, 1)
|
respCh := make(chan *agentIO.OutputEvent, 1)
|
||||||
@ -1706,8 +1709,8 @@ func (h *Handler) handleOpenAICompletions(w http.ResponseWriter, r *http.Request
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
// 带超时的上下文
|
// 带超时的上下文(同 handleChat:客户端断开立即取消;300s 约束长生成)
|
||||||
ctx, cancel := context.WithTimeout(r.Context(), 60*time.Second)
|
ctx, cancel := context.WithTimeout(r.Context(), 300*time.Second)
|
||||||
defer cancel()
|
defer cancel()
|
||||||
|
|
||||||
respCh := make(chan *agentIO.OutputEvent, 1)
|
respCh := make(chan *agentIO.OutputEvent, 1)
|
||||||
@ -1719,7 +1722,7 @@ func (h *Handler) handleOpenAICompletions(w http.ResponseWriter, r *http.Request
|
|||||||
select {
|
select {
|
||||||
case response = <-respCh:
|
case response = <-respCh:
|
||||||
case <-ctx.Done():
|
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
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user