fix(a2a/acp): 回复闭环 + 会话延续 + 同步注入(不再抢占打断)

【问题】
1. 入站请求用 InjectInterruptText 抢占打断当前对话,立即回 202 submitted,
   agent 的回复 emit 到未注册的 channel(a2a/acp)→ 请求方永远拿不到回复文本,
   只能干等超时。(acp 的 session.Replying 从未被填真回复 → SSE 永远 "(未收到回复)")
2. 无法指定/延续 session:a2a 无 session 概念;acp session/new 每次新建、
   不接受调用方 session_id,多轮上下文断裂。
3. a2a/acp 通道未注册为输出设备 → emitResponse 的回复无落点,
   output_list_channels 也不可见 → agent 困惑"回复该发到哪"。

【修复】
- 入站改用 SDK InjectInputSync 同步注入:阻塞等待 agent 处理完成,
  直接把最终回复文本返回给 HTTP 请求方(不再回 202)。
  这是 A2A/ACP 协议的合理形态——客户端控制超时,服务端同步返回。
- session 支持:a2a tasks.send 与 acp session/new 均接受 params.session_id,
  指定则延续已有会话(拼上下文前缀),不指定则新建并返回 session_id。
  单会话保留最多 10 轮历史防膨胀;30min GC 清理 2h 未用会话。
- a2a/acp 注册为输出通道(RegisterOutputChannel):回复有落点,
  output_list_channels 可见,agent 可主动 output_send 推消息。
- 注入提示词明确"直接以文本回复即可,无需调 output_send"——
  agent 不再把回复走 queued 入队而返回干净文本。

端到端验证:
  单轮:status=completed, reply=真实回复文本(非 submitted)
  多轮:同 session_id 第二轮准确复述第一轮问题(上下文生效)
  a2a v1.2.0 / acp v1.1.0 安装 config_kept=true
This commit is contained in:
JianFeeeee
2026-08-26 19:16:46 +08:00
parent a82b5b1626
commit 14aa0c880b
4 changed files with 705 additions and 8 deletions

View File

@ -2,7 +2,7 @@
"name": "a2a",
"name_zh": "A2A 代理通信",
"name_en": "A2A Agent Communication",
"version": "1.1.0",
"version": "1.2.0",
"description": "Agent-to-Agent 协议通信插件,支持双向 A2A 通信:可查询其他 Agent 并回复其请求。提供 HTTP 服务端暴露本 Agent 能力。",
"author": "HomeAgent",
"entry": "plugin.so",

View File

