fix(scheduler)!: D1 更正为「中断从上一个任务之前的完整状态开始」,并实现现场合回

用户明确语义(我此前对 D1 的解析就是错的——当时回答里的“A”指的是 git 选项,
D1 实际要的是方案 B):

  中断打断时,上个任务到达以来的所有上下文现场被保护(含 toolcall),
  然后中断在「上个任务前的那个完整状态」上开始运行;
  中断结束后再把被挂起的任务与其上下文现场加载回中断任务之上,并继续运行。

实现:
- 删除 SeedMsgs 与 D1=A 的“只读前缀”机制:中断任务不再继承被打断任务的任何内容,
  它就是普通新任务,正常走完整 prepare(system prompt + timeline + 自己的输入)
- TaskFrame 新增 PrefixLen(基础前缀长度)与 InputBlocks;
  stepPrepare 在 buildMessages 之后记录 PrefixLen
- 新增 rebaseFramePrefix:恢复时重建基础前缀(中断已提交进 a.context,
  重建的 timeline 含中断效果=“加载回中断之上”),再把本任务自己的尾部
  (Stage 上下文 + 工具轮产物 + 占位)接回;并补回 IsInterrupt 标记与多模态块
- resumeTask 在 runTaskSteps 之前调用 rebaseFramePrefix
- 设计稿 §5.3 改写为「已定:D1=B」并写明实现对应;§6.2 补“重建前缀→接回尾部”;
  §12 的 D1 行更新

测试:
- TestPreempt_HigherPreemptsAndResumes 改为断言「中断不继承、恢复后看得见中断内容」
- 新增 TestPreempt_ResumeRebaseRestoresTailDecorations(前缀重建后尾部装饰补回)
- 原 TestPreempt_SeedPathDoesNotLeakInterruptFlag 随之删除(机制已不存在)

验收:agent 全量 + -race;全仓 build/vet 通过
This commit is contained in:
JianFeeeee
2026-09-13 06:24:20 +08:00
parent bad29c00dc
commit bf0d2510cc
8 changed files with 159 additions and 80 deletions

View File

@ -24,7 +24,6 @@ import (
"sync/atomic"
"time"
agentAPI "gitcode.com/JianFeeeee/HomeAgent/internal/agent/api"
agentIO "gitcode.com/JianFeeeee/HomeAgent/internal/agent/io"
"gitcode.com/JianFeeeee/HomeAgent/internal/events"
sdk "gitcode.com/JianFeeeee/HomeAgent/internal/sdk"
@ -107,10 +106,6 @@ type Task struct {
Event *agentIO.InputEvent // Kind == TaskKindInput
Self selfInputMsg // Kind == TaskKindSelf
// SeedMsgs 是抢占式中断任务的只读前缀(D1=A):由被打断的任务在挂起时
// 附上,使中断任务看得见「进行到哪一步」,但其产出不合并回原任务。
SeedMsgs []agentAPI.Message
// PreemptCount 是本任务被抢占的次数,用于饥饿防护:
// effectiveLevel = min(L4, Level + min(PreemptCount, 2))。
PreemptCount int
@ -425,23 +420,8 @@ func (s *scheduler) suspend(t *Task, f *TaskFrame) {
s.running = nil
}
// D1=A:把被抢占任务的只读前缀交给造成本次抢占的中断任务。
// 选最高优先级的待处理中断;若它已有前缀(嵌套抢占)则不覆盖。
if s.preemptArmed {
var victim *Task
for _, it := range s.pendingInterrupts {
if it.Level < s.preemptLevel {
continue
}
if victim == nil || taskBefore(victim, it) {
victim = it
}
}
if victim != nil && len(victim.SeedMsgs) == 0 {
victim.SeedMsgs = append([]agentAPI.Message(nil), f.Msgs...)
}
}
// D1=B:中断任务在上一个任务之前的完整状态上开始运行,
// 因此这里**不**把被打断任务的任何内容交给它。
s.preemptArmed = false
s.preemptLevel = 0
}
@ -645,9 +625,9 @@ func (a *Agent) executeNewTask(t *Task) {
}()
switch t.Kind {
case TaskKindInput:
f, out = a.runInputTask(t.Event, t.SeedMsgs)
f, out = a.runInputTask(t.Event)
case TaskKindSelf:
f, out = a.runInputTask(selfEvent(t.Self), nil)
f, out = a.runInputTask(selfEvent(t.Self))
}
}()
@ -664,6 +644,10 @@ func (a *Agent) executeNewTask(t *Task) {
// resumeTask 从保存的现场继续一个被抢占的任务。
//
// 关键:不重建帧、不重跑 prepare 段——否则会重复提交上下文与事件。
// resumeTask 从保存的现场继续一个被抢占的任务。
//
// 关键:不重跑 prepare 段(否则会重复提交上下文与事件),而是先把基础前缀
// 重建到「中断任务之上」,再把本任务自己的现场接回去(见 rebaseFramePrefix)。
func (a *Agent) resumeTask(t *Task, f *TaskFrame) {
a.sched.mu.Lock()
a.sched.stats.Resumed++
@ -671,6 +655,7 @@ func (a *Agent) resumeTask(t *Task, f *TaskFrame) {
a.publishEvent(events.EventScheduler, map[string]interface{}{
"action": "resume", "task": t.ID, "level": int(t.Level),
})
a.rebaseFramePrefix(f)
defer func() {
if r := recover(); r != nil {
log.Printf("[agent] resume task#%d panic recovered: %v\n%s",