mirror of
https://gitcode.com/JianFeeeee/HomeAgent.git
synced 2026-09-21 01:18:08 +00:00
feat(scheduler): M7 可观测性 + 压力测试 + 端到端测试
设计依据 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 通过
This commit is contained in:
@ -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",
|
||||
|
||||
216
internal/agent/core/scheduler_e2e_test.go
Normal file
216
internal/agent/core/scheduler_e2e_test.go
Normal file
@ -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())
|
||||
}
|
||||
}
|
||||
@ -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
|
||||
}
|
||||
|
||||
@ -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"
|
||||
|
||||
|
||||
@ -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 模型)的启用状态与身份。
|
||||
|
||||
Reference in New Issue
Block a user