Files
HomeAgent/internal/agent/core/process.go
JianFeeeee ccc2ac2d4d fix(stop): 停止按钮真正生效——停止 ≠ 空中断;鸿蒙 screensue 支持 HTML
两处鸿蒙端缺陷 + 一个跨端(WebUI/GUI/鸿蒙)的停止语义缺陷。

## 症状(实测取证)

1. **鸿蒙终止按钮按下没反应**。POST /chat/interrupt 带空 body,接口回 200
   `{"status":"interrupted"}`,但 journalctl 零中断日志、生成继续跑到自然结束。
2. **鸿蒙 screensue 不解析 HTML**,把标签当普通字符串显示。

## 根因

停止按钮走的是「空内容中断」,而 interceptLoop 有一行
`if text == "" { continue }` —— 空内容被判为「无事发生」直接丢弃。
所以停止指令从未到达调度器;接口那个 200 是不诚实的。

另查明两条会放大症状的既有问题(停止后仍在跑):
- `chatStreamWithFallback`:流式连接失败时无条件回退非流式 `Chat`。
  上下文已取消时这等于**再发一次完整请求**(停止后模型继续生成)。
- `stepLLM`:`context.Canceled` 一律 `outcomeContinue` 重跑本步。
  这是给「被更高中断抢占」用的(现场要交出去、稍后继续),
  但用户按停止是「不要了」,重跑就是停止没生效。

## 修法(按用户明确的设计)

停止 = ①立即结束当前 LLM 推理(不重试、不恢复);
②对**停止那一刻已排队**的 x 条消息,后续在 pre-action 阶段依次短路。

- scheduler:新增 `armStop`(登记快照配额并返回当时排队深度)/`takeStop`/
  `consumeCancel`。配额取快照值(停止后新到的输入不受影响),
  重复按停止取 max 不累加(两个客户端同时按不该翻倍)。
- `interceptLoop`:读 `stop` 标记。停止时 armStop + cancelCurrentLLM;
  **纯停止不再进中断队列**(旧实现把它当空中断入队,所以停完还会活)。
  带注释的停止(`/stop 换个话题`)仍走中断路径。
- `stepLLM`:取消 + `takeStop()` → 直接 `outcomeDone`(不再重跑)。
- `stepPrepare`:`consumeCancel()` 命中即在 pre-action 短路收尾。
- `chatStreamWithFallback`:以 **ctx.Err()** 为判据拒绝回退(不是「错误是不是
  Canceled」——很多 provider 用 Canceled 表示「不支持流式」,那种必须继续回退,
  否则会把探测误判成取消;这条区分是跑全量测试时才暴露的)。
- WebUI handler / CLI `/stop`:空消息时带 `stop:true`。

## 鸿蒙端

- `BridgeCaps.ets`:新增 `looksLikeHtml`(首字符 '<' + 字母开头标签名,
  避免误判 "<3" 这类文本)、`screensueHtml`、`escapeHtmlText`。
- `ScreensuePage.ets`:HTML 走 **RichText**(只解析 HTML 子集、无脚本无网络),
  纯文本仍走 Text。不用 Web 组件:agent 下发的是第三方内容,
  Web 默认带 javaScriptAccess/fileAccess,等于让远端内容在客户端执行脚本。
  注入主题前景色,避免 RichText 用系统默认色导致深色主题下黑字不可见。
- `ChatSession.ets`:`interruptChat` 改发 `{stop:true}`(含类型声明,
  ArkTS 禁止无类型对象字面量),并在本地即时复位忙态 + 提示「已停止」。

## 验证

- 新增 `stop_semantics_test.go`:停止终结任务不重试(provider 调用次数恒为 1)、
  配额是快照(x 条短路、随后新到的不受影响)、重复 arm 取 max。
- `go test ./internal/... ./cmd/...` 全绿。
- 鸿蒙 HAP 构建通过;unsigned 包已装进模拟器(signed 包受
  READ_PASTEBOARD 授权限制装不上,与既有记录一致)。
2026-09-18 11:26:12 +08:00

424 lines
15 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

