mirror of
https://gitcode.com/JianFeeeee/HomeAgent.git
synced 2026-10-04 00:03:59 +00:00
feat(scheduler): M3a+M3b 任务生命周期重构 + 四级优先级抢占
设计依据 docs/zh/input-scheduler-design.md §3–§8、§14。 M3a(行为等价的所有权重构): - processInput 拆为 prepareInputTask / runTaskSteps / finishInputTask, 帧覆盖 prepare→step…→finish;提交与回执只在 finish 段发生一次, 为安全点挂起做准备(挂起不重复提交) - process() 不再持 a.mu(挂起不能持锁),a.mu 字段随之移除 - TaskFrame 增加任务层现场(Evt/CleanInput/IsInterrupt/StartedAt/Terminal/ Level/SeedMsgs)与 taskTerminal / outcomeSuspended - 新增 task_lifecycle_test.go 5 项:正常恰好一次终态、去重 skipped、 on_input 短路、错误终态、consolidation 路由 M3b(优先级与抢占): - interceptLoop 重写:只做「收中断 → 定级 → requestPreempt → 必要时取消 LLM」,绝不触碰帧(不变量 I2);三条降级路径与 interceptCh 兜底退场, 改为统一的 pendingInterrupts - scheduler:pendingInterrupts / suspendPool / 让位信号,nextRef 在三集合上 按统一排序键取值;深度上限 4(canSuspend 在安全点拦下) - 抢占判据 incoming.level > running.level;相等与更低只入队 - 安全点只在 step 之间;执行中的 step(工具 RPC/ONNX/CAS)天然不可抢占; _consolidation_ 整任务视为临界区 - 恢复走 resumeTask:从 frame.Step 继续,不重跑 prepare - D1=A:suspend 把被打断任务的只读前缀交给抢占比它的中断任务(SeedMsgs) - 新增 scheduler_preempt_test.go 5 项:抢占-挂起-恢复(含 R1/R5)、同级更低 不抢占、深度上限、空闲中断不丢、seed 路径不污染标志位 验收:agent 全量 + -race 通过;全仓 build/vet 通过
This commit is contained in:
@ -23,6 +23,7 @@ import (
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
agentAPI "gitcode.com/JianFeeeee/HomeAgent/internal/agent/api"
|
||||
agentIO "gitcode.com/JianFeeeee/HomeAgent/internal/agent/io"
|
||||
)
|
||||
|
||||
@ -84,6 +85,10 @@ type Task struct {
|
||||
|
||||
Event *agentIO.InputEvent // Kind == TaskKindInput
|
||||
Self selfInputMsg // Kind == TaskKindSelf
|
||||
|
||||
// SeedMsgs 是抢占式中断任务的只读前缀(D1=A):由被打断的任务在挂起时
|
||||
// 附上,使中断任务看得见「进行到哪一步」,但其产出不合并回原任务。
|
||||
SeedMsgs []agentAPI.Message
|
||||
}
|
||||
|
||||
// SchedulerStats 是调度器的累计计数(可观测性,设计文档 §11 O2)。
|
||||
@ -97,9 +102,11 @@ type SchedulerStats struct {
|
||||
|
||||
// SchedulerSnapshot 是调度器的原子快照。
|
||||
type SchedulerSnapshot struct {
|
||||
Running *Task
|
||||
Queue []*Task
|
||||
Stats SchedulerStats
|
||||
Running *Task
|
||||
Queue []*Task
|
||||
PendingInterrupts []*Task
|
||||
SuspendPool []*suspendedTask
|
||||
Stats SchedulerStats
|
||||
}
|
||||
|
||||
type scheduler struct {
|
||||
@ -109,13 +116,41 @@ type scheduler struct {
|
||||
seq uint64
|
||||
stats SchedulerStats
|
||||
maxQueue int
|
||||
|
||||
// pendingInterrupts:因优先级不足(或运行任务在临界区)而未立即抢占的中断请求。
|
||||
// 与 readyQueue 分离:取出时以中断语义启动(设计文档 D3)。
|
||||
pendingInterrupts []*Task
|
||||
// suspendPool:被抢占后保存了现场、等待恢复的任务(**不是栈**,按优先级取)。
|
||||
suspendPool []*suspendedTask
|
||||
// preemptArmed/preemptLevel:运行任务的“让位信号”。
|
||||
// interruptLoop 只写这两个字段与 pendingInterrupts;帧永远只由调度器读写。
|
||||
preemptArmed bool
|
||||
preemptLevel Level
|
||||
// maxSuspendDepth:suspendPool 深度上限(设计文档 §6.3,默认 4)。
|
||||
maxSuspendDepth int
|
||||
}
|
||||
|
||||
// suspendedTask 是一个被抢占任务的现场。
|
||||
type suspendedTask struct {
|
||||
Task *Task
|
||||
Frame *TaskFrame
|
||||
}
|
||||
|
||||
// nextSelection 标识 nextRef 从哪个集合取出任务。
|
||||
type nextSelection int
|
||||
|
||||
const (
|
||||
nextNone nextSelection = iota
|
||||
nextReady
|
||||
nextPending
|
||||
nextSuspended
|
||||
)
|
||||
|
||||
func newScheduler(maxQueue int) *scheduler {
|
||||
if maxQueue <= 0 {
|
||||
maxQueue = 256
|
||||
}
|
||||
return &scheduler{maxQueue: maxQueue}
|
||||
return &scheduler{maxQueue: maxQueue, maxSuspendDepth: 4}
|
||||
}
|
||||
|
||||
// hasRoom 报告就绪队列是否还能接收任务。泵入侧据此节流:
|
||||
@ -145,17 +180,164 @@ func (s *scheduler) enqueue(t *Task) bool {
|
||||
}
|
||||
|
||||
// next 取出下一个要执行的任务;队列空返回 nil。
|
||||
//
|
||||
// 保留该签名供已有测试使用;调度器自用 nextRef(需要区分是否携带现场)。
|
||||
func (s *scheduler) next() *Task {
|
||||
t, _, _ := s.nextRef()
|
||||
return t
|
||||
}
|
||||
|
||||
// nextRef 从三个集合中按统一排序键取出下一个任务。
|
||||
//
|
||||
// 设计文档 §4.1:高有效级先;同级先到先服务。挂起任务保留其**原始**入队时刻,
|
||||
// 因此同级时天然倾向“先把旧任务做完”,抑制饥饿。
|
||||
func (s *scheduler) nextRef() (*Task, *TaskFrame, nextSelection) {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
if len(s.queue) == 0 {
|
||||
return nil
|
||||
|
||||
var bestTask *Task
|
||||
var bestFrame *TaskFrame
|
||||
bestKind := nextNone
|
||||
consider := func(t *Task, k nextSelection, fr *TaskFrame) {
|
||||
if bestTask == nil || taskBefore(t, bestTask) {
|
||||
bestTask, bestKind, bestFrame = t, k, fr
|
||||
}
|
||||
}
|
||||
i := pickTaskIndex(s.queue)
|
||||
t := s.queue[i]
|
||||
s.queue = append(s.queue[:i], s.queue[i+1:]...)
|
||||
s.running = t
|
||||
return t
|
||||
for _, t := range s.queue {
|
||||
consider(t, nextReady, nil)
|
||||
}
|
||||
for _, t := range s.pendingInterrupts {
|
||||
consider(t, nextPending, nil)
|
||||
}
|
||||
for _, st := range s.suspendPool {
|
||||
consider(st.Task, nextSuspended, st.Frame)
|
||||
}
|
||||
if bestTask == nil {
|
||||
return nil, nil, nextNone
|
||||
}
|
||||
|
||||
switch bestKind {
|
||||
case nextReady:
|
||||
s.queue = removeTask(s.queue, bestTask)
|
||||
case nextPending:
|
||||
s.pendingInterrupts = removeTask(s.pendingInterrupts, bestTask)
|
||||
case nextSuspended:
|
||||
for i, st := range s.suspendPool {
|
||||
if st.Task == bestTask {
|
||||
s.suspendPool = append(s.suspendPool[:i], s.suspendPool[i+1:]...)
|
||||
break
|
||||
}
|
||||
}
|
||||
}
|
||||
s.running = bestTask
|
||||
return bestTask, bestFrame, bestKind
|
||||
}
|
||||
|
||||
func removeTask(list []*Task, target *Task) []*Task {
|
||||
for i, t := range list {
|
||||
if t == target {
|
||||
return append(list[:i], list[i+1:]...)
|
||||
}
|
||||
}
|
||||
return list
|
||||
}
|
||||
|
||||
// enqueueInterrupt 把一个未立即抢占的中断请求放进 pendingInterrupts。
|
||||
//
|
||||
// 有界:满了丢**最老**的一条并计数(中断是提示性输入,宁可丢旧保新)。
|
||||
func (s *scheduler) enqueueInterrupt(t *Task) {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
s.seq++
|
||||
t.ID = s.seq
|
||||
if t.EnqueuedAt.IsZero() {
|
||||
t.EnqueuedAt = time.Now()
|
||||
}
|
||||
if len(s.pendingInterrupts) >= s.maxQueue {
|
||||
s.pendingInterrupts = s.pendingInterrupts[1:]
|
||||
s.stats.Rejected++
|
||||
}
|
||||
s.pendingInterrupts = append(s.pendingInterrupts, t)
|
||||
}
|
||||
|
||||
// requestPreempt 登记一次中断请求。
|
||||
//
|
||||
// 返回 true 表示“应该尝试取消运行任务正在进行的可取消步骤(LLM 流式)”。
|
||||
//
|
||||
// 无论能否抢占,中断请求都进 pendingInterrupts——这样即使运行任务在抢占生效前
|
||||
// 就正常结束,中断也不会丢(它会被 nextRef 按优先级选出)。
|
||||
func (s *scheduler) requestPreempt(evt *agentIO.InputEvent, level Level) bool {
|
||||
s.mu.Lock()
|
||||
running := s.running
|
||||
s.mu.Unlock()
|
||||
|
||||
s.enqueueInterrupt(newInterruptTask(evt, level))
|
||||
|
||||
if running == nil || level <= running.Level {
|
||||
return false
|
||||
}
|
||||
s.mu.Lock()
|
||||
s.preemptArmed = true
|
||||
s.preemptLevel = level
|
||||
s.mu.Unlock()
|
||||
return true
|
||||
}
|
||||
|
||||
// preemptGrantedFor 报告级别为 level 的运行任务是否应在当前安全点让位。
|
||||
func (s *scheduler) preemptGrantedFor(level Level) bool {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
return s.preemptArmed && s.preemptLevel > level
|
||||
}
|
||||
|
||||
func (s *scheduler) clearPreempt() {
|
||||
s.mu.Lock()
|
||||
s.preemptArmed = false
|
||||
s.preemptLevel = 0
|
||||
s.mu.Unlock()
|
||||
}
|
||||
|
||||
// suspend 保存现场。
|
||||
//
|
||||
// 深度上限(设计文档 §6.3):安全点上的 canSuspend 已提前拦下超限情况,
|
||||
// 此处仅在竞态下兜底计数——绝不丢弃帧(帧丢了会丢副作用记录)。
|
||||
func (s *scheduler) suspend(t *Task, f *TaskFrame) {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
if len(s.suspendPool) >= s.maxSuspendDepth {
|
||||
s.stats.Rejected++
|
||||
}
|
||||
s.suspendPool = append(s.suspendPool, &suspendedTask{Task: t, Frame: f})
|
||||
if s.running == t {
|
||||
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...)
|
||||
}
|
||||
}
|
||||
|
||||
s.preemptArmed = false
|
||||
s.preemptLevel = 0
|
||||
}
|
||||
|
||||
// canSuspend 报告还有下潜余量(安全点用它决定是否真的让位)。
|
||||
func (s *scheduler) canSuspend() bool {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
return len(s.suspendPool) < s.maxSuspendDepth
|
||||
}
|
||||
|
||||
// done 标记任务执行结束。
|
||||
@ -165,9 +347,59 @@ func (s *scheduler) done(t *Task) {
|
||||
if s.running == t {
|
||||
s.running = nil
|
||||
}
|
||||
// 任务正常结束:让位信号不再有意义(中断已在 pendingInterrupts 里)。
|
||||
s.preemptArmed = false
|
||||
s.preemptLevel = 0
|
||||
s.stats.Executed++
|
||||
}
|
||||
|
||||
// currentLevel 返回当前正在执行任务的级别;无 running 时为默认级。
|
||||
//
|
||||
// 用于在 prepare 段把级别写进帧(抢占比较的基准)。
|
||||
func (s *scheduler) currentLevel() Level {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
if s.running != nil {
|
||||
return s.running.Level
|
||||
}
|
||||
return DefaultLevel
|
||||
}
|
||||
|
||||
// taskLevel 是内核的优先级策略:四级的来源(设计文档 §3.2)。
|
||||
//
|
||||
// 优先走注入的查找函数(配置表);未命中则用通道名兜底:
|
||||
// cli/webui/http 为人机交互(L3),system/_consolidation_ 为后台(L1),
|
||||
// 其余一律默认级(L1)。显式才是特权:没有策略就不给抢占权。
|
||||
func (a *Agent) taskLevel(source, channel string) Level {
|
||||
if a.priorityLookup != nil {
|
||||
if l := a.priorityLookup(source, channel); l >= LevelBackground && l <= LevelCritical {
|
||||
return l
|
||||
}
|
||||
}
|
||||
switch channel {
|
||||
case "cli", "webui", "http":
|
||||
return LevelInteractive
|
||||
case channelConsolidation, "system":
|
||||
return LevelBackground
|
||||
}
|
||||
switch source {
|
||||
case "cli", "webui":
|
||||
return LevelInteractive
|
||||
case "system":
|
||||
return LevelBackground
|
||||
}
|
||||
return DefaultLevel
|
||||
}
|
||||
|
||||
// inCriticalSection 报告运行任务是否处于不可抢占区。
|
||||
//
|
||||
// M3b 只处理「整个任务不可抢占」的情形(记忆整理)。工具执行、ONNX、
|
||||
// CAS 落盘属于**单步**临界区——它们由「只在 step 之间检查让位」天然保护,
|
||||
// 不需要在这里列(M4 会把清单显式化)。
|
||||
func (a *Agent) inCriticalSection() bool {
|
||||
return a.currentOutputChannel == channelConsolidation
|
||||
}
|
||||
|
||||
// pickTaskIndex 返回下一个要执行的任务下标(设计文档 §4.1 的选择函数)。
|
||||
//
|
||||
// 排序键:优先级降序 → 入队时刻升序 → ID 升序。
|
||||
@ -197,6 +429,11 @@ func newInputTask(evt *agentIO.InputEvent) *Task {
|
||||
return &Task{Kind: TaskKindInput, Level: DefaultLevel, Event: evt, EnqueuedAt: time.Now()}
|
||||
}
|
||||
|
||||
// newInterruptTask 把一个中断请求包装成任务。
|
||||
func newInterruptTask(evt *agentIO.InputEvent, level Level) *Task {
|
||||
return &Task{Kind: TaskKindInput, Level: level, Event: evt, EnqueuedAt: time.Now()}
|
||||
}
|
||||
|
||||
func newSelfTask(msg selfInputMsg) *Task {
|
||||
return &Task{Kind: TaskKindSelf, Level: DefaultLevel, Self: msg, EnqueuedAt: time.Now()}
|
||||
}
|
||||
@ -210,6 +447,8 @@ func (a *Agent) DumpScheduler() SchedulerSnapshot {
|
||||
defer a.sched.mu.Unlock()
|
||||
snap := SchedulerSnapshot{Running: a.sched.running, Stats: a.sched.stats}
|
||||
snap.Queue = append(snap.Queue, a.sched.queue...)
|
||||
snap.PendingInterrupts = append(snap.PendingInterrupts, a.sched.pendingInterrupts...)
|
||||
snap.SuspendPool = append(snap.SuspendPool, a.sched.suspendPool...)
|
||||
return snap
|
||||
}
|
||||
|
||||
@ -226,9 +465,9 @@ func (a *Agent) schedulerLoop() {
|
||||
for {
|
||||
a.pumpInbox()
|
||||
|
||||
t := a.sched.next()
|
||||
if t == nil {
|
||||
// 就绪队列空:阻塞等新输入或退出。
|
||||
t, f, kind := a.sched.nextRef()
|
||||
if kind == nextNone {
|
||||
// 无待办:阻塞等新输入或退出。
|
||||
select {
|
||||
case evt := <-a.io.InputChan():
|
||||
a.sched.enqueue(newInputTask(evt))
|
||||
@ -239,7 +478,12 @@ func (a *Agent) schedulerLoop() {
|
||||
}
|
||||
continue
|
||||
}
|
||||
a.executeTask(t)
|
||||
|
||||
if kind == nextSuspended {
|
||||
a.resumeTask(t, f)
|
||||
continue
|
||||
}
|
||||
a.executeNewTask(t)
|
||||
}
|
||||
}
|
||||
|
||||
@ -264,11 +508,19 @@ func (a *Agent) pumpInbox() {
|
||||
}
|
||||
}
|
||||
|
||||
// executeTask 执行一个任务,并做**任务级 panic 隔离**(不变量 I6)。
|
||||
// executeTask 执行一个任务(测试与旧调用方的入口);见 executeNewTask。
|
||||
func (a *Agent) executeTask(t *Task) {
|
||||
a.executeNewTask(t)
|
||||
}
|
||||
|
||||
// executeNewTask 执行一个**新建**任务,并做任务级 panic 隔离(不变量 I6)。
|
||||
//
|
||||
// 与改造前的差异(有意):原 eventLoop 在 panic 后重启整个循环,
|
||||
// 现在一个任务的 panic 只丢弃该任务,调度器与其它任务不受影响。
|
||||
func (a *Agent) executeTask(t *Task) {
|
||||
func (a *Agent) executeNewTask(t *Task) {
|
||||
var f *TaskFrame
|
||||
var out stepOutcome = outcomeDone
|
||||
|
||||
func() {
|
||||
defer func() {
|
||||
if r := recover(); r != nil {
|
||||
@ -278,10 +530,36 @@ func (a *Agent) executeTask(t *Task) {
|
||||
}()
|
||||
switch t.Kind {
|
||||
case TaskKindInput:
|
||||
a.handleInput(t.Event)
|
||||
f, out = a.runInputTask(t.Event, t.SeedMsgs)
|
||||
case TaskKindSelf:
|
||||
a.handleSelfInput(t.Self)
|
||||
f, out = a.runInputTask(selfEvent(t.Self), nil)
|
||||
}
|
||||
}()
|
||||
|
||||
if out == outcomeSuspended && f != nil {
|
||||
a.sched.suspend(t, f)
|
||||
return
|
||||
}
|
||||
a.sched.done(t)
|
||||
}
|
||||
|
||||
// resumeTask 从保存的现场继续一个被抢占的任务。
|
||||
//
|
||||
// 关键:不重建帧、不重跑 prepare 段——否则会重复提交上下文与事件。
|
||||
func (a *Agent) resumeTask(t *Task, f *TaskFrame) {
|
||||
defer func() {
|
||||
if r := recover(); r != nil {
|
||||
log.Printf("[agent] resume task#%d panic recovered: %v\n%s",
|
||||
t.ID, r, debug.Stack())
|
||||
a.sched.done(t)
|
||||
}
|
||||
}()
|
||||
|
||||
out := a.runTaskSteps(f)
|
||||
if out == outcomeSuspended {
|
||||
a.sched.suspend(t, f)
|
||||
return
|
||||
}
|
||||
a.finishInputTask(f, out)
|
||||
a.sched.done(t)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user