已合并
feat: A2A网关接入实现 - VersatileA2AGateway适配器 #377
feat: A2A网关接入实现 - VersatileA2AGateway适配器 #377
已合并
xiongxing创建于 7月13日
共 12 个文件变更+769-34
@@ -30,6 +30,8 @@ from __future__ import annotations
30import time30import time
31from typing import Any31from typing import Any
32 32 
33+from loguru import logger
34+ 
33# ════════════════════════════════════════════════════════════════════35# ════════════════════════════════════════════════════════════════════
34# 事件类型白名单36# 事件类型白名单
35# ════════════════════════════════════════════════════════════════════37# ════════════════════════════════════════════════════════════════════
@@ -148,7 +150,7 @@ def wrap_workflow_event(
148 if event_kind not in ("message", "end"):150 if event_kind not in ("message", "end"):
149 # 防御:只接受观察到的两个值;其他值也允许透传,避免阻塞将来协议扩展151 # 防御:只接受观察到的两个值;其他值也允许透传,避免阻塞将来协议扩展
150 pass152 pass
151- return {153+ result = {
152 "success": True,154 "success": True,
153 "agent_id": agent_id,155 "agent_id": agent_id,
154 "conversation_id": conversation_id,156 "conversation_id": conversation_id,
@@ -158,6 +160,11 @@ def wrap_workflow_event(
158 "data": data,160 "data": data,
159 },161 },
160 }162 }
163+ logger.debug(
164+ f"[ResponseWrapper] 输出 workflow_event: event={event_kind}, "
165+ f"conv={conversation_id}, data={str(data)[:200]}"
166+ )
167+ return result
161 168 
162 169 
163# ════════════════════════════════════════════════════════════════════170# ════════════════════════════════════════════════════════════════════
@@ -204,7 +211,7 @@ def wrap_sub_task_event(
204 conversation_id: str,211 conversation_id: str,
205 elapsed: float,212 elapsed: float,
206) -> dict[str, Any]:213) -> dict[str, Any]:
207- return {214+ result = {
208 "success": True,215 "success": True,
209 "agent_id": agent_id,216 "agent_id": agent_id,
210 "conversation_id": conversation_id,217 "conversation_id": conversation_id,
@@ -223,6 +230,11 @@ def wrap_sub_task_event(
223 ),230 ),
224 },231 },
225 }232 }
233+ logger.debug(
234+ f"[ResponseWrapper] 输出 sub_task_event: sub_task_path={sub_task_path}, "
235+ f"node_kind={node_kind}, conv={conversation_id}"
236+ )
237+ return result
226 238 
227 239 
228def wrap_error(240def wrap_error(
@@ -532,6 +532,10 @@ class RemoteAgentHandler:
532 "data": data,532 "data": data,
533 },533 },
534 }534 }
535+ logger.info(
536+ f"[RemoteAgentHandler] SubTaskEvent 盖章: sub_task_path={sub_task_path}, "
537+ f"node_kind={node_kind}, data={str(data)[:200]}"
538+ )
535 await turn_ctx.event_queue.enqueue_event(539 await turn_ctx.event_queue.enqueue_event(
536 dict_to_a2a(event, turn_ctx.task_id, turn_ctx.conv_id)540 dict_to_a2a(event, turn_ctx.task_id, turn_ctx.conv_id)
537 )541 )
@@ -36,6 +36,11 @@ ADAPTER_FASTAPI_WORKERS=1
36# ── Versatile 工作流最终结果节点名 ───────────36# ── Versatile 工作流最终结果节点名 ───────────
37va_workflow_result_node=37va_workflow_result_node=
38 38 
39+# ── A2A 网关模式切换 ───────────────────────────
40+# a2a_gateway(默认):走 A2A 1.0 网关协议(VersatileA2AGateway)
41+# workflow:直连低码工作流(VersatileWorkflow)
42+VA_WORKFLOW_ADAPTER_TYPE=a2a_gateway
43+ 
39# ── Log ───────────────────────────────────────────────────────────────────────44# ── Log ───────────────────────────────────────────────────────────────────────
40ADAPTER_LOG_LEVEL=DEBUG45ADAPTER_LOG_LEVEL=DEBUG
41ADAPTER_LOG_FILE=logs/versatile_adapter.log46ADAPTER_LOG_FILE=logs/versatile_adapter.log
@@ -74,6 +74,11 @@ class A2aVersatileExecutor(AgentExecutor):
74 logger.info(74 logger.info(
75 f"[A2aVA] execute: conv_id={context.context_id}, task_id={context.task_id}"75 f"[A2aVA] execute: conv_id={context.context_id}, task_id={context.task_id}"
76 )76 )
77+ logger.info(
78+ f"[A2aVA] 收到请求: conv_id={context.context_id}, task_id={context.task_id}, "
79+ f"body={str(runner_kw.get('body', {}))[:200]}, headers={runner_kw.get('headers', {})}, "
80+ f"trace_id={runner_kw.get('trace_id', '')}"
81+ )
atomgit-bot
atomgit-botatomgit-bot7月13日

🟡 Medium Priority

executor.py 第 77-81 行新增的 logger.info(...) 将完整的 runner_kw.get('headers', {}) 以 INFO 级别写入日志,未对 token、authorization、cookie 等敏感 header 做任何脱敏处理。body 字段有 [:200] 截断,但 headers 完全没有截断或脱敏。

对比 versatile_proxy.py 中 _log_request 使用 _SENSITIVE_HEADERS 集合对 curl 命令中的敏感 header 做 *** 替换,且请求头日志使用 DEBUG 级别。此处以 INFO 级别(生产环境通常开启)记录完整 headers,若上游 a2a_service 传入的 headers 中包含 [REDACTED]、session cookie 等,将直接写入日志文件。

风险:日志被非授权人员/系统访问时泄露凭据。

建议:将日志级别降为 DEBUG,或对 headers 做脱敏处理(过滤 token/authorization/cookie 等敏感 key),或只记录 header keys 而非 values。

