fix(scheduler): 安全点重新求值中断队列 + 抢占/背压计数修正 + 停机补终态

对照 docs/zh/input-scheduler-design.md 原文修四处(前两处是真缺陷,后两处是
观测面与设计承诺不一致),均配回归用例:

1. §4.3/§5.2「临界区结束后的第一个安全点重新求值」此前**没有实现**:
   全仓唯一的武装点是 registerInterrupt,凡被拦成「入队」的中断只能等当前任务
   自然结束。可达症状:WebUI 终止按钮连按两次,第二次落在 2s 抢占冷却窗内 →
   入队 → 再也不会被求值。修:runTaskSteps 的安全点先 rearmPending()——
   判据与 registerInterrupt 完全同一套(canPreempt + 冷却 + 临界区闸门)。

2. PreemptsByLevel 的语义是「进入 immediate 槽的次数」,但计数发生在
   setImmediateLocked 之前:immediate 是单槽,同一安全点前到达的两条同级中断里
   被降级的那条也被计成抢占。修:setImmediateLocked 只在真占住槽时返回 true,
   计数随之为真;同时把「降级入队」的责任收归调用方,消除同一任务被入队两次的
   隐患(实测该隐患会让中断任务执行两次、Executed 虚高)。

3. 状态面 Preempted 此前拿 Stats.Suspended 顶替,与 preempts_by_level 自相矛盾。
   修:Preempted = Σ PreemptsByLevel[1..4]。

4. §4.4/Q4「满时阻塞发送方 + 计数并打日志」只做了阻塞:pumpInbox 满时直接返回,
   一个字都不计。修:新增 Stats.Backpressure(+DTO 字段) 与只报一次的状态翻转日志;
   同时显式处理 enqueue 返回值(静默丢弃会让同步调用方永久挂起)。

另:Stop() 停机前排空待办——给从未运行与已挂起的、带 ResponseCh 的任务补
skipped 终态,否则 cli/clawhubadapter 这类无超时同步注入方永久挂起(§7 I5、§11.3 X4)。
emitResponse 的 ResponseCh 写入改为非阻塞 + 告警,避免一行写错就卡死调度器 goroutine。

验证:go build/vet 干净;go test -count=1 ./internal/agent/... ./internal/plugin/...
./internal/sdk/... ./cmd/... 全绿;go test -race ./internal/agent/core/ ./internal/sdk/ 干净。
新增 scheduler_rearm_test.go 六个用例(冷却期满重新求值/同级降级不计数/Preempted 求和/
停机补终态/背压计数与翻转/pumpInbox 满计数)。
This commit is contained in:
JianFeeeee
2026-09-14 10:39:25 +08:00
parent 1ce3a917a5
commit c321388a21
6 changed files with 383 additions and 14 deletions

View File

