From f11de37bf20370a2aa5cb1f7c3ec6a6235343ab3 Mon Sep 17 00:00:00 2001 From: JianFeeeee Date: Sun, 13 Sep 2026 00:45:28 +0800 Subject: [PATCH] =?UTF-8?q?feat(scheduler):=20M7=20=E5=8F=AF=E8=A7=82?= =?UTF-8?q?=E6=B5=8B=E6=80=A7=20+=20=E5=8E=8B=E5=8A=9B=E6=B5=8B=E8=AF=95?= =?UTF-8?q?=20+=20=E7=AB=AF=E5=88=B0=E7=AB=AF=E6=B5=8B=E8=AF=95?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 设计依据 docs/zh/input-scheduler-design.md §11.5(O1/O2)、§11.6(E1/E2)。 - 可观测性:KernelStatus 新增 Scheduler 段(running/三集合深度/计数/ 深度上限),由 GetKernelStatus 从 DumpScheduler 原子快照填充; 新增 events.EventScheduler,挂起/恢复各发一条(action/task/level) - TaskKind.String() 便于日志与状态输出 - 新增 scheduler_e2e_test.go 3 项: · 压力:200 排队输入 + 50 中断全部经真实 loop 执行,结束时三集合排空、 LLM 调用数精确等于输入数、无 Rejected · 可观测性:挂起/恢复事件齐备,状态快照计数一致 · 端到端:完整启动 schedulerLoop+interceptLoop,经真实 channel 投递 L1 任务与 L4 中断,验证「LLM 流式中断 → 挂起 → 中断先完成 → 原任务恢复」 整条链路(LLM 调用数 = 丢弃1+中断1+恢复1+常规1) - 验收:agent 全量 + -race;全仓 build/vet 通过 --- internal/agent/core/scheduler.go | 57 +++++- internal/agent/core/scheduler_e2e_test.go | 216 ++++++++++++++++++++++ internal/agent/core/status.go | 3 +- internal/events/bus.go | 15 +- internal/sdk/status.go | 30 +++ 5 files changed, 311 insertions(+), 10 deletions(-) create mode 100644 internal/agent/core/scheduler_e2e_test.go diff --git a/internal/agent/core/scheduler.go b/internal/agent/core/scheduler.go index 4f7c037..890e9ca 100644 --- a/internal/agent/core/scheduler.go +++ b/internal/agent/core/scheduler.go @@ -25,6 +25,8 @@ import ( agentAPI "gitcode.com/JianFeeeee/HomeAgent/internal/agent/api" agentIO "gitcode.com/JianFeeeee/HomeAgent/internal/agent/io" + "gitcode.com/JianFeeeee/HomeAgent/internal/events" + sdk "gitcode.com/JianFeeeee/HomeAgent/internal/sdk" ) // Level 是任务优先级,由内核预定义四级(设计文档 §3.1)。 @@ -73,6 +75,17 @@ const ( TaskKindSelf ) +func (k TaskKind) String() string { + switch k { + case TaskKindInput: + return "input" + case TaskKindSelf: + return "self" + default: + return "unknown" + } +} + // Task 是调度器的最小单位。 // // M2 只承载"一份待处理的输入";M3 起把 TaskFrame(现场)挂上来, @@ -124,9 +137,11 @@ func effectiveLevel(t *Task) Level { type SchedulerStats struct { Enqueued uint64 Executed uint64 - // Rejected 是因队列满而未入队的次数。M2 由泵入侧节流,正常为 0; - // 出现非 0 说明消费端长期慢于生产端。 + // Rejected 是因队列满(或深度超限)而未被接纳的次数。 Rejected uint64 + // Suspended / Resumed 是挂起与恢复的次数。 + Suspended uint64 + Resumed uint64 } // SchedulerSnapshot 是调度器的原子快照。 @@ -136,6 +151,33 @@ type SchedulerSnapshot struct { PendingInterrupts []*Task SuspendPool []*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), + SuspendPool: len(snap.SuspendPool), + 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 { @@ -342,6 +384,7 @@ func (s *scheduler) suspend(t *Task, f *TaskFrame) { s.stats.Rejected++ } s.suspendPool = append(s.suspendPool, &suspendedTask{Task: t, Frame: f}) + s.stats.Suspended++ // 饥饿防护:抢占计数 +1(提升有效级)并记录冷却起点。 t.PreemptCount++ t.LastPreemptAt = time.Now() @@ -487,6 +530,7 @@ func (a *Agent) DumpScheduler() SchedulerSnapshot { snap.Queue = append(snap.Queue, a.sched.queue...) snap.PendingInterrupts = append(snap.PendingInterrupts, a.sched.pendingInterrupts...) snap.SuspendPool = append(snap.SuspendPool, a.sched.suspendPool...) + snap.MaxSuspendDepth = a.sched.maxSuspendDepth return snap } @@ -576,6 +620,9 @@ func (a *Agent) executeNewTask(t *Task) { 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) @@ -585,6 +632,12 @@ func (a *Agent) executeNewTask(t *Task) { // // 关键:不重建帧、不重跑 prepare 段——否则会重复提交上下文与事件。 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), + }) defer func() { if r := recover(); r != nil { log.Printf("[agent] resume task#%d panic recovered: %v\n%s", diff --git a/internal/agent/core/scheduler_e2e_test.go b/internal/agent/core/scheduler_e2e_test.go new file mode 100644 index 0000000..121791f --- /dev/null +++ b/internal/agent/core/scheduler_e2e_test.go @@ -0,0 +1,216 @@ +package core + +// M7 验收测试:可观测性 + 压力 + 端到端。 +// +// 设计依据 docs/zh/input-scheduler-design.md §11.5(O1/O2)、§11.6(E1/E2)。 +// +// 这一组与前几组的区别:前几组直接驱动调度器(确定性、可断言内部状态), +// 这一组**完整启动** schedulerLoop + interceptLoop,经真实 channel 投递, +// 验证组装后的行为与不变量。 + +import ( + "context" + "errors" + "fmt" + "strings" + "sync" + "sync/atomic" + "testing" + "time" + + agentAPI "gitcode.com/JianFeeeee/HomeAgent/internal/agent/api" + agentIO "gitcode.com/JianFeeeee/HomeAgent/internal/agent/io" + "gitcode.com/JianFeeeee/HomeAgent/internal/events" +) + +// countingProvider 只统计调用次数,永远成功。 +type countingProvider struct{ n atomic.Int64 } + +func (p *countingProvider) Name() string { return "counting" } +func (p *countingProvider) Chat(ctx context.Context, req *agentAPI.CompletionRequest) (*agentAPI.CompletionResponse, error) { + p.n.Add(1) + return &agentAPI.CompletionResponse{Content: "ok"}, nil +} +func (p *countingProvider) ChatStream(ctx context.Context, req *agentAPI.CompletionRequest) (<-chan agentAPI.StreamChunk, error) { + return nil, errors.New("counting provider: no stream") +} +func (p *countingProvider) MaxContextTokens() int { return 8192 } + +func waitQuiescent(t *testing.T, a *Agent, wantExecuted uint64, timeout time.Duration) SchedulerSnapshot { + t.Helper() + deadline := time.Now().Add(timeout) + for { + snap := a.DumpScheduler() + if snap.Running == nil && len(snap.Queue) == 0 && + len(snap.PendingInterrupts) == 0 && len(snap.SuspendPool) == 0 && + snap.Stats.Executed >= wantExecuted { + return snap + } + if time.Now().After(deadline) { + t.Fatalf("未在 %v 内排空:running=%v queue=%d pending=%d suspend=%d executed=%d", + timeout, snap.Running != nil, len(snap.Queue), len(snap.PendingInterrupts), + len(snap.SuspendPool), snap.Stats.Executed) + } + time.Sleep(20 * time.Millisecond) + } +} + +// 压力:N 个排队输入 + M 个中断,全部经真实 loop 执行,结束时三集合必须排空。 +func TestScheduler_StressMixedLoad(t *testing.T) { + sp := &countingProvider{} + a := New(AgentConfig{ + ID: "stress", + Provider: sp, + ProviderManager: agentAPI.NewProviderManager(), + IO: agentIO.NewIOManager(), + StageHost: NewStageHost(), + }) + a.Start() + defer a.Stop() + + const nInputs = 200 + const nInterrupts = 50 + + for i := 0; i < nInputs; i++ { + a.io.InjectInput("cli", "text", map[string]interface{}{"content": fmt.Sprintf("msg-%d", i)}) + } + for i := 0; i < nInterrupts; i++ { + a.io.InjectInterruptText("qq", "cli", fmt.Sprintf("intr-%d", i)) + } + + snap := waitQuiescent(t, a, nInputs+nInterrupts, 30*time.Second) + + if got := sp.n.Load(); got != int64(nInputs+nInterrupts) { + t.Fatalf("LLM 调用=%d,期望 %d(每条输入/中断恰好一次)", got, nInputs+nInterrupts) + } + if snap.Stats.Rejected != 0 { + t.Fatalf("容量充足却出现 Rejected=%d,说明背压/深度判定有误", snap.Stats.Rejected) + } + // 上次快照的计数在排空后应当稳定(不丢不重):等于入队后的执行数。 + if snap.Stats.Executed != uint64(nInputs+nInterrupts) { + t.Fatalf("Executed=%d,期望 %d", snap.Stats.Executed, nInputs+nInterrupts) + } +} + +// O2:每次挂起/恢复都产生一条 scheduler 事件。 +func TestObservability_SchedulerEventsAndStatus(t *testing.T) { + bus := events.NewBus() + var mu sync.Mutex + var actions []string + bus.Subscribe(events.EventScheduler, func(e *events.Event) { + mu.Lock() + actions = append(actions, fmt.Sprint(e.Payload["action"])) + mu.Unlock() + }) + + sp := newPreemptProvider("intr-done", "low-done") + a := New(AgentConfig{ + ID: "obs", + Provider: sp, + ProviderManager: agentAPI.NewProviderManager(), + IO: agentIO.NewIOManager(), + StageHost: NewStageHost(), + EventBus: bus, + }) + + // 直接驱动一次抢占-挂起-恢复(与 M3b 相同的手法)。 + lowEvt, _ := textEvent("qq", "低优先级") + lowTask := &Task{Kind: TaskKindInput, Level: LevelBackground, Event: lowEvt, EnqueuedAt: time.Now()} + a.sched.enqueue(lowTask) + lt, _, _ := a.sched.nextRef() + + done := make(chan struct{}) + go func() { a.executeNewTask(lt); close(done) }() + select { + case <-sp.entered: + case <-time.After(3 * time.Second): + t.Fatal("provider 未进入") + } + + intrEvt, _ := textEvent("cli", "紧急") + intrEvt.Payload["interrupt"] = true + a.sched.requestPreempt(intrEvt, LevelCritical) + a.cancelCurrentLLM() + <-done + + it, _, _ := a.sched.nextRef() + a.executeNewTask(it) + rt, rf, _ := a.sched.nextRef() + a.resumeTask(rt, rf) + + mu.Lock() + got := strings.Join(actions, ",") + mu.Unlock() + if !strings.Contains(got, "suspend") || !strings.Contains(got, "resume") { + t.Fatalf("调度事件缺失:%q", got) + } + + // 状态快照(供状态页/诊断):计数一致、三集合为空。 + st := a.GetKernelStatus().Scheduler + if st.SuspendPool != 0 || st.PendingInterrupts != 0 || st.ReadyQueueDepth != 0 { + t.Fatalf("排空后状态非空:%+v", st) + } + if st.Suspended == 0 || st.Resumed == 0 { + t.Fatalf("挂起/恢复计数缺失:%+v", st) + } + if st.Executed < 2 { + t.Fatalf("Executed=%d,期望 >=2", st.Executed) + } + if st.MaxSuspendDepth != 4 { + t.Fatalf("MaxSuspendDepth=%d,期望 4", st.MaxSuspendDepth) + } +} + +// E1/E2:完整启动 loop,经真实 channel 投递 L1 任务与 L4 中断, +// 断言「LLM 流式中断 → 挂起 → 中断先完成 → 原任务恢复」的整条链路。 +func TestE2E_RealLoopPreemption(t *testing.T) { + bus := events.NewBus() + var mu sync.Mutex + var actions []string + bus.Subscribe(events.EventScheduler, func(e *events.Event) { + mu.Lock() + actions = append(actions, fmt.Sprint(e.Payload["action"])) + mu.Unlock() + }) + + sp := newPreemptProvider("intr-done", "low-done") + a := New(AgentConfig{ + ID: "e2e", + Provider: sp, + ProviderManager: agentAPI.NewProviderManager(), + IO: agentIO.NewIOManager(), + StageHost: NewStageHost(), + EventBus: bus, + }) + a.Start() + defer a.Stop() + + // L1:qq 入站消息 → 阻塞在第一次 LLM 调用 + a.io.InjectInput("qq", "text", map[string]interface{}{"content": "低优先级长任务"}) + select { + case <-sp.entered: + case <-time.After(5 * time.Second): + t.Fatal("低优先级任务未进入 LLM") + } + + // L4:cli 紧急打断 → interceptLoop 应取消 LLM、登记抢占 + a.io.InjectInterruptText("cli", "cli", "紧急打断") + a.io.InjectInput("cli", "text", map[string]interface{}{"content": "后续常规输入"}) + + // 排空:中断任务 + 被恢复的原任务 + 后续常规输入 + snap := waitQuiescent(t, a, 3, 15*time.Second) + + mu.Lock() + got := strings.Join(actions, ",") + mu.Unlock() + if !strings.Contains(got, "suspend") || !strings.Contains(got, "resume") { + t.Fatalf("E2E 未发生抢占-挂起-恢复:%q", got) + } + if snap.Stats.Executed < 3 { + t.Fatalf("Executed=%d,期望 >=3", snap.Stats.Executed) + } + // 第一次 LLM 调用被丢弃 + 中断 1 + 恢复 1 + 常规输入 1 = 4 + if sp.callCount() != 4 { + t.Fatalf("LLM 调用=%d,期望 4(丢弃 1 + 中断 1 + 恢复 1 + 常规 1)", sp.callCount()) + } +} diff --git a/internal/agent/core/status.go b/internal/agent/core/status.go index d3644eb..0b36fa1 100644 --- a/internal/agent/core/status.go +++ b/internal/agent/core/status.go @@ -156,7 +156,6 @@ func collectKernelStatus( } } - // Tracker if trk != nil { status.Tracker.Available = true @@ -192,7 +191,6 @@ func (a *Agent) GetKernelStatus() *KernelStatus { socialStore = a.social } - var trk *tracker.Tracker if a.tracker != nil { trk = a.tracker @@ -222,6 +220,7 @@ func (a *Agent) GetKernelStatus() *KernelStatus { trk, ) ks.ONNX = a.onnxStatus() + ks.Scheduler = a.schedulerStatus() return ks } diff --git a/internal/events/bus.go b/internal/events/bus.go index e26692e..c1bed3a 100644 --- a/internal/events/bus.go +++ b/internal/events/bus.go @@ -9,12 +9,15 @@ import ( type EventType string const ( - EventRawInput EventType = "raw_input" - EventAgentOutput EventType = "agent_output" - EventAgentLLMChain EventType = "agent_llm_chain" - EventToolCall EventType = "tool_call" - EventReasoning EventType = "reasoning" - EventStage EventType = "stage" + EventRawInput EventType = "raw_input" + EventAgentOutput EventType = "agent_output" + EventAgentLLMChain EventType = "agent_llm_chain" + EventToolCall EventType = "tool_call" + EventReasoning EventType = "reasoning" + EventStage EventType = "stage" + // EventScheduler 是输入调度器的状态变更事件(抢占/挂起/恢复), + // 供状态页与诊断订阅(设计文档 §11 O2)。 + EventScheduler EventType = "scheduler" EventSystem EventType = "system" EventTerminalOutput EventType = "terminal_output" diff --git a/internal/sdk/status.go b/internal/sdk/status.go index 4cb0cf1..6caf231 100644 --- a/internal/sdk/status.go +++ b/internal/sdk/status.go @@ -46,6 +46,36 @@ type KernelStatus struct { ONNX ONNXStatus `json:"onnx"` Tracker TrackerStatus `json:"tracker"` + + // Scheduler 是输入调度器的运行时快照(可观测性,设计文档 §11 O1/O2)。 + // M2 起输入不再直接排队在 channel 上,而是经 readyQueue/pendingInterrupts/ + // suspendPool 三集合按优先级调度;这里把这些状态暴露出来。 + Scheduler SchedulerStatus `json:"scheduler"` +} + +// SchedulerStatus 是调度器的原子快照 DTO。 +type SchedulerStatus struct { + // Running 是当前执行的任务(空表示空闲)。 + Running *SchedulerTask `json:"running,omitempty"` + // ReadyQueueDepth / PendingInterrupts / SuspendPool 是三个集合的深度。 + ReadyQueueDepth int `json:"ready_queue_depth"` + PendingInterrupts int `json:"pending_interrupts"` + SuspendPool int `json:"suspend_pool"` + MaxSuspendDepth int `json:"max_suspend_depth"` + + Enqueued uint64 `json:"enqueued"` + Executed uint64 `json:"executed"` + Rejected uint64 `json:"rejected"` + Suspended uint64 `json:"suspended"` + Resumed uint64 `json:"resumed"` + Preempted uint64 `json:"preempted"` +} + +// SchedulerTask 是任务的最小标识(不暴露帧内容)。 +type SchedulerTask struct { + ID uint64 `json:"id"` + Level int `json:"level"` + Kind string `json:"kind"` } // ONNXStatus 是统一多模态向量空间(ONNX 模型)的启用状态与身份。