mirror of
https://gitcode.com/JianFeeeee/HomeAgent.git
synced 2026-10-03 15:53:56 +00:00
fix(offload): 积压可能全堵在 io 输入 channel,不在就绪队列
线上实测发现上一版**永不触发**:主 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),修后通过 —— 是真回归测试。
This commit is contained in:
@ -94,6 +94,32 @@ type offloadCandidate struct {
|
|||||||
Event *agentIO.InputEvent
|
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 名。
|
// residentInputChannel 返回"父给某个子投递输入"用的 inputch 名。
|
||||||
//
|
//
|
||||||
// 与 SendToResident 用的是同一个(sub/<id>):子是**不配插件 inputch** 的
|
// 与 SendToResident 用的是同一个(sub/<id>):子是**不配插件 inputch** 的
|
||||||
@ -124,7 +150,19 @@ func (a *Agent) offloadPendingTasks(opts OffloadOptions) int {
|
|||||||
return 0
|
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)
|
cands := a.sched.takeQueuedInputs(opts.MinPending)
|
||||||
if len(cands) == 0 {
|
if len(cands) == 0 {
|
||||||
return 0
|
return 0
|
||||||
|
|||||||
@ -299,3 +299,49 @@ func TestOffloadLoopRunsWhileBusy(t *testing.T) {
|
|||||||
}
|
}
|
||||||
t.Fatal("offloadLoop 在忙时没有转投:检查没有跑在独立 goroutine 里?")
|
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()))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user