feat(streaming): token-level delta events + interrupt for CLI/WebUI/GUI

Expose the LLM token-level streaming deltas (EventReasoningDelta /
EventContentDelta) to every client channel and add user-initiated
interrupt (cancel generation / send interrupt message) to all three
frontends, preserving the existing interrupt-injection semantics.

SDK/events:
  - EventReasoningDelta, EventContentDelta constants exported in the
    public/internal SDK event alias tables.

CLI plugin:
  - handleChat subscribes to both delta events and forwards
    reasoning_delta / content_delta JSON frames (channel-filtered);
    aggregated reasoning/tool_call/response frames still fire as before.
  - New /stop (alias /interrupt) builtin injects an interrupt via
    InjectInterrupt(cliSource, cliChannel) - matches interceptLoop
    semantics: cancels an active stream and re-injects the message as
    a [中断消息] for a restarted turn; with no active LLM it behaves
    as a plain input.

Waiter client (line mode + TUI):
  - streamRender accumulates delta chunks and redraws the current line;
    a reset frame (stream abandoned, e.g. user interrupt) flushes the
    partial buffer so the next turn does not concatenate onto stale
    content. Aggregated frames terminate the delta line and render the
    final text (old servers without deltas behave exactly as before).
  - TUI merges content_delta into the in-flight agent message and seals
    it (final flag) on response/tool_call/error so subsequent deltas
    never append to a finished message.

WebUI:
  - SSE handler subscribes to the two delta events but does NOT record
    them into the replay ring - reconnection replays only aggregated
    events (the final truth), avoiding duplicate delta accumulation.
  - POST /api/v1/chat/interrupt calls InjectInterrupt(webui, webui)
    with optional message; fronted by a Stop button shown only while
    a generation is in flight.

dashboard.html / GUI app.js:
  - Stop button next to Send (hidden until chatLoading); interruptChat
    POSTs /chat/interrupt. Delta listeners append incrementally;
    agent_output (aggregated) now replaces (not appends) the in-flight
    content and marks _final; reset frames finalize the partial message.

process.go:
  - chatStreamWithFallback preserves the context.Canceled/
    DeadlineExceeded contract: a user interrupt returns the canceled
    error (never a partial-content success) so the existing continue
    branch restarts the turn with the [中断消息]. A reset
    EventContentDelta is published so connected clients drop stale
    partial renderings before the new turn begins.

Verified: /stop 'msg' via waiter triggers 'interrupt from cli/cli' in
interceptLoop; unit TestChatStreamCancelPreservesInterrupt confirms the
canceled error propagates instead of being swallowed.
This commit is contained in:
JianFeeeee
2026-08-25 10:50:37 +08:00
parent bb53fddd8b
commit a51079e7cd
11 changed files with 423 additions and 7 deletions

View File

