From 882b230bf186f3a657ef0faefcba6910608e54e3 Mon Sep 17 00:00:00 2001 From: Rishi Kunnath Date: Thu, 13 Aug 2026 19:46:09 +0530 Subject: [PATCH 1/5] fix a2a Message serialization for a2a-sdk >= 1.0 Replace model_dump/Message(**kwargs) with protobuf-compatible MessageToDict/ParseDict and bump a2a-sdk lower bound to >=1.0.0 --- pyproject.toml | 2 +- src/sap_cloud_sdk/extensibility/client.py | 24 ++++++++++++++++++++--- 2 files changed, 22 insertions(+), 4 deletions(-) diff --git a/pyproject.toml b/pyproject.toml index 9baab7b9..fbb119e2 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -75,7 +75,7 @@ dev = [ "sqlalchemy>=2.0.0", "django>=4.0", "flask>=3.0", - "a2a-sdk>=0.2.0", + "a2a-sdk>=1.0.0", "langchain-core>=1.2.7", "langgraph>=0.2.0", "langchain-community>=0.3.0", diff --git a/src/sap_cloud_sdk/extensibility/client.py b/src/sap_cloud_sdk/extensibility/client.py index 3e317236..e2b1306d 100644 --- a/src/sap_cloud_sdk/extensibility/client.py +++ b/src/sap_cloud_sdk/extensibility/client.py @@ -12,6 +12,10 @@ import httpx from a2a.types import Message from opentelemetry.propagate import inject + +from google.protobuf.json_format import MessageToDict as _proto_message_to_dict +from google.protobuf.json_format import ParseDict as _proto_parse_dict + from pydantic_core import ValidationError from sap_cloud_sdk.core.telemetry import Module, Operation @@ -199,6 +203,7 @@ def get_extension_capability_implementation( tenant="1d2e1a41-a28b-431f-9e3f-42e9704bfa75", ) """ + logger.info("Fetching extension capabilities for tenant=%s", tenant) try: return self._transport.get_extension_capability_implementation( capability_id=capability_id, @@ -376,7 +381,7 @@ def call_hook( .get("main", [[{}]])[0][0] .get("json", {}) ) - return Message(**response_json) + return _proto_parse_dict(response_json, Message()) except (KeyError, IndexError, TypeError, ValidationError) as exc: raise ExtensibilityError( f"Failed to extract response from last executed node: {exc}" @@ -400,6 +405,7 @@ async def _discover_n8n_tools( self, agw_client: Any, user_token: Optional[str] ) -> tuple[Any, Any]: tools = await agw_client.list_mcp_tools(user_token=user_token or None) + logger.info("Listed %d MCP tools from Agent Gateway", len(tools)) execute_tool = next( ( @@ -415,6 +421,7 @@ async def _discover_n8n_tools( f"MCP tool '{_EXECUTE_WORKFLOW_TOOL_NAME}' on server '{_N8N_MCP_SERVER_NAME}' " "not found via Agent Gateway." ) + logger.info("Fetched n8n execute_tool: %s", _EXECUTE_WORKFLOW_TOOL_NAME) get_exec_tool = next( ( @@ -430,6 +437,7 @@ async def _discover_n8n_tools( f"MCP tool '{_GET_EXECUTION_TOOL_NAME}' on server '{_N8N_MCP_SERVER_NAME}' " "not found via Agent Gateway." ) + logger.info("Fetched n8n get_exec_tool: %s", _GET_EXECUTION_TOOL_NAME) return execute_tool, get_exec_tool @@ -442,7 +450,10 @@ async def _execute_workflow_via_agw( message: Optional[Any], headers: Optional[dict], ) -> tuple[str, Any]: - message_body = message.model_dump(mode="json") if message is not None else {} + if message is None: + message_body: dict = {} + else: + message_body = _proto_message_to_dict(message, preserving_proto_field_name=True) execute_arguments = { "workflowId": hook.n8n_workflow_config.workflow_id, "inputs": { @@ -455,6 +466,7 @@ async def _execute_workflow_via_agw( }, }, } + logger.info("Executing workflow id=%s", hook.n8n_workflow_config.workflow_id) try: result_str = await agw_client.call_mcp_tool( execute_tool, @@ -480,6 +492,7 @@ async def _execute_workflow_via_agw( ) execution_id = data.get("executionId") + logger.info("Workflow execution complete: execution_id=%s, status=%s", execution_id, status) return str(execution_id), status @staticmethod @@ -494,7 +507,7 @@ def _extract_message(data: dict) -> Message: .get("main", [[{}]])[0][0] .get("json", {}) ) - return Message(**response_json) + return _proto_parse_dict(response_json, Message()) except (KeyError, IndexError, TypeError, ValidationError) as exc: raise TransportError( f"Failed to extract response from last executed node: {exc}" @@ -511,6 +524,7 @@ async def _poll_hook_execution( ) -> Optional[Message]: deadline = time.monotonic() + hook.timeout last_status = initial_status + logger.info("Polling for workflow %s execution result (timeout=%ss)", hook.n8n_workflow_config.workflow_id, hook.timeout) while time.monotonic() < deadline: await asyncio.sleep(_HOOK_POLL_INTERVAL) @@ -541,6 +555,7 @@ async def _poll_hook_execution( ) if last_status == "success": + logger.info("Execution %s completed successfully", execution_id) return self._extract_message(data) if last_status in _EXECUTION_TERMINAL_STATUSES: @@ -615,12 +630,15 @@ async def call_hook_agw( agw_client = create_agw_client( tenant_subdomain, _telemetry_source=Module.EXTENSIBILITY ) + logger.info("AGW client created successfully for tenant_subdomain=%s", tenant_subdomain) execute_tool, get_exec_tool = await self._discover_n8n_tools( agw_client, user_token ) + logger.info("Discovered n8n tools") execution_id, status = await self._execute_workflow_via_agw( agw_client, execute_tool, hook, user_token, message, headers ) + logger.info("Workflow triggered: execution_id=%s, initial_status=%s", execution_id, status) return await self._poll_hook_execution( agw_client, get_exec_tool, hook, execution_id, user_token, status ) From e715394d0fdad232c93b1ca1b25ea7b600a7e843 Mon Sep 17 00:00:00 2001 From: Rishi Kunnath Date: Thu, 13 Aug 2026 19:54:46 +0530 Subject: [PATCH 2/5] update sdk version --- pyproject.toml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pyproject.toml b/pyproject.toml index fbb119e2..8f2bad5f 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -1,6 +1,6 @@ [project] name = "sap-cloud-sdk" -version = "0.43.2" +version = "0.43.3" description = "SAP Cloud SDK for Python" readme = "README.md" license = "Apache-2.0" From f6565753eb555ac76ea5da15a978089caef16afb Mon Sep 17 00:00:00 2001 From: Rishi Kunnath Date: Mon, 17 Aug 2026 22:00:34 +0530 Subject: [PATCH 3/5] removing changes related to a2a sdk issue --- pyproject.toml | 2 +- src/sap_cloud_sdk/extensibility/client.py | 13 +++---------- 2 files changed, 4 insertions(+), 11 deletions(-) diff --git a/pyproject.toml b/pyproject.toml index 8f2bad5f..4daa9f21 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -75,7 +75,7 @@ dev = [ "sqlalchemy>=2.0.0", "django>=4.0", "flask>=3.0", - "a2a-sdk>=1.0.0", + "a2a-sdk>=0.2.0", "langchain-core>=1.2.7", "langgraph>=0.2.0", "langchain-community>=0.3.0", diff --git a/src/sap_cloud_sdk/extensibility/client.py b/src/sap_cloud_sdk/extensibility/client.py index e2b1306d..d90a6dc5 100644 --- a/src/sap_cloud_sdk/extensibility/client.py +++ b/src/sap_cloud_sdk/extensibility/client.py @@ -12,10 +12,6 @@ import httpx from a2a.types import Message from opentelemetry.propagate import inject - -from google.protobuf.json_format import MessageToDict as _proto_message_to_dict -from google.protobuf.json_format import ParseDict as _proto_parse_dict - from pydantic_core import ValidationError from sap_cloud_sdk.core.telemetry import Module, Operation @@ -381,7 +377,7 @@ def call_hook( .get("main", [[{}]])[0][0] .get("json", {}) ) - return _proto_parse_dict(response_json, Message()) + return Message(**response_json) except (KeyError, IndexError, TypeError, ValidationError) as exc: raise ExtensibilityError( f"Failed to extract response from last executed node: {exc}" @@ -450,10 +446,7 @@ async def _execute_workflow_via_agw( message: Optional[Any], headers: Optional[dict], ) -> tuple[str, Any]: - if message is None: - message_body: dict = {} - else: - message_body = _proto_message_to_dict(message, preserving_proto_field_name=True) + message_body = message.model_dump(mode="json") if message is not None else {} execute_arguments = { "workflowId": hook.n8n_workflow_config.workflow_id, "inputs": { @@ -507,7 +500,7 @@ def _extract_message(data: dict) -> Message: .get("main", [[{}]])[0][0] .get("json", {}) ) - return _proto_parse_dict(response_json, Message()) + return Message(**response_json) except (KeyError, IndexError, TypeError, ValidationError) as exc: raise TransportError( f"Failed to extract response from last executed node: {exc}" From 4247889472912f116d0479624003be34d7b688ad Mon Sep 17 00:00:00 2001 From: Rishi Kunnath Date: Mon, 17 Aug 2026 22:16:21 +0530 Subject: [PATCH 4/5] fix formatting --- src/sap_cloud_sdk/extensibility/client.py | 22 ++++++++++++++++++---- 1 file changed, 18 insertions(+), 4 deletions(-) diff --git a/src/sap_cloud_sdk/extensibility/client.py b/src/sap_cloud_sdk/extensibility/client.py index d90a6dc5..20c49419 100644 --- a/src/sap_cloud_sdk/extensibility/client.py +++ b/src/sap_cloud_sdk/extensibility/client.py @@ -485,7 +485,11 @@ async def _execute_workflow_via_agw( ) execution_id = data.get("executionId") - logger.info("Workflow execution complete: execution_id=%s, status=%s", execution_id, status) + logger.info( + "Workflow execution complete: execution_id=%s, status=%s", + execution_id, + status, + ) return str(execution_id), status @staticmethod @@ -517,7 +521,11 @@ async def _poll_hook_execution( ) -> Optional[Message]: deadline = time.monotonic() + hook.timeout last_status = initial_status - logger.info("Polling for workflow %s execution result (timeout=%ss)", hook.n8n_workflow_config.workflow_id, hook.timeout) + logger.info( + "Polling for workflow %s execution result (timeout=%ss)", + hook.n8n_workflow_config.workflow_id, + hook.timeout, + ) while time.monotonic() < deadline: await asyncio.sleep(_HOOK_POLL_INTERVAL) @@ -623,7 +631,9 @@ async def call_hook_agw( agw_client = create_agw_client( tenant_subdomain, _telemetry_source=Module.EXTENSIBILITY ) - logger.info("AGW client created successfully for tenant_subdomain=%s", tenant_subdomain) + logger.info( + "AGW client created successfully for tenant_subdomain=%s", tenant_subdomain + ) execute_tool, get_exec_tool = await self._discover_n8n_tools( agw_client, user_token ) @@ -631,7 +641,11 @@ async def call_hook_agw( execution_id, status = await self._execute_workflow_via_agw( agw_client, execute_tool, hook, user_token, message, headers ) - logger.info("Workflow triggered: execution_id=%s, initial_status=%s", execution_id, status) + logger.info( + "Workflow triggered: execution_id=%s, initial_status=%s", + execution_id, + status, + ) return await self._poll_hook_execution( agw_client, get_exec_tool, hook, execution_id, user_token, status ) From cd266bd3ad08d4febdf57c3ce9a16a8167868a9a Mon Sep 17 00:00:00 2001 From: Rishi Kunnath Date: Wed, 19 Aug 2026 09:22:38 +0530 Subject: [PATCH 5/5] review comments --- src/sap_cloud_sdk/extensibility/client.py | 1 - 1 file changed, 1 deletion(-) diff --git a/src/sap_cloud_sdk/extensibility/client.py b/src/sap_cloud_sdk/extensibility/client.py index 20c49419..d6deda28 100644 --- a/src/sap_cloud_sdk/extensibility/client.py +++ b/src/sap_cloud_sdk/extensibility/client.py @@ -401,7 +401,6 @@ async def _discover_n8n_tools( self, agw_client: Any, user_token: Optional[str] ) -> tuple[Any, Any]: tools = await agw_client.list_mcp_tools(user_token=user_token or None) - logger.info("Listed %d MCP tools from Agent Gateway", len(tools)) execute_tool = next( (