@ -21,15 +21,47 @@ type Plugin struct {
srvMu sync.Mutex
server *http.Server
serverAddr string
// 会话表session_id → 上下文前缀。A2A 无状态协议下由插件侧维护
// 多轮上下文:同 session 的后续请求会把之前的对话拼进注入文本。
sessMu sync.Mutex
sessions map[string]*a2aSession
}
// a2aSession 记录一个会话的轮次历史,用于延续上下文。
type a2aSession struct {
ID string
History []string // 轮次文本 [user1, agent1, user2, agent2, ...]
LastUsed time.Time
}
// maxSessionTurns 单会话保留的最大轮次对数(防上下文无限膨胀)。
const maxSessionTurns = 10
// sessionGCPeriod 会话过期清理周期;超过 2 小时未用的会话回收。
const sessionGCPeriod = 30 * time.Minute
func (p *Plugin) Name() string { return p.name }
func (p *Plugin) Start(s *sdk.PluginSDK) error {
s.SetAutoRestart(true)
p.sdk = s
p.sessions = make(map[string]*a2aSession)
tp := p.name + "_"
// 注册自身为输出通道agent 回复 emit 到本通道时有落点,
// 且 output_list_channels 可见agent 能主动向 a2a 会话推送消息)。
if err := s.RegisterOutputChannel(p.name, 1, "A2A Agent 互联通道(外部 agent 查询的回复由此返回)", sdk.ChannelDef{}, func(args map[string]interface{}) (interface{}, error) {
payload, _ := args["payload"].(string)
log.Printf("[%s] channel output: %s", p.name, truncateRunes(payload, 120))
return map[string]interface{}{"status": "ok"}, nil
}); err != nil {
log.Printf("[%s] register output channel: %v", p.name, err)
}
// 会话 GC后台周期回收长期不用的会话
go p.sessionGCLoop()
s.Settings().RegisterDef(sdk.ConfigDef{
Key: "listen", Default: "127.0.0.1:12000",
Type: "string", DisplayName: "监听地址",
@ -114,6 +146,29 @@ func (p *Plugin) Stop() error {
return nil
}
// sessionGCLoop 周期清理超时会话。
func (p *Plugin) sessionGCLoop() {
ticker := time.NewTicker(sessionGCPeriod)
defer ticker.Stop()
for range ticker.C {
p.sessMu.Lock()
for id, sess := range p.sessions {
if time.Since(sess.LastUsed) > 2*time.Hour {
delete(p.sessions, id)
}
}
p.sessMu.Unlock()
}
}
func truncateRunes(s string, n int) string {
r := []rune(s)
if len(r) <= n {
return s
}
return string(r[:n]) + "..."
}
func (p *Plugin) stopServer() {
p.srvMu.Lock()
defer p.srvMu.Unlock()
@ -191,7 +246,8 @@ func (p *Plugin) handleIncomingA2A(w http.ResponseWriter, r *http.Request) {
ID string `json:"id"`
Method string `json:"method"`
Params struct {
Query string `json:"query,omitempty"`
Query string `json:"query,omitempty"`
SessionID string `json:"session_id,omitempty"`
Message *struct {
Role string `json:"role"`
Parts []struct {
@ -215,19 +271,60 @@ func (p *Plugin) handleIncomingA2A(w http.ResponseWriter, r *http.Request) {
}
queryText = strings.TrimSpace(queryText)
}
// Inject into agent pipeline via interrupt (preempt current processing) or direct input
if queryText != "" {
p.sdk.InjectInterruptText("a2a", "webui", fmt.Sprintf("[来自A2A Agent的查询]\n%s", queryText))
if queryText == "" {
http.Error(w, "query/message.text required", http.StatusBadRequest)
return
}
// Respond with task accepted
// 会话:调用方可指定 session_id 延续多轮上下文;不指定则新建。
sessionID := strings.TrimSpace(req.Params.SessionID)
injectText := queryText
p.sessMu.Lock()
if sessionID != "" {
sess := p.sessions[sessionID]
if sess == nil {
sess = &a2aSession{ID: sessionID, LastUsed: time.Now()}
p.sessions[sessionID] = sess
}
sess.LastUsed = time.Now()
// 有历史则把上下文拼在前面(截尾防爆量)
if len(sess.History) > 0 {
ctxText := strings.Join(sess.History, "\n")
injectText = "[对话上下文]\n" + ctxText + "\n[本轮输入]\n" + queryText
}
} else {
sessionID = fmt.Sprintf("a2a_%d", time.Now().UnixNano())
p.sessions[sessionID] = &a2aSession{ID: sessionID, LastUsed: time.Now()}
}
p.sessMu.Unlock()
// 同步注入:阻塞等待 agent 处理完成拿回复(不再抢占打断、
// 也不再回 202 让请求方永远等不到结果。HTTP 超时由调用方控制。
reply := p.sdk.InjectInputSync(p.name, p.name,
fmt.Sprintf("[来自A2A Agent的查询 session=%s]\n%s\n[注意] 请直接以文本回复本查询,不要调用 output_send__%s——你的最终文本回复会被系统自动返回给请求方。", sessionID, injectText, p.name))
// 回复写回会话历史(下一轮作为上下文)
p.sessMu.Lock()
if sess := p.sessions[sessionID]; sess != nil {
sess.History = append(sess.History, "用户: "+queryText, "助手: "+reply)
if len(sess.History) > maxSessionTurns*2 {
sess.History = sess.History[len(sess.History)-maxSessionTurns*2 :]
}
sess.LastUsed = time.Now()
}
p.sessMu.Unlock()
resp := map[string]interface{}{
"jsonrpc": "2.0",
"id": req.ID,
"result": map[string]interface{}{
"id": fmt.Sprintf("task_%d", time.Now().UnixNano()),
"status": "submitted",
"status": "completed",
"session_id": sessionID,
"message": map[string]interface{}{
"role": "agent",
"parts": []map[string]string{{"type": "text", "text": reply}},
},
},
}
w.Header().Set("Content-Type", "application/json")