| @@ -30,6 +30,8 @@ from __future__ import annotations | |||
| 30 | import time | 30 | import time |
| 31 | from typing import Any | 31 | from 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 | pass | 152 | 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 | ||
| 228 | def wrap_error( | 240 | def 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 工作流最终结果节点名 ─────────── |
| 37 | va_workflow_result_node= | 37 | va_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 ─────────────────────────────────────────────────────────────────────── |
| 40 | ADAPTER_LOG_LEVEL=DEBUG | 45 | ADAPTER_LOG_LEVEL=DEBUG |
| 41 | ADAPTER_LOG_FILE=logs/versatile_adapter.log | 46 | ADAPTER_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 | + ) | ||
| 77 | updater = await self._setup_task(context, event_queue) | 82 | updater = await self._setup_task(context, event_queue) |
| 78 | terminal_sent = False | 83 | 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_id | 155 | conv_id = context.context_id |
| 151 | task_id = context.task_id | 156 | 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: | ||||||||||
🟡 Medium Priority
此时 触发条件:A2A Gateway 返回的 artifact 中 建议:将 改动建议
![]() ![]() 不准确? 🟡 Medium Priority
触发条件:A2A Gateway 返回 后果:适配器层不产出任何终态事件,完全依赖 executor 的 修复方向:在条件中增加对 建议:将条件改为 ![]() ![]() 不准确? | |||||||||||
| 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 | + | ||||||||||
| 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( | ||||||||||
🟡 Medium Priority 在 相比之下, 触发条件:A2A Gateway 返回的 建议:在迭代 改动建议
![]() ![]() 不准确? | |||||||||||
| 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 | ||
| 25 | class VersatileStreamCtx: | 25 | class 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 = False | 34 | self.completed: bool = False |
| 32 | self.is_failed: bool = False | 35 | self.is_failed: bool = False |
| 33 | self.execution_result: str | None = None | 36 | 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 | ||
| 37 | class VersatileProxy(BaseAdapter): | 44 | class 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 【openlibing.ci】识别到代码检查告警抑制注释,匹配工具:pylint,请Committer检视其合理性。 ![]() ![]() | |||
| 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 | 162 | ||
| 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] = None | 46 | versatile_timeout: Optional[int] = None |
| 47 | versatile_headers_template: Json[Dict[str, Any]] = _DEFAULT_VERSATILE_HEADERS_TEMPLATE | 47 | 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 | |||
| 18 | import yaml | 18 | import yaml |
| 19 | from loguru import logger | 19 | from loguru import logger |
| 20 | 20 | ||
| 21 | +from adapters.versatile_a2a_gateway import VersatileA2AGateway | ||
| 21 | from adapters.versatile_controller import VersatileController | 22 | from adapters.versatile_controller import VersatileController |
| 22 | from adapters.versatile_workflow import VersatileWorkflow | 23 | from adapters.versatile_workflow import VersatileWorkflow |
| 23 | from config import get_settings | 24 | from config import get_settings |
| @@ -54,6 +55,18 @@ def _merge_workflow_defaults(raw: dict, workflow_defaults: dict) -> dict: | |||
| 54 | return merged | 55 | 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 | + | ||
| 57 | class _VersatileAdapterConfig: | 70 | class _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 | ||
| 81 | class VersatileAdapterRunner: | 99 | class 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_raw | 126 | + 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 | 131 | ||
| 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 | continue | 157 | 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 b | 159 | 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 【openlibing.ci】识别到代码检查告警抑制注释,匹配工具:pylint,请Committer检视其合理性。 ![]() ![]() | |||
| 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 event | 217 | 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 | + | ||
| 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 | 25 | ||
| 17 | def va_root() -> Path: | 26 | def va_root() -> Path: |
| 18 | """VA 进程根目录。""" | 27 | """VA 进程根目录。""" |
| @@ -114,18 +114,16 @@ class TestMakeTextPart: | |||
| 114 | 114 | ||
| 115 | class TestExtractLoggingContext: | 115 | class TestExtractLoggingContext: |
| 116 | 116 | ||
| 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 | 123 | ||
| 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 【openlibing.ci】识别到代码检查告警抑制注释,匹配工具:pylint,请Committer检视其合理性。 ![]() ![]() | |||
| 5 | + | ||
| 6 | +import pytest | ||
| 7 | + | ||
| 8 | +from adapters.versatile_a2a_gateway import VersatileA2AGateway | ||
| 9 | +from adapters.versatile_proxy import VersatileStreamCtx | ||
| 10 | + | ||
| 11 | + | ||
| 12 | + | ||
| 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 配置 → VersatileWorkflow | 25 | +# - 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 流式请求超时(秒),默认 600 | 33 | +# {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 | ||
| 36 | workflow_defaults: | 46 | workflow_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-userid | 56 | - cust-userid |
| 47 | - cookie | 57 | - 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 | ||
🟡 Medium Priority 在 证据链:
修复方向:从 建议:从 a2a_gateway_defaults.forward_header_whitelist 中移除 ![]() ![]() 不准确? | |||
| 70 | + - x-b3-parentspanid | ||
| 71 | + - x-b3-sampled | ||
| 72 | + - x-biz-tag | ||
| 73 | + - x-user-id | ||
| 74 | + | ||
| 49 | adapters: | 75 | adapters: |
| 50 | - name: default_controller | 76 | - name: default_controller |
| 51 | type: controller | 77 | type: controller |
| @@ -80,3 +106,12 @@ adapters: | |||
| 80 | - authorization | 106 | - authorization |
| 81 | - x-tenant-id | 107 | - 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" | ||


🟡 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。