diff --git a/internal/agent/core/agent.go b/internal/agent/core/agent.go index d2afea0..288d6d9 100644 --- a/internal/agent/core/agent.go +++ b/internal/agent/core/agent.go @@ -109,6 +109,10 @@ type Agent struct { // 高优先级打断通道:interceptLoop 注入,process() 在工具循环轮次间非阻塞读取 interceptCh chan *agentIO.InputEvent + // 输入调度器:就绪队列、任务抽象与快照(见 scheduler.go)。 + // M2 起取代 eventLoop 的隐式 channel 排队。 + sched *scheduler + // 进行中的 LLM 请求取消函数,interceptLoop 可调用以在请求中打断 cancelLLM context.CancelFunc llmMu sync.Mutex @@ -298,6 +302,7 @@ func New(cfg AgentConfig) *Agent { selfInputCh: make(chan selfInputMsg, 64), childTasks: make(map[string]*childTaskState), interceptCh: make(chan *agentIO.InputEvent, 64), + sched: newScheduler(256), pluginHealth: newPluginHealthTracker(), thinkingEnabled: cfg.ThinkingEnabled, inputCfg: cfg.InputProcessing, @@ -315,7 +320,7 @@ func New(cfg AgentConfig) *Agent { func (a *Agent) SetSkillIndexProvider(p SkillIndexProvider) { a.skillIndex = p } func (a *Agent) Start() { - go a.eventLoop() + go a.schedulerLoop() go a.interceptLoop() go a.distillLoop() go a.archiveLoop() diff --git a/internal/agent/core/eventloop.go b/internal/agent/core/eventloop.go index d7b54b4..fc84196 100644 --- a/internal/agent/core/eventloop.go +++ b/internal/agent/core/eventloop.go @@ -13,25 +13,11 @@ import ( pubsdk "gitcode.com/JianFeeeee/homeagent-sdk/sdk" ) -func (a *Agent) eventLoop() { - defer func() { - if r := recover(); r != nil { - log.Printf("[agent] eventLoop panic recovered: %v\n%s", r, debug.Stack()) - time.Sleep(time.Second) - go a.eventLoop() - } - }() - for { - select { - case evt := <-a.io.InputChan(): - a.handleInput(evt) - case msg := <-a.selfInputCh: - a.handleSelfInput(msg) - case <-a.ctx.Done(): - return - } - } -} +// eventLoop 已由 scheduler.go 的 schedulerLoop 取代(M2)。 +// +// 原实现直接在 select 里处理 inputCh/selfInputCh,没有任何可枚举的队列、 +// 无法承载优先级与抢占;现在任务先入就绪队列,由选择函数 pickTaskIndex 决定下一个。 +// 兼容性说明:M2 全部任务为 LevelBackground,因此行为等价于原先的 FIFO。 func (a *Agent) interceptLoop() { defer func() { diff --git a/internal/agent/core/scheduler.go b/internal/agent/core/scheduler.go new file mode 100644 index 0000000..1d0db20 --- /dev/null +++ b/internal/agent/core/scheduler.go @@ -0,0 +1,287 @@ +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 语义逐条一致。抢占、suspendPool、pendingInterrupts、 +// 任务级回执在 M3–M6 加入。 +// +// 并发模型(不变量 I2):readyQueue/running/stats 只由 schedulerLoop 写, +// 外部只读——读取一律经 DumpScheduler() 加锁取快照。 + +import ( + "log" + "runtime/debug" + "sync" + "time" + + agentIO "gitcode.com/JianFeeeee/HomeAgent/internal/agent/io" +) + +// 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" + } +} + +// TaskKind 区分任务来源。 +type TaskKind int + +const ( + // TaskKindInput 来自 io.InputChan(外部/插件注入的输入)。 + TaskKindInput TaskKind = iota + // TaskKindSelf 来自 selfInputCh(内核自循环:记忆整理、子任务通知)。 + TaskKindSelf +) + +// 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 +} + +// SchedulerStats 是调度器的累计计数(可观测性,设计文档 §11 O2)。 +type SchedulerStats struct { + Enqueued uint64 + Executed uint64 + // Rejected 是因队列满而未入队的次数。M2 由泵入侧节流,正常为 0; + // 出现非 0 说明消费端长期慢于生产端。 + Rejected uint64 +} + +// SchedulerSnapshot 是调度器的原子快照。 +type SchedulerSnapshot struct { + Running *Task + Queue []*Task + Stats SchedulerStats +} + +type scheduler struct { + mu sync.Mutex + queue []*Task + running *Task + seq uint64 + stats SchedulerStats + maxQueue int +} + +func newScheduler(maxQueue int) *scheduler { + if maxQueue <= 0 { + maxQueue = 256 + } + return &scheduler{maxQueue: maxQueue} +} + +// 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。 +func (s *scheduler) next() *Task { + s.mu.Lock() + defer s.mu.Unlock() + if len(s.queue) == 0 { + return nil + } + i := pickTaskIndex(s.queue) + t := s.queue[i] + s.queue = append(s.queue[:i], s.queue[i+1:]...) + s.running = t + return t +} + +// done 标记任务执行结束。 +func (s *scheduler) done(t *Task) { + s.mu.Lock() + defer s.mu.Unlock() + if s.running == t { + s.running = nil + } + s.stats.Executed++ +} + +// 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 { + if x.Level != y.Level { + return x.Level > y.Level + } + 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()} +} + +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...) + 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 := a.sched.next() + if t == nil { + // 就绪队列空:阻塞等新输入或退出。 + 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 + } + continue + } + a.executeTask(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 执行一个任务,并做**任务级 panic 隔离**(不变量 I6)。 +// +// 与改造前的差异(有意):原 eventLoop 在 panic 后重启整个循环, +// 现在一个任务的 panic 只丢弃该任务,调度器与其它任务不受影响。 +func (a *Agent) executeTask(t *Task) { + 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: + a.handleInput(t.Event) + case TaskKindSelf: + a.handleSelfInput(t.Self) + } + }() + a.sched.done(t) +} diff --git a/internal/agent/core/scheduler_test.go b/internal/agent/core/scheduler_test.go new file mode 100644 index 0000000..03b57af --- /dev/null +++ b/internal/agent/core/scheduler_test.go @@ -0,0 +1,216 @@ +package core + +// M2 验收测试:调度器骨架(就绪队列、选择函数、快照、panic 隔离)。 +// +// 设计依据 docs/zh/input-scheduler-design.md §11.4(Q1/Q4)与 §11.5(O1/K1)。 + +import ( + "testing" + "time" + + agentAPI "gitcode.com/JianFeeeee/HomeAgent/internal/agent/api" + agentIO "gitcode.com/JianFeeeee/HomeAgent/internal/agent/io" +) + +func mkTask(id uint64, level Level, at time.Time) *Task { + return &Task{ID: id, Level: level, EnqueuedAt: at} +} + +// Q1:选择函数的排序键是 (-Level, EnqueuedAt, ID)。 +func TestPickTaskIndex_Ordering(t *testing.T) { + base := time.Date(2026, 9, 12, 12, 0, 0, 0, time.UTC) + + cases := []struct { + name string + queue []*Task + want []uint64 + }{ + { + name: "高优先级先执行,与入队先后无关", + queue: []*Task{ + mkTask(1, LevelBackground, base), + mkTask(2, LevelCritical, base.Add(time.Second)), + mkTask(3, LevelMessage, base.Add(2*time.Second)), + }, + want: []uint64{2, 3, 1}, + }, + { + name: "同优先级先到先服务", + queue: []*Task{ + mkTask(1, LevelInteractive, base.Add(3*time.Second)), + mkTask(2, LevelInteractive, base.Add(time.Second)), + mkTask(3, LevelInteractive, base.Add(2*time.Second)), + }, + want: []uint64{2, 3, 1}, + }, + { + name: "同优先级同入队时刻用 ID 兜底(保证确定性)", + queue: []*Task{ + mkTask(7, LevelMessage, base), + mkTask(3, LevelMessage, base), + mkTask(5, LevelMessage, base), + }, + want: []uint64{3, 5, 7}, + }, + { + name: "四级全覆盖", + queue: []*Task{ + mkTask(1, LevelBackground, base), + mkTask(2, LevelMessage, base), + mkTask(3, LevelInteractive, base), + mkTask(4, LevelCritical, base), + }, + want: []uint64{4, 3, 2, 1}, + }, + } + + for _, c := range cases { + t.Run(c.name, func(t *testing.T) { + q := append([]*Task(nil), c.queue...) + var got []uint64 + for len(q) > 0 { + i := pickTaskIndex(q) + got = append(got, q[i].ID) + q = append(q[:i], q[i+1:]...) + } + if len(got) != len(c.want) { + t.Fatalf("取出的任务数=%d,期望 %d", len(got), len(c.want)) + } + for i := range got { + if got[i] != c.want[i] { + t.Fatalf("执行顺序=%v,期望 %v", got, c.want) + } + } + }) + } +} + +// Q4:队列有界;满了必须拒绝并计数,而不是静默丢弃或无界增长。 +func TestScheduler_EnqueueBackpressure(t *testing.T) { + s := newScheduler(2) + if !s.enqueue(newSelfTask(selfInputMsg{text: "a"})) { + t.Fatal("第 1 个任务应入队成功") + } + if !s.enqueue(newSelfTask(selfInputMsg{text: "b"})) { + t.Fatal("第 2 个任务应入队成功") + } + if s.hasRoom() { + t.Fatal("队列已满,hasRoom 应为 false") + } + if s.enqueue(newSelfTask(selfInputMsg{text: "c"})) { + t.Fatal("队列满时第 3 个任务必须被拒绝") + } + if s.stats.Rejected != 1 { + t.Fatalf("Rejected=%d,期望 1", s.stats.Rejected) + } + if s.stats.Enqueued != 2 { + t.Fatalf("Enqueued=%d,期望 2", s.stats.Enqueued) + } +} + +// 生命周期:next 置 running 并移出队列;done 清 running 并累加计数。 +func TestScheduler_Lifecycle(t *testing.T) { + s := newScheduler(4) + s.enqueue(newSelfTask(selfInputMsg{text: "a"})) + s.enqueue(newSelfTask(selfInputMsg{text: "b"})) + + t1 := s.next() + if t1 == nil || s.running != t1 { + t.Fatal("next 应取出任务并置为 running") + } + if len(s.queue) != 1 { + t.Fatalf("取出后队列长度=%d,期望 1", len(s.queue)) + } + // 队列内不得同时出现 running(O1:三集合互不重叠)。 + for _, q := range s.queue { + if q == t1 { + t.Fatal("running 任务不得同时留在就绪队列") + } + } + + s.done(t1) + if s.running != nil { + t.Fatal("done 后 running 应为 nil") + } + if s.stats.Executed != 1 { + t.Fatalf("Executed=%d,期望 1", s.stats.Executed) + } + if s.next() == nil { + t.Fatal("队列里还有 b,next 不应为 nil") + } + if s.next() != nil { + t.Fatal("队列已空,next 应返回 nil") + } +} + +// K1:任务 panic 必须被隔离——调度器统计仍然推进,且不向外抛出。 +func TestScheduler_PanicIsolationOnExecuteTask(t *testing.T) { + sp := &scriptProvider{} + a := New(AgentConfig{ + ID: "sched-panic", + Provider: sp, + ProviderManager: agentAPI.NewProviderManager(), + IO: agentIO.NewIOManager(), + }) + + // Event 为 nil:handleInput 解引用即 panic,用来验证 recover 生效。 + task := &Task{Kind: TaskKindInput, Level: DefaultLevel, Event: nil} + a.executeTask(task) // 若未隔离,这里会 panic 冒泡使测试失败 + + if a.sched.stats.Executed != 1 { + t.Fatalf("panic 后 Executed=%d,期望 1(任务失败但调度器存活)", a.sched.stats.Executed) + } + if a.sched.running != nil { + t.Fatal("panic 后 running 必须被清空") + } +} + +// O1 轻量版:快照与内部状态一致,且 running 不出现在 queue 里。 +func TestScheduler_SnapshotConsistency(t *testing.T) { + sp := &scriptProvider{} + a := New(AgentConfig{ + ID: "sched-snap", + Provider: sp, + ProviderManager: agentAPI.NewProviderManager(), + IO: agentIO.NewIOManager(), + }) + + a.sched.enqueue(newSelfTask(selfInputMsg{text: "a"})) + a.sched.enqueue(newSelfTask(selfInputMsg{text: "b"})) + snap := a.DumpScheduler() + if snap.Running != nil { + t.Fatal("尚未 next,快照的 running 应为 nil") + } + if len(snap.Queue) != 2 || snap.Stats.Enqueued != 2 { + t.Fatalf("快照不一致:queue=%d enqueued=%d", len(snap.Queue), snap.Stats.Enqueued) + } + + r := a.sched.next() + a.executeTask(&Task{Kind: TaskKindSelf, Level: DefaultLevel, Self: selfInputMsg{text: "a"}}) + snap = a.DumpScheduler() + if snap.Running != r { + t.Fatal("执行完成后 running 应仍指向未 done 的任务") + } + for _, q := range snap.Queue { + if q == r { + t.Fatal("快照中 running 与 queue 不得重叠") + } + } + + // 队列快照必须是副本:改快照不得影响调度器。 + snap.Queue = append(snap.Queue, &Task{}) + if len(a.DumpScheduler().Queue) != 1 { + t.Fatal("DumpScheduler 必须返回队列副本") + } +} + +// Level 的字面量是持久化/日志契约,改值必须是有意的。 +func TestLevelContract(t *testing.T) { + if LevelBackground != 1 || LevelMessage != 2 || LevelInteractive != 3 || LevelCritical != 4 { + t.Fatalf("四级取值被改动:%d/%d/%d/%d", + LevelBackground, LevelMessage, LevelInteractive, LevelCritical) + } + if DefaultLevel != LevelBackground { + t.Fatalf("默认级必须是 L1(显式才是特权),实际 %v", DefaultLevel) + } +}