Files
HomeAgent/internal/agent/core/offload.go
JianFeeeee 69446a2649 feat(scheduler): 主 agent 忙时把积压任务自动转投给驻留子
问题(2026-09-19 线上实测):主 agent 被长任务占住时(现场:12 分 8 秒、69 次
工具调用),后来到达的消息全部以 level insufficient 排进中断队列干等 —— 同级
中断不能抢占同级运行任务(canPreempt),只能等前一个跑完。而内核本有驻留子
(独立 agent + 独立调度器)可并行干活。

行为(用户 2026-09-19 明确要求):
- 触发:运行任务持续 > offload_busy_after(5m) 且积压 >= offload_min_pending(3)
- 拉起/复用「转投专用」驻留子,把积压的纯排队输入转投过去
- 在原队列位置留下说明「[系统] N 条积压任务已转投给驻留子 agent X 处理…」

通道配置(按用户口径,与人工创建的子刻意不同):
- 不配 inputch(内核的干活 agent,不接收插件用户输入)
- 持有全部输出通道(结果要能发回 qq/webui 等正确通道)

三个设计要点(都是实测撞出来的,写进代码注释与设计文档 §7.1):
1. 检查必须在**独立 goroutine**:schedulerLoop 同步执行任务,放它里面在
   「正忙」期间根本回不到循环顶部 ⇒ 永不触发(我第一版就写错了,测试才发现)。
2. 只转投 TaskQueued 纯排队输入:中断任务带级别语义、self 任务与父的记忆面绑定。
3. 转投失败/关闭时必须把任务**放回队列前端**:吞一条输入比多处理一条更糟。

这是设计 §7「决策在父的模型手里」的**刻意例外**(父正忙、物理上无法决策,
而积压任务本来就是空的),已在文档中显式记录,且默认关闭、由部署方显式打开。

测试 11 条:只取排队输入 / 不足量不取 / 放回不丢任务 / 说明自解释 / 默认关闭 /
空闲不触发 / 端到端转投 / 上限不增殖 / 独立 goroutine 确实会触发。
2026-09-19 16:59:47 +08:00

334 lines
12 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

