feat(scheduler): M5 饥饿防护(抢占计数提升有效级 + 抢占冷却)

设计依据 docs/zh/input-scheduler-design.md §9、§11.5(G1/G2)。

- Task 增加 PreemptCount / LastPreemptAt;effectiveLevel(t) =
  min(L4, Level + min(PreemptCount, 2)):被抢占越多越“值钱”,
  逐步追上抢占它的流,但封顶 L4 因而抢不过真正的紧急输入
- 选择函数 taskBefore 改用有效级;requestPreempt 与 preemptGrantedFor
  同样以有效级比较
- 抢占冷却 preemptCooldown=2s:刚被抢占的任务期内不再被抢占,
  避免同一任务被反复打断到永不完结
- 新增 scheduler_starvation_test.go 4 项:提升与封顶、冷却期内不得再抢占、
  提升后同级不得抢占而更高可、选择函数确实用有效级
- 验收:agent 全量 + -race;全仓 build/vet 通过
This commit is contained in:
JianFeeeee
2026-09-13 00:38:30 +08:00
parent 0f78d7094d
commit 1b7c3fa53e
5 changed files with 180 additions and 25 deletions

View File

@ -89,6 +89,35 @@ type Task struct {
// SeedMsgs 是抢占式中断任务的只读前缀(D1=A):由被打断的任务在挂起时
// 附上,使中断任务看得见「进行到哪一步」,但其产出不合并回原任务。
SeedMsgs []agentAPI.Message
// PreemptCount 是本任务被抢占的次数,用于饥饿防护:
// effectiveLevel = min(L4, Level + min(PreemptCount, 2))。
PreemptCount int
// LastPreemptAt 是上次被抢占的时刻,用于抢占冷却。
LastPreemptAt time.Time
}
// preemptPromotionCap 是抢占计数能带来的最大提升档数。
const preemptPromotionCap = 2
// preemptCooldown 是“刚被抢占过”的冷却期:期内不再被抢占,
// 避免高优先级流把同一任务反复打断到永不完结。
const preemptCooldown = 2 * time.Second
// effectiveLevel 返回任务的**有效**优先级(设计文档 §9 饥饿防护)。
//
// 被抢占越多的任务越“值钱”,从而逐步追上抢占它的流;封顶 L4,
// 因此它永远不会反过来抢占真正的紧急输入。
func effectiveLevel(t *Task) Level {
p := t.PreemptCount
if p > preemptPromotionCap {
p = preemptPromotionCap
}
l := t.Level + Level(p)
if l > LevelCritical {
l = LevelCritical
}
return l
}
// SchedulerStats 是调度器的累计计数(可观测性,设计文档 §11 O2)。
@ -266,28 +295,33 @@ func (s *scheduler) enqueueInterrupt(t *Task) {
//
// 无论能否抢占,中断请求都进 pendingInterrupts——这样即使运行任务在抢占生效前
// 就正常结束,中断也不会丢(它会被 nextRef 按优先级选出)。
//
// 判据用**有效**优先级(饥饿防护),并受抢占冷却约束。
func (s *scheduler) requestPreempt(evt *agentIO.InputEvent, level Level) bool {
s.mu.Lock()
running := s.running
canPreempt := false
if running != nil && level > effectiveLevel(running) {
if running.LastPreemptAt.IsZero() || time.Since(running.LastPreemptAt) >= preemptCooldown {
canPreempt = true
s.preemptArmed = true
s.preemptLevel = level
}
}
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
return canPreempt
}
// preemptGrantedFor 报告级别为 level 的运行任务是否应在当前安全点让位。
func (s *scheduler) preemptGrantedFor(level Level) bool {
// preemptGrantedFor 报告运行任务是否应在当前安全点让位。
func (s *scheduler) preemptGrantedFor() bool {
s.mu.Lock()
defer s.mu.Unlock()
return s.preemptArmed && s.preemptLevel > level
if !s.preemptArmed || s.running == nil {
return false
}
return s.preemptLevel > effectiveLevel(s.running)
}
func (s *scheduler) clearPreempt() {
@ -308,6 +342,9 @@ func (s *scheduler) suspend(t *Task, f *TaskFrame) {
s.stats.Rejected++
}
s.suspendPool = append(s.suspendPool, &suspendedTask{Task: t, Frame: f})
// 饥饿防护:抢占计数 +1(提升有效级)并记录冷却起点。
t.PreemptCount++
t.LastPreemptAt = time.Now()
if s.running == t {
s.running = nil
}
@ -414,10 +451,11 @@ func pickTaskIndex(q []*Task) int {
return best
}
// taskBefore 报告 x 是否应先于 y 执行。
// taskBefore 报告 x 是否应先于 y 执行(按**有效**优先级)。
func taskBefore(x, y *Task) bool {
if x.Level != y.Level {
return x.Level > y.Level
lx, ly := effectiveLevel(x), effectiveLevel(y)
if lx != ly {
return lx > ly
}
if !x.EnqueuedAt.Equal(y.EnqueuedAt) {
return x.EnqueuedAt.Before(y.EnqueuedAt)