@ -2774,6 +2774,9 @@ background:
'<input id="chat-input" placeholder="' +
__("输入消息...", "Type a message...") +
'" onkeydown="if(event.key==\'Enter\')sendChat()">' +
'<button class="btn" onclick="interruptChat()" id="chat-stop-btn" style="display:none;background:var(--danger, #d1383d);color:#fff">' +
__("停止", "Stop") +
"</button>" +
'<button class="btn btn-primary" onclick="sendChat()" id="chat-send-btn">' +
__("发送", "Send") +
"</button>" +
@ -3633,6 +3636,7 @@ background:
async function sendChat() {
var inp = document.getElementById("chat-input");
var btn = document.getElementById("chat-send-btn");
var stopBtn = document.getElementById("chat-stop-btn");
var text = inp.value.trim();
if (!text || state.chatLoading) return;
state.chatStick = true;
@ -3644,6 +3648,7 @@ background:
state.chatStage = __("等待AI回复...", "Waiting for AI...");
btn.disabled = true;
btn.textContent = "";
if (stopBtn) stopBtn.style.display = ""; // 生成期间可停止
rerenderChat(true);
// 触发式 POST:短超时仅确认受理;回复靠 SSE 流式渲染(对齐 GUI 行为)。
try {
@ -3715,10 +3720,26 @@ background:
state.chatStage = "";
btn.disabled = false;
btn.textContent = __("发送", "Send");
var sb2 = document.getElementById("chat-stop-btn");
if (sb2) sb2.style.display = "none"; // 回复完成/失败,隐藏停止按钮
rerenderChat(true);
}
}
// 停止生成 / 发送中断消息。核心拦截语义:有 LLM 在跑则取消当前
// 请求并以 [中断消息] 重启轮次;无则在跑则作为普通消息处理。
async function interruptChat() {
try {
await api("/chat/interrupt", {
method: "POST",
body: JSON.stringify({}),
});
toast(__("已发送中断信号", "Interrupt signal sent"));
} catch (e) {
toast(__("中断失败: ", "Interrupt failed: ") + e.message, true);
}
}
async function queryMemoryChat() {
var q = document.getElementById("mem-query")?.value;
var r = document.getElementById("mem-result-chat");
@ -4119,8 +4140,11 @@ background:
lastM2.role === "assistant" &&
!lastM2._final
) {
// 聚合最终响应:覆盖 delta 累积的中间内容(以聚合为准,
// 含 stage 插件改写后的最终文本),并置 final 结束本轮流式。
lastM2._grow = true;
lastM2.content += p.content;
lastM2.content = p.content;
lastM2._final = true;
rerenderChat();
return;
}
@ -4137,12 +4161,50 @@ background:
content: p.content,
_streaming: true,
_grow: true,
_final: true,
});
rerenderChat();
} catch (ex) {
console.error("[SSE] agent_output error", ex);
}
});
// token 级流式增量:逐块追加到当前回复内容(流式生成中);
// reset 帧表示轮次作废(用户中断):定格已显示的部分内容,置 final。
es.addEventListener("content_delta", function (e) {
try {
var ev = JSON.parse(e.data);
var p = ev.payload || {};
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;
rerenderChat();
}
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;
rerenderChat();
} catch (ex) {}
});
es.addEventListener("terminal_output", function (e) {
try {
var ev = JSON.parse(e.data);
@ -4170,6 +4232,7 @@ background:
try {
var ev = JSON.parse(e.data);
var p = ev.payload || {};
if (p.channel === "_consolidation_") return;
if (p.content) {
state.chatStage = __("AI 思考中...", "AI thinking...");
var last =
@ -4186,12 +4249,39 @@ background:
});
last = state.messages[state.messages.length - 1];
}
last.reasoning_content =
(last.reasoning_content || "") + p.content;
// 聚合 reasoning 帧携带全文:直接覆盖(若已有 delta 累积则等价)
last.reasoning_content = p.content;
rerenderChat();
}
} catch (ex) {}
});
// token 级流式增量:逐块追加到当前思考内容;reset 帧表示轮次作废
es.addEventListener("reasoning_delta", function (e) {
try {
var ev = JSON.parse(e.data);
var p = ev.payload || {};
if (p.channel === "_consolidation_") return;
if (p.reset) 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;
rerenderChat();
} catch (ex) {}
});
es.addEventListener("tool_call", function (e) {
try {
var ev = JSON.parse(e.data);

View File

@ -694,6 +694,7 @@ func (h *Handler) RegisterRoutes(mux *http.ServeMux) {
mux.HandleFunc("/api/v1/tracker/", h.requireAPI(h.handleTracker))
mux.HandleFunc("/api/v1/chat", h.requireAPI(h.handleChat))
mux.HandleFunc("/api/v1/chat/history", h.requireAPI(h.handleChatHistory))
mux.HandleFunc("/api/v1/chat/interrupt", h.requireAPI(h.handleChatInterrupt))
mux.HandleFunc("/api/v1/chat/events", h.requireAPI(h.handleChatEvents))
mux.HandleFunc("/api/v1/terminals", h.requireAPI(h.handleTerminals))
mux.HandleFunc("/api/v1/cmd/history", h.requireAPI(h.handleCmdHistory))
@ -1221,6 +1222,39 @@ func (h *Handler) handleChatHistory(w http.ResponseWriter, r *http.Request) {
writeJSON(w, http.StatusOK, map[string]interface{}{"messages": result})
}
// handleChatInterrupt 注入用户中断:取消正在进行的 LLM 生成并/或发送打断消息。
// 核心拦截语义(interceptLoop):
// - 有 LLM 在跑:cancelLLM 取消当前请求 + 中断入队,process() 以
// [中断消息] 重启轮次,模型看到被打断的上下文和用户新输入;
// - 无 LLM 在跑:作为普通输入处理(等同发了一条消息)。
// message 可选:空则纯取消(仍会注入空内容中断触发取消)。
func (h *Handler) handleChatInterrupt(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodPost {
http.Error(w, "method not allowed", http.StatusMethodNotAllowed)
return
}
if h.sdk == nil {
writeJSON(w, http.StatusServiceUnavailable, map[string]string{"error": "agent unavailable"})
return
}
var body struct {
Message string `json:"message"`
DeviceID string `json:"device_id"`
}
if r.Body != nil {
_ = json.NewDecoder(r.Body).Decode(&body) // body 可选
}
source := "webui"
if body.DeviceID != "" {
source = "webui/" + body.DeviceID
}
h.sdk.InjectInterrupt(source, "webui", "text", map[string]interface{}{
"content": body.Message,
})
writeJSON(w, http.StatusOK, map[string]string{"status": "interrupted"})
}
func (h *Handler) handleTerminals(w http.ResponseWriter, r *http.Request) {
h.termMu.Lock()
terms := make([]*termState, 0, len(h.termStates))
@ -1413,8 +1447,27 @@ func (h *Handler) handleChatEvents(w http.ResponseWriter, r *http.Request) {
}
subTypes := []string{"agent_output", "reasoning", "agent_error", "tool_call", "stage", "agent_llm_chain", "terminal_output"}
// token 级流式增量事件:实时转发给浏览器做逐 token 渲染。
// 不进 sseEventRing —— 断线重连只重放聚合事件(最终真相),
// 避免重放 delta 与聚合内容重复追加。
var unsubs []func()
var seq int64
appendDeltaSub := func(evtType sdk.EventType) {
unsub := h.sdk.Subscribe(evtType, func(evt *sdk.Event) {
data, _ := json.Marshal(evt)
seq++
id := fmt.Sprintf("%d-%d", evt.Timestamp, seq)
select {
case writeCh <- fmt.Sprintf("id: %s\nevent: %s\ndata: %s\n", id, evt.Type, string(data)):
default:
log.Printf("[SSE] DROPPED %s (writeCh full, len=%d)", evt.Type, len(writeCh))
}
})
unsubs = append(unsubs, unsub)
}
appendDeltaSub(sdk.EventReasoningDelta)
appendDeltaSub(sdk.EventContentDelta)
for _, t := range subTypes {
t2 := t
unsub := h.sdk.Subscribe(sdk.EventType(t2), func(evt *sdk.Event) {