From 34c0df27050cecb805e4ffd5217df8e6e49ab099 Mon Sep 17 00:00:00 2001 From: JianFeeeee Date: Sun, 27 Sep 2026 11:27:24 +0800 Subject: [PATCH] =?UTF-8?q?feat(toolcall):=20=E6=89=B9=E6=AC=A1=E5=B9=B6?= =?UTF-8?q?=E5=8F=91=E8=B0=83=E5=BA=A6=E4=B8=8E=E5=90=8C=E9=80=9A=E9=81=93?= =?UTF-8?q?=E4=BF=9D=E5=BA=8F=EF=BC=88=E9=98=B6=E6=AE=B5=202d=EF=BC=8C?= =?UTF-8?q?=E9=97=AD=E5=90=88=202c=20=E5=88=A4=E6=8D=AE=E7=BC=BA=E5=8F=A3?= =?UTF-8?q?=EF=BC=89?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 规则(三条全满足才并发): 1. 批内 >1 个工具 2. **全部**工具声明 ParallelSafe —— 一个不声明就整批降级,不做部分并发 3. 不含需保序的同通道输出发送 SDK: · ToolDef 加 ParallelSafe bool。⚠️ 零值 false 是刻意的:存量插件不改一行 就得到**保守**行为(整批串行),不会因升级被意外并发。声明它是责任 而非特权。纯新增字段,无签名变更。 · io.ToolDef 同步加该字段(设备/通道工具走 io 路径,只查 StageHost 会漏)。 core: · 新增 StepToolBatch —— runTaskSteps 是单线程驱动状态机的, 「每步一个工具」的游标模型无法表达「一批同时跑」,故需独立 step。 · stepToolBatch:fan-out(每工具一 goroutine,各写自己的 toolCtxs[i]) → join → **按索引顺序**串行收尾(after_toolcall / 落消息 / 事件)。 收尾必须串行且按索引:f.Msgs 是共享切片,且按索引落才能让模型读到的 上下文顺序与它自己发出的顺序一致。 · runOneTool 抽出「before_toolcall + 执行」的单工具逻辑,串行/并发两条路共用。 · toolParallelSafe / batchRunnable 判据函数。 ★ 修掉一个我自己引入的竞争:resolveTurnScenes 会把结果记进**共享**的 f.sceneDone / f.turnScene(memorypass.go:289)。最初在每个 goroutine 里 各调一次 —— 既是数据竞争,又会各自触发一次 EnterSceneWithHint, 重复计入场景强度(正是 sceneDone 注释警告过的问题)。改为在 fan-out **之前**解析一次,goroutine 内只读。 判据(parallelsched_test.go,4 条): · 全批 ParallelSafe ⇒ 并发峰值 >= 2(用阻塞设备观察真实并发) · 一个非 ParallelSafe ⇒ 整批串行,但**仍全部执行** · 同 output_send__<通道> 连发 3 条 ⇒ 严格按声明顺序到达 · ★ 并发下每个工具的 ctx 只带自己的 ToolCalls、after 读到自己结果 ★ 并关闭了 2c 的判据缺口:此前两条 2c 判据在**串行**下无法区分 per-tool 与单槽(变体验证后仍全绿)。新增的并发版判据在退回单槽时 触发 **6 处 DATA RACE 报告 + 串味断言失败**(dup:k_a 与 k_a 撞名)。 至此 2c 可记为已验证。 过程中三次自伤: · resolveTurnScenes 竞争(上述); · 我的 harness 用 StageHost 注册 handler 遮蔽了设备工具, slowDevice 根本没被调用("实际 0")——改为在 io.ToolDef 上声明; · 批内并发峰值判据最初用 StageHost 声明 ParallelSafe,掩盖了 「设备工具也需要该字段」这一真实缺口。 回归:internal/agent/... internal/sdk/... internal/plugin/... internal/plugins/... 全绿(17 包);core 包 -race 全绿。 --- internal/agent/core/argvalidate.go | 55 ++++ internal/agent/core/parallelsched_test.go | 318 ++++++++++++++++++++++ internal/agent/core/task.go | 129 ++++++++- internal/agent/io/channel.go | 3 + third_party/homeagent-sdk/sdk/plugin.go | 10 + 5 files changed, 514 insertions(+), 1 deletion(-) create mode 100644 internal/agent/core/parallelsched_test.go diff --git a/internal/agent/core/argvalidate.go b/internal/agent/core/argvalidate.go index 7925c84..3925f72 100644 --- a/internal/agent/core/argvalidate.go +++ b/internal/agent/core/argvalidate.go @@ -241,3 +241,58 @@ func (a *Agent) validateArgsAgainstSchema(tc agentAPI.ToolCall) *sdk.ToolError { } return nil } + +// toolParallelSafe 报告工具是否可被**并发执行**。 +// +// 两条来源都要查(与 validateArgsAgainstSchema 同理):插件工具走 StageHost, +// 设备/通道工具走 IOManager。查不到 ⇒ 保守返回 false(不可并发)。 +// +// 为什么保守:新语义下并发会改变工具的行为前提,让存量插件意外并发 +// 比慢一点危险得多——判不出就该按串行走。 +func (a *Agent) toolParallelSafe(name string) bool { + if a == nil { + return false + } + if a.stageHost != nil { + if def := a.stageHost.ToolDef(name); def != nil { + return def.ParallelSafe + } + } + if a.io != nil { + if def, ok := a.io.ToolDefOf(name); ok { + return def.ParallelSafe + } + } + return false +} + +// batchRunnable 并发执行本批工具。 +// +// 何时并发(三条全满足): +// 1. 批内 >1 个工具 +// 2. **全部**工具都声明 ParallelSafe —— 一个不声明就整批降级, +// 不做"部分并发":部分并发收益不抵其不可预测性 +// 3. 不含需要保序的同通道输出发送(同 output_send__<通道> 多次发送) +// +// 保序为什么不用 ParallelSafe 表达:那属于**批内**约束而非工具属性, +// 且同一工具在不同批里的通道可能不同(output_send__qq 两次、一次 qq 一次 cli)。 +func (f *TaskFrame) batchRunnable(a *Agent) bool { + if f == nil || len(f.PendingTools) <= 1 { + return false + } + chans := map[string]bool{} + for _, tc := range f.PendingTools { + if !a.toolParallelSafe(tc.Name) { + return false + } + // 同通道多次发送必须保序 —— 用户可见消息顺序敏感 + if isOutputDeliveryTool(tc.Name) { + ch := strings.TrimPrefix(tc.Name, "output_send__") + if chans[ch] { + return false + } + chans[ch] = true + } + } + return true +} diff --git a/internal/agent/core/parallelsched_test.go b/internal/agent/core/parallelsched_test.go new file mode 100644 index 0000000..8b57498 --- /dev/null +++ b/internal/agent/core/parallelsched_test.go @@ -0,0 +1,318 @@ +package core + +import ( + "fmt" + "strings" + "sync" + "sync/atomic" + "testing" + "time" + + agentAPI "gitcode.com/JianFeeeee/HomeAgent/internal/agent/api" + agentIO "gitcode.com/JianFeeeee/HomeAgent/internal/agent/io" + sdk "gitcode.com/JianFeeeee/HomeAgent/internal/sdk" +) + +// 阶段 2d:批次调度与保序。 +// +// 三条规则: +// 1. 全批 ParallelSafe ⇒ 并发;否则**整批**降级串行(不做部分并发—— +// 收益不抵不可预测性)。 +// 2. 同一 output_send__<通道> 的多次发送**保序**(用户可见消息顺序敏感)。 +// 3. 并发时每个工具只写自己的 per-tool ctx(这同时闭合 2c 的判据缺口)。 + +// slowDevice 是一个"可并发"的测试设备:每个 Execute 阻塞到被显式放行, +// 用来观察多个工具是否**同时**在执行中。 +type slowDevice struct { + name string + toolNames []string + safe []string // 声明为 ParallelSafe 的工具名 + entered *int32 + active *int32 + maxActive *int32 + release chan struct{} + delay time.Duration +} + +func (d *slowDevice) Name() string { return d.name } +func (d *slowDevice) Type() agentIO.DeviceType { return agentIO.DeviceOutput } +func (d *slowDevice) Description() string { return "slow test device" } +func (d *slowDevice) Tools() []agentIO.ToolDef { + safe := map[string]bool{} + for _, n := range d.safe { + safe[n] = true + } + out := make([]agentIO.ToolDef, 0, len(d.toolNames)) + for _, n := range d.toolNames { + out = append(out, agentIO.ToolDef{Name: n, ParallelSafe: safe[n]}) + } + return out +} + +func (d *slowDevice) Execute(tool string, args map[string]interface{}) (interface{}, error) { + atomic.AddInt32(d.entered, 1) + cur := atomic.AddInt32(d.active, 1) + // 记录并发峰值 + for { + old := atomic.LoadInt32(d.maxActive) + if cur <= old || atomic.CompareAndSwapInt32(d.maxActive, old, cur) { + break + } + } + if d.release != nil { + <-d.release // 阻塞,直到测试放行 + } else if d.delay > 0 { + time.Sleep(d.delay) + } + atomic.AddInt32(d.active, -1) + return "ran:" + tool, nil +} +func (d *slowDevice) Start() error { return nil } +func (d *slowDevice) Stop() error { return nil } +func (d *slowDevice) OutputCapabilities() agentIO.OutputCapability { return agentIO.CapText } +func (d *slowDevice) ChannelDef() agentIO.ChannelDef { return agentIO.ChannelDef{} } + +// ① 并发:全批 ParallelSafe ⇒ 同一时刻有多个工具在执行。 +func TestBatchParallelWhenAllToolsParallelSafe(t *testing.T) { + var entered, active, maxActive int32 + release := make(chan struct{}) + names := []string{"p_a", "p_b", "p_c"} + + sp := &batchProvider{responses: []*agentAPI.CompletionResponse{ + {ToolCalls: []agentAPI.ToolCall{ + {ID: "c1", Name: "p_a", Arguments: map[string]interface{}{}}, + {ID: "c2", Name: "p_b", Arguments: map[string]interface{}{}}, + {ID: "c3", Name: "p_c", Arguments: map[string]interface{}{}}, + }}, + {Content: "final"}, + }} + a := newParallelAgent(t, sp, names, names, &entered, &active, &maxActive, release) + + done := make(chan struct{}) + go func() { + defer close(done) + if out := a.runTaskSteps(a.newTaskFrame("go", a.stageCtxFromInput("go", "", ""))); out != outcomeDone { + t.Errorf("runTaskSteps=%v", out) + } + }() + + // 等到**至少两个**同时进入执行;若串行则永远只有一个,会超时。 + deadline := time.Now().Add(3 * time.Second) + for atomic.LoadInt32(&maxActive) < 2 && time.Now().Before(deadline) { + time.Sleep(5 * time.Millisecond) + } + close(release) + <-done + + if got := atomic.LoadInt32(&maxActive); got < 2 { + t.Errorf("全批 ParallelSafe 却未并发(并发峰值=%d,应 >=2)", got) + } + if entered != int32(len(names)) { + t.Errorf("应执行 %d 个工具,实际 %d", len(names), entered) + } +} + +// ② 整批降级:只要有一个**非** ParallelSafe ⇒ 整批串行(不做部分并发)。 +func TestBatchFallsBackToSerialIfAnyToolNotParallelSafe(t *testing.T) { + var entered, active, maxActive int32 + release := make(chan struct{}) + names := []string{"s_a", "s_b", "s_c"} + + sp := &batchProvider{responses: []*agentAPI.CompletionResponse{ + {ToolCalls: []agentAPI.ToolCall{ + {ID: "c1", Name: "s_a", Arguments: map[string]interface{}{}}, + {ID: "c2", Name: "s_b", Arguments: map[string]interface{}{}}, + {ID: "c3", Name: "s_c", Arguments: map[string]interface{}{}}, + }}, + {Content: "final"}, + }} + // 只声明前两个可并发 —— 第三个不声明 ⇒ 整批必须串行 + a := newParallelAgent(t, sp, names, names[:2], &entered, &active, &maxActive, release) + + done := make(chan struct{}) + go func() { + defer close(done) + a.runTaskSteps(a.newTaskFrame("go", a.stageCtxFromInput("go", "", ""))) + }() + // 给串行留出充分时间:让第一个工具走完并进入第二个 + time.Sleep(150 * time.Millisecond) + observed := atomic.LoadInt32(&maxActive) + close(release) + <-done + + if observed > 1 { + t.Errorf("存在非 ParallelSafe 工具时不应部分并发(并发峰值=%d)", observed) + } + if atomic.LoadInt32(&entered) != int32(len(names)) { + t.Errorf("整批仍应全部执行,实际 %d", entered) + } +} + +// ③ 保序:同一 output_send__<通道> 的多次发送**必须**按声明顺序到达。 +// +// 这是用户可见的语义:同一条通道连发 3 条消息,顺序颠倒用户就读错了。 +func TestBatchSameOutputChannelKeepsOrder(t *testing.T) { + sp := &batchProvider{responses: []*agentAPI.CompletionResponse{ + {ToolCalls: []agentAPI.ToolCall{ + {ID: "c1", Name: "output_send__testch", Arguments: map[string]interface{}{"payload": "first"}}, + {ID: "c2", Name: "output_send__testch", Arguments: map[string]interface{}{"payload": "second"}}, + {ID: "c3", Name: "output_send__testch", Arguments: map[string]interface{}{"payload": "third"}}, + }}, + {Content: "final"}, + }} + a := newPreemptAgent(t, sp) + + var mu sync.Mutex + var got []string + dev := &recordingDevice{name: "testch", caps: agentIO.CapText, + tools: []agentIO.ToolDef{{Name: "output_send__testch"}}} + dev.onExec = func(payload string) { + mu.Lock() + got = append(got, payload) + mu.Unlock() + // 每条之间加延迟:若并发执行,顺序会被打乱 + time.Sleep(20 * time.Millisecond) + } + if err := a.io.RegisterDevice(dev); err != nil { + t.Fatalf("注册设备失败: %v", err) + } + + if out := a.runTaskSteps(a.newTaskFrame("go", a.stageCtxFromInput("go", "", ""))); out != outcomeDone { + t.Fatalf("runTaskSteps=%v", out) + } + want := []string{"first", "second", "third"} + if len(got) != len(want) { + t.Fatalf("应发送 %d 条,实际 %d(%v)", len(want), len(got), got) + } + for i := range want { + if got[i] != want[i] { + t.Errorf("第 %d 条应是 %q,实际 %q —— 同通道发送未保序(完整 %v)", i, want[i], got[i], got) + } + } +} + +// recordingDevice 记录 output 工具的 payload 顺序。 +type recordingDevice struct { + name string + caps agentIO.OutputCapability + tools []agentIO.ToolDef + onExec func(payload string) +} + +func (d *recordingDevice) Name() string { return d.name } +func (d *recordingDevice) Type() agentIO.DeviceType { return agentIO.DeviceOutput } +func (d *recordingDevice) Description() string { return "recording test device" } +func (d *recordingDevice) Tools() []agentIO.ToolDef { return d.tools } +func (d *recordingDevice) Execute(tool string, args map[string]interface{}) (interface{}, error) { + if d.onExec != nil { + p, _ := args["payload"].(string) + d.onExec(p) + } + return map[string]interface{}{"status": "sent"}, nil +} +func (d *recordingDevice) Start() error { return nil } +func (d *recordingDevice) Stop() error { return nil } +func (d *recordingDevice) OutputCapabilities() agentIO.OutputCapability { return d.caps } +func (d *recordingDevice) ChannelDef() agentIO.ChannelDef { return agentIO.ChannelDef{} } + +// newParallelAgent 建一个带 slowDevice 的 agent。 +func newParallelAgent(t *testing.T, sp agentAPI.Provider, names, safe []string, + entered, active, maxActive *int32, release chan struct{}) *Agent { + t.Helper() + a := New(AgentConfig{ + ID: "paragent", + Provider: sp, + ProviderManager: agentAPI.NewProviderManager(), + IO: agentIO.NewIOManager(), + StageHost: NewStageHost(), + }) + if err := a.io.RegisterDevice(&slowDevice{ + name: "slowdev", toolNames: names, safe: safe, + entered: entered, active: active, maxActive: maxActive, release: release, + }); err != nil { + t.Fatalf("注册设备失败: %v", err) + } + return a +} + +// 阶段 2c 判据的**并发版**(闭合此前"串行下测不出差别"的缺口)。 +// +// 此前两条 2c 判据在串行路径下无法区分「per-tool ctx」与「单槽」—— +// 变体验证(toolCtxFor 退回单槽)后仍然全绿。差别只在并发下显现。 +// 本判据在**真实并发批次**下断言: +// +// · 每个工具的 before_toolcall ctx 只带自己的 ToolCalls[0].Name +// · after_toolcall 读到的 Result 属于当前工具,不是批内另一个的 +// · 全程 -race 无数据竞争 +func TestBatchConcurrentEachToolSeesOwnContext(t *testing.T) { + names := []string{"k_a", "k_b", "k_c", "k_d"} + var entered, active, maxActive int32 + release := make(chan struct{}) + + sp := &batchProvider{responses: []*agentAPI.CompletionResponse{ + {ToolCalls: []agentAPI.ToolCall{ + {ID: "c1", Name: names[0], Arguments: map[string]interface{}{}}, + {ID: "c2", Name: names[1], Arguments: map[string]interface{}{}}, + {ID: "c3", Name: names[2], Arguments: map[string]interface{}{}}, + {ID: "c4", Name: names[3], Arguments: map[string]interface{}{}}, + }}, + {Content: "final"}, + }} + a := newParallelAgent(t, sp, names, names, &entered, &active, &maxActive, release) + + var mu sync.Mutex + beforeNames := map[string]string{} + crossTalk := map[string]string{} + + a.stageHost.RegisterStage(sdk.StageBeforeToolcall, func(ctx *sdk.StageContext) error { + if len(ctx.ToolCalls) != 1 { + t.Errorf("并发下 before_toolcall 的 ctx 应只带 1 个 ToolCall,实际 %d", len(ctx.ToolCalls)) + } + name := "" + if len(ctx.ToolCalls) > 0 { + name = ctx.ToolCalls[0].Name + } + mu.Lock() + if prev, dup := beforeNames[name]; dup { + crossTalk["dup:"+name] = "与 " + prev + " 撞名" + } + beforeNames[name] = name + mu.Unlock() + return nil + }) + a.stageHost.RegisterStage(sdk.StageAfterToolcall, func(ctx *sdk.StageContext) error { + if len(ctx.ToolResults) == 0 { + return nil + } + name := ctx.ToolResults[0].Name + res := fmt.Sprint(ctx.ToolResults[0].Result) + // 结果必须含**自己**的名字("ran:k_a"),否则读到的是别人的 + if !strings.Contains(res, name) { + mu.Lock() + crossTalk["after:"+name] = res + mu.Unlock() + } + return nil + }) + + done := make(chan struct{}) + go func() { + defer close(done) + if out := a.runTaskSteps(a.newTaskFrame("go", a.stageCtxFromInput("go", "", ""))); out != outcomeDone { + t.Errorf("runTaskSteps=%v", out) + } + }() + deadline := time.Now().Add(3 * time.Second) + for atomic.LoadInt32(&maxActive) < 2 && time.Now().Before(deadline) { + time.Sleep(5 * time.Millisecond) + } + close(release) + <-done + + if len(crossTalk) > 0 { + t.Errorf("并发下 StageContext 串味: %v", crossTalk) + } + if len(beforeNames) != len(names) { + t.Errorf("每个工具应各看到自己的 ctx,实际看到 %v(期望 %v)", beforeNames, names) + } +} diff --git a/internal/agent/core/task.go b/internal/agent/core/task.go index 04cf4c0..cb24351 100644 --- a/internal/agent/core/task.go +++ b/internal/agent/core/task.go @@ -25,6 +25,7 @@ import ( "fmt" "log" "strings" + "sync" "time" agentAPI "gitcode.com/JianFeeeee/HomeAgent/internal/agent/api" @@ -43,6 +44,13 @@ const ( StepPrepare Step = iota // StepLLM 轮次顶部(中断/占位)+ LLM 调用(含 provider 回退与重试)+ post_action。 StepLLM + // StepToolBatch 并发执行整批工具(阶段 2d)。仅当全批可并发时使用; + // 否则走 StepToolBegin/Exec/After 的串行路径。 + // + // 为何要有独立 step:runTaskSteps 是**单线程**驱动状态机的, + // 并发必须在一个 step 内 fan-out 并 join,否则"每步一个工具"的游标 + // 推进模型无法表达"一批同时跑"。 + StepToolBatch // StepToolBegin 取本批下一个工具,跑 before_toolcall;被拒/插件不健康则跳过。 StepToolBegin // StepToolExec 执行工具。**临界区**:副作用不可回滚,执行中不是安全点。 @@ -504,6 +512,8 @@ func (a *Agent) step(f *TaskFrame) stepOutcome { return a.stepPrepare(f) case StepLLM: return a.stepLLM(f) + case StepToolBatch: + return a.stepToolBatch(f) case StepToolBegin: return a.stepToolBegin(f) case StepToolExec: @@ -705,10 +715,127 @@ func (a *Agent) stepLLM(f *TaskFrame) stepOutcome { // 若某插件将来要改写 args,需在这里改为「回填后重写该条 assistant」。 f.assistantMsgIdx = -1 f.ToolIdx = 0 - f.Step = StepToolBegin + // 阶段 2d:全批可并发(且无需保序)⇒ 走并发 step。 + if f.batchRunnable(a) { + f.Step = StepToolBatch + } else { + f.Step = StepToolBegin + } return outcomeContinue } +// stepToolBatch **并发**执行整批工具(阶段 2d)。 +// +// 结构上是「fan-out → join → 顺序落消息」三段: +// +// 1. fan-out:每个工具一个 goroutine,各自跑 before_toolcall + 执行。 +// 每个 goroutine 只写**自己那份** toolCtxs[i](阶段 2c 的拆分正为此), +// 共享的 f.Msgs / f.ToolResults 在此期间**一律不碰**。 +// 2. join:等全部完成。 +// 3. 落消息:按 **索引顺序**(不是完成顺序)逐个跑 after_toolcall 与落消息。 +// 这一步必须串行——f.Msgs 是共享切片;而且按索引落能让模型读到的 +// 上下文顺序与模型自己发出的顺序一致。 +// +// 落消息为何不放进 goroutine:那样完成顺序不确定 ⇒ 同一批 tool 消息 +// 顺序随机 ⇒ 模型读到的因果关系与实际执行不符。 +func (a *Agent) stepToolBatch(f *TaskFrame) stepOutcome { + n := len(f.PendingTools) + results := make([]batchItemResult, n) + + // 保证 assistant 消息**先于**任何 tool 消息存在(协议要求)。 + f.ensureBatchAssistant() + + // ⚠️ 场景必须**在 fan-out 之前**解析一次:resolveTurnScenes 会把结果 + // 记进共享的 f.sceneDone / f.turnScene(memorypass.go:289), + // 在 N 个 goroutine 里各调一次既是数据竞争,也会各自触发一次 + // EnterSceneWithHint(重复计入场景强度——正是 task.go 里 + // sceneDone 注释警告的「多解析一次就多给场景加一次强度」)。 + for i := 0; i < n; i++ { + a.resolveTurnScenes(f, f.PendingTools[i].Name) + } + + var wg sync.WaitGroup + for i := 0; i < n; i++ { + wg.Add(1) + go func(idx int) { + defer wg.Done() + results[idx] = a.runOneTool(f, idx) + }(i) + } + wg.Wait() + + // 顺序收尾:after_toolcall / 裁剪 / 落消息 / 事件。 + for i := 0; i < n; i++ { + r := results[i] + f.CurTool = f.PendingTools[i] + f.CurToolPlugin = r.plugin + f.CurResult = r.text + f.CurRaw = r.raw + f.ToolIdx = i + f.ToolResults = append(f.ToolResults, ToolResultItem{Name: r.name, Output: r.text}) + if out := a.stepToolAfter(f); out != outcomeContinue { + return out + } + } + f.ToolIdx = n + f.Step = StepTurnEnd + return outcomeContinue +} + +// batchItemResult 是单个工具在并发阶段产出的结果。 +type batchItemResult struct { + name string + plugin string + text string + raw interface{} + // denied 表示被 before_toolcall 拒绝或插件不健康而未执行(已落 tool 消息)。 + denied bool +} + +// runOneTool 执行**单个**工具的「before_toolcall + 实际执行」,不碰共享状态。 +// +// 只写 toolCtxs[idx] 与返回值:f.Msgs / f.ToolResults / f.Cur* 全部由调用方 +// (stepToolBatch 的顺序收尾段,或串行路径的 stepToolBegin/Exec)负责。 +func (a *Agent) runOneTool(f *TaskFrame, idx int) batchItemResult { + tc := f.PendingTools[idx] + pluginName := a.resolveToolPlugin(tc.Name) + res := batchItemResult{name: tc.Name, plugin: pluginName} + + tctx := f.toolCtxFor(idx) + tctx.ToolCalls = []sdk.ToolCall{{ID: tc.ID, Name: tc.Name, Plugin: pluginName, Arguments: tc.Arguments}} + tctx.ToolResults = nil + + // before_toolcall(拒绝则不执行) + if a.runStage(sdk.StageBeforeToolcall, tctx) { + res.text = denialResultText(tctx, tc.Name) + res.denied = true + return res + } + tc.Arguments = tctx.ToolCalls[0].Arguments + + // 插件崩溃态:不执行 + if pluginName != "" && !a.pluginHealth.isHealthy(pluginName) { + res.text = fmt.Sprintf("插件 %s 处于崩溃状态,已跳过执行,等待自动恢复重载", pluginName) + res.denied = true + return res + } + + // 场景已在 stepToolBatch 的 fan-out **之前**解析完毕(避免竞争与重复计强度), + // 这里只读取结果。 + var turn memory.TurnScene + if f != nil && f.sceneDone { + turn = f.turnScene + } + outcome := a.executeToolCallOutcome(tc, f.OutputChannel, turn.Keys...) + res.text = outcome.Text + res.raw = outcome.Raw + tctx.ToolResults = []sdk.ToolResult{{ + CallID: tc.ID, Name: tc.Name, Plugin: pluginName, + Success: !isToolError(outcome.Raw), Result: outcome.Raw, + }} + return res +} + // stepToolBegin 取本批下一个工具;批已耗尽或发生中断则进入收尾。 func (a *Agent) stepToolBegin(f *TaskFrame) stepOutcome { if f.ToolIdx >= len(f.PendingTools) { diff --git a/internal/agent/io/channel.go b/internal/agent/io/channel.go index f0d8bd7..0005f4e 100644 --- a/internal/agent/io/channel.go +++ b/internal/agent/io/channel.go @@ -94,6 +94,9 @@ type ToolDef struct { Description string `json:"description"` Parameters map[string]interface{} `json:"parameters"` Handler ToolHandler `json:"-"` // 可选:插件工具的直接处理器,Device 通过 Execute() 分发 + // ParallelSafe 与 SDK 的 ToolDef.ParallelSafe 同义:声明此设备工具可被 + // **并发执行**。零值 false = 不可并发(保守默认,见 SDK 注释)。 + ParallelSafe bool `json:"parallel_safe,omitempty"` } type InputEvent struct { diff --git a/third_party/homeagent-sdk/sdk/plugin.go b/third_party/homeagent-sdk/sdk/plugin.go index 31832a9..f4f01ea 100644 --- a/third_party/homeagent-sdk/sdk/plugin.go +++ b/third_party/homeagent-sdk/sdk/plugin.go @@ -288,6 +288,16 @@ type ToolDef struct { // ""(默认 none) / RecallPolicyNone / RecallPolicyAuto。 // 默认 none:多数工具输出是噪声;需要「取回真实内容后据它召回」的工具(如 qq_get_message)应显式声明 auto。 RecallPolicy string `json:"recall_policy,omitempty"` + // ParallelSafe 声明此工具**可以被并发执行**(同一批多个 tool_call 同时跑)。 + // + // ⚠️ 零值 false 是刻意的:存量插件不改一行就得到**保守**行为 + //(整批串行),不会因升级被意外并发。声明它是**责任**而非特权。 + // + // 判据(三者皆满足才可并发): + // · handler 自身线程安全(不持有跨调用的可变状态) + // · 不与同批其它工具争抢同一资源(SQLite 写、设备、同一输出通道) + // · 执行顺序无关(顺序敏感的工具应留 false,由内核保序) + ParallelSafe bool `json:"parallel_safe,omitempty"` } // IOInjector provides methods for injecting input and interrupts into the agent pipeline.