From 4e4e0ad656871d67a77e1d2ce2f97936d3a3596e Mon Sep 17 00:00:00 2001 From: JianFeeeee Date: Sun, 13 Sep 2026 00:39:56 +0800 Subject: [PATCH] =?UTF-8?q?feat(scheduler):=20M6=20=E4=BB=BB=E5=8A=A1?= =?UTF-8?q?=E7=BA=A7=E5=9B=9E=E6=89=A7=20+=20=E6=96=AD=E9=93=BE=E7=82=B9?= =?UTF-8?q?=E7=BB=9F=E4=B8=80=E4=B8=BA=E7=BB=88=E6=80=81=E4=BA=8B=E4=BB=B6?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 设计依据 docs/zh/input-scheduler-design.md §7(I5)、§11.3(X1/X2/X4)。 - 新增 emitSkippedReply:跳过路径(去重命中、空输入)给同步调用方一个 skipped 终态,但**不发 agent_output 事件**(避免 WebUI 聊天记录凭空多出 空消息)。修复前 cli/clawhubadapter 这类无超时的同步注入在去重命中时永久挂起 - 回执按任务归属:中断任务的回执只写自己的 ResponseCh,绝不误投给被挂起的等待者; 被挂起任务恢复并结束后才拿到自己的回执 - 新增 task_terminal_test.go 3 项:X1 回执不误投、X2 空输入有 skipped 终态、 X4 无超时同步调用方 1s 内拿到终态(回归判据) - 更新 M3a 去重用例:由「不得回执」改为「必须有 skipped 终态」 - 验收:agent 全量 + -race;全仓 build/vet 通过 --- internal/agent/core/eventloop.go | 34 ++++++ internal/agent/core/task.go | 5 + internal/agent/core/task_lifecycle_test.go | 11 +- internal/agent/core/task_terminal_test.go | 132 +++++++++++++++++++++ 4 files changed, 180 insertions(+), 2 deletions(-) create mode 100644 internal/agent/core/task_terminal_test.go diff --git a/internal/agent/core/eventloop.go b/internal/agent/core/eventloop.go index aa03fb2..ccfefba 100644 --- a/internal/agent/core/eventloop.go +++ b/internal/agent/core/eventloop.go @@ -271,6 +271,40 @@ func (a *Agent) mediaToBlocks(payload map[string]interface{}, mediaType string, return blocks, alt } +// emitSkippedReply 给被跳过任务的**同步**调用方一个终态。 +// +// 为什么要单独一条路径而不是复用 emitResponse:跳过意味着“我们没有处理这条输入”, +// 不应对外发 agent_output 事件(否则 WebUI 聊天记录会凭空多出一条空消息), +// 但必须写 ResponseCh——否则 cli/clawhub 这类无超时的同步注入会永久挂起。 +// +// 非阻塞写:ResponseCh 由同步调用方以 cap=1 创建,调用方超时离开后仍可写入。 +func (a *Agent) emitSkippedReply(evt *agentIO.InputEvent, reason string) { + if evt == nil || evt.ResponseCh == nil { + return + } + ch := evt.OutputChannel + if ch == "" { + ch = evt.Source + } + payload := map[string]interface{}{ + "content": "", + "request_id": evt.RequestID, + "skipped": true, + "reason": reason, + } + select { + case evt.ResponseCh <- &agentIO.OutputEvent{ + RequestID: evt.RequestID, + Target: evt.Source, + Type: "text", + Payload: payload, + Done: true, + OutputChannel: ch, + }: + default: + } +} + func (a *Agent) emitResponse(evt *agentIO.InputEvent, response string) { stageCtx := &sdk.StageContext{ FinalText: response, diff --git a/internal/agent/core/task.go b/internal/agent/core/task.go index 8c448a0..79149a5 100644 --- a/internal/agent/core/task.go +++ b/internal/agent/core/task.go @@ -214,6 +214,8 @@ func (a *Agent) prepareInputTask(evt *agentIO.InputEvent) (*TaskFrame, taskTermi in, ok := a.resolveInput(evt) if !ok { + // 空输入(文本与媒体都空):没有可处理内容,但同步调用方仍在等回执。 + a.emitSkippedReply(evt, "empty_input") return nil, terminalSkipped } @@ -222,6 +224,9 @@ func (a *Agent) prepareInputTask(evt *agentIO.InputEvent) (*TaskFrame, taskTermi // 是同一句,拿它去重会把连发的两张图误判成重复。 if len(in.blocks) == 0 && a.isDuplicateInput(evt.Source, in.text) { log.Printf("[agent] dropped duplicate input from %s: %s", evt.Source, truncateStr(in.text, 60)) + // 去重是「不处理」而不是「不回」,否则同步调用方(cli/clawhub 无超时) + // 会永久挂起(设计文档 §7 不变量 I5、§11.3 X2/X4)。 + a.emitSkippedReply(evt, "duplicate") return nil, terminalSkipped } diff --git a/internal/agent/core/task_lifecycle_test.go b/internal/agent/core/task_lifecycle_test.go index 747f605..743952b 100644 --- a/internal/agent/core/task_lifecycle_test.go +++ b/internal/agent/core/task_lifecycle_test.go @@ -102,8 +102,15 @@ func TestLifecycle_DuplicateSkippedHasTerminal(t *testing.T) { if a.context.Len() != after1 { t.Fatalf("去重命中不得提交上下文:%d → %d", after1, a.context.Len()) } - if len(ch2) != 0 { - t.Fatal("去重命中不得回执(原实现静默 return)") + if len(ch2) != 1 { + t.Fatal("去重命中必须回一个 skipped 终态,否则同步调用方永久挂起") + } + r := <-ch2 + if skipped, _ := r.Payload["skipped"].(bool); !skipped { + t.Fatalf("去重回执必须带 skipped=true,实际 %+v", r.Payload) + } + if reason, _ := r.Payload["reason"].(string); reason != "duplicate" { + t.Fatalf("reason=%q,期望 duplicate", reason) } if outputs != outputs1 { t.Fatalf("去重命中不得发输出事件:%d → %d", outputs1, outputs) diff --git a/internal/agent/core/task_terminal_test.go b/internal/agent/core/task_terminal_test.go new file mode 100644 index 0000000..a74fa76 --- /dev/null +++ b/internal/agent/core/task_terminal_test.go @@ -0,0 +1,132 @@ +package core + +// M6 验收测试:任务级回执与断链点统一为终态事件。 +// +// 设计依据 docs/zh/input-scheduler-design.md §7(不变量 I5)、§11.3(X1–X4)。 +// +// 问题背景:回执原先由全局 emitResponse 写(无任务归属),且 processInput 有多条 +// 「提前 return 而不 emit」的路径(解析失败、去重、consolidation)——同步调用方 +// 若不自带超时(cli、clawhubadapter)就会永久挂起。 + +import ( + "testing" + "time" + + agentAPI "gitcode.com/JianFeeeee/HomeAgent/internal/agent/api" + agentIO "gitcode.com/JianFeeeee/HomeAgent/internal/agent/io" +) + +// X1:回执按任务归属,中断的回执绝不投给被挂起的等待者。 +func TestTerminal_TaskScopedReplyNotMisrouted(t *testing.T) { + sp := newPreemptProvider("intr-done", "low-done") + a := newPreemptAgent(t, sp) + + lowEvt, lowCh := textEvent("qq", "低优先级任务") + lowTask := &Task{Kind: TaskKindInput, Level: LevelBackground, Event: lowEvt, EnqueuedAt: time.Now()} + if !a.sched.enqueue(lowTask) { + t.Fatal("入队失败") + } + lt, _, _ := a.sched.nextRef() + + done := make(chan struct{}) + go func() { a.executeNewTask(lt); close(done) }() + select { + case <-sp.entered: + case <-time.After(3 * time.Second): + t.Fatal("provider 未被调用") + } + + intrEvt, intrCh := textEvent("cli", "紧急打断") + intrEvt.Payload["interrupt"] = true + if !a.sched.requestPreempt(intrEvt, LevelCritical) { + t.Fatal("L4 应抢占 L1") + } + a.cancelCurrentLLM() + select { + case <-done: + case <-time.After(3 * time.Second): + t.Fatal("未挂起") + } + + // 执行中断任务 → 只应写它自己的回执通道。 + it, _, k := a.sched.nextRef() + if k != nextPending { + t.Fatalf("应取到 pending 中断,kind=%v", k) + } + a.executeNewTask(it) + if len(intrCh) != 1 { + t.Fatalf("中断任务应回执到自己的通道,实际 %d", len(intrCh)) + } + if len(lowCh) != 0 { + t.Fatal("中断的回执绝不能被投给被挂起的等待者") + } + + // 恢复并结束后,原任务才拿到自己的回执。 + rt, rf, k2 := a.sched.nextRef() + if k2 != nextSuspended { + t.Fatalf("应恢复被抢占任务,kind=%v", k2) + } + a.resumeTask(rt, rf) + if len(lowCh) != 1 { + t.Fatalf("恢复任务结束后应恰好回执一次,实际 %d", len(lowCh)) + } + if got, _ := (<-lowCh).Payload["content"].(string); got != "low-done" { + t.Fatalf("原任务回执内容=%q,期望 low-done", got) + } +} + +// X2:空输入(解析失败)也必须有终态回执。 +func TestTerminal_EmptyInputGetsSkippedReply(t *testing.T) { + a := newLifecycleAgent(t, &scriptProvider{}, nil, NewStageHost()) + ch := make(chan *agentIO.OutputEvent, 1) + evt := &agentIO.InputEvent{ + RequestID: "r-empty", + Source: "cli", + Type: "text", + Payload: map[string]interface{}{}, // 无 content,无媒体块 + OutputChannel: "cli", + ResponseCh: ch, + } + if _, out := a.runInputTask(evt, nil); out != outcomeDone { + t.Fatalf("空输入应正常返回,实际 %v", out) + } + if len(ch) != 1 { + t.Fatal("空输入必须回 skipped 终态") + } + if reason, _ := (<-ch).Payload["reason"].(string); reason != "empty_input" { + t.Fatalf("reason=%q,期望 empty_input", reason) + } + if a.context.Len() != 0 { + t.Fatal("空输入不得写入上下文") + } +} + +// X4:无超时的同步调用方(cli / clawhubadapter)在断链路径上不再永久挂起。 +// +// 这是回归判据:修复前 `InjectTextSync` 遇到去重命中会永久阻塞。 +func TestTerminal_NoTimeoutSyncCallerDoesNotHang(t *testing.T) { + a := newLifecycleAgent(t, &scriptProvider{script: []*agentAPI.CompletionResponse{ + {Content: "第一次"}, {Content: "第二次"}, + }}, nil, NewStageHost()) + + // 第一次成功 + e1, ch1 := textEvent("cli", "重复内容") + if _, out := a.runInputTask(e1, nil); out != outcomeDone { + t.Fatalf("首次=%v", out) + } + if len(ch1) != 1 { + t.Fatal("首次应有回执") + } + + // 第二次(去重命中):模拟同步调用方阻塞等待——必须在 1s 内拿到终态。 + e2, ch2 := textEvent("cli", "重复内容") + go a.runInputTask(e2, nil) + select { + case r := <-ch2: + if skipped, _ := r.Payload["skipped"].(bool); !skipped { + t.Fatalf("应为 skipped 终态,实际 %+v", r.Payload) + } + case <-time.After(time.Second): + t.Fatal("去重命中让同步调用方永久挂起(X4 回归)") + } +}