From fe1d2672d8e4949897f4111ade0f41f70378c212 Mon Sep 17 00:00:00 2001 From: JianFeeeee Date: Sat, 19 Sep 2026 17:05:38 +0800 Subject: [PATCH] =?UTF-8?q?fix(offload):=20=E7=A7=AF=E5=8E=8B=E5=8F=AF?= =?UTF-8?q?=E8=83=BD=E5=85=A8=E5=A0=B5=E5=9C=A8=20io=20=E8=BE=93=E5=85=A5?= =?UTF-8?q?=20channel=EF=BC=8C=E4=B8=8D=E5=9C=A8=E5=B0=B1=E7=BB=AA?= =?UTF-8?q?=E9=98=9F=E5=88=97?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 线上实测发现上一版**永不触发**:主 agent 跑着 6×45s 的长任务、我连发 4 条消息, scheduler 始终显示 queue=0、residents=0,转投一次都没发生。 根因:schedulerLoop 是**同步执行**任务的,所以「正忙」期间它根本回不到循环顶部 去调 pumpInbox —— 后到的输入全堆在 io.inputCh(容量 256)里,压根没进 sched.queue。 而 takeQueuedInputs 只看 s.queue ⇒ 恒取不到东西。 ★ 仓库里早记过同一个坑:armStop 的注释写着「pending 是还没被 pumpInbox 搬进队列 的那一段……只数 s.queue 会得到 0(实测),配额随之失效」。我重犯了它。 修法:转投前先 drainInboxToQueue() 把 channel 里的输入搬进队列。 与 pumpInbox 的区别是**不要求 hasRoom** —— pumpInbox 满时会停下保留背压, 而转投场景恰恰是「队列空、输入堵在 channel」(调度器回不到 pumpInbox)。 队列上限仍由 enqueue 把关,放不下的给同步调用方 skipped 终态(不丢、不阻塞)。 回归测试 TestOffloadSeesInputsStuckInChannel 精确复现该现场状态: 已实测它对着修复前的逻辑**会失败**(期望 3 实际 0),修后通过 —— 是真回归测试。 --- internal/agent/core/offload.go | 40 ++++++++++++++++++++++++- internal/agent/core/offload_test.go | 46 +++++++++++++++++++++++++++++ 2 files changed, 85 insertions(+), 1 deletion(-) diff --git a/internal/agent/core/offload.go b/internal/agent/core/offload.go index 96523c5..ffeef00 100644 --- a/internal/agent/core/offload.go +++ b/internal/agent/core/offload.go @@ -94,6 +94,32 @@ type offloadCandidate struct { Event *agentIO.InputEvent } +// drainInboxToQueue 把 io 输入 channel 里**已经到达但尚未被搬运**的输入 +// 搬进就绪队列(非阻塞;取空为止)。 +// +// 它与 pumpInbox 做的事一样,但**不要求 hasRoom**:pumpInbox 在队列满时 +// 会停下以保留背压,而转投场景恰恰是「队列空/不满、但输入堵在 channel 里」 +// (因为调度器正忙于执行任务、根本回不到 pumpInbox)。 +// +// 队列上限仍由 enqueue 把关:满了就停下,超出的输入留在 channel 里。 +func (a *Agent) drainInboxToQueue() { + for { + select { + case evt := <-a.io.InputChan(): + if !a.sched.enqueue(newInputTask(evt)) { + // 队列满:放不进去。不能丢,也不能阻塞(我们是后台 goroutine, + // 阻塞会把这个循环永远卡住)——给同步调用方一个终态后丢弃, + // 与 pumpInbox 的 queue_full 处置一致。 + a.sched.noteBackpressure() + a.emitSkippedReply(evt, "queue_full") + return + } + default: + return + } + } +} + // residentInputChannel 返回"父给某个子投递输入"用的 inputch 名。 // // 与 SendToResident 用的是同一个(sub/):子是**不配插件 inputch** 的 @@ -124,7 +150,19 @@ func (a *Agent) offloadPendingTasks(opts OffloadOptions) int { return 0 } - // ② 收集可转投的排队任务;不够量就不值得拉起一个 agent。 + // ② 先把积压从 io 的输入 channel **搬进就绪队列**。 + // + // !!这是本特性最容易写错的一步(我第一版就错了,写完后线上实测永不触发): + // schedulerLoop 是**同步执行**任务的,所以「正忙」期间它根本不会回到循环顶部 + // 去调 pumpInbox —— 这时后到的输入全部堆在 io.inputCh(容量 256)里, + // **压根没进 sched.queue**。只数 s.queue 会得到 0,转投就永远不触发。 + // + // 仓库里早记过同一个坑:armStop 的注释写着「pending 是还没被 pumpInbox 搬进 + // 队列的那一段……只数 s.queue 会得到 0(实测),配额随之失效」。 + // 这里必须在同一层把这件事做对,而不是重犯。 + a.drainInboxToQueue() + + // ③ 收集可转投的排队任务;不够量就不值得拉起一个 agent。 cands := a.sched.takeQueuedInputs(opts.MinPending) if len(cands) == 0 { return 0 diff --git a/internal/agent/core/offload_test.go b/internal/agent/core/offload_test.go index 02f64ab..7592c40 100644 --- a/internal/agent/core/offload_test.go +++ b/internal/agent/core/offload_test.go @@ -299,3 +299,49 @@ func TestOffloadLoopRunsWhileBusy(t *testing.T) { } t.Fatal("offloadLoop 在忙时没有转投:检查没有跑在独立 goroutine 里?") } + +// ★★ 回归:积压可能**全在 io 输入 channel 里**,不在 sched.queue。 +// +// 这是本特性最容易写错、而且我在线上真踩了的一步:schedulerLoop 是同步执行 +// 任务的,所以「正忙」期间它根本回不到 pumpInbox —— 后到的输入全堆在 +// io.inputCh(容量 256)里,sched.queue 恒为 0。 +// +// 只数 s.queue 的实现在线上**永不触发**(实测:主 agent 跑着 6×45s 的任务、 +// 我连发 4 条消息,队列始终显示 0、residents 始终 0)。 +// 仓库里 armStop 早记过同一个坑("只数 s.queue 会得到 0"),这里钉死不重犯。 +func TestOffloadSeesInputsStuckInChannel(t *testing.T) { + root, main := newRootWithoutSchedulerLoop(t) + defer main.Close() + + opts := DefaultOffloadOptions() + opts.Enabled = true + opts.BusyAfter = time.Nanosecond + opts.MinPending = 3 + opts.MaxResidents = 1 + + // 伪造"正忙" + root.sched.enqueue(makeQueuedInput(1)) + root.sched.nextRef() + + // 关键:把 3 条消息注入 **io 输入 channel**,不碰 sched.queue。 + // 这精确复现"调度器忙于执行任务、pumpInbox 没被调用"的现场状态。 + for i := 0; i < 3; i++ { + root.io.InjectInputTo("webui", "webui", "text", + map[string]interface{}{"content": "stuck"}) + } + if len(root.sched.queue) != 0 { + t.Fatalf("前置条件:此时 sched.queue 应为 0(输入还没被搬运),实际 %d", len(root.sched.queue)) + } + if root.io.PendingInputs() != 3 { + t.Fatalf("前置条件:输入应堆在 channel 里,实际 %d", root.io.PendingInputs()) + } + + // 转投必须能看到它们(先搬进队列再取) + moved := root.offloadPendingTasks(opts) + if moved != 3 { + t.Fatalf("★ 堆在 channel 里的积压必须被看见并转投:期望 3,实际 %d", moved) + } + if len(root.Residents()) != 1 { + t.Fatalf("应拉起 1 个驻留子,实际 %d", len(root.Residents())) + } +}