likedislike
不准确?
77 updater = await self._setup_task(context, event_queue)82 updater = await self._setup_task(context, event_queue)
78 terminal_sent = False83 terminal_sent = False
79 84 
@@ -149,13 +154,12 @@ class A2aVersatileExecutor(AgentExecutor):
149 async def cancel(self, context: RequestContext, event_queue: EventQueue) -> None:154 async def cancel(self, context: RequestContext, event_queue: EventQueue) -> None:
150 conv_id = context.context_id155 conv_id = context.context_id
151 task_id = context.task_id156 task_id = context.task_id
152- input_data = self._build_first_input(context.message)157+ self._build_first_input(context.message)
153- agent_id = input_data.get("agent_id", "")
154 158 
155 updater = TaskUpdater(event_queue, task_id, conv_id)159 updater = TaskUpdater(event_queue, task_id, conv_id)
156 await updater.cancel()160 await updater.cancel()
157 logger.info(161 logger.info(
158- f"[A2aVA] cancel: conv_id={conv_id}, task_id={task_id}, agent_id={agent_id}"162+ f"[A2aVA] cancel: conv_id={conv_id}, task_id={task_id}"
159 )163 )
160 164 
161 async def _emit_proxy_artifact(165 async def _emit_proxy_artifact(
@@ -304,7 +308,6 @@ class A2aVersatileExecutor(AgentExecutor):
304 def _extract_logging_context(input_data: dict, conv_id: str) -> dict:308 def _extract_logging_context(input_data: dict, conv_id: str) -> dict:
305 return {309 return {
306 "trace_id": input_data.get("trace_id", ""),310 "trace_id": input_data.get("trace_id", ""),
307- "agent_id": input_data.get("agent_id", ""),
308 "conv_id": conv_id,311 "conv_id": conv_id,
309 }312 }
310 313 
@@ -318,4 +321,5 @@ class A2aVersatileExecutor(AgentExecutor):
318 "body": input_data.get("body", {}),321 "body": input_data.get("body", {}),
319 "headers": input_data.get("headers", {}),322 "headers": input_data.get("headers", {}),
320 "params": input_data.get("params", {}),323 "params": input_data.get("params", {}),
324+ "trace_id": input_data.get("trace_id", ""),
321 }325 }
@@ -0,0 +1,333 @@
1+# coding: utf-8
2+# Copyright (c) Huawei Technologies Co., Ltd. 2026-2026. All rights reserved
3+ 
4+"""
5+VersatileA2AGateway - A2A 网关协议适配器。
6+ 
7+将 VA 内部统一入参转换为 A2A 1.0 JSON-RPC SendStreamingMessage 请求,
8+解析 A2A 1.0 SSE 响应(taskArtifactUpdate / taskStatusUpdate),
9+映射为 VA 统一 AdapterEvent。
10+ 
11+data_proxy 输出格式与低码工作流一致({"event":"message","data":{"text":"..."}}),
12+上游 a2a_service 和前端无需适配。
13+"""
14+from __future__ import annotations
15+ 
16+import json
17+import uuid
18+from typing import Optional
19+ 
20+from loguru import logger
21+ 
22+from adapters.versatile_proxy import VersatileProxy, VersatileStreamCtx
23+from event.events import (
24+ AdapterEvent,
25+ DataProxyContent,
26+ ExecutionCompletedContent,
27+ ExecutionInputRequiredContent,
28+)
29+ 
30+ 
31+class VersatileA2AGateway(VersatileProxy):
32+ """A2A 网关协议适配器。
33+ 
34+ 职责:
35+ 1. 拼 URL:{a2a_gateway_base}/a2a/{agent_card_name}
36+ 2. 构造 Header:token(必选) + userId(必选) + B3/X-Biz-Tag(可选透传)
37+ 3. 构造 A2A 1.0 SendStreamingMessage 请求体(configuration 由 VA 固定生成)
38+ 4. 解析 A2A 1.0 SSE:taskArtifactUpdate / taskStatusUpdate
39+ 5. 映射为 VA 统一 AdapterEvent,data_proxy 格式与低码工作流一致
40+ """
41+ 
42+ def __init__(
43+ self,
44+ a2a_gateway_base: str,
45+ agent_card_name: str,
46+ token: str,
47+ url_template: str,
48+ timeout: int = 600,
49+ headers_template: Optional[dict] = None,
50+ forward_header_whitelist: Optional[set[str]] = None,
51+ workflow_result_node: Optional[str] = None,
52+ ) -> None:
53+ super().__init__(url_template, timeout, headers_template, forward_header_whitelist)
54+ self._a2a_gateway_base = a2a_gateway_base
55+ self._agent_card_name = agent_card_name
56+ self._token = token
57+ # workflow_result_node 保留但不使用,A2A Gateway 模式按方案 D 提取结果
58+ self._workflow_result_node = workflow_result_node
59+ # dispatch_stream 入口设置,提前初始化以支持直接调用钩子方法测试
60+ self._conv_id: str = ""
61+ self._trace_id: str = ""
62+ self._passed_headers: dict = {}
63+ self._cached_user_id: str = ""
64+ 
65+ # ── 钩子方法覆盖 ──────────────────────────────────────────
66+ 
67+ def _build_url(self, conv_id: str) -> str:
68+ url = self._url_template.format(
69+ a2a_gateway_base=self._a2a_gateway_base,
70+ agent_card_name=self._agent_card_name,
71+ )
72+ logger.info(f"[VersatileA2AGateway] URL: {url}, conv_id={conv_id}")
73+ return url
74+ 
75+ def _build_headers(self, headers: Optional[dict] = None) -> dict:
76+ merged = super()._build_headers(headers)
77+ 
78+ # 必选 Header
79+ merged["token"] = self._token
80+ self._cached_user_id = self._extract_user_id()
81+ merged["userId"] = self._cached_user_id
82+ logger.debug(f"[VersatileA2AGateway] 必选 Header: token=***, userId={merged['userId']}")
83+ 
84+ # 可选 Header(上游有则透传,无则不写)
85+ if self._trace_id or self._passed_headers.get("x-b3-traceid"):
86+ merged["X-B3-TraceId"] = self._trace_id or self._passed_headers.get("x-b3-traceid", "")
87+ if self._passed_headers.get("x-b3-parentspanid"):
88+ merged["X-B3-ParentSpanId"] = self._passed_headers["x-b3-parentspanid"]
89+ merged["X-B3-SpanId"] = uuid.uuid4().hex[:16]
90+ if self._passed_headers.get("x-b3-sampled"):
91+ merged["X-B3-Sampled"] = self._passed_headers["x-b3-sampled"]
92+ if self._passed_headers.get("x-biz-tag"):
93+ merged["X-Biz-Tag"] = self._passed_headers["x-biz-tag"]
94+ 
95+ return merged
96+ 
97+ def _build_request_body(self, body: dict) -> dict:
98+ inputs = (body.get("custom_data") or {}).get("inputs") or body.get("input") or {}
99+ query = inputs.get("query", "")
100+ 
101+ request_body = {
102+ "jsonrpc": "2.0",
103+ "id": f"call-versatile-{uuid.uuid4()}",
104+ "method": "SendStreamingMessage",
105+ "params": {
106+ "metadata": {
107+ "userId": self._cached_user_id,
108+ "traceId": self._trace_id,
109+ "versatile": {
110+ "inputs": inputs,
111+ },
112+ },
113+ "payload": {
114+ "message": {
115+ "role": "ROLE_USER",
116+ "messageId": f"msg-{uuid.uuid4()}",
117+ "conversationId": self._conv_id,
118+ "parts": [
119+ {"type": "text", "text": query},
120+ ],
121+ "metadata": {},
122+ },
123+ "configuration": {
124+ "blocking": True,
125+ "acceptedOutputModes": ["text/plain"],
126+ },
127+ },
128+ },
129+ }
130+ logger.info(
131+ f"[VersatileA2AGateway] 构造请求体: method=SendStreamingMessage, "
132+ f"conversationId={self._conv_id}, query={query!r:.60}"
133+ )
134+ logger.debug(
135+ f"[VersatileA2AGateway] 完整请求体: {json.dumps(request_body, ensure_ascii=False)[:500]}"
136+ )
137+ return request_body
138+ 
139+ def _process_chunk(self, chunk: str, ctx: VersatileStreamCtx) -> list[AdapterEvent]:
140+ try:
141+ parsed = json.loads(chunk)
142+ except json.JSONDecodeError:
143+ logger.warning(f"[VersatileA2AGateway] 无法解析 SSE 行: {chunk!r:.80}")
144+ return [AdapterEvent(data_proxy=DataProxyContent(raw_data=chunk))]
145+ 
146+ # JSON-RPC error
147+ if "error" in parsed:
148+ error = parsed["error"]
149+ ctx.completed = True
150+ ctx.is_failed = True
151+ ctx.error_message = json.dumps(error, ensure_ascii=False)
152+ logger.error(f"[VersatileA2AGateway] JSON-RPC error: {error}")
153+ return [AdapterEvent(data_proxy=DataProxyContent(raw_data=chunk))]
154+ 
155+ result = parsed.get("result", {})
156+ if not isinstance(result, dict):
157+ logger.warning(f"[VersatileA2AGateway] result 非 dict,透传: {chunk!r:.80}")
158+ return [AdapterEvent(data_proxy=DataProxyContent(raw_data=chunk))]
159+ payload = result.get("payload", {})
160+ 
161+ if "taskArtifactUpdate" in payload:
162+ return self._handle_artifact_update(payload["taskArtifactUpdate"], ctx)
163+ 
164+ if "taskStatusUpdate" in payload:
165+ return self._handle_status_update(payload["taskStatusUpdate"], ctx)
166+ 
167+ # 未知类型 - 透传
168+ logger.debug(f"[VersatileA2AGateway] 未知 payload 类型,透传: {chunk!r:.80}")
169+ return [AdapterEvent(data_proxy=DataProxyContent(raw_data=chunk))]
170+ 
171+ def _on_stream_end(self, ctx: VersatileStreamCtx) -> list[AdapterEvent]:
172+ if not ctx.completed:
173+ logger.info(f"[VersatileA2AGateway] 流关闭未收到 status-update, 兜底 INPUT_REQUIRED, conv_id={self._conv_id}")
174+ return [AdapterEvent(execution_input_required=ExecutionInputRequiredContent())]
175+ 
176+ if ctx.input_required:
177+ logger.info(f"[VersatileA2AGateway] INPUT_REQUIRED, conv_id={self._conv_id}")
178+ return [AdapterEvent(execution_input_required=ExecutionInputRequiredContent())]
179+ 
180+ if ctx.is_failed or ctx.execution_result:
atomgit-botatomgit-bot
atomgit-botatomgit-bot7月13日