package core
// 积压任务的**自动转投**:主 agent 长时间忙时,把排队中的任务改投给内核拉起的
// 驻留子,并在原队列位置留一条"已转投"提示。
//
// 为什么要这个(用户 2026-09-19 提出的实际需求):
// 实测主 agent 被一条长任务占住时(当天现场:12 分 8 秒、69 次工具调用),
// 后来的 QQ 消息全部以 "level insufficient" 排进中断队列干等 —— 同级中断
// 不能抢占同级运行任务(scheduler.canPreempt),只能等前一个跑完。
// 而内核明明有驻留子(独立 agent + 独立调度器 + 共享输出通道视图)可以并行干活。
//
// 与设计文档 §7 的关系(**这是刻意的例外,必须显式记录**):
// 设计原文写「创建/销毁/回收/查看/发送是父可调用的原语;**决策在父的模型手里**——
// 内核不替父决定」。本特性让**内核**主动创建并使用驻留子,属于对该原则的例外。
// 之所以可接受:父此刻正忙(无法做决策),而积压任务**本来就是空的**——
// 转投只是把"排队干等"换成"有人在做",不改变任何已提交决策的语义。
// 若不做例外,这个能力就只能由父的模型发起,而它恰恰是忙不过来的那个。
import (
"fmt"
"log"
"strings"
"time"
"runtime/debug"
agentIO "gitcode.com/JianFeeeee/HomeAgent/internal/agent/io"
)
// OffloadOptions 是自动转投的判定与执行参数。
type OffloadOptions struct {
// BusyAfter:运行任务已持续多久算"长时间工作"(0 = 用默认)。
BusyAfter time.Duration
// MinPending:至少要积压多少条才值得拉起驻留子(0 = 用默认)。
MinPending int
// MaxResidents:为转投而拉起的驻留子上限(0 = 用默认)。
MaxResidents int
// Enabled 为 false 时完全关闭(默认关:见 DefaultOffloadOptions 的说明)。
Enabled bool
}
// 默认参数。
//
// 为什么默认**关闭**:自动拉起是"内核替父做决策",改变的是系统行为而非修 bug;
// 且它会让日志/账单里凭空多出一个 agent 在干活。默认关闭、由部署方显式打开,
// 与「显式才是特权」(scheduler.DefaultLevel 的同一条理由)一致。
const (
defaultOffloadBusyAfter = 5 * time.Minute
defaultOffloadMinPending = 3
defaultOffloadMaxResident = 2
)
// DefaultOffloadOptions 返回默认参数(Enabled=false)。
func DefaultOffloadOptions() OffloadOptions {
return OffloadOptions{
BusyAfter: defaultOffloadBusyAfter,
MinPending: defaultOffloadMinPending,
MaxResidents: defaultOffloadMaxResident,
Enabled: false,
}
}
func (o OffloadOptions) normalized() OffloadOptions {
if o.BusyAfter <= 0 {
o.BusyAfter = defaultOffloadBusyAfter
}
if o.MinPending <= 0 {
o.MinPending = defaultOffloadMinPending
}
if o.MaxResidents <= 0 {
o.MaxResidents = defaultOffloadMaxResident
}
return o
}
// offloadNotice 是替换被转投任务的那条提示的正文。
//
// 它必须**自己说清是系统做的**:用户看到队列里出现一条没人发过的消息时,
// 唯一能解释这件事的就是这句话本身。
func offloadNotice(count int, residentID string) string {
return fmt.Sprintf(
"[系统] %d 条积压任务已转投给驻留子 agent %s 处理(主 agent 正忙于长任务,"+
"内核为它们拉起了独立 agent 并行执行)。它们的回复会由 %s 直接发到对应通道;"+
"本提示仅用于说明「那几条消息不会再由你处理」,无需为它们采取任何行动。",
count, residentID, residentID)
}
// offloadCandidate 是一条可被转投的排队任务。
//
// 只有**纯排队输入**(TaskQueued + KindInput)可转投:
// - 中断任务带级别语义(可能正在等待抢占时机),转投会打乱中断阶梯;
// - self 任务是内核内部记账(记忆整理等),与父的记忆面绑定,不能换 agent。
type offloadCandidate struct {
Event *agentIO.InputEvent
}
// residentInputChannel 返回"父给某个子投递输入"用的 inputch 名。
//
// 与 SendToResident 用的是同一个(sub/<id>):子是**不配插件 inputch** 的
// (opts.InputChs 为空),它的入站口就只有父给它的这一条,因此必须与
// SendToResident 保持一致,否则转投的消息会落到一个父不知道的通道名上。
func residentInputChannel(id string) string { return "sub/" + id }
// offloadPendingTasks 检查是否需要转投,需要则拉起/复用一个驻留子并搬运任务。
//
// 返回实际转投的任务条数(0 = 未触发/未转投)。
//
// ❗并发前提:本函数会被 offloadLoop 在**任务执行期间**调用(那正是它的意义),
// 因此它与 schedulerLoop 是并发跑的。所有对队列的读写都经 scheduler 的锁,
// 而"取走哪些任务"与"放回什么"都在同一次锁内完成,不存在丢任务的窗口。
func (a *Agent) offloadPendingTasks(opts OffloadOptions) int {
opts = opts.normalized()
if !opts.Enabled {
return 0
}
// ① 判定:运行任务是否已忙够久。运行任务为空说明压根不忙,不做。
running := a.sched.runningTask()
if running == nil {
return 0
}
busyFor := a.sched.runningFor()
if busyFor < opts.BusyAfter {
return 0
}
// ② 收集可转投的排队任务;不够量就不值得拉起一个 agent。
cands := a.sched.takeQueuedInputs(opts.MinPending)
if len(cands) == 0 {
return 0
}
// ③ 找或拉起一个"转投专用"驻留子。
residentID, err := a.ensureOffloadResident(opts)
if err != nil {
// 拉不起来就把任务**放回队列**,绝不能丢:丢一条输入比多处理一条更糟
// (与 routeInputByOwner 的兜底同一条理由)。
a.sched.requeueFront(cands)
log.Printf("[offload] 无法为 %d 条积压任务准备驻留子,已放回队列: %v", len(cands), err)
return 0
}
// ④ 搬运:逐条投进子的 inputch。
//
// 注意这里**逐条转发原文**而不是打包成一条:任务本身带 Source/OutputChannel
// 等路由信息,打包会让子无法把回复发回正确的通道(qq 私聊 vs 群聊不同)。
var moved int
for _, c := range cands {
if c.Event == nil {
continue
}
if err := a.forwardInputToResident(residentID, c.Event); err != nil {
// 某条投不进去:放回原队列,其余继续(部分成功好过全部回滚)。
a.sched.requeueFront([]offloadCandidate{c})
log.Printf("[offload] 转发任务给驻留子 %s 失败,已放回队列: %v", residentID, err)
continue
}
moved++
}
if moved == 0 {
return 0
}
// ⑤ 在**原队列位置**留下提示(用户要求的那条说明)。
//
// 为什么留在队列里而不是只记日志:队列顺序就是主 agent 接下来要处理的事;
// 用户看会话记录时,需要在这里就看到"那几条去哪儿了",而不是去翻内核日志。
a.sched.requeueFront([]offloadCandidate{
{Event: a.syntheticEvent(offloadNotice(moved, residentID))},
})
log.Printf("[offload] 主 agent 已忙 %s,把 %d 条积压任务转投给驻留子 %s(队列留 1 条说明)",
busyFor.Truncate(time.Second), moved, residentID)
return moved
}
// forwardInputToResident 把一条输入原文投给指定驻留子的 inputch。
//
// 走 InjectInputTo(排队输入,非中断):转投的是"待办工作",不是"打断子"。
// 子的 io 上有 inputRouter(routeInputByOwner),但投递目标是**它自己的** inputch
// 且 Owner 就是它,因此不会被再次路由走。
func (a *Agent) forwardInputToResident(residentID string, evt *agentIO.InputEvent) error {
a.residentMu.Lock()
rc := a.residents[residentID]
a.residentMu.Unlock()
if rc == nil || rc.agent == nil || rc.agent.io == nil {
return fmt.Errorf("驻留子 %s 不存在或不可用", residentID)
}
payload := map[string]interface{}{}
for k, v := range evt.Payload {
payload[k] = v
}
// 带上来源线索,让子知道这条不是父当前任务的续接,而是转投的独立请求。
payload["offloaded_from"] = string(a.id)
payload["offloaded_at"] = time.Now().Format(time.RFC3339)
ch := residentInputChannel(residentID)
rc.agent.io.InjectInputToOpts(evt.Source, ch, evt.Type, payload, agentIO.InjectOptions{})
return nil
}
// ensureOffloadResident 返回一个可用于转投的驻留子 id,必要时拉起一个新的。
//
// 复用规则:优先复用"内核为转投而建"且仍 running、还没满的驻留子;
// 都不可用时(在 MaxResidents 内)新建一个。
func (a *Agent) ensureOffloadResident(opts OffloadOptions) (string, error) {
a.residentMu.Lock()
var reusable []string
for id, rc := range a.residents {
if rc == nil || !rc.offloadOwned {
continue
}
rc.mu.Lock()
state := rc.state
rc.mu.Unlock()
if state == "running" {
reusable = append(reusable, id)
}
}
a.residentMu.Unlock()
// 复用:按 id 稳定排序后取第一个,避免每次挑到不同的子(可预测性)。
if len(reusable) > 0 {
sortStrings(reusable)
return reusable[0], nil
}
// 计数:只为转投而建的子是否已达上限(人工建的子不计入)。
a.residentMu.Lock()
owned := 0
for _, rc := range a.residents {
if rc != nil && rc.offloadOwned {
owned++
}
}
a.residentMu.Unlock()
if owned >= opts.MaxResidents {
return "", fmt.Errorf("转投专用驻留子已达上限 %d", opts.MaxResidents)
}
id := fmt.Sprintf("offload-%d", time.Now().Unix())
if a.dataDir == "" {
return "", fmt.Errorf("未配置 DataDir,无法为驻留子分配 temp 图库路径")
}
info, err := a.SpawnResident(ResidentOptions{
ID: id,
// 不配 inputch:它是内核的**干活** agent,不接收任何插件的用户输入
// (用户要求"不配输入通道")。它只由父经 sub/<id> 投喂任务。
InputChs: nil,
// 全部输出通道:它要能把结果发回 qq/webui 等正确通道
// (用户要求"持有全部输出通道")。nil = 完整授权。
AllowedOutputs: nil,
TempPath: a.residentTempPath(id),
OffloadOwned: true,
})
if err != nil {
return "", err
}
return info.ID, nil
}
// residentTempPath 计算某个驻留子的 temp 图记忆路径(与既有约定一致)。
func (a *Agent) residentTempPath(id string) string {
return strings.TrimRight(a.dataDir, "/") + "/residents/" + id + "/graph.db"
}
// syntheticEvent 造一条"内核自己发的"输入事件(用于队列里的转投说明)。
//
// Source 取 kernel:这条消息不是任何用户发来的,日志与用户界面里都应看得出。
// 不带 ResponseCh:没有同步调用方在等它(它只是给主 agent 看的一句说明)。
func (a *Agent) syntheticEvent(text string) *agentIO.InputEvent {
return &agentIO.InputEvent{
Source: "kernel",
Type: "text",
OutputChannel: "kernel",
Payload: map[string]interface{}{"content": text},
}
}
// sortStrings 是一个不引入 sort 依赖的小排序(候选集极小,插入排序足够)。
func sortStrings(s []string) {
for i := 1; i < len(s); i++ {
for j := i; j > 0 && s[j] < s[j-1]; j-- {
s[j], s[j-1] = s[j-1], s[j]
}
}
}
// offloadLoop 周期性检查「主 agent 是否被长任务占住 + 是否有积压」。
//
// 为什么必须是**独立 goroutine**而不是 schedulerLoop 里的一步:
//
// schedulerLoop 是**同步执行**任务的(executeNewTask 会一直阻塞到任务结束),
// 所以「正忙」期间它根本不会回到循环顶部 —— 把检查放在那里等于永不触发。
// 这正是本特性存在的理由(主 agent 忙时无人处理积压),不能在实现上重犯。
//
// 检查间隔取 BusyAfter 的 1/5(不低于 1 秒):保证在跨过阈值后能在合理时间内
// 触发,又不至于空转打日志。
func (a *Agent) offloadLoop() {
defer func() {
if r := recover(); r != nil {
log.Printf("[agent] offloadLoop panic recovered: %v\n%s", r, debug.Stack())
time.Sleep(time.Second)
go a.offloadLoop()
}
}()
if !a.offload.Enabled {
return // 未启用:不占 goroutine,也不打日志(默认关闭是常态)
}
// 只让**根 agent** 做转投:驻留子自己也可能忙,但让子再去拉孙子会形成
// 无界增殖(每层都能拉 MAX 个),而积压的源头是根那一条调度链。
if a.parentID != "" {
return
}
interval := a.offload.BusyAfter / 5
if interval < time.Second {
interval = time.Second
}
ticker := time.NewTicker(interval)
defer ticker.Stop()
for {
select {
case <-ticker.C:
a.offloadPendingTasks(a.offload)
case <-a.ctx.Done():
return
}
}
}