mirror of
https://gitcode.com/JianFeeeee/HomeAgent.git
synced 2026-09-21 09:28:14 +00:00
设计:docs/zh/resident-subagent-design.md §6/§7/§8/§9/§10。
## N3 生命周期(resident.go)
- `SpawnResident`:主库**受限句柄** + 自己的 temp 实例(`LightMemory`)⇒ 子的轻量内核;
划入 inputch(登记归属)、授权输出通道、注入任务提示词;为父登记 `child/<id>` 入站 inputch;
建独立 `IOManager`(共享通道登记表);启动子。
- `DestroyResident`:停子内核、归还划入的 inputch(回到未分配)、丢弃 temp 目录、出登记表。
- `Stop()` → `StopResidents()`:**父退出必须销毁全部子、不留孤儿**(设计 §10 硬约束)。
- `Residents()` 登记表快照;`ResidentTable(id)` 父 pull 子的处理表(不打断)。
## N4 跨 agent 投递
- 子→父:`notify_parent` → 投进父的 `child/<id>` inputch,优先级 **L3**。
- 父→子:`SendToResident` → 投进子的 inputch,优先级 **L4**;
`isKernelLevelSource` 泛化为"该 agent 的上级"(`AgentConfig.KernelSource`)⇒
只有父能在子的阶梯上产生 L4(子内部一律 ≤L3)。
- 子的 contextfull → 父侧 `raiseKernelInterrupt`(内核级事件,带子标识,父侧 L4)。
## N5 inputch 处理表
- 子持有;`inputch_note` 主动写**优先**,轮末 `autoRecordInputch` 兜底 ⇒ 每轮必有记录。
- 压缩时清表(表记的是被压掉那段窗口的逐轮处理)。
## N6 contextfull(判据修正)
❗初版判据是"拼好的 `f.Msgs` 估算 > 90% 窗口",**结构上永不成立**:
`buildMessages` 拿到的 `budget.ContextTokens` 由 `targetUsage = 0.8 × 窗口` 推出,
时间线**在拼进消息之前就被预算裁过**,`f.Msgs` 封顶在 ~80% 窗口。
(初版测试用一个比系统提示词还小的窗口才勉强越线 —— 那等于什么都没测。)
现判据 = **未裁剪的积累上下文**(`a.context.Recent(0)`)超过窗口 90%:
它超过就说明下一轮必须丢事件,这正是"上下文满"。
三处置:`CompressResident`(保留语义:`TrimKeepRecent` 保留最近 N 条 + 清表)/
`ReclaimResident`(取消语义:`ExportTriples` 读 temp → 父选出要保留的 → `Commit` 进 main → 取消该子)/
`DestroyResident`(立刻销毁并移除)。
## 工具面
`resident_agents`(父,单工具多动作:list/create/send/inspect/compress/reclaim/destroy)、
`notify_parent` + `inputch_note`(子)。声明条件式:父才有前者,子才有后两者。
## 验收
`resident_test.go` 五项:生命周期与不留孤儿、双向投递(含"子的主动消息不得以 L4 出现")、
处理表(自动写 vs 主动写优先)、contextfull + 三处置、
**压力 8 子 × 12 轮(父→子 L4 与普通输入各半)+ 双向汇报 + 父退出清理**。
全仓 go test ./... 37 包 ok / 0 FAIL;`-race`(agent/memory/plugin)干净;
压力 `-race -count=3` 通过。
938 lines
31 KiB
Go
938 lines
31 KiB
Go
package core
|
||
|
||
// 输入调度器:四级中断优先级 · 可抢占 · 现场保存/恢复。
|
||
//
|
||
// 设计依据 docs/zh/input-scheduler-design.md。
|
||
//
|
||
// # 模型(两类别 + 四级)
|
||
//
|
||
// 类别由**用哪个注入 API**决定,与通道名无关:
|
||
// - TaskInterrupt:InjectInterrupt* 注入。带级别 L1..L4,可抢占,
|
||
// 可被更高级中断打断(被打断的现场压入**中断栈**)。
|
||
// - TaskQueued:InjectText*/InjectInputSync* 与内核自循环。**无级别**,
|
||
// 用于“不需及时处理”的场景,可被**任何**中断打断。
|
||
//
|
||
// 级别只属于中断:
|
||
// - L1..L3 由插件在 InjectOptions.Priority 里声明(见 clampPluginLevel);
|
||
// - L4 给“立即打断”能力:内核自身(raiseKernelInterrupt:panic / selfip)
|
||
// 与**内核级插件**(编译期内置插件,如 WebUI 终止按钮)可声明;
|
||
// 外部插件经 proc 桥被夹到 L3,core 里也再判一次来源。
|
||
//
|
||
// # 选择顺序
|
||
//
|
||
// 1. immediate —— 刚抢占成功的中断(抢占必须立即生效)
|
||
// 2. 中断队列 L4→L1(同级 FIFO)
|
||
// 3. 中断栈顶(与 2 的队头比级别,取高者;栈顶无级别时中断必胜)
|
||
// 4. 排队队列(FIFO)
|
||
//
|
||
// # 并发模型(不变量 I2)
|
||
//
|
||
// queue/running/栈/stats 只由 schedulerLoop 与调度 goroutine 写;
|
||
// interruptLoop 只写中断登记与让位信号,**从不碰帧**。外部读取一律经
|
||
// DumpScheduler() 加锁取快照。
|
||
|
||
import (
|
||
"fmt"
|
||
"log"
|
||
"runtime/debug"
|
||
"strings"
|
||
"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 是**中断**的优先级,由内核预定义四级。
|
||
//
|
||
// 语义:它衡量“这项工作有多不能等”,与具体通道名无关。
|
||
// 插件在中断注入时通过 InjectOptions.Priority 声明 L1..L3;
|
||
// **L4 由内核独占**(panic、内核事件 selfip),插件声明 L4 会被夹到 L3。
|
||
//
|
||
// 排队输入(InjectText* / InjectInputSync*)**没有级别**:它们本就是
|
||
// “不需及时处理”的那一类,可被任何中断打断(见 TaskClass)。
|
||
type Level int
|
||
|
||
const (
|
||
// LevelBackground L1:完全可等。例:QQ/微信这类异步消息、批量通知。
|
||
LevelBackground Level = 1
|
||
// LevelMessage L2:一般提醒。例:插件希望尽快看到、但不紧急的提示。
|
||
LevelMessage Level = 2
|
||
// LevelInteractive L3:需及时处理。例:时钟/定时器到达、终端输出、交互输入。
|
||
LevelInteractive Level = 3
|
||
// LevelCritical L4:**内核独占**。panic 中断、内核事件中断(selfip)。
|
||
// 插件不得声明此级。
|
||
LevelCritical Level = 4
|
||
)
|
||
|
||
// DefaultLevel 是未显式声明时的中断级别。
|
||
//
|
||
// 取最低级是刻意的:**显式才是特权**,新插件不会默认拿到抢占权。
|
||
const DefaultLevel = LevelBackground
|
||
|
||
// clampPluginLevel 把**非内核级**来源声明的级别夹到 L1..L3。
|
||
//
|
||
// L4 是“立即打断”能力(panic / 内核事件 / 内核级插件的终止按钮),
|
||
// 只给内核与编译期内置插件;外部插件声明 L4 会被夹到 L3。
|
||
func clampPluginLevel(l Level) Level {
|
||
if l < LevelBackground {
|
||
return DefaultLevel
|
||
}
|
||
if l > LevelInteractive {
|
||
return LevelInteractive
|
||
}
|
||
return l
|
||
}
|
||
|
||
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 已删除。
|
||
//
|
||
// 为何不保留:优先级是**内核内部属性**,不是配置项。
|
||
// 曾一度做成 `core.agent.priority.<channel>`(配置中心可见),
|
||
// 那等于把内核的调度内部属性外化成运维配置,与设计意图相反。
|
||
//
|
||
// 现在级别的来源只有两个(见 Task/Level 注释):
|
||
// - 插件在中断注入时声明(InjectOptions.Priority,L1..L3);
|
||
// - 内核内部产生 L4(panic / selfip)。
|
||
|
||
// TaskClass 是任务的两大类别——**由“用哪个注入 API”决定,与通道名无关**。
|
||
//
|
||
// 这是模型的核心区分:
|
||
// - InjectInterrupt* → TaskInterrupt:带级别,可抢占,可被更高级中断打断(→ 中断栈)
|
||
// - InjectText* / InjectInputSync* / 内核自循环 → TaskQueued:无级别,
|
||
// 可被**任何**中断打断(“用于不需要及时处理的场景”)
|
||
type TaskClass int
|
||
|
||
const (
|
||
TaskQueued TaskClass = iota
|
||
TaskInterrupt
|
||
)
|
||
|
||
func (c TaskClass) String() string {
|
||
switch c {
|
||
case TaskQueued:
|
||
return "queued"
|
||
case TaskInterrupt:
|
||
return "interrupt"
|
||
default:
|
||
return "unknown"
|
||
}
|
||
}
|
||
|
||
// 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 是调度器的最小单位。
|
||
type Task struct {
|
||
ID uint64
|
||
Class TaskClass
|
||
// Level 仅对 TaskInterrupt 有意义;TaskQueued 恒为 0(无级别)。
|
||
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 返回任务的**有效**级别。
|
||
//
|
||
// 排队输入恒为 0(无级别):任何中断(≥ L1)都大于它——这正好实现
|
||
// “排队输入可被任何中断打断”。
|
||
//
|
||
// 中断则叠加饥饿防护:被抢占越多的中断越“值钱”,逐步追上抢占它的流;
|
||
// 封顶 L4,因此它永远不会反过来抢占内核紧急中断。
|
||
func effectiveLevel(t *Task) Level {
|
||
if t.Class != TaskInterrupt {
|
||
return 0
|
||
}
|
||
p := t.PreemptCount
|
||
if p > preemptPromotionCap {
|
||
p = preemptPromotionCap
|
||
}
|
||
l := t.Level + Level(p)
|
||
if l > LevelCritical {
|
||
l = LevelCritical
|
||
}
|
||
return l
|
||
}
|
||
|
||
// canPreempt 是唯一的抢占判据。
|
||
//
|
||
// 由于 effectiveLevel(排队)=0,这一个比较同时覆盖两条规则:
|
||
// - running 是排队任务 → 任何中断(≥L1)都能抢占;
|
||
// - running 是中断 Li → 只有 Lj > Li 的中断能抢占(严格大于)。
|
||
func canPreempt(incoming, running *Task) bool {
|
||
if incoming == nil || running == nil {
|
||
return false
|
||
}
|
||
if incoming.Class != TaskInterrupt {
|
||
return false // 排队输入从不抢占
|
||
}
|
||
return effectiveLevel(incoming) > effectiveLevel(running)
|
||
}
|
||
|
||
// SchedulerStats 是调度器的累计计数(可观测性,设计文档 §11 O2)。
|
||
//
|
||
// InterruptsByLevel / PreemptsByLevel 按**中断级别**分桶(下标 1..4):
|
||
// “各级中断各登记了多少、各真正抢断了多少次”。按级别验收(而不是只看总数)
|
||
// 是这套调度器的核心判据——总数相同、级别分布不同,行为完全不同。
|
||
type SchedulerStats struct {
|
||
Enqueued uint64
|
||
Executed uint64
|
||
// Rejected 是因队列满(或深度超限)而未被接纳的次数。
|
||
Rejected uint64
|
||
// Suspended / Resumed 是挂起与恢复的次数。
|
||
// 不变量:系统排空后 Suspended == Resumed(挂起必然被恢复),
|
||
// 因此两者各自只在**一处**计数(suspend / resumeTask)。
|
||
Suspended uint64
|
||
Resumed uint64
|
||
// InterruptsByLevel[1..4]:各级中断被**登记**的次数(含未抢占成功的)。
|
||
InterruptsByLevel [5]uint64
|
||
// PreemptsByLevel[1..4]:各级中断**判定为可抢占并进入 immediate**的次数。
|
||
// 注意它与 Suspended 不等价:受害者可能在让位信号生效前就自行结束,
|
||
// 此时抢占者仍然"下一个运行",但没有挂起发生。
|
||
PreemptsByLevel [5]uint64
|
||
}
|
||
|
||
// bumpInterruptLevel 按级别累加(级别必须落在 1..4,否则忽略——
|
||
// 排队任务没有级别,不该出现在中断计数里)。
|
||
func (st *SchedulerStats) bumpInterruptLevel(dst *[5]uint64, lv Level) {
|
||
if lv >= LevelBackground && lv <= LevelCritical {
|
||
dst[lv]++
|
||
}
|
||
}
|
||
|
||
// SchedulerSnapshot 是调度器的原子快照。
|
||
type SchedulerSnapshot struct {
|
||
Running *Task
|
||
// Queue 是排队输入队列(无级别,FIFO)。
|
||
Queue []*Task
|
||
// InterruptQueues[level] 是四条中断队列(下标 1..4,同级 FIFO)。
|
||
InterruptQueues [5][]*Task
|
||
// Immediate 是刚抢占成功、将在下一个安全点立即运行的中断(最多一个)。
|
||
Immediate *Task
|
||
// PendingInterrupts = 四条中断队列 + Immediate(对外的待处理中断总数视图)。
|
||
PendingInterrupts []*Task
|
||
// SuspendStack:中断栈(含嵌套抢占的多个现场),**栈顶**优先恢复。
|
||
SuspendStack []*suspendedTask
|
||
Stats SchedulerStats
|
||
// MaxInterruptFrames 是中断栈帧数的结构上界(= 中断级数,不是配置项)。
|
||
MaxInterruptFrames 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.MaxInterruptFrames,
|
||
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
|
||
|
||
// interruptQueues[level]:四条**中断队列**(level 1..4),同级 FIFO。
|
||
// 未能立即抢占的中断(级别不足,或运行任务在临界区)按级别入队,
|
||
// nextRef 从 L4 到 L1 依次扫描。
|
||
interruptQueues [5][]*Task
|
||
// immediate:刚抢占成功的中断。抢占必须**立即生效**,所以它不经队列,
|
||
// 在下一个安全点直接运行。这也消除了“抢占者与被抢占者同级”的比较问题——
|
||
// 抢占者根本不需要和栈顶比。
|
||
immediate *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{}
|
||
// maxInterruptFrames:中断栈帧数的**结构上界**,不是配置项。
|
||
//
|
||
// 链条 = 排队(L0) ← I(L1) ← I(L2) ← I(L3) ← I(L4 运行中),
|
||
// 被挂起 4 帧;L4 之上没有更高级别,链到此为止。超限只可能是内核 bug,
|
||
// 因此这里只做防御性计数,**不降级、不丢弃帧**。
|
||
maxInterruptFrames int
|
||
}
|
||
|
||
// suspendedTask 是一个被抢占任务的现场。
|
||
type suspendedTask struct {
|
||
Task *Task
|
||
Frame *TaskFrame
|
||
}
|
||
|
||
// nextSelection 标识 nextRef 从哪个集合取出任务。
|
||
type nextSelection int
|
||
|
||
const (
|
||
nextNone nextSelection = iota
|
||
nextReady
|
||
nextInterrupt
|
||
nextImmediate
|
||
nextSuspended
|
||
)
|
||
|
||
func newScheduler(maxQueue int) *scheduler {
|
||
if maxQueue <= 0 {
|
||
maxQueue = 256
|
||
}
|
||
return &scheduler{
|
||
maxQueue: maxQueue,
|
||
maxInterruptFrames: int(LevelCritical), // 结构推论:= 中断级数
|
||
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
|
||
}
|
||
|
||
// allocateIDLocked 分配任务 ID 与入队时刻(调用方持锁)。
|
||
func (s *scheduler) allocateIDLocked(t *Task) {
|
||
s.seq++
|
||
t.ID = s.seq
|
||
if t.EnqueuedAt.IsZero() {
|
||
t.EnqueuedAt = time.Now()
|
||
}
|
||
}
|
||
|
||
// 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.allocateIDLocked(t)
|
||
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 选出下一个任务。优先顺序:
|
||
//
|
||
// 1. immediate —— 刚抢占成功的中断(抢占必须立即生效)
|
||
// 2. 中断队列 L4→L1(同级 FIFO)
|
||
// 3. 中断栈顶(与 2 比级别取高者;栈顶是排队任务时视为最低)
|
||
// 4. 排队队列(FIFO)
|
||
func (s *scheduler) nextRef() (*Task, *TaskFrame, nextSelection) {
|
||
s.mu.Lock()
|
||
defer s.mu.Unlock()
|
||
|
||
if s.immediate != nil {
|
||
t := s.immediate
|
||
s.immediate = nil
|
||
s.running = t
|
||
return t, nil, nextImmediate
|
||
}
|
||
|
||
qTask, qLevel := s.highestInterruptLocked()
|
||
|
||
// 中断栈:只比**栈顶**(严格 LIFO)。栈内不做优先级重排——
|
||
// 嵌套抢占天然使栈自底向上级别递增,且“后被打断的先恢复”才是栈语义。
|
||
if n := len(s.suspendStack); n > 0 {
|
||
top := s.suspendStack[n-1]
|
||
// 栈顶 vs 最高级待处理中断:取高者(持平归栈顶,维持 LIFO 与公平)。
|
||
if qTask == nil || effectiveLevel(top.Task) >= qLevel {
|
||
// 这里只负责“选出”;Resumed 由 resumeTask 计一次(否则会双计,
|
||
// 使“排空后 Suspended == Resumed”这条不变量失真)。
|
||
s.suspendStack = s.suspendStack[:n-1]
|
||
s.running = top.Task
|
||
return top.Task, top.Frame, nextSuspended
|
||
}
|
||
}
|
||
|
||
if qTask != nil {
|
||
s.popInterruptLocked(qLevel)
|
||
s.running = qTask
|
||
return qTask, nil, nextInterrupt
|
||
}
|
||
|
||
if len(s.queue) > 0 {
|
||
t := s.queue[0]
|
||
s.queue = s.queue[1:]
|
||
s.running = t
|
||
return t, nil, nextReady
|
||
}
|
||
return nil, nil, nextNone
|
||
}
|
||
|
||
// highestInterruptLocked 返回当前最高级非空中断队列的队头及其级别。
|
||
func (s *scheduler) highestInterruptLocked() (*Task, Level) {
|
||
for lv := LevelCritical; lv >= LevelBackground; lv-- {
|
||
if q := s.interruptQueues[lv]; len(q) > 0 {
|
||
return q[0], lv
|
||
}
|
||
}
|
||
return nil, 0
|
||
}
|
||
|
||
// popInterruptLocked 弹出某级别中断队列的队头(调用方已确认非空)。
|
||
func (s *scheduler) popInterruptLocked(lv Level) {
|
||
s.interruptQueues[lv] = s.interruptQueues[lv][1:]
|
||
}
|
||
|
||
// interruptCountLocked 统计所有待处理中断(含 immediate 槽)。
|
||
func (s *scheduler) interruptCountLocked() int {
|
||
n := 0
|
||
for lv := LevelBackground; lv <= LevelCritical; lv++ {
|
||
n += len(s.interruptQueues[lv])
|
||
}
|
||
if s.immediate != nil {
|
||
n++
|
||
}
|
||
return n
|
||
}
|
||
|
||
// setImmediateLocked 登记一个应“立即运行”的抢占者。
|
||
//
|
||
// 槽只有一格:若已有抢占者且新的级别更高,旧的降级入队;否则新的入队。
|
||
func (s *scheduler) setImmediateLocked(t *Task) {
|
||
if s.immediate != nil && effectiveLevel(t) <= effectiveLevel(s.immediate) {
|
||
s.enqueueInterruptLocked(t)
|
||
return
|
||
}
|
||
if s.immediate != nil {
|
||
s.enqueueInterruptLocked(s.immediate)
|
||
}
|
||
s.allocateIDLocked(t)
|
||
s.stats.Enqueued++
|
||
s.immediate = t
|
||
}
|
||
|
||
func removeTask(list []*Task, target *Task) []*Task {
|
||
for i, t := range list {
|
||
if t == target {
|
||
return append(list[:i], list[i+1:]...)
|
||
}
|
||
}
|
||
return list
|
||
}
|
||
|
||
// enqueueInterruptLocked 把一个未立即抢占的中断按其级别入队(调用方持锁)。
|
||
//
|
||
// 有界:满了丢**最老**的一条并计数(中断是提示性输入,宁可丢旧保新)。
|
||
func (s *scheduler) enqueueInterruptLocked(t *Task) {
|
||
s.allocateIDLocked(t)
|
||
if s.interruptCountLocked() >= s.maxQueue {
|
||
for lv := LevelBackground; lv <= LevelCritical; lv++ {
|
||
if len(s.interruptQueues[lv]) > 0 {
|
||
s.popInterruptLocked(lv)
|
||
s.stats.Rejected++
|
||
break
|
||
}
|
||
}
|
||
}
|
||
lv := t.Level
|
||
if lv < LevelBackground || lv > LevelCritical {
|
||
lv = DefaultLevel
|
||
}
|
||
s.interruptQueues[lv] = append(s.interruptQueues[lv], t)
|
||
s.stats.Enqueued++
|
||
}
|
||
|
||
// requestPreempt 登记一次中断请求(class=TaskInterrupt)。
|
||
//
|
||
// 返回 true 表示“应该尝试取消运行任务正在进行的可取消步骤(LLM 流式)”。
|
||
//
|
||
// 判据是 canPreempt(由优先级级别系统一承担),并受抢占冷却约束:
|
||
// - running 是排队任务 → 任何中断都抢占;
|
||
// - running 是中断 Li → 仅 Lj > Li 的中断抢占。
|
||
//
|
||
// 能抢占时把中断放进 immediate(立即生效);否则按其级别入队,等当前任务
|
||
// 结束或下一个安全点再处理——无论哪种,中断都不会丢。
|
||
//
|
||
// 临界区(如记忆整理)内不 arm、不取消:中断只入队,等临界区结束后的安全点处理,
|
||
// 这是设计 §4.3 的硬要求——那个位置的“不抢占”不能只是不让位,还必须不取消。
|
||
// level 必须是**已解析好**的中断级别(含特权判定):
|
||
// 生产路径只有 interruptLoop,它用 (*Agent).interruptLevel 得出 level;
|
||
// 内核自身用 requestKernelPreempt(固定 L4)。本函数不再夹取,
|
||
// 否则内核级插件的 L4 会被无辜削掉。
|
||
func (s *scheduler) requestPreempt(evt *agentIO.InputEvent, level Level) bool {
|
||
return s.registerInterrupt(newInterruptTask(evt, level))
|
||
}
|
||
|
||
// requestKernelPreempt 是**内核**中断入口(panic / 内核事件 selfip)。
|
||
//
|
||
// 级别固定 L4,且**不夹取**——这是 L4 的唯一来源,插件永远够不到。
|
||
func (s *scheduler) requestKernelPreempt(evt *agentIO.InputEvent) bool {
|
||
return s.registerInterrupt(newKernelInterruptTask(evt))
|
||
}
|
||
|
||
// registerInterrupt 是登记中断的公共实现(任务已带好 Class/Level)。
|
||
func (s *scheduler) registerInterrupt(t *Task) bool {
|
||
s.mu.Lock()
|
||
running := s.running
|
||
critical := s.critical.Load()
|
||
s.stats.bumpInterruptLevel(&s.stats.InterruptsByLevel, t.Level)
|
||
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)
|
||
}
|
||
}
|
||
if !arm {
|
||
s.enqueueInterruptLocked(t)
|
||
}
|
||
s.mu.Unlock()
|
||
|
||
if !arm {
|
||
s.signalWake()
|
||
}
|
||
return arm
|
||
}
|
||
|
||
// 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 保存现场。
|
||
//
|
||
// 深度上界是**结构推论**(= 中断级数),不是配置项:安全点上的 canSuspend 已提前
|
||
// 拦下超限情况,此处仅在竞态下兜底计数——绝不丢弃帧(帧丢了会丢副作用记录)。
|
||
func (s *scheduler) suspend(t *Task, f *TaskFrame) {
|
||
s.mu.Lock()
|
||
defer s.mu.Unlock()
|
||
if len(s.suspendStack) >= s.maxInterruptFrames {
|
||
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.maxInterruptFrames
|
||
}
|
||
|
||
// done 标记任务执行结束。
|
||
func (s *scheduler) done(t *Task) {
|
||
s.mu.Lock()
|
||
defer s.mu.Unlock()
|
||
if s.running == t {
|
||
s.running = nil
|
||
}
|
||
// 任务正常结束:让位信号不再有意义(中断已在中断队列/immediate 里)。
|
||
s.preemptArmed = false
|
||
s.preemptLevel = 0
|
||
s.stats.Executed++
|
||
}
|
||
|
||
// currentLevel 返回当前正在执行任务的级别;无 running 时为默认级。
|
||
//
|
||
// 用于在 prepare 段把级别写进帧(抢占比较的基准)。
|
||
// interruptLevel 返回一次**中断注入**的级别。
|
||
//
|
||
// 级别是“这项工作有多不能等”,由来源在 InjectOptions.Priority 里声明
|
||
// (排队注入没有级别,它们的 TaskClass 是 TaskQueued)。
|
||
//
|
||
// privileged 表示来源是**内核级插件**(编译期内置插件,见 isKernelLevelSource):
|
||
// - privileged=true → 可用到 L4(实现“立即打断”,如 WebUI 终止按钮)
|
||
// - privileged=false → 夹到 L1..L3;空/非法一律降级为 DefaultLevel(L1)
|
||
//
|
||
// 另有完全绕过本函数的 L4 来源:内核自身的 raiseKernelInterrupt(panic / selfip)。
|
||
func interruptLevel(evt *agentIO.InputEvent, privileged bool) Level {
|
||
if evt == nil || evt.Payload == nil {
|
||
return DefaultLevel
|
||
}
|
||
raw, _ := evt.Payload["priority"].(string)
|
||
l, ok := parseInterruptLevel(raw)
|
||
if !ok {
|
||
return DefaultLevel
|
||
}
|
||
if privileged {
|
||
return l
|
||
}
|
||
return clampPluginLevel(l)
|
||
}
|
||
|
||
// isKernelLevelSource 报告某来源是否是**内核级插件**(编译期内置插件)。
|
||
//
|
||
// 只有它们能声明 L4(见 interruptLevel)。判据是插件注册表里的“内置工厂”,
|
||
// 而不是插件自报的名字本身——外部插件经 proc 桥时已被夹到 L3,这里是第二道闸。
|
||
//
|
||
// source 的约定是 `插件名` 或 `插件名/实例`(如 webui/<deviceID>),故取第一段。
|
||
func (a *Agent) isKernelLevelSource(source string) bool {
|
||
if source == "" {
|
||
return false
|
||
}
|
||
name := source
|
||
if i := strings.IndexByte(name, '/'); i > 0 {
|
||
name = name[:i]
|
||
}
|
||
// ① 编译期内置插件(根 agent 的 L4 来源之一)。
|
||
if a.pluginReg != nil && a.pluginReg.IsBuiltinPlugin(name) {
|
||
return true
|
||
}
|
||
// ② **本 agent 的上级**(驻留子的父)—— 设计 §6.1 的 L4 通则:
|
||
// 子的阶梯上只有父能产生 L4,所以父的"发送消息"一定能打断子。
|
||
return a.kernelSource != "" && name == a.kernelSource
|
||
}
|
||
|
||
// parseInterruptLevel 解析插件声明的级别字符串("L1".."L3")。
|
||
// 只认字面量:拼写错误必须降级成默认级而不是被静默当成别的级别。
|
||
func parseInterruptLevel(s string) (Level, bool) {
|
||
switch s {
|
||
case "L1", "l1":
|
||
return LevelBackground, true
|
||
case "L2", "l2":
|
||
return LevelMessage, true
|
||
case "L3", "l3":
|
||
return LevelInteractive, true
|
||
case "L4", "l4":
|
||
// 内核级:解析出来但会被 clamp 夹到 L3。
|
||
return LevelCritical, true
|
||
default:
|
||
return 0, false
|
||
}
|
||
}
|
||
|
||
// newInputTask 把一个**排队输入**包装成任务(无级别)。
|
||
func newInputTask(evt *agentIO.InputEvent) *Task {
|
||
return &Task{Class: TaskQueued, Kind: TaskKindInput, Event: evt, EnqueuedAt: time.Now()}
|
||
}
|
||
|
||
// newInterruptTask 把一个中断请求包装成任务(带级别)。
|
||
func newInterruptTask(evt *agentIO.InputEvent, level Level) *Task {
|
||
return &Task{Class: TaskInterrupt, Kind: TaskKindInput, Level: level, Event: evt, EnqueuedAt: time.Now()}
|
||
}
|
||
|
||
// newSelfTask 包装内核自循环输入——它是**排队任务**:记忆整理/子任务通知
|
||
// 不需要及时处理,可被任何中断打断。
|
||
func newSelfTask(msg selfInputMsg) *Task {
|
||
return &Task{Class: TaskQueued, Kind: TaskKindSelf, Self: msg, EnqueuedAt: time.Now()}
|
||
}
|
||
|
||
// newKernelInterruptTask 构造一个**内核级中断**(L4)。
|
||
//
|
||
// 这是 L4 的唯一来源:panic 中断、内核事件中断(selfip)。
|
||
// 插件永远拿不到这个入口——它不经 InjectOptions,也不经 proc 桥。
|
||
func newKernelInterruptTask(evt *agentIO.InputEvent) *Task {
|
||
return &Task{Class: TaskInterrupt, Kind: TaskKindInput, Level: LevelCritical, Event: evt, 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.Immediate = a.sched.immediate
|
||
for lv := LevelBackground; lv <= LevelCritical; lv++ {
|
||
snap.InterruptQueues[lv] = append(snap.InterruptQueues[lv], a.sched.interruptQueues[lv]...)
|
||
snap.PendingInterrupts = append(snap.PendingInterrupts, a.sched.interruptQueues[lv]...)
|
||
}
|
||
if a.sched.immediate != nil {
|
||
snap.PendingInterrupts = append(snap.PendingInterrupts, a.sched.immediate)
|
||
}
|
||
snap.SuspendStack = append(snap.SuspendStack, a.sched.suspendStack...)
|
||
snap.MaxInterruptFrames = a.sched.maxInterruptFrames
|
||
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。
|
||
// raiseKernelInterrupt 是 **L4 的唯一入口**:panic 中断与内核事件中断(selfip)。
|
||
//
|
||
// 它不经 io.InputChan(那是外部/插件输入),而是直接向调度器登记一条内核中断:
|
||
// 级别固定 L4、不夹取、不受插件声明影响。这正是“L4 只有内核持有”的落点。
|
||
//
|
||
// 能否抢占由调度器按统一判据决定;若会抢占,则顺手取消可取消的 LLM 流式步骤
|
||
// (与 interceptLoop 对插件中断的处理完全一致)。
|
||
func (a *Agent) raiseKernelInterrupt(source, channel, text string) {
|
||
if a.sched == nil {
|
||
return
|
||
}
|
||
evt := &agentIO.InputEvent{
|
||
Source: source,
|
||
Type: "interrupt",
|
||
OutputChannel: channel,
|
||
Payload: map[string]interface{}{
|
||
"content": text,
|
||
"interrupt": true,
|
||
"interrupt_source": source,
|
||
"interrupt_channel": channel,
|
||
"kernel": true,
|
||
},
|
||
}
|
||
if a.sched.requestKernelPreempt(evt) {
|
||
a.cancelCurrentLLM()
|
||
}
|
||
}
|
||
|
||
// reportTaskPanic 把一个任务 panic 报告成内核 L4 中断。
|
||
//
|
||
// 递归保护是**结构性**的:若 panic 的任务本身就是 L4 内核中断,则不再产生新的
|
||
// L4——否则同一个 panic 会自我放大成中断风暴,与“内核事件”应有的语义相反。
|
||
func (a *Agent) reportTaskPanic(t *Task, r interface{}) {
|
||
if t.Class == TaskInterrupt && t.Level >= LevelCritical {
|
||
return
|
||
}
|
||
a.raiseKernelInterrupt("kernel", "kernel",
|
||
fmt.Sprintf("内核事件:任务 #%d 发生 panic:%v(该任务已被丢弃,调度器存活)", t.ID, r))
|
||
}
|
||
|
||
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())
|
||
a.reportTaskPanic(t, r)
|
||
}
|
||
}()
|
||
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.reportTaskPanic(t, r)
|
||
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)
|
||
}
|