🟡 Medium Priority

VersatileA2AGateway._on_stream_end 第 182 行使用真值判断 ctx.execution_result 来决定是否产出 execution_completed 终态事件。但 _handle_artifact_update 在处理非文本 parts(如纯文件 parts)时,仍会将空字符串 "" 存入 ctx.artifact_texts(第 230 行),并设置 ctx.last_artifact_id(第 231 行)。当 TASK_STATE_COMPLETED 到来时,_handle_status_update 会从 artifact_texts 中取出空字符串赋给 ctx.execution_result(第 259 行)。

此时 ctx.is_failed or ctx.execution_result → False or "" → ""(falsy),_on_stream_end 返回 [],不产出 execution_completed。executor 的 for 循环结束后触发兜底逻辑(executor.py 第 133 行),以 result=None 完成任务——丢失了原本可能从 status.message 降级提取到的有效文本。

触发条件:A2A Gateway 返回的 artifact 中 parts 只有非 text 类型(如 {"type":"file", ...}),同时 TASK_STATE_COMPLETED 的 message 中包含可用于降级提取的文本。

建议:将 ctx.execution_result 的真值检查改为 ctx.execution_result is not None,以区分"已设置空结果"与"未设置结果"两种语义。ctx.execution_result 初始为 None(VersatileStreamCtx 第 36 行),只有被明确赋值后才应视为有结果。

改动建议
180
- if ctx.is_failed or ctx.execution_result:
180
+ if ctx.is_failed or ctx.execution_result is not None:
应用建议
likedislike
不准确?
atomgit-botatomgit-bot7月13日

🟡 Medium Priority

_on_stream_end 第 182 行的条件 if ctx.is_failed or ctx.execution_result: 在 ctx.execution_result 为空字符串 "" 时求值为 False or "" → ""(falsy),导致整个方法返回空列表 []。

触发条件:A2A Gateway 返回 TASK_STATE_COMPLETED 但无 artifact 且 status message 也为空,或 artifact 累积的文本恰为空串。此时 ctx.completed=True、ctx.is_failed=False、ctx.execution_result=""。

后果:适配器层不产出任何终态事件,完全依赖 executor 的 if not terminal_sent 兜底逻辑发送 result=None 的 completed 事件。虽然 executor 有兜底,但适配器自身的契约被打破(completed 状态不产生事件),且当 executor 兜底逻辑未来变更时可能引入回归。

修复方向:在条件中增加对 ctx.completed 的判断,确保 completed 状态下始终产出 ExecutionCompletedContent 事件。

建议:将条件改为 if ctx.is_failed or ctx.execution_result or ctx.completed:,或在条件体内明确处理 completed-but-no-result 情况,确保产出 ExecutionCompletedContent(is_failed=False, result="", error_message="") 事件。