@ -229,6 +229,10 @@ type SchedulerStats struct {
Executed uint64
// Rejected 是因队列满(或深度超限)而未被接纳的次数。
Rejected uint64
// Backpressure 是就绪队列满、输入被挡回 channel 的次数
// (设计 §4.4 / §11.4 Q4:满时阻塞发送方,**必须计数并打日志**)。
// 与 Rejected 的区别:Rejected 是「丢了」,Backpressure 是「暂时不收、发送方在等」。
Backpressure uint64
// Suspended / Resumed 是挂起与恢复的次数。
// 不变量:系统排空后 Suspended == Resumed(挂起必然被恢复),
// 因此两者各自只在**一处**计数(suspend / resumeTask)。
@ -286,7 +290,13 @@ func (a *Agent) schedulerStatus() sdk.SchedulerStatus {
Rejected: snap.Stats.Rejected,
Suspended: snap.Stats.Suspended,
Resumed: snap.Stats.Resumed,
Preempted: snap.Stats.Suspended,
Backpressure: snap.Stats.Backpressure,
}
// Preempted 是「各级抢占成功次数之和」,**不是** Suspended:受害者可能在
// 让位信号生效前就自行结束,此时有抢占而没有挂起(见 PreemptsByLevel 注释)。
// 此前这里直接拿 Suspended 顶替,导致 DTO 里 preempted 与 preempts_by_level 自相矛盾。
for lv := LevelBackground; lv <= LevelCritical; lv++ {
out.Preempted += snap.Stats.PreemptsByLevel[lv]
}
if snap.Running != nil {
out.Running = &sdk.SchedulerTask{
@ -342,6 +352,9 @@ type scheduler struct {
// critical 报告运行任务是否在不可抢占临界区(如记忆整理)。
// 由于 interceptLoop 要读它,必须是原子的:帧仍只由调度器读写。
critical atomic.Bool
// backpressured 记录「就绪队列满」这一状态的翻转,用于只打一次日志。
// 满着的时候 pumpInbox 每轮都会走到,逐轮打日志会把日志刷爆。
backpressured bool
// wake 用于把空闲的调度器叫醒:pendingInterrupts 不是 channel,
// 没有这个信号时“空闲时到达的中断”会一直等下一次输入(设计 §5.1 ③)。
wake chan struct{}
@ -403,6 +416,28 @@ func (s *scheduler) hasRoom() bool {
return len(s.queue) < s.maxQueue
}
// noteBackpressure 记一次背压,并报告这是否是「从有空间 → 满」的翻转。
//
// 为什么需要翻转信息:满的时候每轮泵入都会调用本函数,逐轮打日志会刷爆;
// 而设计 §4.4 要求「必须计数并打日志」——两者靠这个布尔量同时满足。
func (s *scheduler) noteBackpressure() bool {
s.mu.Lock()
defer s.mu.Unlock()
s.stats.Backpressure++
if s.backpressured {
return false
}
s.backpressured = true
return true
}
// clearBackpressure 在就绪队列重新可收(泵空)时复位翻转标记。
func (s *scheduler) clearBackpressure() {
s.mu.Lock()
defer s.mu.Unlock()
s.backpressured = false
}
// allocateIDLocked 分配任务 ID 与入队时刻(调用方持锁)。
func (s *scheduler) allocateIDLocked(t *Task) {
s.seq++
@ -512,10 +547,22 @@ func (s *scheduler) interruptCountLocked() int {
// setImmediateLocked 登记一个应“立即运行”的抢占者。
//
// 槽只有一格:若已有抢占者且新的级别更高,旧的降级入队;否则新的入队。
func (s *scheduler) setImmediateLocked(t *Task) {
// 返回 true 表示 t **确实占住了 immediate 槽**;false 表示它被降级进了自己的
// 级别队列(immediate 是单槽,这是设计要求的降级分支,见设计 §2「至多一个」)。
//
// 调用方必须用返回值决定是否计入 PreemptsByLevel:那条计数器的语义是
// 「进入 immediate 的次数」,被降级的中断从未进过 immediate。
// setImmediateLocked 尝试把 t 放进 immediate 槽。
//
// 返回 true:t 已占住 immediate(若原有抢占者被顶掉,它**已被**降级入队)。
// 返回 false:t 没有进 immediate,且本函数**未动 t** —— 调用方负责按级别入队。
//
// 把「降级入队」的责任留给调用方,是为了让「到底入队了几次」只有一个出口:
// 早先由本函数在返回 false 前自行入队,调用方又照着 false 再入一次,
// 同一任务就会在队列里出现两份(实测:中断任务被执行两次、Executed 虚高)。
func (s *scheduler) setImmediateLocked(t *Task) bool {
if s.immediate != nil && effectiveLevel(t) <= effectiveLevel(s.immediate) {
s.enqueueInterruptLocked(t)
return
return false
}
if s.immediate != nil {
s.enqueueInterruptLocked(s.immediate)
@ -523,6 +570,7 @@ func (s *scheduler) setImmediateLocked(t *Task) {
s.allocateIDLocked(t)
s.stats.Enqueued++
s.immediate = t
return true
}
func removeTask(list []*Task, target *Task) []*Task {
@ -593,11 +641,16 @@ func (s *scheduler) registerInterrupt(t *Task) bool {
arm := false
if !critical && canPreempt(t, running) {
if running.LastPreemptAt.IsZero() || time.Since(running.LastPreemptAt) >= preemptCooldown {
arm = true
s.preemptArmed = true
s.preemptLevel = t.Level
s.stats.bumpInterruptLevel(&s.stats.PreemptsByLevel, t.Level)
s.setImmediateLocked(t)
// 只有**真的占住 immediate 槽**才算一次抢占,才计入 PreemptsByLevel:
// immediate 是单槽,若它被另一个更高级的抢占者占着,t 会走上而下的
// 「否则入队」分支——那种情况 t 从未进入 immediate(否则同一安全点前
// 到达两条同级中断时该计数会高估)。
if s.setImmediateLocked(t) {
arm = true
s.preemptArmed = true
s.preemptLevel = t.Level
s.stats.bumpInterruptLevel(&s.stats.PreemptsByLevel, t.Level)
}
}
}
if !arm {
@ -628,6 +681,54 @@ func (s *scheduler) clearPreempt() {
s.mu.Unlock()
}
// rearmPending 在**安全点重新求值**中断队列(设计 §4.3 / §5.2)。
//
// 为什么必须有这一步:中断只在 registerInterrupt 里被武装一次,而那一刻运行任务
// 可能正在临界区(S_TOOL_EXEC / ONNX / CAS)或处于抢占冷却期,于是请求只能入队。
// 若安全点不再回头看队列,它就永远等不到执行——只能等当前任务**自然结束**,
// 这违背设计承诺的「临界区期间到达的抢占请求……在临界区结束后的第一个安全点
// 重新求值」。可复现症状:WebUI 终止按钮连按两次,第二次(落在 2s 冷却窗内)
// 入队后再也不会被求值,「终止」看起来没反应。
//
// 判据与 registerInterrupt **完全同一套**(canPreempt + 冷却 + 临界区闸门),
// 因此不会凭空制造设计之外的抢占。
func (s *scheduler) rearmPending() {
s.mu.Lock()
defer s.mu.Unlock()
// 已有让位信号、或 immediate 槽已被占用:下一个安全点的选择已经在路上,
// 不必(也不该)重复武装。
if s.preemptArmed || s.immediate != nil || s.running == nil || s.critical.Load() {
return
}
// 冷却期内不武装:与 registerInterrupt 同一判据(抗饥饿)。
if !s.running.LastPreemptAt.IsZero() && time.Since(s.running.LastPreemptAt) < preemptCooldown {
return
}
// 中断队列本就按级别组织:从最高级往下找第一条能抢占的队头。
// (队列里的任务有效级恒等于基础级,故「第一条能抢」= 最高级可抢占者。)
for lv := LevelCritical; lv >= LevelBackground; lv-- {
q := s.interruptQueues[lv]
if len(q) == 0 {
continue
}
t := q[0]
if !canPreempt(t, s.running) {
continue
}
s.popInterruptLocked(lv)
if s.setImmediateLocked(t) {
s.preemptArmed = true
s.preemptLevel = t.Level
s.stats.bumpInterruptLevel(&s.stats.PreemptsByLevel, t.Level)
} else {
// immediate 槽没拿到(理论上进不来,顶部已判 immediate == nil):放回队列,
// 否则任务会凭空消失。
s.enqueueInterruptLocked(t)
}
return
}
}
// suspend 保存现场。
//
// 深度上界是**结构推论**(= 中断级数),不是配置项:安全点上的 canSuspend 已提前
@ -660,6 +761,37 @@ func (s *scheduler) canSuspend() bool {
return len(s.suspendStack) < s.maxInterruptFrames
}
// pendingEvents 收集**尚未执行**(排队队列 / 四条中断队列 / immediate)与
// **已挂起**(中断栈)任务所携带的、且带同步回执通道的输入事件。
//
// 用途只有一个:停机收尾。这些任务不会再被调度,若不给它们补终态,
// 无超时的同步注入方(cli / clawhubadapter)会永久挂起(设计 §7 I5、§11.3 X2/X4)。
func (s *scheduler) pendingEvents() []*agentIO.InputEvent {
s.mu.Lock()
defer s.mu.Unlock()
var out []*agentIO.InputEvent
add := func(t *Task) {
if t != nil && t.Event != nil && t.Event.ResponseCh != nil {
out = append(out, t.Event)
}
}
for _, t := range s.queue {
add(t)
}
for lv := LevelBackground; lv <= LevelCritical; lv++ {
for _, t := range s.interruptQueues[lv] {
add(t)
}
}
add(s.immediate)
for _, f := range s.suspendStack {
if f != nil {
add(f.Task)
}
}
return out
}
// done 标记任务执行结束。
func (s *scheduler) done(t *Task) {
s.mu.Lock()
@ -818,9 +950,11 @@ func (a *Agent) schedulerLoop() {
// 无待办:阻塞等新输入、新中断(wake)或退出。
select {
case evt := <-a.io.InputChan():
a.sched.enqueue(newInputTask(evt))
if !a.sched.enqueue(newInputTask(evt)) {
a.emitSkippedReply(evt, "queue_full")
}
case msg := <-a.selfInputCh:
a.sched.enqueue(newSelfTask(msg))
_ = a.sched.enqueue(newSelfTask(msg))
case <-a.sched.wake:
// 中断已入 pendingInterrupts,回到循环顶部重新挑选。
case <-a.ctx.Done():
@ -847,15 +981,28 @@ func (a *Agent) pumpInbox() {
for a.sched.hasRoom() {
select {
case evt := <-a.io.InputChan():
a.sched.enqueue(newInputTask(evt))
// 返回值必须处理:静默丢弃会让同步调用方永久挂起(回执路径 E)。
if !a.sched.enqueue(newInputTask(evt)) {
a.sched.noteBackpressure()
a.emitSkippedReply(evt, "queue_full")
}
case msg := <-a.selfInputCh:
a.sched.enqueue(newSelfTask(msg))
// 自循环输入没有同步调用方,满时记一次背压即可。
if !a.sched.enqueue(newSelfTask(msg)) {
a.sched.noteBackpressure()
}
case <-a.ctx.Done():
return
default:
a.sched.clearBackpressure()
return
}
}
// 队列满:输入留在 channel 里,发送方阻塞(设计 §4.4「阻塞发送方」)。
// 必须计数并打日志——否则运维看到 Rejected=0 会以为没背压,而输入正卡在 channel。
if a.sched.noteBackpressure() {
log.Printf("[agent] ready queue full (%d), input channel backpressured", a.sched.maxQueue)
}
}
// executeTask 执行一个任务(测试与旧调用方的入口);见 executeNewTask。