package core
import (
"context"
"encoding/json"
"errors"
"fmt"
"log"
"strings"
agentAPI "gitcode.com/JianFeeeee/HomeAgent/internal/agent/api"
"gitcode.com/JianFeeeee/HomeAgent/internal/events"
sdk "gitcode.com/JianFeeeee/HomeAgent/internal/sdk"
)
// continuationPlaceholder 是工具轮之后补的 user 占位内容。
//
// zen 兼容网关要求请求最后一条必须是 userthinking 续写模式校验),工具轮
// 产出 assistant/tool 结尾会被 400 拒绝;首轮 system 结尾不补,否则会覆盖
// 真实用户输入。
//
// 用独立常量 + 精确等值判定,是因为这条消息是**核心自己插入的**、不是用户输入,
// 所以可以安全地按内容识别并在补位前移除上一条,保证至多一条。
const continuationPlaceholder = "请根据以上工具结果继续。"
// replyDeliveredPlaceholder 是「本批工具调用全部是输出通道发送」之后补的占位。
//
// 为何不能继续用通用的「请继续」异步通道qq/wechat的回复**只能**经
// output_send__* 交付(纯文本不送达,见 buildSystemPrompt 的输出规则)。于是
// 模型「已经回复完了」的表达形式就是一个工具调用,而紧随其后的
// 「请根据以上工具结果继续。」会被读成「还要再做一步」——能做的「一步」恰好
// 还是再发一条消息。两者叠加成自我强化的发送循环:生产实测单轮 34 次
// output_send__qq、持续 514 秒,直到 QQ 插件自己的循环保险拒绝发送才停下。
//
// 所以这里换成一条明确的终止许可:已回复完就直接返回纯文本收尾。
const replyDeliveredPlaceholder = "若你的回复已完成,直接返回纯文本即可结束本轮,无需再调用任何工具。"
// continuationFor 选择工具轮之后补位的 user 占位文案。
// replyOnly 表示上一批工具调用全部是输出通道发送(即模型刚交付了回复)。
func continuationFor(replyOnly bool) string {
if replyOnly {
return replyDeliveredPlaceholder
}
return continuationPlaceholder
}
// isOutputDeliveryTool 判断工具是否是「向输出通道交付内容」。
// output_send__{channel}_help 只是查询用法,不算交付。
func isOutputDeliveryTool(name string) bool {
return strings.HasPrefix(name, "output_send__") && !strings.HasSuffix(name, "_help")
}
// isContinuationPlaceholder 判断一条 user 消息是否是本机制插入的占位。
// 只按两个常量精确匹配,不碰任何真实用户消息。
func isContinuationPlaceholder(m agentAPI.Message) bool {
return m.Role == "user" &&
(m.Content == continuationPlaceholder || m.Content == replyDeliveredPlaceholder)
}
// toolOutputForQuery 返回用于相关性计算的工具输出**有效内容**。
//
// 为什么要过 Cleaner 而不是直接用原始 resultContextPolicy=prune 的入参是
// **相关性查询向量**——它决定保留/归档哪些上下文事件。原始工具输出里混着
// ANSI 转义、base64、JSON 包装等噪声,直接拿去向量化会让打分失真。
// 而 ToolDef.Cleaner 的契约本就写着“仅在向量化/jieba/蒸馏时调用”,裁剪正是
// 在向量化,所以这里必须过它(此前只在构建事件向量时用了,裁剪查询漏了)。
//
// Cleaner 未注册或 RPC 失败时回退原文(清洗是计算层优化,不能因此丢内容);
// 返回空串时也回退——空串会让查询向量退化成零向量,裁剪就失去判据。
func (a *Agent) toolOutputForQuery(toolName, raw string) string {
if a.stageHost == nil {
return raw
}
cleaner := a.stageHost.ToolDefCleaner(toolName)
if cleaner == nil {
return raw
}
if cleaned := cleaner(raw); cleaned != "" {
return cleaned
}
return raw
}
// recallMsgMarker 是工具触发召回时注入的 system 消息前缀。
// 用它做去重与替换的识别标(与用户/中断的 system 消息区分开)。
const recallMsgMarker = "【记忆召回】"
// appendOrReplaceRecall 把一段召回文本作为 system 消息挂到消息末尾。
//
// 同一任务内多次触发(如模型多次调用 qq_get_message时**替换**上一条召回,
// 而不是累加:否则召回会线性叠进 prompt把上下文与 token 预算越挤越紧。
// 替换位置固定在末尾,不影响 tool/assistant 消息的配对。
func appendOrReplaceRecall(msgs []agentAPI.Message, recallText string) []agentAPI.Message {
if recallText == "" {
return msgs
}
full := recallMsgMarker + "\n" + recallText
for i := len(msgs) - 1; i >= 0; i-- {
if msgs[i].Role == "system" && strings.HasPrefix(msgs[i].Content, recallMsgMarker) {
msgs[i].Content = full
return msgs
}
}
return append(msgs, agentAPI.Message{Role: "system", Content: full})
}
// dropContinuationPlaceholders 移除此前由本机制插入的 user 占位。
//
// 为什么必须移除而不仅仅是“不再追加”:`msgs` 在循环外创建、循环内只增不减,
// 占位是核心自己插的、不是用户说的话。不移除的话prompt 里就会线性叠上
// N 条一模一样的“继续”,把前缀上下文(含记忆注入)往后挤。
func dropContinuationPlaceholders(msgs []agentAPI.Message) []agentAPI.Message {
out := msgs[:0]
for _, m := range msgs {
if isContinuationPlaceholder(m) {
continue
}
out = append(out, m)
}
return out
}
// chatStreamWithFallback 优先流式调用 provider失败时回退非流式 Chat()。
//
// 流式路径ChatStream 拿到 chunk channel逐块累积 content/reasoning_content
// 并发布 EventReasoningDelta / EventContentDelta 增量事件(新订阅者可选订,
// 旧订阅者不认识自然忽略)。流结束后拼出与 Chat() 等价的 CompletionResponse
// 返回——process() 的后续逻辑stageCtx/聚合事件/工具循环)完全不变。
//
// 回退条件ChatStream 返回错误连接失败、provider 不支持流式)。
// 已收到部分 chunk 后出错则不回退(避免重复生成),直接返回已累积内容。
//
// 超时收益:首包 ~1-3s 到达即建立活性,后续只要 token 在流动就不会触发
// 空闲超时;总生成时长不再受限於 180s 整体超时。
func chatStreamWithFallback(ctx context.Context, p agentAPI.Provider, req *agentAPI.CompletionRequest, a *Agent, channel string) (*agentAPI.CompletionResponse, error) {
ch, err := p.ChatStream(ctx, req)
if err != nil {
// 上下文已取消/超时:**绝不能**回退到非流式 Chat。
//
// 回退意味着再发一次完整请求,而这时用户已经按了停止(或请求已超时),
// 结果是“按了停止又跑了一遍”——停止按钮看起来毫无反应的一个真实成因。
//
// 判据取 **ctx.Err()** 而不是“错误是不是 context.Canceled”很多 provider
// 不支持流式时也回 Canceled 表示“请走非流式”(仓里大量假 provider 即如此),
// 那种情况必须继续回退,否则会把“不支持流式”误当成“已被取消”。
if ctx.Err() != nil {
if err == nil {
err = ctx.Err()
}
return nil, err
}
log.Printf("[agent] stream connect failed (%v), falling back to non-stream chat", err)
return p.Chat(ctx, req)
}
resp, accErr := accumulateStream(ctx, ch, a, channel)
// 中断/超时取消必须保持取消语义传给调用方(与原 Chat() 行为一致:
// 被 cancel 时丢弃已收内容返回 err让 process() 的 continue 分支
// 重启轮次并以 [中断消息] 注入打断内容。绝不能把部分内容当成功返回,
// 否则用户打断会被无视、继续执行工具/输出。
if errors.Is(accErr, context.Canceled) || errors.Is(accErr, context.DeadlineExceeded) {
// 通知客户端:本轮流式作废,清空 delta 累积并定格已显示内容
if a != nil {
a.publishEvent(events.EventContentDelta, map[string]interface{}{
"content": "",
"channel": channel,
"reset": true,
})
}
return resp, accErr
}
if accErr == nil {
return resp, nil
}
// 其他错误(网络中断等):已累积到实质内容则返回部分结果,否则回退非流式。
if resp != nil && (resp.Content != "" || len(resp.ToolCalls) > 0) {
log.Printf("[agent] stream interrupted mid-way (%v), returning partial result", accErr)
return resp, nil
}
// 取消类错误不能回退(否则等于再跑一遍完整的非流式请求)。
// 同样以 ctx.Err() 为准真取消才拦provider 探活不算。
if ctx.Err() != nil {
return resp, accErr
}
log.Printf("[agent] stream failed before content (%v), falling back to non-stream chat", accErr)
return p.Chat(ctx, req)
}
// toolCallAcc 累积流式 tool call 的各个分片。OpenAI 风格:每个 index 的
// id/name/arguments 跨多个 chunk 增量到达arguments 是 JSON 字符串分片。
type toolCallAcc struct {
id string
name string
argsRaw strings.Builder
}
// accumulateStream 消费 chunk channel累积为完整 CompletionResponse
// 同时发布增量事件。返回的 response 与非流式 Chat() 的返回等价。
func accumulateStream(ctx context.Context, ch <-chan agentAPI.StreamChunk, a *Agent, channel string) (*agentAPI.CompletionResponse, error) {
resp := &agentAPI.CompletionResponse{
ToolCalls: make([]agentAPI.ToolCall, 0),
}
accs := make(map[int]*toolCallAcc) // index → 累积中的 tool call
var lastFinish string
flushToolCall := func(idx int) {
acc := accs[idx]
if acc == nil {
return
}
if acc.name == "" {
log.Printf("[agent] stream tool_call idx=%d flushed with EMPTY name (args=%q) — dropped", idx, truncateStr(acc.argsRaw.String(), 120))
delete(accs, idx)
return
}
args, argsOK := parseToolArgsJSON(acc.argsRaw.String())
raw := strings.TrimSpace(acc.argsRaw.String())
// 空参诊断:区分「上游没发分片」(raw="")、「混拼污染」(解析失败) 与「合法空对象」({})。
if !argsOK {
log.Printf("[agent] stream tool_call %s (idx=%d) argument fragments invalid JSON: %q", acc.name, idx, truncateStr(raw, 200))
} else if raw == "" {
log.Printf("[agent] stream tool_call %s (idx=%d) received NO argument fragments", acc.name, idx)
}
tc := agentAPI.ToolCall{
ID: acc.id,
Name: acc.name,
Arguments: args,
}
resp.ToolCalls = append(resp.ToolCalls, tc)
delete(accs, idx)
}
for {
select {
case ck, ok := <-ch:
if !ok {
for idx := range accs {
flushToolCall(idx)
}
if lastFinish != "" {
resp.FinishReason = lastFinish
}
return resp, nil
}
if ck.ReasoningContent != "" {
resp.ReasoningContent += ck.ReasoningContent
if a != nil {
a.publishEvent(events.EventReasoningDelta, map[string]interface{}{
"content": ck.ReasoningContent,
"channel": channel,
})
}
}
if ck.Content != "" {
resp.Content += ck.Content
if a != nil {
a.publishEvent(events.EventContentDelta, map[string]interface{}{
"content": ck.Content,
"channel": channel,
})
}
}
// 增量 tool call 分片OpenAI 风格按 index 字段拼接 id/name/arguments。
// 注意必须用分片自带的 StreamIndex上游 JSON "index"),不能用 Go
// range 序号:每个 SSE chunk 通常只含一个 tool_call 元素slice 序号
// 恒为 0并行多工具调用index=0,1,2...)的分片会全部污染到同一个桶,
// 导致 name 相互覆盖、args 碎片混拼解析失败(空参数工具调用)。
for _, tc := range ck.ToolCalls {
idx := tc.StreamIndex
if idx == 0 && tc.Name == "" && tc.RawArguments == "" {
continue
}
acc := accs[idx]
if acc == nil {
acc = &toolCallAcc{}
accs[idx] = acc
}
if tc.ID != "" {
acc.id = tc.ID
}
if tc.Name != "" {
acc.name = tc.Name
}
// arguments 以 JSON 字符串分片到达OpenAI 标准),拼接后最终解析
if tc.RawArguments != "" {
acc.argsRaw.WriteString(tc.RawArguments)
}
}
if ck.Done && ck.FinishReason != "" {
lastFinish = ck.FinishReason
}
if ck.Usage != nil {
resp.TokenUsage = *ck.Usage
}
case <-ctx.Done():
for idx := range accs {
flushToolCall(idx)
}
return resp, ctx.Err()
}
}
}
// parseToolArgsJSON 将经过完整拼接的 tool call arguments JSON 字符串解析为 map。
// 第二个返回值 ok=false 表示分片拼接结果不是合法 JSON分片污染/丢失),
// 与「合法的空对象 {}」相区分。
func parseToolArgsJSON(s string) (map[string]interface{}, bool) {
if strings.TrimSpace(s) == "" {
return map[string]interface{}{}, true
}
var m map[string]interface{}
if err := json.Unmarshal([]byte(s), &m); err == nil && m != nil {
return m, true
}
return map[string]interface{}{}, false
}
func convertToolCalls(tcs []agentAPI.ToolCall) []sdk.ToolCall {
if tcs == nil {
return nil
}
result := make([]sdk.ToolCall, len(tcs))
for i, tc := range tcs {
result[i] = sdk.ToolCall{ID: tc.ID, Name: tc.Name, Arguments: tc.Arguments}
}
return result
}
func convertBackToolCalls(tcs []sdk.ToolCall) []agentAPI.ToolCall {
if tcs == nil {
return nil
}
result := make([]agentAPI.ToolCall, len(tcs))
for i, tc := range tcs {
result[i] = agentAPI.ToolCall{ID: tc.ID, Name: tc.Name, Arguments: tc.Arguments}
}
return result
}
func (a *Agent) docStoreSize() int {
if a.docStore == nil {
return 0
}
s := a.docStore.Stats()
if n, ok := s["doc_count"]; ok {
if ni, ok := n.(int); ok {
return ni
}
}
return 0
}
func (a *Agent) formatMergedTimeline(maxTokens int) string {
a.context.mu.Lock()
events := make([]*ContextEvent, len(a.context.events))
copy(events, a.context.events)
a.context.mu.Unlock()
if len(events) == 0 {
return ""
}
// 第一轮:从最新到最旧,计算在预算内能放多少条
headerTokens := EstimateTokens("【对话时序】\n")
remaining := maxTokens - headerTokens
include := 0
for i := len(events) - 1; i >= 0; i-- {
e := events[i]
// ❗单位必须与 EstimateTokens 一致rune×2。这里曾用 `len()`**字节**)再 ×2
// CJK 一字 3 字节 ⇒ 中文事件被高估 3 倍,窗口还有余量也会提前 break
// 把更早的事件整段丢掉实测2384 字的中文事件被估成 14398 token > 8192
estTokens := EstimateTokens(e.Source) + EstimateTokens(e.Input) + 40
if e.Response != "" {
estTokens += 120
}
if remaining-estTokens < 0 && include > 0 {
break
}
remaining -= estTokens
include++
}
if include == 0 && len(events) > 0 {
include = 1
}
// 第二轮:按时间正序渲染
start := len(events) - include
if start < 0 {
start = 0
}
var sb strings.Builder
sb.WriteString("【对话时序】\n")
for _, e := range events[start:] {
sb.WriteString(fmt.Sprintf("[%s] %s: %s",
e.Timestamp.Format("15:04:05"), e.Source, e.Input))
if len(e.ToolsUsed) > 0 {
sb.WriteString(fmt.Sprintf(" → 调用工具: %s", strings.Join(e.ToolsUsed, ", ")))
}
if e.Response != "" {
sb.WriteString(fmt.Sprintf(" → %s", truncateStr(e.Response, 120)))
}
sb.WriteString("\n")
}
return sb.String()
}
func (a *Agent) buildMessages(sysPrompt, input string, ctxTokens int) []agentAPI.Message {
msgs := []agentAPI.Message{{Role: "system", Content: sysPrompt}}
if ctxTok := a.formatMergedTimeline(ctxTokens); ctxTok != "" {
msgs = append(msgs, agentAPI.Message{Role: "system", Content: ctxTok})
}
msgs = append(msgs, agentAPI.Message{Role: "user", Content: input})
return msgs
}