likedislike
不准确?
181+ logger.info(
182+ f"[VersatileA2AGateway] 终态: is_failed={ctx.is_failed}, "
183+ f"result={ctx.execution_result!r:.60}, conv_id={self._conv_id}"
184+ )
185+ return [AdapterEvent(execution_completed=ExecutionCompletedContent(
186+ is_failed=ctx.is_failed,
187+ result=ctx.execution_result or "",
188+ error_message=ctx.error_message,
189+ ))]
190+ 
191+ return []
192+ 
193+ # ── 内部方法 ──────────────────────────────────────────────
194+ 
195+ def _extract_user_id(self) -> str:
196+ """从 passed_headers 中大小写不敏感地提取 userId。
197+ 
198+ a2a_service 通过 Starlette 的 dict(request.headers) 传入 headers,
199+ 所有 key 会被转为小写(如 x-user-id / userid / cust-userid)。
200+ 此处按优先级依次查找:userid > x-user-id > cust-userid。
201+ """
202+ for key, value in self._passed_headers.items():
203+ key_lower = key.lower()
204+ if key_lower in ("userid", "x-user-id", "cust-userid") and value:
205+ return value
206+ return ""
207+ 
208+ def _handle_artifact_update(self, update: dict, ctx: VersatileStreamCtx) -> list[AdapterEvent]:
209+ artifact = update.get("artifact", {})
210+ artifact_id = artifact.get("artifactId", "_default")
211+ append = update.get("append", False)
212+ last_chunk = update.get("lastChunk", False)
213+ 
214+ # 提取所有 text part 的文本
215+ parts = artifact.get("parts", [])
216+ if not isinstance(parts, list):
217+ parts = []
218+ text = "".join(
219+ part.get("text", "")
220+ for part in parts
221+ if isinstance(part, dict) and part.get("type") == "text"
222+ )
223+ 
224+ # 按 append 累积(方案 D)
225+ if append and artifact_id in ctx.artifact_texts:
226+ ctx.artifact_texts[artifact_id] += text
227+ else:
228+ ctx.artifact_texts[artifact_id] = text
229+ ctx.last_artifact_id = artifact_id
230+ 
231+ logger.debug(
232+ f"[VersatileA2AGateway] artifact: id={artifact_id}, append={append}, "
233+ f"last_chunk={last_chunk}, text={text!r:.60}"
234+ )
235+ 
236+ # data_proxy 转换为低码工作流统一格式(空文本跳过,避免前端收到空消息)
237+ if not text:
238+ return []
239+ frame = {"event": "message", "data": {"text": text}}
240+ raw = json.dumps(frame, ensure_ascii=False)
241+ logger.debug(f"[VersatileA2AGateway] 输出 AdapterEvent(data_proxy): {raw[:200]}")
242+ return [AdapterEvent(data_proxy=DataProxyContent(raw_data=raw))]
243+ 
244+ def _handle_status_update(self, update: dict, ctx: VersatileStreamCtx) -> list[AdapterEvent]:
245+ status = update.get("status", update) # 兼容 status 嵌套或直接在 update 上
246+ state = status.get("state", update.get("state", ""))
247+ if not isinstance(state, str):
248+ state = str(state) if state is not None else ""
249+ state_lower = state.replace("TASK_STATE_", "").lower()
250+ 
251+ logger.info(f"[VersatileA2AGateway] status-update: state={state} ({state_lower}), conv_id={self._conv_id}")
252+ 
253+ if state_lower == "completed":
254+ ctx.completed = True
255+ # 取最后一个 artifactId 的累积文本
256+ if ctx.last_artifact_id and ctx.last_artifact_id in ctx.artifact_texts:
257+ ctx.execution_result = ctx.artifact_texts[ctx.last_artifact_id]
258+ logger.info(
259+ f"[VersatileA2AGateway] 结果提取: artifact_id={ctx.last_artifact_id}, "
260+ f"result={ctx.execution_result!r:.60}"
261+ )
262+ else:
263+ # 降级:从 status.message 提取
264+ ctx.execution_result = self._extract_text_from_status_message(status)
265+ if ctx.execution_result:
266+ logger.info(f"[VersatileA2AGateway] 结果提取(降级message): {ctx.execution_result!r:.60}")
267+ 
268+ # 补发 end 帧(与低码工作流格式一致)
269+ logger.debug(f"[VersatileA2AGateway] 补发 end 帧")
270+ logger.debug(f"[VersatileA2AGateway] 输出 AdapterEvent(data_proxy): {{\"event\":\"end\"}}")
271+ return [AdapterEvent(data_proxy=DataProxyContent(raw_data='{"event":"end"}'))]
272+ 
273+ if state_lower == "failed":
274+ ctx.completed = True
275+ ctx.is_failed = True
276+ ctx.error_message = self._extract_text_from_status_message(status) or "A2A Gateway 返回 FAILED"
277+ logger.error(f"[VersatileA2AGateway] FAILED: error={ctx.error_message!r:.100}")
278+ 
279+ # 转换为低码 error 帧格式
280+ error_frame = json.dumps({
281+ "event": "error",
282+ "data": {"code": "", "message": ctx.error_message},
283+ }, ensure_ascii=False)
284+ logger.debug(f"[VersatileA2AGateway] 输出 AdapterEvent(data_proxy): {error_frame[:200]}")
285+ return [AdapterEvent(data_proxy=DataProxyContent(raw_data=error_frame))]
286+ 
287+ if state_lower == "canceled":
288+ ctx.completed = True
289+ ctx.is_failed = True
290+ ctx.error_message = self._extract_text_from_status_message(status) or "A2A Gateway 返回 CANCELED"
291+ logger.warning(f"[VersatileA2AGateway] CANCELED: error={ctx.error_message!r:.100}")
292+ 
293+ error_frame = json.dumps({
294+ "event": "error",
295+ "data": {"code": "", "message": ctx.error_message},
296+ }, ensure_ascii=False)
297+ logger.debug(f"[VersatileA2AGateway] 输出 AdapterEvent(data_proxy): {error_frame[:200]}")
298+ return [AdapterEvent(data_proxy=DataProxyContent(raw_data=error_frame))]
299+ 
300+ if state_lower == "rejected":
301+ ctx.completed = True
302+ ctx.is_failed = True
303+ ctx.error_message = self._extract_text_from_status_message(status) or "A2A Gateway 返回 REJECTED"
304+ logger.warning(f"[VersatileA2AGateway] REJECTED: error={ctx.error_message!r:.100}")
305+ 
306+ error_frame = json.dumps({
307+ "event": "error",
308+ "data": {"code": "", "message": ctx.error_message},
309+ }, ensure_ascii=False)
310+ logger.debug(f"[VersatileA2AGateway] 输出 AdapterEvent(data_proxy): {error_frame[:200]}")
311+ return [AdapterEvent(data_proxy=DataProxyContent(raw_data=error_frame))]
312+ 
313+ if state_lower == "input_required":
314+ ctx.completed = True
315+ ctx.input_required = True
316+ return []
317+ 
318+ # working / submitted 等非终态 - 忽略
319+ logger.debug(f"[VersatileA2AGateway] 非终态 status: {state}")
320+ return []
321+ 
322+ @staticmethod
323+ def _extract_text_from_status_message(status: dict) -> str:
324+ """从 taskStatusUpdate 的 message.parts 中提取文本。"""
325+ message = status.get("message")
326+ if not message or not isinstance(message, dict):
327+ return ""
328+ parts = message.get("parts", [])
329+ return "".join(
atomgit-bot
atomgit-botatomgit-bot7月13日

🟡 Medium Priority

在 VersatileA2AGateway._extract_text_from_status_message 方法中(第 330-334 行),parts = message.get("parts", []) 获得的 parts 可能不是列表(例如上游返回 {"parts": {"unexpected": "dict"}}),后续的生成器表达式 for part in parts 将抛出 TypeError: 'dict' object is not iterable。

相比之下,_handle_artifact_update 方法(第 218 行)已对此做了防御:if not isinstance(parts, list): parts = []。同样的防护应在 _extract_text_from_status_message 中添加,保持一致性并避免运行时崩溃。

触发条件:A2A Gateway 返回的 taskStatusUpdate.message.parts 字段为非列表值(如 dict、字符串)。

建议:在迭代 parts 前增加 isinstance(parts, list) 检查,与 _handle_artifact_update 保持一致。

改动建议
329
+ parts = message.get("parts", [])
330
+ if not isinstance(parts, list):
331
+ parts = []
329
332
  return "".join(
应用建议
likedislike
不准确?
330+ part.get("text", "")
331+ for part in parts
332+ if part.get("type") == "text"
333+ )
@@ -23,15 +23,22 @@ from adapters.base_adapter import BaseAdapter
23 23 
24 24 
25class VersatileStreamCtx:25class VersatileStreamCtx:
26- """SSE 行循环中累积的可变状态,贯穿 _process_line → _process_chunk → _on_stream_end。"""26+ """SSE 行循环中累积的可变状态,贯穿 _process_line -> _process_chunk -> _on_stream_end。"""
27 27 
28- __slots__ = ("completed", "is_failed", "execution_result", "error_message")28+ __slots__ = (
29+ "completed", "is_failed", "execution_result", "error_message",
30+ "artifact_texts", "last_artifact_id", "input_required",
31+ )
29 32 
30 def __init__(self) -> None:33 def __init__(self) -> None:
31 self.completed: bool = False34 self.completed: bool = False
32 self.is_failed: bool = False35 self.is_failed: bool = False
33 self.execution_result: str | None = None36 self.execution_result: str | None = None
34 self.error_message: str = ""37 self.error_message: str = ""
38+ # A2A Gateway 专用:artifact 累积
39+ self.artifact_texts: dict[str, str] = {}
40+ self.last_artifact_id: str = ""
41+ self.input_required: bool = False
35 42 
36 43 
37class VersatileProxy(BaseAdapter):44class VersatileProxy(BaseAdapter):
@@ -98,13 +105,19 @@ class VersatileProxy(BaseAdapter):
98 105 
99 # ── 主流程 ────────────────────────────────────────────────106 # ── 主流程 ────────────────────────────────────────────────
100 107 
101- async def dispatch_stream(108+ async def dispatch_stream( # pylint: disable=too-many-arguments
102 self,109 self,
O
OopenJiuwen-bot7月14日

此条代码评论区间+108至109

【openlibing.ci】识别到代码检查告警抑制注释,匹配工具:pylint,请Committer检视其合理性。

likedislike
103 conv_id: str,110 conv_id: str,
104 headers: Optional[dict] = None,111 headers: Optional[dict] = None,
105 params: Optional[dict] = None,112 params: Optional[dict] = None,
106 body: Optional[dict] = None,113 body: Optional[dict] = None,
114+ trace_id: str = "",
107 ) -> AsyncGenerator[AdapterEvent, None]:115 ) -> AsyncGenerator[AdapterEvent, None]:
116+ # 存储入口上下文,供子类钩子方法使用
117+ self._conv_id = conv_id
118+ self._trace_id = trace_id
119+ self._passed_headers = headers or {}
120+ 
108 url = self._build_url(conv_id)121 url = self._build_url(conv_id)
109 req_headers = self._build_headers(headers)122 req_headers = self._build_headers(headers)
110 request_body = self._build_request_body(body or {})123 request_body = self._build_request_body(body or {})
@@ -149,10 +162,14 @@ class VersatileProxy(BaseAdapter):
149 @staticmethod162 @staticmethod
150 async def _log_request(request: httpx.Request) -> None:163 async def _log_request(request: httpx.Request) -> None:
151 """记录请求日志(生成 curl 命令)。"""164 """记录请求日志(生成 curl 命令)。"""
165+ _sensitive_headers = {"token", "authorization", "cookie", "set-cookie"}
152 body = await request.aread()166 body = await request.aread()
153 cmd = f"curl -X {request.method} '{request.url}'"167 cmd = f"curl -X {request.method} '{request.url}'"
154 for key, value in request.headers.items():168 for key, value in request.headers.items():
155- cmd += f" -H '{key}: {value}'"169+ if key.lower() in _sensitive_headers:
170+ cmd += f" -H '{key}: ***'"
171+ else:
172+ cmd += f" -H '{key}: {value}'"
156 if body:173 if body:
157 cmd += f" -d '{body.decode('utf-8', errors='replace')}'"174 cmd += f" -d '{body.decode('utf-8', errors='replace')}'"
158 banner_start = f"{'='*20} Proxy Request (Stream) Start {'='*20}"175 banner_start = f"{'='*20} Proxy Request (Stream) Start {'='*20}"
@@ -46,6 +46,9 @@ class Settings(BaseSettings):
46 versatile_timeout: Optional[int] = None46 versatile_timeout: Optional[int] = None
47 versatile_headers_template: Json[Dict[str, Any]] = _DEFAULT_VERSATILE_HEADERS_TEMPLATE47 versatile_headers_template: Json[Dict[str, Any]] = _DEFAULT_VERSATILE_HEADERS_TEMPLATE
48 versatile_adapter_type: str = "controller" # "controller" 或 "workflow"48 versatile_adapter_type: str = "controller" # "controller" 或 "workflow"
49+ versatile_workflow_adapter_type: str = Field(
50+ default="a2a_gateway", alias="va_workflow_adapter_type",
51+ ) # 环境变量 VA_WORKFLOW_ADAPTER_TYPE,默认走 a2a_gateway 网关模式,设为 "workflow" 切换回直连低码工作流
49 versatile_workflow_result_node: Optional[str] = Field(52 versatile_workflow_result_node: Optional[str] = Field(
50 default=None, alias="va_workflow_result_node",53 default=None, alias="va_workflow_result_node",
51 ) # 环境变量 VA_WORKFLOW_RESULT_NODE,代码中用 versatile_ 前缀访问54 ) # 环境变量 VA_WORKFLOW_RESULT_NODE,代码中用 versatile_ 前缀访问
@@ -18,6 +18,7 @@ from typing import AsyncGenerator, Optional
18import yaml18import yaml
19from loguru import logger19from loguru import logger
20 20 
21+from adapters.versatile_a2a_gateway import VersatileA2AGateway
21from adapters.versatile_controller import VersatileController22from adapters.versatile_controller import VersatileController
22from adapters.versatile_workflow import VersatileWorkflow23from adapters.versatile_workflow import VersatileWorkflow
23from config import get_settings24from config import get_settings
@@ -54,6 +55,18 @@ def _merge_workflow_defaults(raw: dict, workflow_defaults: dict) -> dict:
54 return merged55 return merged
55 56 
56 57 
58+def _merge_a2a_gateway_defaults(raw: dict, a2a_gateway_defaults: dict) -> dict:
59+ if raw.get("type") != "a2a_gateway":
60+ return raw
61+ 
62+ merged = {**a2a_gateway_defaults, **raw}
63+ default_headers = a2a_gateway_defaults.get("headers_template") or {}
64+ adapter_headers = raw.get("headers_template") or {}
65+ if default_headers or adapter_headers:
66+ merged["headers_template"] = {**default_headers, **adapter_headers}
67+ return merged
68+ 
69+ 
57class _VersatileAdapterConfig:70class _VersatileAdapterConfig:
58 """单个 adapter 的配置。"""71 """单个 adapter 的配置。"""
59 72 
@@ -61,6 +74,7 @@ class _VersatileAdapterConfig:
61 "name", "type", "url_template", "timeout",74 "name", "type", "url_template", "timeout",
62 "headers_template", "forward_header_whitelist",75 "headers_template", "forward_header_whitelist",
63 "workflow_result_node", "workflow_id", "intent",76 "workflow_result_node", "workflow_id", "intent",
77+ "a2a_gateway_base", "agent_card_name", "token",
64 )78 )
65 79 
66 def __init__(self, raw: dict) -> None:80 def __init__(self, raw: dict) -> None:
@@ -76,12 +90,17 @@ class _VersatileAdapterConfig:
76 self.workflow_result_node = raw.get("workflow_result_node")90 self.workflow_result_node = raw.get("workflow_result_node")
77 self.workflow_id = raw.get("workflow_id")91 self.workflow_id = raw.get("workflow_id")
78 self.intent = raw.get("intent")92 self.intent = raw.get("intent")
93+ # A2A Gateway 专用
94+ self.a2a_gateway_base = raw.get("a2a_gateway_base", "")
95+ self.agent_card_name = raw.get("agent_card_name", "")
96+ self.token = raw.get("token", "")
79 97 
80 98 
81class VersatileAdapterRunner:99class VersatileAdapterRunner:
82 """配置驱动的动态路由 Runner。"""100 """配置驱动的动态路由 Runner。"""
83 101 
84 def __init__(self, config_path: Optional[Path] = None) -> None:102 def __init__(self, config_path: Optional[Path] = None) -> None:
103+ self._workflow_adapter_type = os.getenv("VA_WORKFLOW_ADAPTER_TYPE", "a2a_gateway")
85 path = config_path or _resolve_default_config_path()104 path = config_path or _resolve_default_config_path()
86 self._adapters = self._load_config(path)105 self._adapters = self._load_config(path)
87 if not self._adapters:106 if not self._adapters:
@@ -100,11 +119,14 @@ class VersatileAdapterRunner:
100 with open(path, encoding="utf-8") as f:119 with open(path, encoding="utf-8") as f:
101 raw = yaml.safe_load(f) or {}120 raw = yaml.safe_load(f) or {}
102 workflow_defaults = raw.get("workflow_defaults", {})121 workflow_defaults = raw.get("workflow_defaults", {})
122+ a2a_gateway_defaults = raw.get("a2a_gateway_defaults", {})
103 adapters_raw = raw.get("adapters", [])123 adapters_raw = raw.get("adapters", [])
104- return [124+ configs: list[_VersatileAdapterConfig] = []
105- _VersatileAdapterConfig(_merge_workflow_defaults(b, workflow_defaults))125+ for b in adapters_raw:
106- for b in adapters_raw126+ merged = _merge_workflow_defaults(b, workflow_defaults)
107- ]127+ merged = _merge_a2a_gateway_defaults(merged, a2a_gateway_defaults)
128+ configs.append(_VersatileAdapterConfig(merged))
129+ return configs
108 130 
109 @staticmethod131 @staticmethod
110 def _build_from_settings() -> list[_VersatileAdapterConfig]:132 def _build_from_settings() -> list[_VersatileAdapterConfig]:
@@ -131,7 +153,7 @@ class VersatileAdapterRunner:
131 def _match_workflow(self, target: dict) -> Optional[_VersatileAdapterConfig]:153 def _match_workflow(self, target: dict) -> Optional[_VersatileAdapterConfig]:
132 """根据 target 中的 workflow_id 或 intent 匹配 workflow 配置。"""154 """根据 target 中的 workflow_id 或 intent 匹配 workflow 配置。"""
133 for b in self._adapters:155 for b in self._adapters:
134- if b.type != "workflow":156+ if b.type != self._workflow_adapter_type:
135 continue157 continue
136 if b.workflow_id and target.get("workflow_id") == b.workflow_id:158 if b.workflow_id and target.get("workflow_id") == b.workflow_id:
137 return b159 return b
@@ -152,6 +174,17 @@ class VersatileAdapterRunner:
152 forward_header_whitelist=whitelist,174 forward_header_whitelist=whitelist,
153 workflow_result_node=cfg.workflow_result_node,175 workflow_result_node=cfg.workflow_result_node,
154 )176 )
177+ if cfg.type == "a2a_gateway":
178+ return VersatileA2AGateway(
179+ a2a_gateway_base=cfg.a2a_gateway_base,
180+ agent_card_name=cfg.agent_card_name,
181+ token=cfg.token,
182+ url_template=cfg.url_template,
183+ timeout=cfg.timeout,
184+ headers_template=cfg.headers_template,
185+ forward_header_whitelist=whitelist,
186+ workflow_result_node=cfg.workflow_result_node,
187+ )
155 return VersatileController(188 return VersatileController(
156 url_template=cfg.url_template,189 url_template=cfg.url_template,
157 timeout=cfg.timeout,190 timeout=cfg.timeout,
@@ -160,8 +193,9 @@ class VersatileAdapterRunner:
160 workflow_result_node=cfg.workflow_result_node,193 workflow_result_node=cfg.workflow_result_node,
161 )194 )
162 195 
163- async def run_async(self, target: dict, headers: dict, params: dict,196+ async def run_async( # pylint: disable=too-many-arguments
164- body: dict) -> AsyncGenerator[AdapterEvent, None]:197+ self, target: dict, headers: dict, params: dict,
O
OopenJiuwen-bot7月14日

此条代码评论区间+196至+197

【openlibing.ci】识别到代码检查告警抑制注释,匹配工具:pylint,请Committer检视其合理性。

likedislike
198+ body: dict, trace_id: str = "") -> AsyncGenerator[AdapterEvent, None]:
165 """根据 target 动态匹配配置并创建适配器,驱动流。"""199 """根据 target 动态匹配配置并创建适配器,驱动流。"""
166 conv_id = target.get("conversation_id", "")200 conv_id = target.get("conversation_id", "")
167 wf_cfg = self._match_workflow(target)201 wf_cfg = self._match_workflow(target)
@@ -170,7 +204,14 @@ class VersatileAdapterRunner:
170 raise ValueError(f"无法匹配 adapter 配置: target={target}")204 raise ValueError(f"无法匹配 adapter 配置: target={target}")
171 205 
172 logger.debug(f"[VersatileAdapterRunner] target 匹配 adapter={cfg.name} ({cfg.type})")206 logger.debug(f"[VersatileAdapterRunner] target 匹配 adapter={cfg.name} ({cfg.type})")
207+ logger.info(
208+ f"[VersatileAdapterRunner] 路由决策: intent={target.get('intent', '')}, "
209+ f"adapter={cfg.name} (type={cfg.type}), conv_id={conv_id}"
210+ )
173 adapter = self._create_adapter(cfg)211 adapter = self._create_adapter(cfg)
174 212 
175- async for event in adapter.dispatch_stream(conv_id, headers=headers, params=params, body=body):213+ async for event in adapter.dispatch_stream(
214+ conv_id, headers=headers, params=params, body=body,
215+ trace_id=trace_id,
216+ ):
176 yield event217 yield event
@@ -13,6 +13,15 @@ if str(_VA_ROOT) not in sys.path:
13 sys.path.append(str(_VA_ROOT))13 sys.path.append(str(_VA_ROOT))
14 14 
15 15 
16+@pytest.fixture(autouse=True)
17+def _va_workflow_adapter_type_workflow(monkeypatch):
18+ """现有测试的 YAML 配置均使用 type: workflow 适配器,
19+ 默认设置 VA_WORKFLOW_ADAPTER_TYPE=workflow 使 _match_workflow 匹配 workflow 类型。
20+ 需测试 a2a_gateway 路由的用例可覆盖此环境变量。
21+ """
22+ monkeypatch.setenv("VA_WORKFLOW_ADAPTER_TYPE", "workflow")
23+ 
24+ 
16@pytest.fixture25@pytest.fixture
17def va_root() -> Path:26def va_root() -> Path:
18 """VA 进程根目录。"""27 """VA 进程根目录。"""
@@ -114,18 +114,16 @@ class TestMakeTextPart:
114 114 
115class TestExtractLoggingContext:115class TestExtractLoggingContext:
116 @staticmethod116 @staticmethod
117- def test_extracts_trace_and_agent_id():117+ def test_extracts_trace_id():
118- input_data = {"trace_id": "t-1", "agent_id": "a-1", "other": "x"}118+ input_data = {"trace_id": "t-1", "other": "x"}
119 ctx = A2aVersatileExecutor._extract_logging_context(input_data, "conv-1")119 ctx = A2aVersatileExecutor._extract_logging_context(input_data, "conv-1")
120 assert ctx["trace_id"] == "t-1"120 assert ctx["trace_id"] == "t-1"
121- assert ctx["agent_id"] == "a-1"
122 assert ctx["conv_id"] == "conv-1"121 assert ctx["conv_id"] == "conv-1"
123 122 
124 @staticmethod123 @staticmethod
125 def test_missing_fields_default_empty():124 def test_missing_fields_default_empty():
126 ctx = A2aVersatileExecutor._extract_logging_context({}, "conv-1")125 ctx = A2aVersatileExecutor._extract_logging_context({}, "conv-1")
127 assert ctx["trace_id"] == ""126 assert ctx["trace_id"] == ""
128- assert ctx["agent_id"] == ""
129 127 
130 128 
131# ════════════════════════════════════════════════════════════════════129# ════════════════════════════════════════════════════════════════════
@@ -0,0 +1,274 @@
1+# coding: utf-8
2+"""VersatileA2AGateway 单元测试。"""
3+# pylint: disable=protected-access,add-staticmethod-or-classmethod-decorator
4+import json
O
OopenJiuwen-bot7月14日

