Files
HomeAgent/internal/agent/io/inputroute_test.go
JianFeeeee 4707b05498 fix(scheduler): inputch 划给子后输入只流向子 —— 补上「进内核之前」的输入路由
用户指出的语义(设计稿 §4.1 早已写明):
**inputch 是可分配资源**,「路由发生在**进内核之前**」—— 划给某个 agent 后,
该通道的输入**只流向那个 agent**;outputch 不同,授权是**非独占**的,
父依旧可以通过它发送内容。

而代码里 inputch 划拨只做了**登记**,没有做**路由**:
- 插件注入输入的 io 是**根 agent 的**(`cmd/homed` 里 `pluginReg.SetIOManager(iom)`);
- 唯一消费输入的是「该 io 自己的调度器」(`scheduler.go` 读 `a.io.InputChan()`);
- `ChannelRegistry.Assign` 只把 Owner 写进登记表,**没有任何转发动作**。

⇒ 现场表现(用户线上联调):子挂 `inputch=[timer]`,**timer 的输入却打在父身上**
(日志 `[agent] interrupt from timer/timer`),子侧 `轮次=0` 永远不动。
登记表里的 Owner 于是沦为标签。

改法(按 §4.1 把路由放回"进内核之前"):
- `IOManager` 增加 `InputRouter`(`SetInputRouter`),并把**五处直接入队**收口到
  `deliverInput`:`InjectInput` / `InjectInputSync` / `InjectInputTo` /
  `InjectInputSyncTo` / `InjectInterrupt`(排队与中断两条路都过路由)。
- 内核注入路由器 `Agent.routeInputByOwner`:查 inputch 的 Owner —— 归自己/未分配 ⇒
  本内核处理;归自己的某个驻留子 ⇒ `DeliverRouted` 交给它(**不再进父的队列**);
  归一个不存在的 agent ⇒ **不吞输入**,父兜底 + 留痕(吞掉输入比多处理一条更糟)。
- `DeliverRouted` 是"已路由"的投递口,不再二次路由(避免成环)。
- 同步输入的 `ResponseCh` 随事件一起走 ⇒ 回答由持有者写回同一回程(§4.3)。

判据(新增 6 条):
- io 层:被接管时排队/中断都**不入本内核队列**(且中断确实经过路由)/ 放行与未设
  路由器时与历史行为一致 / `DeliverRouted` 不再触发路由
- 内核层:划给子的 inputch 输入进**子**(子 Executed>0)且**父 Enqueued 不变** /
  归属到不存在的 agent 时父兜底(不吞)/ 未分配的 inputch 仍归父
2026-09-13 15:35:54 +08:00

75 lines
2.6 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 io
import "testing"
// 输入路由:路由器说"已被别的 agent 接管"时,事件**不得**进本内核队列。
//
// 语义(设计 §4.1inputch 是可分配资源,划给某个 agent 后输入只流向它 ——
// "父也顺便看到一份"是错的。
func TestInputRouter_TakesOverExclusively(t *testing.T) {
m := NewIOManager()
var got []*InputEvent
var sawInterrupt bool
m.SetInputRouter(func(evt *InputEvent, isInterrupt bool) bool {
got = append(got, evt)
sawInterrupt = sawInterrupt || isInterrupt
return true // 全部接管
})
m.InjectInputTo("plugin-x", "sub/in", "text", map[string]interface{}{"content": "a"})
m.InjectInterruptTextOpts("plugin-x", "sub/in", "b", InjectOptions{})
if len(got) < 2 {
t.Fatalf("路由器应被调用(含中断路径),实际 %d 次", len(got))
}
if n := len(m.InputChan()); n != 0 {
t.Fatalf("被接管的排队输入不得进本内核队列,实际 %d 条", n)
}
if !sawInterrupt {
t.Fatal("中断注入也必须经过路由(否则中断会绕过 inputch 归属直投父)")
}
if got[0].OutputChannel != "sub/in" {
t.Fatalf("路由器应拿到事件的 inputch得到 %q", got[0].OutputChannel)
}
}
// 路由器放行(返回 false或未设置时行为与以前完全一致。
func TestInputRouter_PassthroughKeepsOldBehaviour(t *testing.T) {
m := NewIOManager()
calls := 0
m.SetInputRouter(func(evt *InputEvent, isInterrupt bool) bool { calls++; return false })
m.InjectInputTo("plugin-x", "sub/in", "text", map[string]interface{}{"content": "a"})
if calls != 1 {
t.Fatalf("路由器应被调用一次,实际 %d", calls)
}
if n := len(m.InputChan()); n != 1 {
t.Fatalf("放行的输入应进本内核队列,实际 %d 条", n)
}
if _, ok := <-m.InputChan(); !ok {
t.Fatal("队列应可读")
}
// 未设路由器:直接入队(历史行为)
m2 := NewIOManager()
m2.InjectInterruptText("plugin-x", "sub/in", "c")
if n := len(m2.InputInterruptChan()); n != 1 {
t.Fatalf("未设路由器时中断应直接入队,实际 %d 条", n)
}
}
// DeliverRouted 是不再二次路由的投递口(路由器实现把事件交给持有者)。
func TestDeliverRouted_SkipsSecondRouting(t *testing.T) {
m := NewIOManager()
routerCalls := 0
m.SetInputRouter(func(evt *InputEvent, isInterrupt bool) bool { routerCalls++; return true })
m.DeliverRouted(&InputEvent{OutputChannel: "sub/in"}, false)
if routerCalls != 0 {
t.Fatalf("DeliverRouted 不应再触发路由(会成环),实际 %d 次", routerCalls)
}
if n := len(m.InputChan()); n != 1 {
t.Fatalf("应已入队,实际 %d 条", n)
}
}