From 69446a2649af2cd5b8c63918f49b101e65e6033e Mon Sep 17 00:00:00 2001 From: JianFeeeee Date: Sat, 19 Sep 2026 16:59:47 +0800 Subject: [PATCH] =?UTF-8?q?feat(scheduler):=20=E4=B8=BB=20agent=20?= =?UTF-8?q?=E5=BF=99=E6=97=B6=E6=8A=8A=E7=A7=AF=E5=8E=8B=E4=BB=BB=E5=8A=A1?= =?UTF-8?q?=E8=87=AA=E5=8A=A8=E8=BD=AC=E6=8A=95=E7=BB=99=E9=A9=BB=E7=95=99?= =?UTF-8?q?=E5=AD=90?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 问题(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 确实会触发。 --- cmd/homed/bootstrap.go | 6 + docs/zh/resident-subagent-design.md | 34 +++ internal/agent/core/agent.go | 10 + internal/agent/core/offload.go | 333 ++++++++++++++++++++++++++++ internal/agent/core/offload_test.go | 301 +++++++++++++++++++++++++ internal/agent/core/resident.go | 11 +- internal/agent/core/scheduler.go | 95 ++++++++ internal/config/registry.go | 10 + 8 files changed, 799 insertions(+), 1 deletion(-) create mode 100644 internal/agent/core/offload.go create mode 100644 internal/agent/core/offload_test.go diff --git a/cmd/homed/bootstrap.go b/cmd/homed/bootstrap.go index afb9b0f..9dc8a1f 100644 --- a/cmd/homed/bootstrap.go +++ b/cmd/homed/bootstrap.go @@ -487,6 +487,12 @@ func newMainAgent(cfg *types.Config, cfgReg *internalConfig.ConfigRegistry, prov ReviewInterval: cfgReg.GetDuration("core.agent.review_interval", 120*time.Minute), MergeInterval: cfgReg.GetDuration("core.agent.merge_interval", 120*time.Minute), MaxToolTurns: cfgReg.GetInt("core.agent.max_tool_turns", 10), + Offload: agentCore.OffloadOptions{ + Enabled: cfgReg.GetBool("core.agent.offload_enabled", false), + BusyAfter: cfgReg.GetDuration("core.agent.offload_busy_after", 5*time.Minute), + MinPending: cfgReg.GetInt("core.agent.offload_min_pending", 3), + MaxResidents: cfgReg.GetInt("core.agent.offload_max_residents", 2), + }, ContextSavePath: filepath.Join(cfg.Daemon.DataDir, "memory", "context.json"), EmbeddingModelPath: cfgReg.GetString("core.agent.embedding_model_path", ""), Embedder: embedder, diff --git a/docs/zh/resident-subagent-design.md b/docs/zh/resident-subagent-design.md index 168fd6e..8da0e3b 100644 --- a/docs/zh/resident-subagent-design.md +++ b/docs/zh/resident-subagent-design.md @@ -340,6 +340,40 @@ - **创建/销毁/回收/查看/发送**是**父可调用的原语(工具)**;**决策**(压还是收、收哪些) 在父的模型手里 —— 内核不替父决定。 + +### 7.1 例外:积压任务自动转投(内核主动拉起)[已定 · 唯一例外] + +**背景(2026-09-19 线上实测)**:主 agent 被一条长任务占住时(当天现场:12 分 8 秒、 +69 次工具调用),后来的 QQ 消息全部以 `level insufficient` 排进中断队列干等 —— +同级中断不能抢占同级运行任务(`scheduler.canPreempt`),只能等前一个跑完。 +而内核明明有驻留子(独立 agent + 独立调度器 + 共享输出通道视图)可以并行干活。 + +**行为**:当运行任务已持续超过 `core.agent.offload_busy_after`(默认 5m) +**且**排队输入积到 `offload_min_pending`(默认 3)条时,内核: +1. 拉起(或在 `offload_max_residents` 内复用一个)**转投专用驻留子**; +2. 把积压的**纯排队输入**转投给它; +3. 在原队列位置留下一条说明:`[系统] N 条积压任务已转投给驻留子 agent X 处理…`。 + +**转投驻留子的通道配置(刻意与人工创建的子不同)**: +- **不配 inputch**:它是内核的干活 agent,不接收任何插件的用户输入; +- **持有全部输出通道**(`AllowedOutputs` 为空 = 完整授权):它必须能把结果发回 + qq/webui 等正确通道(否则干活结果无处可去)。 + +**为什么这是对上述原则的例外,且可接受**:父此刻正忙(物理上无法做决策), +而积压任务**本来就是空的** —— 转投只是把「排队干等」换成「有人在做」, +不改变任何已提交决策的语义。若不做例外,这个能力就只能由父的模型发起, +而它恰恰是忙不过来的那个。 + +**默认关闭**(`core.agent.offload_enabled=false`):它改变的是系统行为而非修 bug, +按「显式才是特权」(与 `scheduler.DefaultLevel` 同一条理由)由部署方打开。 + +**不做的事(边界)**: +- 只转投 `TaskQueued` 纯排队输入。中断任务带级别语义(转投会打乱中断阶梯)、 + self 任务是内核内部记账(与父的记忆面绑定)—— 两者都不动。 +- 只对**根 agent** 生效:子再去拉孙子会形成无界增殖,而积压的源头是根那条链。 +- 实现上检查跑在**独立 goroutine**:`schedulerLoop` 是同步执行的, + 放在那里在「正忙」期间根本不会回到循环顶部(等于永不触发)。 + - **默认完整授权**[已定]:子默认拿到全部插件与工具(含输出门); 父可在创建时**收窄**(收窄工具子集、收窄可用输出通道集合)。 ⚠️ 默认含输出门意味着**子可以直接对用户通道发消息**;若要默认收窄,改一处默认即可。 diff --git a/internal/agent/core/agent.go b/internal/agent/core/agent.go index a91db44..f77f47a 100644 --- a/internal/agent/core/agent.go +++ b/internal/agent/core/agent.go @@ -142,6 +142,8 @@ type Agent struct { // 工具轮次硬上限(0 = 不限);见 AgentConfig.MaxToolTurns。 maxToolTurns int + // offload 是积压任务自动转投的参数(见 offload.go)。 + offload OffloadOptions // 进行中的 LLM 请求取消函数,interceptLoop 可调用以在请求中打断 cancelLLM context.CancelFunc @@ -279,6 +281,12 @@ type AgentConfig struct { // MaxToolTurns 是单个任务允许的工具轮次上限(0 = 不限)。 // 设计文档 D6:主循环必须有硬上限,否则模型不停调用就永不完结。 MaxToolTurns int + + // Offload 是「积压任务自动转投给驻留子」的参数(见 offload.go)。 + // + // 默认关闭(Enabled=false):它让**内核替父做决策**,是设计 §7 + //「决策在父的模型手里」的刻意例外,因此必须由部署方显式打开。 + Offload OffloadOptions } func New(cfg AgentConfig) *Agent { @@ -369,6 +377,7 @@ func New(cfg AgentConfig) *Agent { childTasks: make(map[string]*childTaskState), sched: newScheduler(256), maxToolTurns: cfg.MaxToolTurns, + offload: cfg.Offload, pluginHealth: newPluginHealthTracker(), thinkingEnabled: cfg.ThinkingEnabled, inputCfg: cfg.InputProcessing, @@ -407,6 +416,7 @@ func (a *Agent) Start() { go a.archiveLoop() go a.mergeLoop() go a.reviewLoop() + go a.offloadLoop() a.subscribeTerminalRegistry() a.reembedStaleMedia() a.migrateLegacyGraphMedia() diff --git a/internal/agent/core/offload.go b/internal/agent/core/offload.go new file mode 100644 index 0000000..96523c5 --- /dev/null +++ b/internal/agent/core/offload.go @@ -0,0 +1,333 @@ +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/):子是**不配插件 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/ 投喂任务。 + 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 + } + } +} diff --git a/internal/agent/core/offload_test.go b/internal/agent/core/offload_test.go new file mode 100644 index 0000000..02f64ab --- /dev/null +++ b/internal/agent/core/offload_test.go @@ -0,0 +1,301 @@ +package core + +import ( + "path/filepath" + "strings" + "testing" + "time" + + agentAPI "gitcode.com/JianFeeeee/HomeAgent/internal/agent/api" + agentIO "gitcode.com/JianFeeeee/HomeAgent/internal/agent/io" + "gitcode.com/JianFeeeee/HomeAgent/internal/memory" +) + +// newRootWithoutSchedulerLoop 造一个**不启动后台循环**的根 agent。 +// +// 为什么测试必须用它:newRootWith 会 a.Start(),于是真实的 schedulerLoop +// 与测试**并发**跑,它会瞬间把测试排进队列的任务执行掉并清空 running +// ⇒ "主 agent 正忙"这个前提会被后台循环消掉,转投判定随机失效 +// (实测:同一测试两次运行结果不同,一个过一个不过)。 +// 本特性测的是**判定 + 搬运**这两步的语义,不需要真的把任务跑起来。 +func newRootWithoutSchedulerLoop(t *testing.T) (*Agent, *memory.GraphDB) { + t.Helper() + dir := t.TempDir() + main, err := memory.NewGraphDB(filepath.Join(dir, "main.db")) + if err != nil { + t.Fatal(err) + } + a := New(AgentConfig{ + ID: "parent", + Provider: &countingProvider{}, + ProviderManager: agentAPI.NewProviderManager(), + IO: agentIO.NewIOManager(), + StageHost: NewStageHost(), + Memory: main, + DataDir: dir, + }) + t.Cleanup(func() { a.Stop(); main.Close() }) + return a, main +} + +// makeQueuedInput 造一条排队输入任务(Event 非空,Class=TaskQueued)。 +func makeQueuedInput(id int) *Task { + return newInputTask(&agentIO.InputEvent{ + RequestID: "req", + Source: "qq", + Type: "text", + OutputChannel: "qq", + Payload: map[string]interface{}{"content": "hello"}, + }) +} + +// 转投只应该动**纯排队输入**:中断任务带级别语义、self 任务是内核内部记账, +// 搬走它们会分别破坏中断阶梯与记忆整理。 +func TestTakeQueuedInputsOnlyTakesQueuedInputs(t *testing.T) { + s := newScheduler(64) + // 混合:1 条排队输入 + 1 条中断 + 1 条 self + 3 条排队输入 + s.enqueue(makeQueuedInput(1)) + s.enqueue(newInterruptTask(&agentIO.InputEvent{Source: "qq", OutputChannel: "qq"}, LevelMessage)) + s.enqueue(newSelfTask(selfInputMsg{text: "distill", channel: "cli"})) + s.enqueue(makeQueuedInput(2)) + s.enqueue(makeQueuedInput(3)) + s.enqueue(makeQueuedInput(4)) + + if len(s.queue) != 6 { + t.Fatalf("就绪队列应有 6 条(4 排队输入 + 1 中断 + 1 self),实际 %d", len(s.queue)) + } + got := s.takeQueuedInputs(3) + if len(got) != 3 { + t.Fatalf("应取走 3 条排队输入,实际 %d", len(got)) + } + for _, c := range got { + if c.Event == nil { + t.Fatal("取出的候选不得为空事件") + } + } + // self 与中断必须还在 + var hasSelf, hasInterrupt bool + for _, tt := range s.queue { + if tt.Kind == TaskKindSelf { + hasSelf = true + } + if tt.Class == TaskInterrupt { + hasInterrupt = true + } + } + if !hasSelf { + t.Error("self 任务被误取(会破坏记忆整理)") + } + if !hasInterrupt { + t.Error("中断任务被误取(会破坏中断阶梯)") + } +} + +// 不够量时**一条都不取**:拉起一个 agent 的成本不该为一条任务付。 +// 这条保证「要么不动、要么成批移动」。 +func TestTakeQueuedInputsIsAllOrNothing(t *testing.T) { + s := newScheduler(64) + s.enqueue(makeQueuedInput(1)) + s.enqueue(makeQueuedInput(2)) + + if got := s.takeQueuedInputs(3); got != nil { + t.Fatalf("不足 3 条时不应取走任何任务,实际取走 %d", len(got)) + } + if len(s.queue) != 2 { + t.Errorf("队列不应被改动,实际剩 %d", len(s.queue)) + } +} + +// ★ 安全不变量:转投失败必须把任务**放回队列**。 +// 吞掉一条输入比多处理一条更糟——用户会看到"消息发出去了却没人理"。 +func TestRequeueFrontKeepsAllTasks(t *testing.T) { + s := newScheduler(64) + s.enqueue(makeQueuedInput(1)) + s.enqueue(makeQueuedInput(2)) + s.enqueue(makeQueuedInput(3)) + + taken := s.takeQueuedInputs(3) + if len(taken) != 3 { + t.Fatalf("应取走 3 条,实际 %d", len(taken)) + } + if len(s.queue) != 0 { + t.Fatalf("取走后队列应空,实际 %d", len(s.queue)) + } + + s.requeueFront(taken) + if len(s.queue) != 3 { + t.Fatalf("★ 放回后必须一条不少:期望 3,实际 %d", len(s.queue)) + } + // 放回的是**前端**:它们比队列里原有的一切都早 + s.enqueue(makeQueuedInput(4)) + if s.queue[len(s.queue)-1].Event.Payload["content"] != "hello" { + t.Error("放回的任务应在队列前端") + } +} + +// 转投说明必须自己说清是系统做的:用户看到队列里出现一条没人发过的消息时, +// 唯一能解释这件事的就是这句话本身。 +func TestOffloadNoticeExplainsItself(t *testing.T) { + msg := offloadNotice(3, "offload-123") + for _, want := range []string{"系统", "3 条", "offload-123", "转投"} { + if !strings.Contains(msg, want) { + t.Errorf("说明缺少 %q:%s", want, msg) + } + } +} + +// 默认必须是**关闭**:自动拉起是内核替父做决策(设计 §7 的例外), +// 不能默默改变系统行为。 +func TestOffloadDisabledByDefault(t *testing.T) { + opts := DefaultOffloadOptions() + if opts.Enabled { + t.Error("默认必须关闭") + } + s := newScheduler(64) + s.enqueue(makeQueuedInput(1)) + s.enqueue(makeQueuedInput(2)) + s.enqueue(makeQueuedInput(3)) + a := &Agent{sched: s} + if n := a.offloadPendingTasks(opts); n != 0 { + t.Errorf("关闭时不得转投,实际转了 %d", n) + } + if len(s.queue) != 3 { + t.Errorf("关闭时队列不得被改动,实际 %d", len(s.queue)) + } +} + +// 不忙(无运行任务)时不转投:没有"长任务占住"这个前提,排队就是正常的。 +func TestOffloadSkippedWhenIdle(t *testing.T) { + opts := DefaultOffloadOptions() + opts.Enabled = true + opts.BusyAfter = time.Nanosecond + opts.MinPending = 1 + + s := newScheduler(64) + s.enqueue(makeQueuedInput(1)) + a := &Agent{sched: s} + if n := a.offloadPendingTasks(opts); n != 0 { + t.Errorf("空闲时不应转投,实际 %d", n) + } +} + +// ★ 端到端:主 agent 忙时,积压任务应真的被搬到驻留子,且队列里留下说明。 +// 这是本特性的核心行为 —— 只测"判定函数返回 0/非 0"不够, +// 必须证明任务**换了 agent 且原队列留下了可读的交代**。 +func TestOffloadMovesTasksToResidentEndToEnd(t *testing.T) { + root, main := newRootWithoutSchedulerLoop(t) + defer main.Close() + + opts := DefaultOffloadOptions() + opts.Enabled = true + opts.BusyAfter = time.Nanosecond // 立即算"忙" + opts.MinPending = 2 + opts.MaxResidents = 1 + + // 伪造"正在跑一条长任务":转投判定要求 running 非空。 + root.sched.enqueue(makeQueuedInput(1)) + root.sched.nextRef() // 把它变成 running + + // 再排 2 条积压 + root.sched.enqueue(makeQueuedInput(2)) + root.sched.enqueue(makeQueuedInput(3)) + + moved := root.offloadPendingTasks(opts) + if moved != 2 { + t.Fatalf("应转投 2 条,实际 %d", moved) + } + + // ① 确实拉起了一个驻留子,且标记为"为转投而建" + list := root.Residents() + if len(list) != 1 { + t.Fatalf("应拉起 1 个驻留子,实际 %d", len(list)) + } + resident := list[0] + if !strings.HasPrefix(resident.ID, "offload-") { + t.Errorf("驻留子应为转投专用命名,实际 %s", resident.ID) + } + // ② 它不配任何插件 inputch(用户要求),但持有全部输出通道(nil=全授权) + if len(resident.InputChs) != 0 { + t.Errorf("转投驻留子不应配 inputch,实际 %v", resident.InputChs) + } + if len(resident.AllowedOutputs) != 0 { + t.Errorf("转投驻留子应持有全部输出通道(空=全授权),实际 %v", resident.AllowedOutputs) + } + t.Logf("驻留子 %s: inputch=%v outputs=%v", resident.ID, resident.InputChs, resident.AllowedOutputs) + + // ③ 原队列里留下说明(且说明是内核发的) + if len(root.sched.queue) != 1 { + t.Fatalf("原队列应只剩 1 条说明,实际 %d", len(root.sched.queue)) + } + notice := root.sched.queue[0] + if notice.Event == nil || notice.Event.Source != "kernel" { + t.Fatalf("留下的应是内核说明,实际 %+v", notice.Event) + } + content, _ := notice.Event.Payload["content"].(string) + for _, want := range []string{"2 条", resident.ID, "转投"} { + if !strings.Contains(content, want) { + t.Errorf("说明缺少 %q:%s", want, content) + } + } + t.Logf("队列说明: %s", content) +} + +// 达到上限后不得无界增殖:每个 tick 都拉一个新子会把机器拖垮。 +func TestOffloadRespectsResidentCap(t *testing.T) { + root, main := newRootWithoutSchedulerLoop(t) + defer main.Close() + + opts := DefaultOffloadOptions() + opts.Enabled = true + opts.BusyAfter = time.Nanosecond + opts.MinPending = 1 + opts.MaxResidents = 1 + + root.sched.enqueue(makeQueuedInput(1)) + root.sched.nextRef() + + // 第一轮:拉起 1 个 + root.sched.enqueue(makeQueuedInput(2)) + if n := root.offloadPendingTasks(opts); n != 1 { + t.Fatalf("第一轮应转 1 条,实际 %d", n) + } + // 第二轮:已达上限,但**会复用**刚建的那个子,所以仍能转投 + root.sched.enqueue(makeQueuedInput(3)) + if n := root.offloadPendingTasks(opts); n != 1 { + t.Fatalf("第二轮应复用已有驻留子,实际转 %d", n) + } + if got := len(root.Residents()); got != 1 { + t.Fatalf("★ 不得越过上限增殖:期望 1 个驻留子,实际 %d", got) + } +} + +// ★ 回归:检查必须发生在**独立 goroutine** 里。 +// +// schedulerLoop 是同步执行任务的(executeNewTask 阻塞到任务结束), +// 所以"正忙"期间它根本不会回到循环顶部 —— 把检查放在那里的实现 +// 永远不会触发(我第一版就是这么写的,测出来才发现)。 +// 本测试钉死:offloadLoop 确实起了自己的 goroutine 并能被唤醒干活。 +func TestOffloadLoopRunsWhileBusy(t *testing.T) { + root, main := newRootWithoutSchedulerLoop(t) + defer main.Close() + + root.offload = OffloadOptions{ + Enabled: true, BusyAfter: 10 * time.Millisecond, + MinPending: 1, MaxResidents: 1, + } + // 伪造"正忙":直接占住 running(不启动真实调度循环,避免它把任务跑掉) + root.sched.enqueue(makeQueuedInput(1)) + root.sched.nextRef() + root.sched.enqueue(makeQueuedInput(2)) + + go root.offloadLoop() // 独立 goroutine,正是被测的点 + + deadline := time.Now().Add(3 * time.Second) + for time.Now().Before(deadline) { + if len(root.Residents()) > 0 { + return // 成功:忙时后台循环把积压转走了 + } + time.Sleep(20 * time.Millisecond) + } + t.Fatal("offloadLoop 在忙时没有转投:检查没有跑在独立 goroutine 里?") +} diff --git a/internal/agent/core/resident.go b/internal/agent/core/resident.go index 156ba94..aa508a0 100644 --- a/internal/agent/core/resident.go +++ b/internal/agent/core/resident.go @@ -48,6 +48,12 @@ type ResidentOptions struct { Capacity int // TempPath 是它 temp 图记忆的存储路径(必填;与子同生共死)。 TempPath string + // OffloadOwned 标记这是**内核为转投而拉起**的驻留子(见 offload.go)。 + // + // 为什么要区分:自动转投会复用、也会按上限限制自己拉起的子, + // 但**绝不能**把人(父的模型)建的子算进去或拿去复用 —— + // 那会把人工安排的工作负载搬到一个本来在做别的事的 agent 上。 + OffloadOwned bool } // ResidentInfo 是父对某个驻留子的可查询状态(登记表条目 + 状态面摘要)。 @@ -83,7 +89,9 @@ type residentChild struct { dir string inputChs []string allowed []string - createdAt time.Time + createdAt time.Time + // offloadOwned 标记这是内核为转投拉起的子(见 offload.go)。 + offloadOwned bool mu sync.Mutex state string // running | contextfull | stopped @@ -175,6 +183,7 @@ func (a *Agent) SpawnResident(opts ResidentOptions) (ResidentInfo, error) { inputChs: append([]string(nil), opts.InputChs...), allowed: append([]string(nil), opts.AllowedOutputs...), createdAt: time.Now(), state: "running", + offloadOwned: opts.OffloadOwned, } // ④ 子 → 父的主动消息(**L3 中断**,带子标识):投进父的 inputch。 diff --git a/internal/agent/core/scheduler.go b/internal/agent/core/scheduler.go index 62eed2e..324ab35 100644 --- a/internal/agent/core/scheduler.go +++ b/internal/agent/core/scheduler.go @@ -412,6 +412,101 @@ func (s *scheduler) signalWake() { } } +// runningTask 返回当前正在运行的任务(无则 nil)。 +// +// 供自动转投判定"主 agent 是否忙"用。本函数只读取指针,不参与调度决策。 +func (s *scheduler) runningTask() *Task { + s.mu.Lock() + defer s.mu.Unlock() + return s.running +} + +// runningFor 返回当前运行任务已持续多久;无运行任务时为 0。 +// +// 为什么用 EnqueuedAt 而不是另开一个"开始时间"字段:Task 已有 EnqueuedAt +// (allocateIDLocked 保证非零),它对外排队输入就是"多久没人处理它", +// 恰好也是转投判定关心的量。新增字段会多一份需要维护的真相。 +func (s *scheduler) runningFor() time.Duration { + s.mu.Lock() + defer s.mu.Unlock() + if s.running == nil || s.running.EnqueuedAt.IsZero() { + return 0 + } + return time.Since(s.running.EnqueuedAt) +} + +// takeQueuedInputs 从就绪队列**前端**取出至多 max 条纯排队输入(用于转投)。 +// +// 只取 TaskQueued + KindInput: +// - 中断任务带级别语义,転投会打乱中断阶梯,不动; +// - self 任务(记忆整理等)与父的记忆面绑定,不能换 agent。 +// +// 只从**前端**取:队列是 FIFO,前端就是"最久没人处理"的那几条; +// 从尾部抽会把后来者先送走,反而拉长前面几等的等待。 +// +// 少于 min 条时**什么都不取**(返回 nil):拉起一个 agent 的成本不该为 +// 一条任务付;宁等下一轮积到量再一起转。这使得本函数要么不动、要么成批移动。 +func (s *scheduler) takeQueuedInputs(min int) []offloadCandidate { + s.mu.Lock() + defer s.mu.Unlock() + var out []offloadCandidate + for _, t := range s.queue { + if t == nil || t.Class != TaskQueued || t.Kind != TaskKindInput || t.Event == nil { + continue + } + out = append(out, offloadCandidate{Event: t.Event}) + if len(out) >= min { + break + } + } + if len(out) < min { + return nil + } + // 真的取走:重建队列,抽掉刚选中的那些事件(按指针同一性判断)。 + selected := make(map[*agentIO.InputEvent]bool, len(out)) + for _, c := range out { + selected[c.Event] = true + } + kept := s.queue[:0] + for _, t := range s.queue { + if t != nil && t.Event != nil && selected[t.Event] { + continue + } + kept = append(kept, t) + } + s.queue = kept + return out +} + +// requeueFront 把任务放回队列**前端**(保持原相对顺序)。 +// +// 转投失败/未生效时用,保证"输入只多不少":吞掉一条输入比多处理一条更糟 +// (与 routeOnInputByOwner 的兜底同一条理由)。 +// self 任务在転投路径上不可能出现(takeQueuedInputs 已排掉),所以这里只需处理 Event。 +func (s *scheduler) requeueFront(cands []offloadCandidate) { + if len(cands) == 0 { + return + } + s.mu.Lock() + defer s.mu.Unlock() + restored := make([]*Task, 0, len(cands)) + for _, c := range cands { + if c.Event == nil { + continue + } + t := newInputTask(c.Event) + // 保持"这些任务比当前队列里的一切都早"的语义。 + if t.EnqueuedAt.IsZero() { + t.EnqueuedAt = time.Now() + } + restored = append(restored, t) + } + if len(restored) == 0 { + return + } + s.queue = append(restored, s.queue...) +} + // setCritical 由调度器 goroutine 在任务进入/离开临界区时设置。 func (s *scheduler) setCritical(v bool) { s.critical.Store(v) } diff --git a/internal/config/registry.go b/internal/config/registry.go index 3a6aa3f..7fd56db 100644 --- a/internal/config/registry.go +++ b/internal/config/registry.go @@ -734,6 +734,12 @@ func (r *ConfigRegistry) seedDBValues(dataDir string) { set("core.agent.max_tool_turns", "10") set("core.agent.max_context_size", "30") + // 积压转投默认**关闭**(false):它让内核替父做决策(设计 §7 的刻意例外), + // 所以必须由部署方显式打开,而不是默默改变系统行为。 + set("core.agent.offload_enabled", "false") + set("core.agent.offload_busy_after", "5m") + set("core.agent.offload_min_pending", "3") + set("core.agent.offload_max_residents", "2") set("core.agent.distill_interval", "30m") set("core.agent.archive_interval", "60m") set("core.agent.review_interval", "120m") @@ -836,6 +842,10 @@ func (r *ConfigRegistry) seedCoreDefs(dataDir string) { reg(ConfigDef{Key: "core.log.path", Default: filepath.Join(dataDir, "log"), Type: "string", DisplayName: "日志目录", Description: "日志文件输出目录", Category: "paths"}) reg(ConfigDef{Key: "core.agent.max_tool_turns", Default: "10", Type: "int", DisplayName: "最大工具轮次", Description: "单次请求允许的最大工具调用轮数", Category: "agent"}) + reg(ConfigDef{Key: "core.agent.offload_enabled", Default: "false", Type: "bool", DisplayName: "积压任务自动转投", Description: "主 agent 长时间忙时,把积压任务自动转投给内核拉起的驻留子(需重启生效)", Category: "agent"}) + reg(ConfigDef{Key: "core.agent.offload_busy_after", Default: "5m", Type: "string", DisplayName: "转投触发忙时长", Description: "运行任务持续超过该时长才考虑转投(如 5m)", Category: "agent"}) + reg(ConfigDef{Key: "core.agent.offload_min_pending", Default: "3", Type: "int", DisplayName: "转投最少积压条数", Description: "积压少于该条数时不值得拉起驻留子", Category: "agent"}) + reg(ConfigDef{Key: "core.agent.offload_max_residents", Default: "2", Type: "int", DisplayName: "转投驻留子上限", Description: "自动转投最多拉起几个驻留子(人工创建的不计)", Category: "agent"}) reg(ConfigDef{Key: "core.agent.max_context_size", Default: "30", Type: "int", DisplayName: "最大上下文", Description: "上下文窗口中保留的最大消息条数", Category: "agent"}) reg(ConfigDef{Key: "core.agent.distill_interval", Default: "30m", Type: "duration", DisplayName: "蒸馏间隔", Description: "记忆蒸馏的执行间隔", Category: "agent"}) reg(ConfigDef{Key: "core.agent.archive_interval", Default: "60m", Type: "duration", DisplayName: "冷文档归档间隔", Description: "冷文档归档(L2→L3)的执行间隔", Category: "agent"})