此条代码评论区间+2至+4

【openlibing.ci】识别到代码检查告警抑制注释,匹配工具:pylint,请Committer检视其合理性。

likedislike
5+ 
6+import pytest
7+ 
8+from adapters.versatile_a2a_gateway import VersatileA2AGateway
9+from adapters.versatile_proxy import VersatileStreamCtx
10+ 
11+ 
12+@pytest.fixture
13+def gw():
14+ return VersatileA2AGateway(
15+ a2a_gateway_base="https://a2a-gateway.example.com",
16+ agent_card_name="KnowledgeAgent",
17+ token="test-token",
18+ url_template="{a2a_gateway_base}/a2a/{agent_card_name}",
19+ timeout=600,
20+ headers_template={"Accept": "text/event-stream"},
21+ )
22+ 
23+ 
24+def _make_artifact_update(artifact_id, text, append=True, last_chunk=False):
25+ return {
26+ "jsonrpc": "2.0",
27+ "id": "req-001",
28+ "result": {
29+ "taskId": "task-001",
30+ "payload": {
31+ "taskArtifactUpdate": {
32+ "conversationId": "conv-001",
33+ "append": append,
34+ "lastChunk": last_chunk,
35+ "artifact": {
36+ "artifactId": artifact_id,
37+ "type": "text",
38+ "parts": [{"type": "text", "text": text}],
39+ },
40+ }
41+ },
42+ },
43+ }
44+ 
45+ 
46+def _make_status_update(state, message_text=None):
47+ payload = {
48+ "taskStatusUpdate": {
49+ "conversationId": "conv-001",
50+ "state": state,
51+ }
52+ }
53+ if message_text:
54+ payload["taskStatusUpdate"]["message"] = {
55+ "parts": [{"type": "text", "text": message_text}]
56+ }
57+ return {
58+ "jsonrpc": "2.0",
59+ "id": "req-001",
60+ "result": {"taskId": "task-001", "payload": payload},
61+ }
62+ 
63+ 
64+class TestBuildUrl:
65+ def test_build_url(self, gw):
66+ gw._conv_id = "conv-001"
67+ url = gw._build_url("conv-001")
68+ assert url == "https://a2a-gateway.example.com/a2a/KnowledgeAgent"
69+ 
70+ 
71+class TestBuildHeaders:
72+ def test_required_headers(self, gw):
73+ gw._passed_headers = {"userId": "user-123"}
74+ gw._cached_user_id = gw._extract_user_id()
75+ gw._trace_id = "trace-001"
76+ headers = gw._build_headers({"userId": "user-123"})
77+ assert headers["token"] == "test-token"
78+ assert headers["userId"] == "user-123"
79+ 
80+ def test_optional_b3_headers_present(self, gw):
81+ gw._passed_headers = {
82+ "userId": "user-123",
83+ "x-b3-traceid": "trace-abc",
84+ "x-b3-parentspanid": "parent-xyz",
85+ "x-b3-sampled": "1",
86+ "x-biz-tag": "finance",
87+ }
88+ gw._trace_id = ""
89+ headers = gw._build_headers(gw._passed_headers)
90+ assert headers["X-B3-TraceId"] == "trace-abc"
91+ assert headers["X-B3-ParentSpanId"] == "parent-xyz"
92+ assert "X-B3-SpanId" in headers
93+ assert headers["X-B3-Sampled"] == "1"
94+ assert headers["X-Biz-Tag"] == "finance"
95+ 
96+ def test_optional_b3_headers_absent(self, gw):
97+ gw._passed_headers = {"userId": "user-123"}
98+ gw._trace_id = ""
99+ headers = gw._build_headers({"userId": "user-123"})
100+ assert "X-B3-TraceId" not in headers
101+ assert "X-B3-ParentSpanId" not in headers
102+ assert "X-Biz-Tag" not in headers
103+ 
104+ 
105+class TestBuildRequestBody:
106+ def test_request_body_structure(self, gw):
107+ gw._conv_id = "conv-001"
108+ gw._trace_id = "trace-001"
109+ gw._passed_headers = {"userId": "user-123"}
110+ gw._cached_user_id = gw._extract_user_id()
111+ body = {
112+ "input": {"query": "请推荐理财", "intent": "knowledge_qa"},
113+ "custom_data": {"inputs": {"query": "请推荐理财", "intent": "knowledge_qa"}},
114+ }
115+ rb = gw._build_request_body(body)
116+ assert rb["method"] == "SendStreamingMessage"
117+ assert rb["params"]["payload"]["message"]["conversationId"] == "conv-001"
118+ assert rb["params"]["payload"]["message"]["parts"][0]["text"] == "请推荐理财"
119+ assert rb["params"]["payload"]["configuration"]["blocking"] is True
120+ assert rb["params"]["payload"]["configuration"]["acceptedOutputModes"] == ["text/plain"]
121+ assert rb["params"]["metadata"]["userId"] == "user-123"
122+ assert rb["params"]["metadata"]["traceId"] == "trace-001"
123+ assert rb["params"]["metadata"]["versatile"]["inputs"]["query"] == "请推荐理财"
124+ 
125+ 
126+class TestArtifactUpdate:
127+ def test_append_true_accumulate(self, gw):
128+ ctx = VersatileStreamCtx()
129+ chunk1 = json.dumps(_make_artifact_update("A1", "低风险理财", append=True, last_chunk=False))
130+ chunk2 = json.dumps(_make_artifact_update("A1", "第一类", append=True, last_chunk=True))
131+ 
132+ events1 = gw._process_chunk(chunk1, ctx)
133+ assert ctx.artifact_texts == {"A1": "低风险理财"}
134+ assert len(events1) == 1
135+ frame1 = json.loads(events1[0].data_proxy.raw_data)
136+ assert frame1 == {"event": "message", "data": {"text": "低风险理财"}}
137+ 
138+ events2 = gw._process_chunk(chunk2, ctx)
139+ assert ctx.artifact_texts == {"A1": "低风险理财第一类"}
140+ assert len(events2) == 1
141+ frame2 = json.loads(events2[0].data_proxy.raw_data)
142+ assert frame2 == {"event": "message", "data": {"text": "第一类"}}
143+ 
144+ def test_append_false_overwrite(self, gw):
145+ ctx = VersatileStreamCtx()
146+ chunk = json.dumps(_make_artifact_update("A1", "完整文本", append=False, last_chunk=True))
147+ events = gw._process_chunk(chunk, ctx)
148+ assert ctx.artifact_texts == {"A1": "完整文本"}
149+ assert len(events) == 1
150+ frame = json.loads(events[0].data_proxy.raw_data)
151+ assert frame == {"event": "message", "data": {"text": "完整文本"}}
152+ 
153+ def test_multiple_artifacts(self, gw):
154+ ctx = VersatileStreamCtx()
155+ gw._process_chunk(json.dumps(_make_artifact_update("A1", "中间过程", append=False, last_chunk=True)), ctx)
156+ gw._process_chunk(json.dumps(_make_artifact_update("A2", "最终回答", append=True, last_chunk=True)), ctx)
157+ assert ctx.artifact_texts == {"A1": "中间过程", "A2": "最终回答"}
158+ assert ctx.last_artifact_id == "A2"
159+ 
160+ 
161+class TestStatusUpdate:
162+ def test_completed_with_result(self, gw):
163+ ctx = VersatileStreamCtx()
164+ ctx.artifact_texts = {"A1": "低风险理财推荐"}
165+ ctx.last_artifact_id = "A1"
166+ chunk = json.dumps(_make_status_update("TASK_STATE_COMPLETED"))
167+ events = gw._process_chunk(chunk, ctx)
168+ assert ctx.completed is True
169+ assert ctx.execution_result == "低风险理财推荐"
170+ # 补发 end 帧
171+ assert len(events) == 1
172+ assert json.loads(events[0].data_proxy.raw_data) == {"event": "end"}
173+ 
174+ def test_completed_without_artifact(self, gw):
175+ ctx = VersatileStreamCtx()
176+ chunk = json.dumps(_make_status_update("TASK_STATE_COMPLETED"))
177+ events = gw._process_chunk(chunk, ctx)
178+ assert ctx.completed is True
179+ assert ctx.execution_result == "" or ctx.execution_result is None
180+ 
181+ def test_failed_with_message(self, gw):
182+ ctx = VersatileStreamCtx()
183+ chunk = json.dumps(_make_status_update("TASK_STATE_FAILED", "Model only support text input"))
184+ events = gw._process_chunk(chunk, ctx)
185+ assert ctx.completed is True
186+ assert ctx.is_failed is True
187+ assert ctx.error_message == "Model only support text input"
188+ assert len(events) == 1
189+ frame = json.loads(events[0].data_proxy.raw_data)
190+ assert frame["event"] == "error"
191+ assert frame["data"]["message"] == "Model only support text input"
192+ 
193+ def test_input_required(self, gw):
194+ ctx = VersatileStreamCtx()
195+ chunk = json.dumps(_make_status_update("TASK_STATE_INPUT_REQUIRED"))
196+ events = gw._process_chunk(chunk, ctx)
197+ assert ctx.completed is True
198+ assert ctx.input_required is True
199+ assert len(events) == 0
200+ 
201+ def test_working_ignored(self, gw):
202+ ctx = VersatileStreamCtx()
203+ chunk = json.dumps(_make_status_update("TASK_STATE_WORKING"))
204+ events = gw._process_chunk(chunk, ctx)
205+ assert ctx.completed is False
206+ assert len(events) == 0
207+ 
208+ 
209+class TestOnStreamEnd:
210+ def test_stream_end_without_status_yields_input_required(self, gw):
211+ """流关闭未收到 status-update 时兜底 INPUT_REQUIRED。"""
212+ ctx = VersatileStreamCtx()
213+ events = gw._on_stream_end(ctx)
214+ assert len(events) == 1
215+ assert events[0].execution_input_required is not None
216+ 
217+ def test_input_required_terminal(self, gw):
218+ ctx = VersatileStreamCtx()
219+ ctx.completed = True
220+ ctx.input_required = True
221+ events = gw._on_stream_end(ctx)
222+ assert len(events) == 1
223+ assert events[0].execution_input_required is not None
224+ 
225+ def test_completed_terminal(self, gw):
226+ ctx = VersatileStreamCtx()
227+ ctx.completed = True
228+ ctx.execution_result = "推荐结果"
229+ events = gw._on_stream_end(ctx)
230+ assert len(events) == 1
231+ assert events[0].execution_completed.is_failed is False
232+ assert events[0].execution_completed.result == "推荐结果"
233+ 
234+ def test_failed_terminal(self, gw):
235+ ctx = VersatileStreamCtx()
236+ ctx.completed = True
237+ ctx.is_failed = True
238+ ctx.error_message = "执行失败"
239+ events = gw._on_stream_end(ctx)
240+ assert len(events) == 1
241+ assert events[0].execution_completed.is_failed is True
242+ assert events[0].execution_completed.error_message == "执行失败"
243+ 
244+ 
245+class TestResultExtraction:
246+ def test_strategy_d_last_artifact(self, gw):
247+ ctx = VersatileStreamCtx()
248+ gw._process_chunk(json.dumps(_make_artifact_update("A1", "中间过程", append=False, last_chunk=True)), ctx)
249+ gw._process_chunk(json.dumps(_make_artifact_update("A2", "最终", append=True, last_chunk=False)), ctx)
250+ gw._process_chunk(json.dumps(_make_artifact_update("A2", "回答", append=True, last_chunk=True)), ctx)
251+ gw._process_chunk(json.dumps(_make_status_update("TASK_STATE_COMPLETED")), ctx)
252+ assert ctx.execution_result == "最终回答"
253+ 
254+ def test_jsonrpc_error(self, gw):
255+ ctx = VersatileStreamCtx()
256+ chunk = json.dumps({"jsonrpc": "2.0", "id": "req-001", "error": {"code": -32600, "message": "Invalid Request"}})
257+ events = gw._process_chunk(chunk, ctx)
258+ assert ctx.completed is True
259+ assert ctx.is_failed is True
260+ assert len(events) == 1
261+ 
262+ def test_workflow_result_node_ignored(self, gw):
263+ """workflow_result_node 配置不影响 A2A Gateway。"""
264+ gw2 = VersatileA2AGateway(
265+ a2a_gateway_base="https://gw.example.com",
266+ agent_card_name="Agent",
267+ token="tok",
268+ url_template="{a2a_gateway_base}/a2a/{agent_card_name}",
269+ workflow_result_node="WorkflowQAResponseNode",
270+ )
271+ ctx = VersatileStreamCtx()
272+ gw2._process_chunk(json.dumps(_make_artifact_update("A1", "文本", append=False, last_chunk=True)), ctx)
273+ gw2._process_chunk(json.dumps(_make_status_update("TASK_STATE_COMPLETED")), ctx)
274+ assert ctx.execution_result == "文本"
@@ -16,22 +16,32 @@
16# - forward_header_whitelist 不合并;adapter 配置后整体覆盖默认白名单16# - forward_header_whitelist 不合并;adapter 配置后整体覆盖默认白名单
17# - controller 通常只有一个,请在 controller adapter 内完整配置17# - controller 通常只有一个,请在 controller adapter 内完整配置
18#18#
19+# a2a_gateway_defaults 仅对 type: a2a_gateway 的 adapter 生效(同上逻辑)。
20+# 通过环境变量 VA_WORKFLOW_ADAPTER_TYPE 切换:
21+# VA_WORKFLOW_ADAPTER_TYPE=workflow -> _match_workflow 匹配 type: workflow
22+# VA_WORKFLOW_ADAPTER_TYPE=a2a_gateway -> _match_workflow 匹配 type: a2a_gateway(默认)
23+#
19# adapters 列表可包含多个 adapter,运行时按 target 动态匹配:24# adapters 列表可包含多个 adapter,运行时按 target 动态匹配:
20-# - target 含 workflow_id / intent 且匹配到 workflow 配置 → VersatileWorkflow25+# - VA_WORKFLOW_ADAPTER_TYPE=workflow 时,target 含 workflow_id / intent 且匹配到 workflow 配置 -> VersatileWorkflow
21-# - 否则 → VersatileController(使用第一个 type=controller 的配置)26+# - VA_WORKFLOW_ADAPTER_TYPE=a2a_gateway 时,target 含 workflow_id / intent 且匹配到 a2a_gateway 配置 -> VersatileA2AGateway
27+# - 否则 -> VersatileController(使用第一个 type=controller 的配置)
22#28#
23# 各字段说明:29# 各字段说明:
24-# name — 适配器名称(日志标识)30+# name - 适配器名称(日志标识)
25-# type — controller(一级控制器)或 workflow(低码工作流)31+# type - controller / workflow / a2a_gateway
26-# url_template — URL 模板,支持占位符 {conversation_id}、{workflow_id}32+# url_template - URL 模板,支持占位符 {conversation_id}、{workflow_id}、
27-# timeout — HTTP 流式请求超时(秒),默认 60033+# {a2a_gateway_base}、{agent_card_name}
28-# headers_template — 调用 Versatile 时注入的默认请求头34+# timeout - HTTP 流式请求超时(秒),默认 600
29-# forward_header_whitelist — 仅转发白名单中的外部请求头(小写);不设置则转发全部35+# headers_template - 调用时注入的默认请求头
30-# workflow_result_node — (controller/workflow 可选)SSE 流中匹配 node_type=="QA" 且36+# forward_header_whitelist - 仅转发白名单中的外部请求头(小写);不设置则转发全部
37+# workflow_result_node - (controller/workflow 可选)SSE 流中匹配 node_type=="QA" 且
31# node_name==此值的帧,提取 text 作为 workflow_result;38# node_name==此值的帧,提取 text 作为 workflow_result;
32-# adapter 级配置,每个 adapter 可不同39+# A2A Gateway 模式不使用此字段(取最后 artifact)
33-# workflow_id — (仅 workflow)target 中 workflow_id 精确匹配时命中40+# workflow_id - (仅 workflow/a2a_gateway)target 中 workflow_id 精确匹配时命中
34-# intent — (仅 workflow)target 中 intent 匹配时命中41+# intent - (仅 workflow/a2a_gateway)target 中 intent 匹配时命中
42+# a2a_gateway_base - (仅 a2a_gateway)网关基础地址
43+# agent_card_name - (仅 a2a_gateway)AgentCard 名称
44+# token - (仅 a2a_gateway)鉴权 token
35 45 
36workflow_defaults:46workflow_defaults:
37 url_template: "https://versatile.example.com/workflows/{workflow_id}/conversations/{conversation_id}"47 url_template: "https://versatile.example.com/workflows/{workflow_id}/conversations/{conversation_id}"
@@ -46,6 +56,22 @@ workflow_defaults:
46 - cust-userid56 - cust-userid
47 - cookie57 - cookie
48 58 
59+# A2A Gateway 默认配置(VA_WORKFLOW_ADAPTER_TYPE=a2a_gateway 时生效)
60+a2a_gateway_defaults:
61+ a2a_gateway_base: "https://a2a-gateway.example.com"
62+ url_template: "{a2a_gateway_base}/a2a/{agent_card_name}"
63+ timeout: 600
64+ headers_template:
65+ Accept: "text/event-stream, application/json"
66+ Content-Type: "application/json"
67+ forward_header_whitelist:
68+ - x-b3-traceid
69+ - x-b3-spanid
atomgit-bot
atomgit-botatomgit-bot7月13日

