package core // 输入调度器(M2:骨架)。 // // 设计依据 docs/zh/input-scheduler-design.md。 // // M2 只建立结构,不引入抢占: // - 显式的 readyQueue 与 Task 抽象(取代 eventLoop 里隐式的 channel 排队); // - 统一的排序键 (-Level, EnqueuedAt, ID)(设计文档 §4.1); // - 原子快照 DumpScheduler() 与计数(可观测性); // - **每任务 panic 隔离**:panic 只使该任务失败,调度器本身存活(不变量 I6)。 // // M2 全部任务都是 LevelBackground(默认级),因此排序结果等价于 FIFO—— // 与改造前的 channel 语义逐条一致。抢占、中断栈、pendingInterrupts、 // 任务级回执在 M3–M6 加入。 // // 并发模型(不变量 I2):readyQueue/running/stats 只由 schedulerLoop 写, // 外部只读——读取一律经 DumpScheduler() 加锁取快照。 import ( "log" "runtime/debug" "sync" "sync/atomic" "time" agentIO "gitcode.com/JianFeeeee/HomeAgent/internal/agent/io" "gitcode.com/JianFeeeee/HomeAgent/internal/events" sdk "gitcode.com/JianFeeeee/HomeAgent/internal/sdk" ) // Level 是任务优先级,由内核预定义四级(设计文档 §3.1)。 // // 取值域刻意只有四档:不引入任意整数,避免"9 级比 4 级大但没人知道怎么排"。 type Level int const ( // LevelBackground 后台维护:心跳蒸馏/归档/合并/复审、子任务、consolidation。 LevelBackground Level = 1 // LevelMessage 异步消息:QQ/微信等入站消息、插件通知。 LevelMessage Level = 2 // LevelInteractive 人机交互:用户在 CLI/WebUI 的直接对话。 LevelInteractive Level = 3 // LevelCritical 紧急打断:显式打断、系统告警、安全类中断。 LevelCritical Level = 4 ) // DefaultLevel 是未显式声明时的优先级。 // // 取最低级是刻意的:**显式才是特权**,新插件不会默认拿到抢占权。 const DefaultLevel = LevelBackground func (l Level) String() string { switch l { case LevelBackground: return "L1-background" case LevelMessage: return "L2-message" case LevelInteractive: return "L3-interactive" case LevelCritical: return "L4-critical" default: return "L?-unknown" } } // ParseLevel 已删除。 // // 为何不保留:优先级是**内核内部属性**,不是配置项——内核预定义四级 // (L1 后台 / L2 消息 / L3 交互 / L4 紧急),由内核按内部规则为任务与中断定级。 // 曾一度做成 `core.agent.priority.` 这种“策略表 + 字符串解析”, // 那等于把内核的内部属性外化成运维配置,与设计意图相反。 // TaskKind 区分任务来源。 type TaskKind int const ( // TaskKindInput 来自 io.InputChan(外部/插件注入的输入)。 TaskKindInput TaskKind = iota // TaskKindSelf 来自 selfInputCh(内核自循环:记忆整理、子任务通知)。 TaskKindSelf ) func (k TaskKind) String() string { switch k { case TaskKindInput: return "input" case TaskKindSelf: return "self" default: return "unknown" } } // Task 是调度器的最小单位。 // // M2 只承载"一份待处理的输入";M3 起把 TaskFrame(现场)挂上来, // 使其成为可挂起/可恢复的执行单元。 type Task struct { ID uint64 Level Level Kind TaskKind EnqueuedAt time.Time Event *agentIO.InputEvent // Kind == TaskKindInput Self selfInputMsg // Kind == TaskKindSelf // 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)。 type SchedulerStats struct { Enqueued uint64 Executed uint64 // Rejected 是因队列满(或深度超限)而未被接纳的次数。 Rejected uint64 // Suspended / Resumed 是挂起与恢复的次数。 Suspended uint64 Resumed uint64 } // SchedulerSnapshot 是调度器的原子快照。 type SchedulerSnapshot struct { Running *Task Queue []*Task PendingInterrupts []*Task // SuspendStack:中断栈(含嵌套抢占的多个现场),**栈顶**优先恢复。 SuspendStack []*suspendedTask Stats SchedulerStats MaxSuspendDepth int } // schedulerStatus 把快照转成对外的状态 DTO(不暴露帧内容)。 func (a *Agent) schedulerStatus() sdk.SchedulerStatus { if a.sched == nil { return sdk.SchedulerStatus{} } snap := a.DumpScheduler() out := sdk.SchedulerStatus{ ReadyQueueDepth: len(snap.Queue), PendingInterrupts: len(snap.PendingInterrupts), SuspendStack: len(snap.SuspendStack), MaxSuspendDepth: snap.MaxSuspendDepth, Enqueued: snap.Stats.Enqueued, Executed: snap.Stats.Executed, Rejected: snap.Stats.Rejected, Suspended: snap.Stats.Suspended, Resumed: snap.Stats.Resumed, Preempted: snap.Stats.Suspended, } if snap.Running != nil { out.Running = &sdk.SchedulerTask{ ID: snap.Running.ID, Level: int(snap.Running.Level), Kind: snap.Running.Kind.String(), } } return out } type scheduler struct { mu sync.Mutex queue []*Task running *Task seq uint64 stats SchedulerStats maxQueue int // pendingInterrupts:因优先级不足(或运行任务在临界区)而未立即抢占的中断请求。 // 与 readyQueue 分离:取出时以中断语义启动(设计文档 D3)。 pendingInterrupts []*Task // suspendStack:**中断栈**。被抢占后保存现场的任务压栈(LIFO), // 用于“中断被中断”的嵌套场景:只有**栈顶**参与恢复选择,栈内不做优先级重排。 suspendStack []*suspendedTask // preemptArmed/preemptLevel:运行任务的“让位信号”。 // interruptLoop 只写这两个字段与 pendingInterrupts;帧永远只由调度器读写。 preemptArmed bool preemptLevel Level // critical 报告运行任务是否在不可抢占临界区(如记忆整理)。 // 由于 interceptLoop 要读它,必须是原子的:帧仍只由调度器读写。 critical atomic.Bool // wake 用于把空闲的调度器叫醒:pendingInterrupts 不是 channel, // 没有这个信号时“空闲时到达的中断”会一直等下一次输入(设计 §5.1 ③)。 wake chan struct{} // maxSuspendDepth:中断栈深度上限(设计文档 §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, maxSuspendDepth: 4, wake: make(chan struct{}, 1)} } // signalWake 非阻塞地唤醒调度器。 func (s *scheduler) signalWake() { select { case s.wake <- struct{}{}: default: } } // setCritical 由调度器 goroutine 在任务进入/离开临界区时设置。 func (s *scheduler) setCritical(v bool) { s.critical.Store(v) } // inCritical 报告运行任务是否在不可抢占临界区。 func (s *scheduler) inCritical() bool { return s.critical.Load() } // hasRoom 报告就绪队列是否还能接收任务。泵入侧据此节流: // 队列满则停止从 channel 取,让背压落回 channel 本身。 func (s *scheduler) hasRoom() bool { s.mu.Lock() defer s.mu.Unlock() return len(s.queue) < s.maxQueue } // enqueue 入队;队列满返回 false(调用方负责计数)。 func (s *scheduler) enqueue(t *Task) bool { s.mu.Lock() defer s.mu.Unlock() if len(s.queue) >= s.maxQueue { s.stats.Rejected++ return false } s.seq++ t.ID = s.seq if t.EnqueuedAt.IsZero() { t.EnqueuedAt = time.Now() } s.stats.Enqueued++ s.queue = append(s.queue, t) return true } // 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() var bestTask *Task var bestFrame *TaskFrame bestKind := nextNone consider := func(t *Task, k nextSelection, fr *TaskFrame) { if bestTask == nil { bestTask, bestKind, bestFrame = t, k, fr return } lt, lb := effectiveLevel(t), effectiveLevel(bestTask) if lt != lb { if lt > lb { bestTask, bestKind, bestFrame = t, k, fr } return } // 同级时 **pending 中断优先**。 // // 为何必需:一次抢占生效后,被挂起的原任务会因饥饿防护提升有效级, // 于是与抢占者同级;若此时按“先到先服务”,原任务(入队更早)会被 // 立刻选回,抢占者永远排不到——抢占变成空转。 if k == nextPending && bestKind != nextPending { bestTask, bestKind, bestFrame = t, k, fr return } if bestKind == nextPending && k != nextPending { return } // 同级同类:先到先服务(ID 兜底保证确定性)。 if t.EnqueuedAt.Before(bestTask.EnqueuedAt) || (t.EnqueuedAt.Equal(bestTask.EnqueuedAt) && t.ID < bestTask.ID) { bestTask, bestKind, bestFrame = t, k, fr } } for _, t := range s.queue { consider(t, nextReady, nil) } for _, t := range s.pendingInterrupts { consider(t, nextPending, nil) } // 中断栈:只比**栈顶**(严格 LIFO)。栈内不做优先级重排—— // 嵌套抢占天然使栈自底向上优先级递增,且“后被打断的先恢复”才是栈语义。 if n := len(s.suspendStack); n > 0 { top := s.suspendStack[n-1] consider(top.Task, nextSuspended, top.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: // 只有栈顶可能被选中,故弹出即截断末位。 s.suspendStack = s.suspendStack[:len(s.suspendStack)-1] } 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) s.signalWake() } // requestPreempt 登记一次中断请求。 // // 返回 true 表示“应该尝试取消运行任务正在进行的可取消步骤(LLM 流式)”。 // // 无论能否抢占,中断请求都进 pendingInterrupts——这样即使运行任务在抢占生效前 // 就正常结束,中断也不会丢(它会被 nextRef 按优先级选出)。 // // 判据用**有效**优先级(饥饿防护),并受抢占冷却约束。 // // 临界区(如记忆整理)内不 arm、不取消:中断只入队,等临界区结束后的安全点处理, // 这是设计 §4.3 的硬要求——那个位置的“不抢占”不能只是不让位,还必须不取消。 func (s *scheduler) requestPreempt(evt *agentIO.InputEvent, level Level) bool { s.mu.Lock() running := s.running critical := s.critical.Load() canPreempt := false if !critical && 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)) return canPreempt } // preemptGrantedFor 报告运行任务是否应在当前安全点让位。 func (s *scheduler) preemptGrantedFor() bool { s.mu.Lock() defer s.mu.Unlock() if !s.preemptArmed || s.running == nil { return false } return s.preemptLevel > effectiveLevel(s.running) } 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.suspendStack) >= s.maxSuspendDepth { s.stats.Rejected++ } s.suspendStack = append(s.suspendStack, &suspendedTask{Task: t, Frame: f}) s.stats.Suspended++ // 饥饿防护:抢占计数 +1(提升有效级)并记录冷却起点。 t.PreemptCount++ t.LastPreemptAt = time.Now() if s.running == t { s.running = nil } // D1=B:中断任务在上一个任务之前的完整状态上开始运行, // 因此这里**不**把被打断任务的任何内容交给它。 s.preemptArmed = false s.preemptLevel = 0 } // canSuspend 报告还有下潜余量(安全点用它决定是否真的让位)。 func (s *scheduler) canSuspend() bool { s.mu.Lock() defer s.mu.Unlock() return len(s.suspendStack) < s.maxSuspendDepth } // done 标记任务执行结束。 func (s *scheduler) done(t *Task) { s.mu.Lock() defer s.mu.Unlock() 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 是内核为任务定级的内部规则。 // // ❗优先级是**内核内部属性**,不做成配置项:内核预定义四级,并按内部规则 // 为任务与中断定级。下方规则只是 v1 的内部缺省值——它决定“谁能让位于谁”, // 属于内核自己的隐私,不对外暴露为运维可调项。 // // 缺省:cli/webui/http 为人机交互(L3),system/_consolidation_ 为后台(L1), // 其余一律默认级(L1)。 func (a *Agent) taskLevel(source, channel string) Level { 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 升序。 // 纯函数:便于对抢占/优先级矩阵做确定性单测。 func pickTaskIndex(q []*Task) int { best := 0 for i := 1; i < len(q); i++ { if taskBefore(q[i], q[best]) { best = i } } return best } // taskBefore 报告 x 是否应先于 y 执行(按**有效**优先级)。 func taskBefore(x, y *Task) bool { lx, ly := effectiveLevel(x), effectiveLevel(y) if lx != ly { return lx > ly } if !x.EnqueuedAt.Equal(y.EnqueuedAt) { return x.EnqueuedAt.Before(y.EnqueuedAt) } return x.ID < y.ID } 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()} } // DumpScheduler 返回调度器的原子快照(供状态页/测试断言)。 func (a *Agent) DumpScheduler() SchedulerSnapshot { if a.sched == nil { return SchedulerSnapshot{} } a.sched.mu.Lock() 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.SuspendStack = append(snap.SuspendStack, a.sched.suspendStack...) snap.MaxSuspendDepth = a.sched.maxSuspendDepth return snap } // schedulerLoop 是唯一的任务执行者(取代原 eventLoop 的输入处理)。 func (a *Agent) schedulerLoop() { defer func() { if r := recover(); r != nil { log.Printf("[agent] schedulerLoop panic recovered: %v\n%s", r, debug.Stack()) time.Sleep(time.Second) go a.schedulerLoop() } }() for { a.pumpInbox() t, f, kind := a.sched.nextRef() if kind == nextNone { // 无待办:阻塞等新输入、新中断(wake)或退出。 select { case evt := <-a.io.InputChan(): a.sched.enqueue(newInputTask(evt)) case msg := <-a.selfInputCh: a.sched.enqueue(newSelfTask(msg)) case <-a.sched.wake: // 中断已入 pendingInterrupts,回到循环顶部重新挑选。 case <-a.ctx.Done(): return } continue } if kind == nextSuspended { a.resumeTask(t, f) continue } a.executeNewTask(t) } } // pumpInbox 把 channel 里**已就绪**的输入搬进就绪队列(非阻塞)。 // // 为什么不直接边收边执行:先把已到达的输入收进队列,选择函数才有意义—— // M3 起抢占必然要看"队列里还压着什么",而 channel 不是可枚举的结构。 // // 队列满即停止泵入(背压落回 channel,语义与设计文档 §4.4 一致)。 func (a *Agent) pumpInbox() { for a.sched.hasRoom() { select { case evt := <-a.io.InputChan(): a.sched.enqueue(newInputTask(evt)) case msg := <-a.selfInputCh: a.sched.enqueue(newSelfTask(msg)) case <-a.ctx.Done(): return default: return } } } // executeTask 执行一个任务(测试与旧调用方的入口);见 executeNewTask。 func (a *Agent) executeTask(t *Task) { a.executeNewTask(t) } // executeNewTask 执行一个**新建**任务,并做任务级 panic 隔离(不变量 I6)。 // // 与改造前的差异(有意):原 eventLoop 在 panic 后重启整个循环, // 现在一个任务的 panic 只丢弃该任务,调度器与其它任务不受影响。 func (a *Agent) executeNewTask(t *Task) { var f *TaskFrame var out stepOutcome = outcomeDone func() { defer func() { if r := recover(); r != nil { log.Printf("[agent] task#%d (%s) panic recovered: %v\n%s", t.ID, t.Level, r, debug.Stack()) } }() switch t.Kind { case TaskKindInput: f, out = a.runInputTask(t.Event) case TaskKindSelf: f, out = a.runInputTask(selfEvent(t.Self)) } }() if out == outcomeSuspended && f != nil { a.sched.suspend(t, f) a.publishEvent(events.EventScheduler, map[string]interface{}{ "action": "suspend", "task": t.ID, "level": int(t.Level), }) return } a.sched.done(t) } // resumeTask 从保存的现场继续一个被抢占的任务。 // // 关键:不重建帧、不重跑 prepare 段——否则会重复提交上下文与事件。 // resumeTask 从保存的现场继续一个被抢占的任务。 // // 关键:不重跑 prepare 段(否则会重复提交上下文与事件),而是先把基础前缀 // 重建到「中断任务之上」,再把本任务自己的现场接回去(见 rebaseFramePrefix)。 func (a *Agent) resumeTask(t *Task, f *TaskFrame) { a.sched.mu.Lock() a.sched.stats.Resumed++ a.sched.mu.Unlock() 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", 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) }