From c321388a2102cb9ad499a52deec3527e5cfd1f4c Mon Sep 17 00:00:00 2001 From: JianFeeeee Date: Mon, 14 Sep 2026 10:39:25 +0800 Subject: [PATCH] =?UTF-8?q?fix(scheduler):=20=E5=AE=89=E5=85=A8=E7=82=B9?= =?UTF-8?q?=E9=87=8D=E6=96=B0=E6=B1=82=E5=80=BC=E4=B8=AD=E6=96=AD=E9=98=9F?= =?UTF-8?q?=E5=88=97=20+=20=E6=8A=A2=E5=8D=A0/=E8=83=8C=E5=8E=8B=E8=AE=A1?= =?UTF-8?q?=E6=95=B0=E4=BF=AE=E6=AD=A3=20+=20=E5=81=9C=E6=9C=BA=E8=A1=A5?= =?UTF-8?q?=E7=BB=88=E6=80=81?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 对照 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 满计数)。 --- internal/agent/core/agent.go | 21 +++ internal/agent/core/eventloop.go | 9 +- internal/agent/core/scheduler.go | 173 ++++++++++++++++-- internal/agent/core/scheduler_rearm_test.go | 186 ++++++++++++++++++++ internal/agent/core/task.go | 4 + internal/sdk/status.go | 4 + 6 files changed, 383 insertions(+), 14 deletions(-) create mode 100644 internal/agent/core/scheduler_rearm_test.go diff --git a/internal/agent/core/agent.go b/internal/agent/core/agent.go index 59c5b23..39c95f2 100644 --- a/internal/agent/core/agent.go +++ b/internal/agent/core/agent.go @@ -404,9 +404,30 @@ func (a *Agent) Start() { func (a *Agent) Stop() { // 父退出**必须**销毁全部驻留子(设计 §10 硬约束:子不得比父活得久、不留孤儿)。 a.StopResidents() + // 停机前给待办任务补终态。运行中的任务会经 cancel → LLM 失败 → emitResponse + // 自然拿到终态,但**从未运行**(排队/待处理)与**已挂起**的任务不会有任何人 + // 回它们;带 ResponseCh 的同步注入方(cli / clawhubadapter 均无超时)会永久挂起 + // (设计 §7 I5、§11.3 X2/X4)。必须在 cancel 之前做:cancel 会让调度器直接 return。 + a.drainPendingInterrupts("agent_stopped") a.cancel() } +// drainPendingInterrupts 给排队/待处理/已挂起任务中带同步回执通道的调用方补一条 +// skipped 终态(复用 emitSkippedReply:非阻塞写,不对外发 agent_output 事件)。 +func (a *Agent) drainPendingInterrupts(reason string) { + if a.sched == nil { + return + } + pending := a.sched.pendingEvents() + if len(pending) == 0 { + return + } + for _, evt := range pending { + a.emitSkippedReply(evt, reason) + } + log.Printf("[agent] %s: 停机,%d 条待办任务已补 skipped 终态", a.id, len(pending)) +} + // graphMemoryOf 决定本 agent 的图记忆共同面实现。 // // - 轻量内核(给了 LightMemory):用 LightMemory,**整理面保持 nil**; diff --git a/internal/agent/core/eventloop.go b/internal/agent/core/eventloop.go index 7579a76..48cbbda 100644 --- a/internal/agent/core/eventloop.go +++ b/internal/agent/core/eventloop.go @@ -331,13 +331,20 @@ func (a *Agent) emitResponse(evt *agentIO.InputEvent, response string) { payload["usage"] = stageCtx.TokenUsage } if evt.ResponseCh != nil { - evt.ResponseCh <- &agentIO.OutputEvent{ + // 非阻塞写:ResponseCh 由同步调用方以 cap=1 创建。按不变量 I5(每任务恰一次 + // 终态)这里永远写得进去;但一旦哪天写出第二次,阻塞会卡死**调度器 goroutine** + // (整个 agent 停摆),而丢弃只是丢一条回执——与 emitSkippedReply 对称。 + select { + case evt.ResponseCh <- &agentIO.OutputEvent{ RequestID: evt.RequestID, Target: evt.Source, Type: "text", Payload: payload, Done: true, OutputChannel: ch, + }: + default: + log.Printf("[agent] ResponseCh 已满,终态回执被丢弃(request=%s,可能违反不变量 I5)", evt.RequestID) } } diff --git a/internal/agent/core/scheduler.go b/internal/agent/core/scheduler.go index 8f1988e..330aaca 100644 --- a/internal/agent/core/scheduler.go +++ b/internal/agent/core/scheduler.go @@ -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。 diff --git a/internal/agent/core/scheduler_rearm_test.go b/internal/agent/core/scheduler_rearm_test.go new file mode 100644 index 0000000..d207e40 --- /dev/null +++ b/internal/agent/core/scheduler_rearm_test.go @@ -0,0 +1,186 @@ +package core + +// 回归测试:安全点「重新求值」、抢占计数语义、停机补终态、背压计数。 +// +// 对照设计稿原文修正的四条: +// +// 1. §4.3/§5.2 —— 临界区期间到达的抢占请求「不丢失:按级别进入中断队列, +// 在**临界区结束后的第一个安全点重新求值**」。实现里此前没有这一步: +// 唯一的武装点是 registerInterrupt,凡被拦成「入队」的中断只能等当前任务 +// **自然结束**。可复现症状:WebUI 终止按钮连按两次,第二次落在 2s 抢占冷却 +// 窗内 → 入队 → 再也不会被求值,「终止」看起来没反应。 +// 2. SchedulerStats.PreemptsByLevel 的语义是「判定可抢占**并进入 immediate** 的 +// 次数」;此前在 setImmediateLocked 之前就计数,于是同一安全点前到达的两条同级 +// 中断里、被降级入队的那条也被计入(immediate 是单槽,降级是设计要求的路径)。 +// 3. 状态面 Preempted 此前直接拿 Stats.Suspended 顶替,与 preempts_by_level 自相矛盾。 +// 4. §4.4 / §11.4 Q4 —— 就绪队列满必须「阻塞发送方 + **计数并打日志**」。 + +import ( + "testing" + "time" + + agentAPI "gitcode.com/JianFeeeee/HomeAgent/internal/agent/api" + agentIO "gitcode.com/JianFeeeee/HomeAgent/internal/agent/io" +) + +// 被冷却拦成入队的中断,在冷却期满后的第一个安全点必须被重新武装。 +// +// 这是「终止按钮连按两次」的最小复现:第一次抢占成功(受害者进入 2s 冷却), +// 第二次在冷却窗内只能入队——修复前它就永远等不到执行了。 +func TestRearm_CooldownExpiryPromotesQueuedInterrupt(t *testing.T) { + a := newPreemptAgent(t, newPreemptProvider()) + + victim := &Task{ID: 1, Class: TaskInterrupt, Level: LevelBackground, EnqueuedAt: time.Now()} + a.sched.immediate = victim + a.sched.nextRef() // running = victim + // 模拟「刚被抢占过」:冷却起点就在此刻,且抢占提升已生效(有效级 L2)。 + victim.PreemptCount = 1 + victim.LastPreemptAt = time.Now() + + evt, _ := textEvent("cli", "第二次终止") + evt.Payload["interrupt"] = true + if a.sched.requestPreempt(evt, LevelCritical) { + t.Fatal("抢占冷却期内不得抢占(应入队)") + } + if got := a.sched.stats.PreemptsByLevel[LevelCritical]; got != 0 { + t.Fatalf("被冷却拦成入队的中断不得计入抢占数,实际 %d", got) + } + if n := len(a.sched.interruptQueues[LevelCritical]); n != 1 { + t.Fatalf("应恰好入队一条,实际 %d(>1 说明入队路径重复)", n) + } + + // 冷却期满 → 安全点的「重新求值」必须把它武装起来(修复前缺失的正是这一步)。 + victim.LastPreemptAt = time.Now().Add(-3 * time.Second) + a.sched.rearmPending() + if !a.sched.preemptGrantedFor() { + t.Fatal("冷却期结束后应重新武装让位信号(设计 §4.3「第一个安全点重新求值」)") + } + snap := a.DumpScheduler() + if snap.Immediate == nil { + t.Fatal("重新求值后应把该中断提升进 immediate 槽") + } + if got := snap.Stats.PreemptsByLevel[LevelCritical]; got != 1 { + t.Fatalf("真正占住 immediate 才能计一次抢占,实际 %d", got) + } + // 重新求值不该把任务复制一份:队列必须空、immediate 恰好一条。 + if n := len(snap.InterruptQueues[LevelCritical]); n != 0 { + t.Fatalf("提升后 L4 队列应空,实际 %d", n) + } +} + +// 同一安全点前到达的两条同级中断:immediate 是单槽,第二条只能降级入队; +// 它**没有**进入 immediate,因此不得计入 PreemptsByLevel,也不得被入队两次。 +func TestRearm_SameLevelSecondPreempterIsQueuedNotCounted(t *testing.T) { + a := newPreemptAgent(t, newPreemptProvider()) + + victim := &Task{ID: 1, Class: TaskQueued, EnqueuedAt: time.Now()} + a.sched.immediate = victim + a.sched.nextRef() // running = 排队任务(有效级 0,任何中断都能抢) + + e1, _ := textEvent("cli", "irq-1") + if !a.sched.requestPreempt(e1, LevelInteractive) { + t.Fatal("第一条 L3 应抢占排队任务") + } + e2, _ := textEvent("cli", "irq-2") + a.sched.requestPreempt(e2, LevelInteractive) // 同级 → 降级入队 + + snap := a.DumpScheduler() + if got := snap.Stats.PreemptsByLevel[LevelInteractive]; got != 1 { + t.Fatalf("被降级的同级第二条不得计入抢占数(期望 1,实际 %d)", got) + } + if n := len(snap.InterruptQueues[LevelInteractive]); n != 1 { + t.Fatalf("被降级的那条应在 L3 队列里**恰好**出现一次,实际 %d", n) + } + if snap.Immediate == nil { + t.Fatal("第一条应留在 immediate 槽,两条都不能丢") + } +} + +// 状态面 Preempted 必须是「各级抢占数之和」,不能拿 Suspended 顶替。 +func TestStatus_PreemptedEqualsSumOfLevels(t *testing.T) { + a := New(AgentConfig{ID: "rt-sum", ProviderManager: agentAPI.NewProviderManager(), IO: agentIO.NewIOManager()}) + if a.sched == nil { + t.Fatal("agent 应带调度器") + } + a.sched.mu.Lock() + a.sched.stats.PreemptsByLevel[LevelBackground] = 2 + a.sched.stats.PreemptsByLevel[LevelInteractive] = 3 + a.sched.stats.Suspended = 99 // 故意与抢占数不等 + a.sched.mu.Unlock() + + got := a.schedulerStatus() + if got.Preempted != 5 { + t.Fatalf("Preempted 应为各级抢占数之和 5,实际 %d(拿 Suspended 顶替会得 99)", got.Preempted) + } + if got.Suspended != 99 { + t.Fatalf("Suspended 应原样透传,实际 %d", got.Suspended) + } +} + +// 停机必须给「从未运行」与「已挂起」的同步任务补终态, +// 否则 cli / clawhubadapter 这类无超时的同步注入方会永久挂起(设计 §7 I5、§11.3 X2/X4)。 +func TestStop_DrainsPendingSyncTasks(t *testing.T) { + a := newPreemptAgent(t, newPreemptProvider()) + + queuedEvt, queuedCh := textEvent("cli", "排队中,永远不会被调度") + if !a.sched.enqueue(newInputTask(queuedEvt)) { + t.Fatal("入队失败") + } + suspEvt, suspCh := textEvent("cli", "已挂起,停机时不会恢复") + a.sched.suspend( + &Task{ID: 2, Class: TaskInterrupt, Level: LevelInteractive, EnqueuedAt: time.Now(), Event: suspEvt}, + a.newTaskFrame("挂起", a.stageCtxFromInput("挂起", "", "")), + ) + + a.Stop() + + for name, ch := range map[string]chan *agentIO.OutputEvent{"排队": queuedCh, "挂起": suspCh} { + select { + case r := <-ch: + if r == nil || r.Payload["skipped"] != true { + t.Fatalf("%s 任务停机时应补 skipped 终态,实际 %+v", name, r) + } + case <-time.After(2 * time.Second): + t.Fatalf("%s 任务停机未补终态(同步调用方会永久挂起)", name) + } + } +} + +// 背压计数:持续满只报一次「翻转」不发生;计数本身每次都要累加。 +func TestBackpressure_CounterAndTransition(t *testing.T) { + s := newScheduler(1) + if !s.noteBackpressure() { + t.Fatal("首次背压应报告「翻转」") + } + if s.noteBackpressure() { + t.Fatal("持续背压不得重复报告翻转(否则日志会被刷爆)") + } + if s.stats.Backpressure != 2 { + t.Fatalf("背压计数应为 2,实际 %d", s.stats.Backpressure) + } + s.clearBackpressure() + if !s.noteBackpressure() { + t.Fatal("队列恢复后再满应再次报告翻转") + } +} + +// 集成:就绪队列满时 pumpInbox 必须计一次背压(Rejected 保持 0——背压不是丢弃)。 +func TestBackpressure_PumpInboxCountsWhenFull(t *testing.T) { + a := newPreemptAgent(t, newPreemptProvider()) + a.sched.maxQueue = 1 + + a.io.InjectInput("cli", "text", map[string]interface{}{"content": "第一条"}) + a.io.InjectInput("cli", "text", map[string]interface{}{"content": "第二条"}) + a.pumpInbox() + + snap := a.DumpScheduler() + if len(snap.Queue) != 1 { + t.Fatalf("maxQueue=1 时队列应恰好 1 条,实际 %d", len(snap.Queue)) + } + if snap.Stats.Backpressure == 0 { + t.Fatal("队列满必须计一次背压(设计 §4.4/Q4:阻塞发送方 + 计数)") + } + if snap.Stats.Rejected != 0 { + t.Fatalf("背压不是丢弃,Rejected 必须保持 0,实际 %d", snap.Stats.Rejected) + } +} diff --git a/internal/agent/core/task.go b/internal/agent/core/task.go index 2d7139a..d46595f 100644 --- a/internal/agent/core/task.go +++ b/internal/agent/core/task.go @@ -178,6 +178,10 @@ func (a *Agent) runTaskSteps(f *TaskFrame) stepOutcome { for i := 0; i < maxSteps; i++ { // 安全点:只在 step 之间检查让位。临界区(StepToolExec)不在此列, // 因为让位信号由 interruptLoop 置位、而本循环是唯一读帧者。 + // + // 先「重新求值」再判让位:临界区(或抢占冷却期)内被拦成入队的中断, + // 必须在这里重新武装——否则它只能等当前任务自然结束(设计 §4.3/§5.2)。 + a.sched.rearmPending() if !isCriticalChannel(f.OutputChannel) && a.sched.preemptGrantedFor() && a.sched.canSuspend() { return outcomeSuspended } diff --git a/internal/sdk/status.go b/internal/sdk/status.go index 3f5044a..8737ab2 100644 --- a/internal/sdk/status.go +++ b/internal/sdk/status.go @@ -89,7 +89,11 @@ type SchedulerStatus struct { Rejected uint64 `json:"rejected"` Suspended uint64 `json:"suspended"` Resumed uint64 `json:"resumed"` + // Preempted = Σ PreemptsByLevel[1..4],即「真正抢占成功」的次数。 + // 它与 Suspended 不等价(受害者可能先自行结束),因此不是 Suspended 的别名。 Preempted uint64 `json:"preempted"` + // Backpressure 是就绪队列满、输入被挡回 channel 的次数(暂时不收,不是丢弃)。 + Backpressure uint64 `json:"backpressure"` } // SchedulerFrame 是中断栈里的一帧(供图形化展示"压了几层现场")。