🟡 Medium Priority

在 a2a_gateway_defaults 的 forward_header_whitelist(第 69 行)中包含了 x-b3-spanid。父类 VersatileProxy._build_headers(versatile_proxy.py 第 79 行)会按白名单原样转发该头(带原始 key x-b3-spanid 和原始值)。而子类 VersatileA2AGateway._build_headers(versatile_a2a_gateway.py 第 90 行)在 x-b3-parentspanid 存在时会通过 uuid.uuid4().hex[:16] 重新生成 X-B3-SpanId。HTTP 头部是大小写不敏感的(RFC 7230),因此最终请求中会同时存在 x-b3-spanid(转发的上游 span-id)和 X-B3-SpanId(新生成的 span-id),两个值不同,导致下游无法确定哪个是正确的 span-id,破坏分布式追踪语义。

证据链:

  1. YAML 第 69 行:- x-b3-spanid 在白名单中
  2. versatile_proxy.py 第 79 行:父类按原始 key 转发匹配白名单的 header
  3. versatile_a2a_gateway.py 第 90 行:子类生成新 X-B3-SpanId = uuid.uuid4().hex[:16]
  4. 两者同时出现在同一次 HTTP 请求中,值不同,HTTP 大小写不敏感

修复方向:从 a2a_gateway_defaults.forward_header_whitelist 中移除 x-b3-spanid(span-id 应由网关重新生成,不应透传上游值)。

