feat(a2a/acp): 会话历史查询 session.get + 客户端 session_id 透传

a2a:
- tasks.get 从空壳改为按 session_id 返回会话内近 N 条消息(默认10);
  新增 session.get 别名同语义
- a2a_query(出站)接受 session_id 参数透传给目标 agent,
  响应回显 session_id + 延续提示
- A2AParams/A2AResult 加 session_id 字段;工具描述补说明

acp:
- 新增 JSON-RPC method session/get:按 session_id 返回近 N 条消息
- acp_query(出站)接受 session_id 透传给 session/new,
  响应回显 + 延续提示
- params 结构体加 limit 字段

端到端验证(回环本机):
  a2a tasks.send→建会话;session.get→返回[user/agent]交替消息列表
  同 session 第二轮延续上下文正确(记数字→答数字)
  acp session/new + session/get 同样通过
This commit is contained in:
JianFeeeee
2026-08-26 19:30:32 +08:00
parent f12521d033
commit 0823c3a1e0
4 changed files with 151 additions and 13 deletions

View File

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

View File

@ -77,6 +77,7 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error {
"properties": map[string]interface{}{ "properties": map[string]interface{}{
"agent_url": map[string]interface{}{"type": "string", "description": "目标 Agent 的 A2A 端点 URL"}, "agent_url": map[string]interface{}{"type": "string", "description": "目标 Agent 的 A2A 端点 URL"},
"query": map[string]interface{}{"type": "string", "description": "发送给目标 Agent 的文本查询"}, "query": map[string]interface{}{"type": "string", "description": "发送给目标 Agent 的文本查询"},
"session_id": map[string]interface{}{"type": "string", "description": "可选。上次调用返回的 session_id,传入可延续与该 agent 的多轮对话上下文"},
"timeout": map[string]interface{}{"type": "integer", "description": "超时时间(秒),默认 60"}, "timeout": map[string]interface{}{"type": "integer", "description": "超时时间(秒),默认 60"},
}, },
"required": []string{"agent_url", "query"}, "required": []string{"agent_url", "query"},
@ -169,6 +170,43 @@ func truncateRunes(s string, n int) string {
return string(r[:n]) + "..." return string(r[:n]) + "..."
} }
// sessionMessages 返回指定会话的近 limit 条消息(时间正序),
// 会话不存在返回 nil。消息格式 [{role, text, ts}]。
func (p *Plugin) sessionMessages(sessionID string, limit int) []map[string]interface{} {
p.sessMu.Lock()
sess := p.sessions[sessionID]
var hist []string
var lastUsed time.Time
if sess != nil {
hist = append([]string{}, sess.History...)
lastUsed = sess.LastUsed
}
p.sessMu.Unlock()
if sess == nil {
return nil
}
_ = lastUsed
// History 交替 [user, agent, user, agent...],取末尾 limit 条,保持时间正序
start := 0
if len(hist) > limit {
start = len(hist) - limit
}
msgs := make([]map[string]interface{}, 0, len(hist)-start)
for i := start; i < len(hist); i++ {
role, text := "user", hist[i]
if after, ok := strings.CutPrefix(text, "用户: "); ok {
role, text = "user", after
} else if after, ok := strings.CutPrefix(text, "助手: "); ok {
role, text = "agent", after
}
msgs = append(msgs, map[string]interface{}{
"role": role,
"text": text,
})
}
return msgs
}
func (p *Plugin) stopServer() { func (p *Plugin) stopServer() {
p.srvMu.Lock() p.srvMu.Lock()
defer p.srvMu.Unlock() defer p.srvMu.Unlock()
@ -248,6 +286,7 @@ func (p *Plugin) handleIncomingA2A(w http.ResponseWriter, r *http.Request) {
Params struct { Params struct {
Query string `json:"query,omitempty"` Query string `json:"query,omitempty"`
SessionID string `json:"session_id,omitempty"` SessionID string `json:"session_id,omitempty"`
Limit int `json:"limit,omitempty"`
Message *struct { Message *struct {
Role string `json:"role"` Role string `json:"role"`
Parts []struct { Parts []struct {
@ -330,11 +369,37 @@ func (p *Plugin) handleIncomingA2A(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "application/json") w.Header().Set("Content-Type", "application/json")
json.NewEncoder(w).Encode(resp) json.NewEncoder(w).Encode(resp)
case "tasks.get": case "tasks.get", "session.get":
// 按 session_id 返回会话内近 N 条消息(默认 10 条)。
sessionID := strings.TrimSpace(req.Params.SessionID)
if sessionID == "" {
sessionID = strings.TrimSpace(req.Params.Query)
}
limit := 10
if req.Params.Limit > 0 && req.Params.Limit <= 100 {
limit = req.Params.Limit
}
msgs := p.sessionMessages(sessionID, limit)
if msgs == nil {
w.Header().Set("Content-Type", "application/json")
json.NewEncoder(w).Encode(map[string]interface{}{
"jsonrpc": "2.0", "id": req.ID,
"result": map[string]interface{}{
"session_id": sessionID,
"status": "not_found",
"messages": []interface{}{},
},
})
return
}
w.Header().Set("Content-Type", "application/json") w.Header().Set("Content-Type", "application/json")
json.NewEncoder(w).Encode(map[string]interface{}{ json.NewEncoder(w).Encode(map[string]interface{}{
"jsonrpc": "2.0", "id": req.ID, "jsonrpc": "2.0", "id": req.ID,
"result": map[string]interface{}{"id": req.Params.Query, "status": "unknown"}, "result": map[string]interface{}{
"session_id": sessionID,
"status": "completed",
"messages": msgs,
},
}) })
default: default:
@ -378,9 +443,10 @@ type A2ARequest struct {
} }
type A2AParams struct { type A2AParams struct {
Query string `json:"query,omitempty"` Query string `json:"query,omitempty"`
Message *A2AMessage `json:"message,omitempty"` SessionID string `json:"session_id,omitempty"`
TaskID string `json:"id,omitempty"` Message *A2AMessage `json:"message,omitempty"`
TaskID string `json:"id,omitempty"`
} }
type A2AResponse struct { type A2AResponse struct {
@ -393,6 +459,7 @@ type A2AResponse struct {
type A2AResult struct { type A2AResult struct {
TaskID string `json:"id,omitempty"` TaskID string `json:"id,omitempty"`
Status string `json:"status,omitempty"` Status string `json:"status,omitempty"`
SessionID string `json:"session_id,omitempty"`
Message *A2AMessage `json:"message,omitempty"` Message *A2AMessage `json:"message,omitempty"`
AgentCard *A2AAgentCard `json:"agent_card,omitempty"` AgentCard *A2AAgentCard `json:"agent_card,omitempty"`
} }
@ -458,6 +525,7 @@ func (p *Plugin) handleA2ADiscover(args map[string]interface{}) (interface{}, er
func (p *Plugin) handleA2AQuery(args map[string]interface{}) (interface{}, error) { func (p *Plugin) handleA2AQuery(args map[string]interface{}) (interface{}, error) {
agentURL, _ := args["agent_url"].(string) agentURL, _ := args["agent_url"].(string)
query, _ := args["query"].(string) query, _ := args["query"].(string)
sessionID, _ := args["session_id"].(string) // 可选:延续对方会话
timeoutSec := 60 timeoutSec := 60
if v, ok := args["timeout"].(float64); ok && v > 0 { if v, ok := args["timeout"].(float64); ok && v > 0 {
timeoutSec = int(v) timeoutSec = int(v)
@ -479,7 +547,8 @@ func (p *Plugin) handleA2AQuery(args map[string]interface{}) (interface{}, error
ID: fmt.Sprintf("a2a_%d", time.Now().UnixNano()), ID: fmt.Sprintf("a2a_%d", time.Now().UnixNano()),
Method: "tasks.send", Method: "tasks.send",
Params: A2AParams{ Params: A2AParams{
Message: &A2AMessage{Role: "user", Parts: []A2APart{{Text: query, Type: "text"}}}, SessionID: sessionID,
Message: &A2AMessage{Role: "user", Parts: []A2APart{{Text: query, Type: "text"}}},
}, },
} }
@ -518,10 +587,18 @@ func (p *Plugin) handleA2AQuery(args map[string]interface{}) (interface{}, error
replyText = strings.TrimSpace(replyText) replyText = strings.TrimSpace(replyText)
} }
return map[string]interface{}{ result := map[string]interface{}{
"task_id": a2aResp.Result.TaskID, "status": a2aResp.Result.Status, "task_id": a2aResp.Result.TaskID, "status": a2aResp.Result.Status,
"response": replyText, "response": replyText,
}, nil }
if a2aResp.Result.SessionID != "" || sessionID != "" {
result["session_id"] = a2aResp.Result.SessionID
if result["session_id"] == "" {
result["session_id"] = sessionID
}
result["note"] = "延续会话:下次调用传此 session_id 可保持上下文"
}
return result, nil
} }
// ---- Management Handlers ---- // ---- Management Handlers ----

View File

@ -2,7 +2,7 @@
"name": "acp", "name": "acp",
"name_zh": "ACP 代理通信", "name_zh": "ACP 代理通信",
"name_en": "ACP Agent Client Protocol", "name_en": "ACP Agent Client Protocol",
"version": "1.1.0", "version": "1.2.0",
"description": "Agent Client Protocol 通信插件:充当 ACP 服务端接受其他 Agent 的任务请求,同时提供客户端工具向远程 ACP Agent(如 opencode)发起会话并读取回复", "description": "Agent Client Protocol 通信插件:充当 ACP 服务端接受其他 Agent 的任务请求,同时提供客户端工具向远程 ACP Agent(如 opencode)发起会话并读取回复",
"author": "HomeAgent", "author": "HomeAgent",
"entry": "plugin.so", "entry": "plugin.so",

View File

@ -71,6 +71,7 @@ func (p *Plugin) Start(s *sdk.PluginSDK) error {
"properties": map[string]interface{}{ "properties": map[string]interface{}{
"server_url": map[string]interface{}{"type": "string", "description": "目标 ACP 服务端地址(如 http://127.0.0.1:13000)"}, "server_url": map[string]interface{}{"type": "string", "description": "目标 ACP 服务端地址(如 http://127.0.0.1:13000)"},
"prompt": map[string]interface{}{"type": "string", "description": "发送给目标 Agent 的任务描述"}, "prompt": map[string]interface{}{"type": "string", "description": "发送给目标 Agent 的任务描述"},
"session_id": map[string]interface{}{"type": "string", "description": "可选。上次调用返回的 session_id,传入可延续与该 agent 的多轮对话上下文"},
"timeout": map[string]interface{}{"type": "integer", "description": "等待回复超时(秒),默认 120"}, "timeout": map[string]interface{}{"type": "integer", "description": "等待回复超时(秒),默认 120"},
}, },
"required": []string{"server_url", "prompt"}, "required": []string{"server_url", "prompt"},
@ -184,6 +185,7 @@ func (p *Plugin) handleSessionPost(w http.ResponseWriter, r *http.Request) {
Text string `json:"text"` Text string `json:"text"`
} `json:"request,omitempty"` } `json:"request,omitempty"`
SessionID string `json:"session_id,omitempty"` SessionID string `json:"session_id,omitempty"`
Limit int `json:"limit,omitempty"`
Final bool `json:"final,omitempty"` Final bool `json:"final,omitempty"`
} `json:"params,omitempty"` } `json:"params,omitempty"`
} }
@ -256,6 +258,59 @@ func (p *Plugin) handleSessionPost(w http.ResponseWriter, r *http.Request) {
}, },
}) })
case "session/get":
// 按 session_id 返回会话内近 N 条消息(默认 10 条,时间正序)
sid := req.Params.SessionID
p.mu.RLock()
st := p.sessions[sid]
var hist []string
if st != nil {
hist = append([]string{}, st.History...)
}
p.mu.RUnlock()
if st == nil {
w.Header().Set("Content-Type", "application/json")
json.NewEncoder(w).Encode(map[string]interface{}{
"jsonrpc": "2.0", "id": req.ID,
"result": map[string]interface{}{
"session_id": sid,
"status": "not_found",
"messages": []interface{}{},
},
})
return
}
limit := 10
if req.Params.Limit > 0 && req.Params.Limit <= 100 {
limit = req.Params.Limit
}
start := 0
if len(hist) > limit {
start = len(hist) - limit
}
msgs := make([]map[string]interface{}, 0, len(hist)-start)
for i := start; i < len(hist); i++ {
role, text := "user", hist[i]
if after, ok := strings.CutPrefix(text, "用户: "); ok {
role, text = "user", after
} else if after, ok := strings.CutPrefix(text, "助手: "); ok {
role, text = "agent", after
}
msgs = append(msgs, map[string]interface{}{
"role": role,
"text": text,
})
}
w.Header().Set("Content-Type", "application/json")
json.NewEncoder(w).Encode(map[string]interface{}{
"jsonrpc": "2.0", "id": req.ID,
"result": map[string]interface{}{
"session_id": sid,
"status": "completed",
"messages": msgs,
},
})
case "session/update": case "session/update":
sid := req.Params.SessionID sid := req.Params.SessionID
p.mu.Lock() p.mu.Lock()
@ -391,6 +446,7 @@ func (p *Plugin) handleAcpQuery(args map[string]interface{}) (interface{}, error
if prompt == "" { if prompt == "" {
return map[string]interface{}{"error": "prompt 不能为空"}, nil return map[string]interface{}{"error": "prompt 不能为空"}, nil
} }
sessionID, _ := args["session_id"].(string) // 可选:延续对方会话
timeoutSec := 120 timeoutSec := 120
if v, ok := args["timeout"].(float64); ok && v > 0 { if v, ok := args["timeout"].(float64); ok && v > 0 {
timeoutSec = int(v) timeoutSec = int(v)
@ -399,12 +455,16 @@ func (p *Plugin) handleAcpQuery(args map[string]interface{}) (interface{}, error
endpoint := serverURL + "/api/session" endpoint := serverURL + "/api/session"
client := &http.Client{Timeout: time.Duration(timeoutSec) * time.Second} client := &http.Client{Timeout: time.Duration(timeoutSec) * time.Second}
params := map[string]interface{}{
"request": map[string]interface{}{"text": prompt},
}
if sessionID != "" {
params["session_id"] = sessionID
}
newBody, _ := json.Marshal(map[string]interface{}{ newBody, _ := json.Marshal(map[string]interface{}{
"jsonrpc": "2.0", "id": "acp-" + fmt.Sprintf("%d", time.Now().UnixNano()), "jsonrpc": "2.0", "id": "acp-" + fmt.Sprintf("%d", time.Now().UnixNano()),
"method": "session/new", "method": "session/new",
"params": map[string]interface{}{ "params": params,
"request": map[string]interface{}{"text": prompt},
},
}) })
req, _ := http.NewRequest("POST", endpoint, bytes.NewReader(newBody)) req, _ := http.NewRequest("POST", endpoint, bytes.NewReader(newBody))
@ -469,6 +529,7 @@ func (p *Plugin) handleAcpQuery(args map[string]interface{}) (interface{}, error
"session_id": sid, "session_id": sid,
"status": "completed", "status": "completed",
"reply": replyText, "reply": replyText,
"note": "延续会话:下次调用传此 session_id 可保持上下文",
}, nil }, nil
} }