建议:从 a2a_gateway_defaults.forward_header_whitelist 中移除 - x-b3-spanid 这一行。SpanId 由 A2A Gateway 的 _build_headers 重新生成(versatile_a2a_gateway.py:90),不应透传上游的 span-id。

likedislike
不准确?
70+ - x-b3-parentspanid
71+ - x-b3-sampled
72+ - x-biz-tag
73+ - x-user-id
74+ 
49adapters:75adapters:
50 - name: default_controller76 - name: default_controller
51 type: controller77 type: controller
@@ -80,3 +106,12 @@ adapters:
80 - authorization106 - authorization
81 - x-tenant-id107 - x-tenant-id
82 workflow_result_node: "SpecialQAResponseNode"108 workflow_result_node: "SpecialQAResponseNode"
109+ 
110+ # A2A Gateway 适配器(VA_WORKFLOW_ADAPTER_TYPE=a2a_gateway 时命中)
111+ # 公共字段继承 a2a_gateway_defaults,只需配置差异项
112+ - name: versatile_gateway_knowledge_qa
113+ type: a2a_gateway
114+ intent: "knowledge_qa"
115+ agent_card_name: "KnowledgeAgent"
116+ token: "workflow-knowledge-token"
117+ workflow_result_node: "WorkflowQAResponseNode"