From 37cba9765295a2d6d2fe340d8b6a21beb2343c21 Mon Sep 17 00:00:00 2001 From: "lixuefei.nice" Date: Thu, 23 Jul 2026 21:00:39 +0800 Subject: [PATCH 1/8] feat(skill-sandbox): pre commit --- tests/tools/builtin_tools/test_agentkit.py | 95 +++++++++ .../builtin_tools/test_run_sandbox_agent.py | 162 ++++++++++++++- veadk/tools/builtin_tools/_agentkit.py | 43 ++++ veadk/tools/builtin_tools/execute_skills.py | 195 +++++++++++++++++- 4 files changed, 491 insertions(+), 4 deletions(-) diff --git a/tests/tools/builtin_tools/test_agentkit.py b/tests/tools/builtin_tools/test_agentkit.py index 3233974c..64ef0614 100644 --- a/tests/tools/builtin_tools/test_agentkit.py +++ b/tests/tools/builtin_tools/test_agentkit.py @@ -203,5 +203,100 @@ def test_builds_exec_bash_invoke_tool_request(self): ) +class TestEnsureAgentkitSessionEndpoint(unittest.TestCase): + @classmethod + def setUpClass(cls): + cls.agentkit_module = _load_agentkit_module() + + def test_creates_session_and_prefers_public_endpoint(self): + captured = {} + + class FakeCreateSessionRequest: + def __init__(self, **kwargs): + captured["create_request"] = kwargs + + class FakeGetSessionRequest: + def __init__(self, **kwargs): + captured["get_request"] = kwargs + + class FakeClient: + def __init__(self, **kwargs): + captured["client"] = kwargs + + def create_session(self, _request): + return types.SimpleNamespace(session_id="session-1") + + def get_session(self, _request): + return types.SimpleNamespace( + endpoint="https://public.example", + internal_endpoint="http://internal.example", + ) + + fake_tools_types = types.ModuleType("agentkit.sdk.tools.types") + fake_tools_types.CreateSessionRequest = FakeCreateSessionRequest + fake_tools_types.GetSessionRequest = FakeGetSessionRequest + fake_tools_client = types.ModuleType("agentkit.sdk.tools.client") + fake_tools_client.AgentkitToolsClient = FakeClient + fake_tools_package = types.ModuleType("agentkit.sdk.tools") + fake_tools_package.types = fake_tools_types + fake_sdk_package = types.ModuleType("agentkit.sdk") + fake_agentkit_package = types.ModuleType("agentkit") + + with patch.dict( + sys.modules, + { + "agentkit": fake_agentkit_package, + "agentkit.sdk": fake_sdk_package, + "agentkit.sdk.tools": fake_tools_package, + "agentkit.sdk.tools.types": fake_tools_types, + "agentkit.sdk.tools.client": fake_tools_client, + }, + ): + with ( + patch.object( + self.agentkit_module, + "get_agentkit_endpoint_config", + return_value=("agentkit", "cn-beijing", "host", "https"), + ), + patch.object( + self.agentkit_module, + "get_agentkit_credentials", + return_value=("ak", "sk", {"X-Security-Token": "token"}), + ), + ): + endpoint = self.agentkit_module.ensure_agentkit_session_endpoint( + tool_id="tool-1", + tool_user_session_id="user-session-1", + tool_state={"state": "value"}, + ttl=900, + ) + + self.assertEqual(endpoint, "https://public.example") + self.assertEqual( + captured["client"], + { + "access_key": "ak", + "secret_key": "sk", + "region": "cn-beijing", + "session_token": "token", + }, + ) + self.assertEqual( + captured["create_request"], + { + "ToolId": "tool-1", + "UserSessionId": "user-session-1", + "Ttl": 900, + }, + ) + self.assertEqual( + captured["get_request"], + { + "ToolId": "tool-1", + "SessionId": "session-1", + }, + ) + + if __name__ == "__main__": unittest.main() diff --git a/tests/tools/builtin_tools/test_run_sandbox_agent.py b/tests/tools/builtin_tools/test_run_sandbox_agent.py index bd7f9d8a..09e39062 100644 --- a/tests/tools/builtin_tools/test_run_sandbox_agent.py +++ b/tests/tools/builtin_tools/test_run_sandbox_agent.py @@ -76,7 +76,10 @@ def _load_run_sandbox_agent_module(): return module -def _load_execute_skills_module(run_sandbox_agent): +def _load_execute_skills_module( + run_sandbox_agent, + ensure_agentkit_session_endpoint=lambda **_kwargs: "", +): module_path = ( Path(__file__).resolve().parents[3] / "veadk" @@ -101,12 +104,17 @@ def _load_execute_skills_module(run_sandbox_agent): fake_agentkit = types.ModuleType("veadk.tools.builtin_tools._agentkit") fake_agentkit.get_agentkit_account_id = lambda _state: "test-account" fake_agentkit.resolve_agentkit_tool_id = lambda _name: "test-tool" + fake_agentkit.ensure_agentkit_session_endpoint = ensure_agentkit_session_endpoint fake_runner = types.ModuleType("veadk.tools.builtin_tools.run_sandbox_agent") fake_runner.run_sandbox_agent = run_sandbox_agent fake_utils = types.ModuleType("veadk.utils") fake_utils.__path__ = [] # type: ignore[attr-defined] fake_logger = types.ModuleType("veadk.utils.logger") - fake_logger.get_logger = lambda _name: object() + fake_logger.get_logger = lambda _name: types.SimpleNamespace( + debug=lambda *_args, **_kwargs: None, + warning=lambda *_args, **_kwargs: None, + error=lambda *_args, **_kwargs: None, + ) stub_modules = { "google": fake_google, @@ -212,5 +220,155 @@ def fake_run_sandbox_agent(**kwargs): ) +class TestExecuteSkillsSkillApi(unittest.TestCase): + def _tool_context(self): + invocation_context = types.SimpleNamespace( + session=types.SimpleNamespace(id="session-1"), + agent=types.SimpleNamespace(name="agent"), + user_id="user", + ) + return types.SimpleNamespace( + state={"TIP_TOKEN_KEY": "tip-from-state"}, + _invocation_context=invocation_context, + ) + + def test_prefers_new_skill_execute_api_when_endpoint_is_available(self): + captured_requests = [] + + class FakeResponse: + def __enter__(self): + return self + + def __exit__(self, *_args): + return None + + def read(self): + return b'{"content": "api result"}' + + def fake_urlopen(request, timeout=None): + captured_requests.append((request, timeout)) + return FakeResponse() + + fallback_calls = [] + + def fake_run_sandbox_agent(**kwargs): + fallback_calls.append(kwargs) + return "fallback" + + module = _load_execute_skills_module( + fake_run_sandbox_agent, + ensure_agentkit_session_endpoint=lambda **_kwargs: "https://sandbox.test", + ) + + with patch.object(module.request, "urlopen", fake_urlopen): + result = module.execute_skills("do work", tool_context=self._tool_context()) + + self.assertEqual(result, "api result") + self.assertEqual([], fallback_calls) + self.assertEqual(1, len(captured_requests)) + request_obj, timeout = captured_requests[0] + self.assertEqual("https://sandbox.test/v1/skills/execute", request_obj.full_url) + self.assertEqual(900, timeout) + self.assertEqual("POST", request_obj.get_method()) + self.assertEqual("tip-from-state", request_obj.headers["X-tip-token-key"]) + self.assertIn(b'"prompt": "do work"', request_obj.data) + + def test_falls_back_to_legacy_runcode_when_skill_api_returns_404(self): + class NotFoundResponse: + def read(self): + return b"not found" + + def fake_urlopen(_request, timeout=None): + raise module.error.HTTPError( + url="https://sandbox.test/v1/skills/execute", + code=404, + msg="Not Found", + hdrs={}, + fp=NotFoundResponse(), + ) + + fallback_calls = [] + + def fake_run_sandbox_agent(**kwargs): + fallback_calls.append(kwargs) + return "legacy result" + + module = _load_execute_skills_module( + fake_run_sandbox_agent, + ensure_agentkit_session_endpoint=lambda **_kwargs: "https://sandbox.test", + ) + + with patch.object(module.request, "urlopen", fake_urlopen): + result = module.execute_skills("do work", tool_context=self._tool_context()) + + self.assertEqual(result, "legacy result") + self.assertEqual(1, len(fallback_calls)) + self.assertEqual("do work", fallback_calls[0]["workflow_prompt"]) + + def test_falls_back_to_legacy_runcode_when_session_endpoint_is_unavailable(self): + fallback_calls = [] + + def fake_run_sandbox_agent(**kwargs): + fallback_calls.append(kwargs) + return "legacy result" + + def raise_endpoint_error(**_kwargs): + raise RuntimeError("session unsupported") + + module = _load_execute_skills_module( + fake_run_sandbox_agent, + ensure_agentkit_session_endpoint=raise_endpoint_error, + ) + + result = module.execute_skills("do work", tool_context=self._tool_context()) + + self.assertEqual(result, "legacy result") + self.assertEqual(1, len(fallback_calls)) + self.assertEqual("do work", fallback_calls[0]["workflow_prompt"]) + + def test_stream_mode_aggregates_text_chunks_from_skill_api_sse(self): + sse_body = ( + "event: chunk\n" + 'data: {"request_id":"req_1","type":"progress","content":"started","metadata":{}}\n\n' + "event: chunk\n" + 'data: {"request_id":"req_1","type":"text","content":"hello ","metadata":{}}\n\n' + "event: chunk\n" + 'data: {"request_id":"req_1","type":"text","content":"world","metadata":{}}\n\n' + "event: done\n" + 'data: {"request_id":"req_1","type":"progress","content":"done","metadata":{}}\n\n' + ).encode() + + class FakeResponse: + def __enter__(self): + return self + + def __exit__(self, *_args): + return None + + def read(self): + return sse_body + + captured_urls = [] + + def fake_urlopen(request, timeout=None): + captured_urls.append(request.full_url) + return FakeResponse() + + module = _load_execute_skills_module( + lambda **_kwargs: "fallback", + ensure_agentkit_session_endpoint=lambda **_kwargs: "https://sandbox.test/", + ) + + with patch.object(module.request, "urlopen", fake_urlopen): + result = module.execute_skills( + "do work", + tool_context=self._tool_context(), + prefer_stream=True, + ) + + self.assertEqual(result, "hello world") + self.assertEqual(["https://sandbox.test/v1/skills/stream"], captured_urls) + + if __name__ == "__main__": unittest.main() diff --git a/veadk/tools/builtin_tools/_agentkit.py b/veadk/tools/builtin_tools/_agentkit.py index c51621f2..90f99b23 100644 --- a/veadk/tools/builtin_tools/_agentkit.py +++ b/veadk/tools/builtin_tools/_agentkit.py @@ -206,3 +206,46 @@ def invoke_agentkit_exec_bash( header=header, scheme=scheme, ) + + +def ensure_agentkit_session_endpoint( + *, + tool_id: str, + tool_user_session_id: str, + tool_state: Optional[dict[str, Any]] = None, + ttl: int = 1800, + prefer_internal_endpoint: bool = False, +) -> str: + """Create or reuse an AgentKit tool session and return its HTTP endpoint.""" + from agentkit.sdk.tools import types as tools_types + from agentkit.sdk.tools.client import AgentkitToolsClient + + _, region, _, _ = get_agentkit_endpoint_config() + ak, sk, header = get_agentkit_credentials(tool_state) + session_token = header.get("X-Security-Token", "") + client = AgentkitToolsClient( + access_key=ak, + secret_key=sk, + region=region, + session_token=session_token, + ) + session = client.create_session( + tools_types.CreateSessionRequest( + ToolId=tool_id, + UserSessionId=tool_user_session_id, + Ttl=ttl, + ) + ) + session_id = session.session_id + if not session_id: + return session.internal_endpoint or session.endpoint or "" + + current_session = client.get_session( + tools_types.GetSessionRequest( + ToolId=tool_id, + SessionId=session_id, + ) + ) + if prefer_internal_endpoint: + return current_session.internal_endpoint or current_session.endpoint or "" + return current_session.endpoint or current_session.internal_endpoint or "" diff --git a/veadk/tools/builtin_tools/execute_skills.py b/veadk/tools/builtin_tools/execute_skills.py index 029e7416..6ea5e2ea 100644 --- a/veadk/tools/builtin_tools/execute_skills.py +++ b/veadk/tools/builtin_tools/execute_skills.py @@ -12,11 +12,16 @@ # See the License for the specific language governing permissions and # limitations under the License. +import json +import os from typing import Optional +from urllib import error, request +from urllib.parse import urljoin from google.adk.tools import ToolContext from veadk.tools.builtin_tools._agentkit import ( + ensure_agentkit_session_endpoint, get_agentkit_account_id, resolve_agentkit_tool_id, ) @@ -26,10 +31,184 @@ logger = get_logger(__name__) +_SKILL_API_COMPATIBILITY_STATUS_CODES = frozenset({404, 405}) +_SKILL_API_TIMEOUT = 900 + + +class _SkillApiCompatibilityMiss(Exception): + """Raised when the sandbox does not expose the new Skill HTTP API.""" + + +def _tool_user_session_id(tool_context: ToolContext) -> str: + invocation_context = tool_context._invocation_context + session_id = invocation_context.session.id + agent_name = invocation_context.agent.name + user_id = invocation_context.user_id + return agent_name + "_" + user_id + "_" + session_id + + +def _tip_token_key(tool_context: ToolContext | None) -> str | None: + if tool_context is None: + return os.getenv("TIP_TOKEN_KEY") or None + state = tool_context.state or {} + return ( + state.get("TIP_TOKEN_KEY") + or state.get("tip_token_key") + or os.getenv("TIP_TOKEN_KEY") + or None + ) + + +def _skill_api_enabled(env_vars: Optional[dict[str, str]]) -> bool: + protocol = os.getenv("VEADK_EXECUTE_SKILLS_PROTOCOL", "auto").strip().lower() + if protocol in {"legacy", "runcode", "run_code"}: + return False + # Per-execution env vars are guaranteed by the legacy RunCode path. The new + # Skill HTTP API does not accept arbitrary per-request env overrides. + return not env_vars + + +def _skill_api_url(endpoint: str, path: str) -> str: + if not endpoint: + raise _SkillApiCompatibilityMiss("AgentKit session endpoint is empty") + return urljoin(endpoint.rstrip("/") + "/", path.lstrip("/")) + + +def _post_skill_api_json( + *, + endpoint: str, + path: str, + payload: dict[str, object], + tip_token_key: str | None, + timeout: int, +) -> bytes: + headers = { + "Content-Type": "application/json", + "Accept": "application/json, text/event-stream", + } + if tip_token_key: + headers["X-Tip-Token-Key"] = tip_token_key + + req = request.Request( + _skill_api_url(endpoint, path), + data=json.dumps(payload).encode("utf-8"), + headers=headers, + method="POST", + ) + try: + with request.urlopen(req, timeout=timeout) as response: + return response.read() + except error.HTTPError as exc: + if exc.code in _SKILL_API_COMPATIBILITY_STATUS_CODES: + raise _SkillApiCompatibilityMiss( + f"Skill HTTP API is not available: {exc.code}" + ) from exc + detail = exc.read().decode("utf-8", errors="replace") + raise RuntimeError( + f"Skill HTTP API request failed with HTTP {exc.code}: {detail}" + ) from exc + except error.URLError as exc: + raise _SkillApiCompatibilityMiss( + f"Skill HTTP API endpoint is not reachable: {exc.reason}" + ) from exc + + +def _parse_skill_execute_response(raw: bytes) -> str: + try: + payload = json.loads(raw.decode("utf-8")) + except json.JSONDecodeError: + return raw.decode("utf-8", errors="replace") + + if isinstance(payload, dict): + if isinstance(payload.get("content"), str): + return payload["content"] + data = payload.get("data") + if isinstance(data, dict) and isinstance(data.get("content"), str): + return data["content"] + return json.dumps(payload, ensure_ascii=False) + + +def _parse_skill_stream_response(raw: bytes) -> str: + chunks: list[str] = [] + event_name = "message" + data_lines: list[str] = [] + + def flush_event() -> None: + nonlocal event_name, data_lines + if not data_lines: + event_name = "message" + return + data = "\n".join(data_lines) + try: + payload = json.loads(data) + except json.JSONDecodeError: + payload = {} + + if event_name == "error": + content = payload.get("content") if isinstance(payload, dict) else None + if isinstance(content, str): + raise RuntimeError(content) + raise RuntimeError(data) + if isinstance(payload, dict) and payload.get("type") == "text": + content = payload.get("content") + if isinstance(content, str): + chunks.append(content) + + event_name = "message" + data_lines = [] + + for line in raw.decode("utf-8", errors="replace").splitlines(): + if not line: + flush_event() + continue + if line.startswith(":"): + continue + if line.startswith("event:"): + event_name = line[len("event:") :].strip() + elif line.startswith("data:"): + data_lines.append(line[len("data:") :].strip()) + + flush_event() + return "".join(chunks) + + +def _execute_skills_via_skill_api( + *, + workflow_prompt: str, + tool_id: str, + tool_context: ToolContext, + prefer_stream: bool, + timeout: int, +) -> str: + try: + endpoint = ensure_agentkit_session_endpoint( + tool_id=tool_id, + tool_user_session_id=_tool_user_session_id(tool_context), + tool_state=tool_context.state if tool_context else None, + ttl=max(timeout, 1800), + ) + except Exception as exc: + raise _SkillApiCompatibilityMiss( + f"AgentKit session endpoint is not available: {exc}" + ) from exc + path = "/v1/skills/stream" if prefer_stream else "/v1/skills/execute" + raw = _post_skill_api_json( + endpoint=endpoint, + path=path, + payload={"prompt": workflow_prompt}, + tip_token_key=_tip_token_key(tool_context), + timeout=timeout, + ) + if prefer_stream: + return _parse_skill_stream_response(raw) + return _parse_skill_execute_response(raw) + + def execute_skills( workflow_prompt: str, tool_context: ToolContext = None, env_vars: Optional[dict[str, str]] = None, + prefer_stream: bool = False, ) -> str: """Execute skills in a sandbox and return the output. @@ -43,10 +222,22 @@ def execute_skills( Returns: str: The output of the code execution. """ - timeout = 900 + timeout = _SKILL_API_TIMEOUT tool_id = resolve_agentkit_tool_id("AGENTKIT_TOOL_ID_SKILLS") - account_id = get_agentkit_account_id(tool_context.state if tool_context else None) + if tool_context is not None and _skill_api_enabled(env_vars): + try: + return _execute_skills_via_skill_api( + workflow_prompt=workflow_prompt, + tool_id=tool_id, + tool_context=tool_context, + prefer_stream=prefer_stream, + timeout=timeout, + ) + except _SkillApiCompatibilityMiss as exc: + logger.debug(f"Falling back to legacy SkillEnv execution: {exc}") + + account_id = get_agentkit_account_id(tool_context.state if tool_context else None) extra_env_vars = {} if account_id: extra_env_vars["TOS_SKILLS_DIR"] = ( From 998407b3e8d1e1a3896837f2cbda593bd40154a6 Mon Sep 17 00:00:00 2001 From: "lixuefei.nice" Date: Thu, 23 Jul 2026 21:28:00 +0800 Subject: [PATCH 2/8] feat(skill-sandbox): fill new sandbox --- .../builtin_tools/test_run_sandbox_agent.py | 176 +++++++++++++++++- veadk/tools/builtin_tools/execute_skills.py | 2 + 2 files changed, 176 insertions(+), 2 deletions(-) diff --git a/tests/tools/builtin_tools/test_run_sandbox_agent.py b/tests/tools/builtin_tools/test_run_sandbox_agent.py index 09e39062..a1e92db8 100644 --- a/tests/tools/builtin_tools/test_run_sandbox_agent.py +++ b/tests/tools/builtin_tools/test_run_sandbox_agent.py @@ -278,7 +278,10 @@ class NotFoundResponse: def read(self): return b"not found" - def fake_urlopen(_request, timeout=None): + def close(self): + return None + + def fake_urlopen(_request, **_kwargs): raise module.error.HTTPError( url="https://sandbox.test/v1/skills/execute", code=404, @@ -305,6 +308,40 @@ def fake_run_sandbox_agent(**kwargs): self.assertEqual(1, len(fallback_calls)) self.assertEqual("do work", fallback_calls[0]["workflow_prompt"]) + def test_falls_back_to_legacy_runcode_when_skill_api_returns_405(self): + class MethodNotAllowedResponse: + def read(self): + return b"method not allowed" + + def close(self): + return None + + def fake_urlopen(_request, **_kwargs): + raise module.error.HTTPError( + url="https://sandbox.test/v1/skills/execute", + code=405, + msg="Method Not Allowed", + hdrs={}, + fp=MethodNotAllowedResponse(), + ) + + fallback_calls = [] + + def fake_run_sandbox_agent(**kwargs): + fallback_calls.append(kwargs) + return "legacy result" + + module = _load_execute_skills_module( + fake_run_sandbox_agent, + ensure_agentkit_session_endpoint=lambda **_kwargs: "https://sandbox.test", + ) + + with patch.object(module.request, "urlopen", fake_urlopen): + result = module.execute_skills("do work", tool_context=self._tool_context()) + + self.assertEqual(result, "legacy result") + self.assertEqual(1, len(fallback_calls)) + def test_falls_back_to_legacy_runcode_when_session_endpoint_is_unavailable(self): fallback_calls = [] @@ -326,6 +363,108 @@ def raise_endpoint_error(**_kwargs): self.assertEqual(1, len(fallback_calls)) self.assertEqual("do work", fallback_calls[0]["workflow_prompt"]) + def test_env_vars_disable_skill_api_and_preserve_legacy_per_request_env(self): + fallback_calls = [] + + def fake_run_sandbox_agent(**kwargs): + fallback_calls.append(kwargs) + return "legacy result" + + module = _load_execute_skills_module( + fake_run_sandbox_agent, + ensure_agentkit_session_endpoint=lambda **_kwargs: self.fail( + "Skill API should be disabled when env_vars are provided" + ), + ) + + result = module.execute_skills( + "do work", + tool_context=self._tool_context(), + env_vars={"CUSTOM_VALUE": "custom"}, + ) + + self.assertEqual(result, "legacy result") + self.assertEqual( + { + "TOS_SKILLS_DIR": "tos://agentkit-platform-test-account/skills/", + "CUSTOM_VALUE": "custom", + }, + fallback_calls[0]["extra_env_vars"], + ) + + def test_legacy_protocol_env_disables_skill_api(self): + fallback_calls = [] + + def fake_run_sandbox_agent(**kwargs): + fallback_calls.append(kwargs) + return "legacy result" + + module = _load_execute_skills_module( + fake_run_sandbox_agent, + ensure_agentkit_session_endpoint=lambda **_kwargs: self.fail( + "Skill API should be disabled by VEADK_EXECUTE_SKILLS_PROTOCOL" + ), + ) + + with patch.dict( + module.os.environ, + {"VEADK_EXECUTE_SKILLS_PROTOCOL": "legacy"}, + clear=False, + ): + result = module.execute_skills("do work", tool_context=self._tool_context()) + + self.assertEqual(result, "legacy result") + self.assertEqual(1, len(fallback_calls)) + + def test_missing_tool_context_keeps_legacy_execution_path(self): + fallback_calls = [] + + def fake_run_sandbox_agent(**kwargs): + fallback_calls.append(kwargs) + return "legacy result" + + module = _load_execute_skills_module( + fake_run_sandbox_agent, + ensure_agentkit_session_endpoint=lambda **_kwargs: self.fail( + "Skill API requires a tool_context" + ), + ) + + result = module.execute_skills("do work", tool_context=None) + + self.assertEqual(result, "legacy result") + self.assertEqual(1, len(fallback_calls)) + self.assertIsNone(fallback_calls[0]["tool_context"]) + + def test_non_compatibility_skill_api_http_error_is_not_swallowed(self): + class ServerErrorResponse: + def read(self): + return b"internal error" + + def close(self): + return None + + def fake_urlopen(_request, **_kwargs): + raise module.error.HTTPError( + url="https://sandbox.test/v1/skills/execute", + code=500, + msg="Internal Server Error", + hdrs={}, + fp=ServerErrorResponse(), + ) + + fallback_calls = [] + module = _load_execute_skills_module( + lambda **kwargs: fallback_calls.append(kwargs) or "legacy result", + ensure_agentkit_session_endpoint=lambda **_kwargs: "https://sandbox.test", + ) + + with patch.object(module.request, "urlopen", fake_urlopen): + with self.assertRaisesRegex(RuntimeError, "HTTP 500: internal error"): + module.execute_skills("do work", tool_context=self._tool_context()) + + self.assertEqual([], fallback_calls) + def test_stream_mode_aggregates_text_chunks_from_skill_api_sse(self): sse_body = ( "event: chunk\n" @@ -350,7 +489,7 @@ def read(self): captured_urls = [] - def fake_urlopen(request, timeout=None): + def fake_urlopen(request, **_kwargs): captured_urls.append(request.full_url) return FakeResponse() @@ -369,6 +508,39 @@ def fake_urlopen(request, timeout=None): self.assertEqual(result, "hello world") self.assertEqual(["https://sandbox.test/v1/skills/stream"], captured_urls) + def test_stream_mode_raises_skill_api_error_event(self): + sse_body = ( + "event: error\n" + 'data: {"request_id":"req_1","type":"text","content":"skill failed","metadata":{}}\n\n' + ).encode() + + class FakeResponse: + def __enter__(self): + return self + + def __exit__(self, *_args): + return None + + def read(self): + return sse_body + + module = _load_execute_skills_module( + lambda **_kwargs: "fallback", + ensure_agentkit_session_endpoint=lambda **_kwargs: "https://sandbox.test", + ) + + with patch.object( + module.request, + "urlopen", + lambda *_args, **_kwargs: FakeResponse(), + ): + with self.assertRaisesRegex(RuntimeError, "skill failed"): + module.execute_skills( + "do work", + tool_context=self._tool_context(), + prefer_stream=True, + ) + if __name__ == "__main__": unittest.main() diff --git a/veadk/tools/builtin_tools/execute_skills.py b/veadk/tools/builtin_tools/execute_skills.py index 6ea5e2ea..de0d21a9 100644 --- a/veadk/tools/builtin_tools/execute_skills.py +++ b/veadk/tools/builtin_tools/execute_skills.py @@ -12,6 +12,8 @@ # See the License for the specific language governing permissions and # limitations under the License. +from __future__ import annotations + import json import os from typing import Optional From 3c5ba5b0c3a245b58ddb71f1fe602b365415d1c3 Mon Sep 17 00:00:00 2001 From: "lixuefei.nice" Date: Tue, 28 Jul 2026 15:41:06 +0800 Subject: [PATCH 3/8] feat(skil sandbox): add 404 error hint --- frontend/src/adk/runSseError.ts | 18 ++++++++++++++++-- frontend/tests/runSseError.test.mjs | 8 ++++++++ 2 files changed, 24 insertions(+), 2 deletions(-) diff --git a/frontend/src/adk/runSseError.ts b/frontend/src/adk/runSseError.ts index c32eb4e1..1e10a3c5 100644 --- a/frontend/src/adk/runSseError.ts +++ b/frontend/src/adk/runSseError.ts @@ -4,11 +4,25 @@ const PERSISTENT_MEMORY_HINT = "提示:该 404 可能是多实例部署使用 in-memory 或 SQLite 短期记忆导致的:" + "请求落到不同实例后,会话无法找到。请改用基于数据库的持久化短期记忆存储。"; +const MISSING_RUN_SSE_HINT = + "提示:如果当前连接的是旧版 Skill 沙箱,该沙箱可能接口不存在。" + + "请升级 Skill 沙箱镜像或切换到支持 /run_sse 的新版沙箱。"; + /** Preserve a run error and append guidance for the common cross-instance 404. */ export function formatRunSseError(error: unknown): string { const message = String(error); - if (!SESSION_NOT_FOUND_PATTERN.test(message) || message.includes(PERSISTENT_MEMORY_HINT)) { + if (!SESSION_NOT_FOUND_PATTERN.test(message)) { + return message; + } + const hints = []; + if (!message.includes(PERSISTENT_MEMORY_HINT)) { + hints.push(PERSISTENT_MEMORY_HINT); + } + if (!message.includes(MISSING_RUN_SSE_HINT)) { + hints.push(MISSING_RUN_SSE_HINT); + } + if (hints.length === 0) { return message; } - return `${message}\n\n${PERSISTENT_MEMORY_HINT}`; + return `${message}\n\n${hints.join("\n\n")}`; } diff --git a/frontend/tests/runSseError.test.mjs b/frontend/tests/runSseError.test.mjs index 253bbc5b..7611c837 100644 --- a/frontend/tests/runSseError.test.mjs +++ b/frontend/tests/runSseError.test.mjs @@ -28,6 +28,14 @@ test("adds persistent-memory guidance to run_sse 404 variants", () => { } }); +test("adds upgrade guidance when run_sse may be missing", () => { + const formatted = formatRunSseError("run_sse failed: 404"); + + assert.match(formatted, /接口不存在/); + assert.match(formatted, /升级/); + assert.match(formatted, /Skill 沙箱/); +}); + test("leaves unrelated errors unchanged", () => { for (const error of ["run_sse failed: 500", "create_session failed: 404"]) { assert.equal(formatRunSseError(error), error); From 4bec5ef2c72474e17b9addb8c32c340975858d93 Mon Sep 17 00:00:00 2001 From: "lixuefei.nice" Date: Tue, 28 Jul 2026 21:14:29 +0800 Subject: [PATCH 4/8] feat(skil sandbox): support new version skill sandbox --- .../builtin_tools/test_run_sandbox_agent.py | 37 ++++++++++++++----- veadk/tools/builtin_tools/execute_skills.py | 36 ++++++++++++++---- 2 files changed, 56 insertions(+), 17 deletions(-) diff --git a/tests/tools/builtin_tools/test_run_sandbox_agent.py b/tests/tools/builtin_tools/test_run_sandbox_agent.py index a1e92db8..4905bc82 100644 --- a/tests/tools/builtin_tools/test_run_sandbox_agent.py +++ b/tests/tools/builtin_tools/test_run_sandbox_agent.py @@ -273,7 +273,21 @@ def fake_run_sandbox_agent(**kwargs): self.assertEqual("tip-from-state", request_obj.headers["X-tip-token-key"]) self.assertIn(b'"prompt": "do work"', request_obj.data) - def test_falls_back_to_legacy_runcode_when_skill_api_returns_404(self): + def test_skill_api_url_preserves_agentkit_endpoint_query_auth(self): + module = _load_execute_skills_module( + lambda **_kwargs: "fallback", + ensure_agentkit_session_endpoint=lambda **_kwargs: "", + ) + + self.assertEqual( + "https://sandbox.test/v1/skills/execute?faasInstanceName=inst&Authorization=key", + module._skill_api_url( + "https://sandbox.test/?faasInstanceName=inst&Authorization=key", + "/v1/skills/execute", + ), + ) + + def test_requires_sandbox_upgrade_when_skill_api_returns_404(self): class NotFoundResponse: def read(self): return b"not found" @@ -302,13 +316,15 @@ def fake_run_sandbox_agent(**kwargs): ) with patch.object(module.request, "urlopen", fake_urlopen): - result = module.execute_skills("do work", tool_context=self._tool_context()) + with self.assertRaisesRegex( + RuntimeError, + r"HTTP 404.*(?:升级|upgrade).*Skill HTTP API", + ): + module.execute_skills("do work", tool_context=self._tool_context()) - self.assertEqual(result, "legacy result") - self.assertEqual(1, len(fallback_calls)) - self.assertEqual("do work", fallback_calls[0]["workflow_prompt"]) + self.assertEqual([], fallback_calls) - def test_falls_back_to_legacy_runcode_when_skill_api_returns_405(self): + def test_requires_sandbox_upgrade_when_skill_api_returns_405(self): class MethodNotAllowedResponse: def read(self): return b"method not allowed" @@ -337,10 +353,13 @@ def fake_run_sandbox_agent(**kwargs): ) with patch.object(module.request, "urlopen", fake_urlopen): - result = module.execute_skills("do work", tool_context=self._tool_context()) + with self.assertRaisesRegex( + RuntimeError, + r"HTTP 405.*(?:升级|upgrade).*Skill HTTP API", + ): + module.execute_skills("do work", tool_context=self._tool_context()) - self.assertEqual(result, "legacy result") - self.assertEqual(1, len(fallback_calls)) + self.assertEqual([], fallback_calls) def test_falls_back_to_legacy_runcode_when_session_endpoint_is_unavailable(self): fallback_calls = [] diff --git a/veadk/tools/builtin_tools/execute_skills.py b/veadk/tools/builtin_tools/execute_skills.py index de0d21a9..68b94a26 100644 --- a/veadk/tools/builtin_tools/execute_skills.py +++ b/veadk/tools/builtin_tools/execute_skills.py @@ -18,7 +18,7 @@ import os from typing import Optional from urllib import error, request -from urllib.parse import urljoin +from urllib.parse import urlsplit, urlunsplit from google.adk.tools import ToolContext @@ -33,12 +33,22 @@ logger = get_logger(__name__) -_SKILL_API_COMPATIBILITY_STATUS_CODES = frozenset({404, 405}) +_SKILL_API_UPGRADE_STATUS_CODES = frozenset({404}) _SKILL_API_TIMEOUT = 900 +_SKILL_API_UPGRADE_HINT = ( + "提示:当前 Skill 沙箱镜像未实现 Skill HTTP API " + "(/v1/skills/execute|stream),可能是旧版沙箱镜像。" + "请升级 Skill 沙箱镜像或切换到支持 Skill HTTP API 的新版沙箱。" +) + class _SkillApiCompatibilityMiss(Exception): - """Raised when the sandbox does not expose the new Skill HTTP API.""" + """Raised when the sandbox endpoint is unreachable so legacy fallback should apply.""" + + +class _SkillApiUpgradeRequired(Exception): + """Raised when the sandbox is reachable but does not expose the Skill HTTP API.""" def _tool_user_session_id(tool_context: ToolContext) -> str: @@ -73,7 +83,13 @@ def _skill_api_enabled(env_vars: Optional[dict[str, str]]) -> bool: def _skill_api_url(endpoint: str, path: str) -> str: if not endpoint: raise _SkillApiCompatibilityMiss("AgentKit session endpoint is empty") - return urljoin(endpoint.rstrip("/") + "/", path.lstrip("/")) + parts = urlsplit(endpoint) + endpoint_path = parts.path.rstrip("/") + skill_path = path.lstrip("/") + joined_path = f"{endpoint_path}/{skill_path}" if endpoint_path else f"/{skill_path}" + return urlunsplit( + (parts.scheme, parts.netloc, joined_path, parts.query, parts.fragment) + ) def _post_skill_api_json( @@ -101,9 +117,9 @@ def _post_skill_api_json( with request.urlopen(req, timeout=timeout) as response: return response.read() except error.HTTPError as exc: - if exc.code in _SKILL_API_COMPATIBILITY_STATUS_CODES: - raise _SkillApiCompatibilityMiss( - f"Skill HTTP API is not available: {exc.code}" + if exc.code in _SKILL_API_UPGRADE_STATUS_CODES: + raise _SkillApiUpgradeRequired( + f"Skill HTTP API returned HTTP {exc.code}. {_SKILL_API_UPGRADE_HINT}" ) from exc detail = exc.read().decode("utf-8", errors="replace") raise RuntimeError( @@ -236,8 +252,12 @@ def execute_skills( prefer_stream=prefer_stream, timeout=timeout, ) + except _SkillApiUpgradeRequired as exc: + raise RuntimeError(str(exc)) from exc except _SkillApiCompatibilityMiss as exc: - logger.debug(f"Falling back to legacy SkillEnv execution: {exc}") + logger.warning( + f"Skill HTTP API endpoint unreachable, falling back to legacy RunCode: {exc}" + ) account_id = get_agentkit_account_id(tool_context.state if tool_context else None) extra_env_vars = {} From a8e3b03146488ce39efa46e483a5fe7c8e691776 Mon Sep 17 00:00:00 2001 From: "lixuefei.nice" Date: Tue, 28 Jul 2026 21:32:58 +0800 Subject: [PATCH 5/8] feat(skil sandbox): 404 hint --- veadk/tools/builtin_tools/execute_skills.py | 12 +++++++++++- 1 file changed, 11 insertions(+), 1 deletion(-) diff --git a/veadk/tools/builtin_tools/execute_skills.py b/veadk/tools/builtin_tools/execute_skills.py index 68b94a26..1947a4d9 100644 --- a/veadk/tools/builtin_tools/execute_skills.py +++ b/veadk/tools/builtin_tools/execute_skills.py @@ -42,6 +42,11 @@ "请升级 Skill 沙箱镜像或切换到支持 Skill HTTP API 的新版沙箱。" ) +_SKILL_STREAM_MISSING_HINT = ( + "提示:当前 Skill 沙箱镜像未实现 /v1/skills/stream 接口(HTTP 404)。" + "请升级 Skill 沙箱镜像到支持 /v1/skills/stream 的新版沙箱。" +) + class _SkillApiCompatibilityMiss(Exception): """Raised when the sandbox endpoint is unreachable so legacy fallback should apply.""" @@ -118,8 +123,13 @@ def _post_skill_api_json( return response.read() except error.HTTPError as exc: if exc.code in _SKILL_API_UPGRADE_STATUS_CODES: + hint = ( + _SKILL_STREAM_MISSING_HINT + if path.rstrip("/").endswith("/v1/skills/stream") + else _SKILL_API_UPGRADE_HINT + ) raise _SkillApiUpgradeRequired( - f"Skill HTTP API returned HTTP {exc.code}. {_SKILL_API_UPGRADE_HINT}" + f"Skill HTTP API returned HTTP {exc.code}. {hint}" ) from exc detail = exc.read().decode("utf-8", errors="replace") raise RuntimeError( From cda9b9579c730521d0805bd09a383d0166346c56 Mon Sep 17 00:00:00 2001 From: "lixuefei.nice" Date: Tue, 28 Jul 2026 22:31:18 +0800 Subject: [PATCH 6/8] feat(skil sandbox): manage codes --- frontend/src/adk/runSseError.ts | 18 +- frontend/tests/runSseError.test.mjs | 8 - tests/tools/builtin_tools/test_agentkit.py | 57 ++++++ .../builtin_tools/test_run_sandbox_agent.py | 185 +++++------------- veadk/tools/builtin_tools/_agentkit.py | 12 +- veadk/tools/builtin_tools/execute_skills.py | 132 +++++-------- 6 files changed, 167 insertions(+), 245 deletions(-) diff --git a/frontend/src/adk/runSseError.ts b/frontend/src/adk/runSseError.ts index 1e10a3c5..c32eb4e1 100644 --- a/frontend/src/adk/runSseError.ts +++ b/frontend/src/adk/runSseError.ts @@ -4,25 +4,11 @@ const PERSISTENT_MEMORY_HINT = "提示:该 404 可能是多实例部署使用 in-memory 或 SQLite 短期记忆导致的:" + "请求落到不同实例后,会话无法找到。请改用基于数据库的持久化短期记忆存储。"; -const MISSING_RUN_SSE_HINT = - "提示:如果当前连接的是旧版 Skill 沙箱,该沙箱可能接口不存在。" + - "请升级 Skill 沙箱镜像或切换到支持 /run_sse 的新版沙箱。"; - /** Preserve a run error and append guidance for the common cross-instance 404. */ export function formatRunSseError(error: unknown): string { const message = String(error); - if (!SESSION_NOT_FOUND_PATTERN.test(message)) { - return message; - } - const hints = []; - if (!message.includes(PERSISTENT_MEMORY_HINT)) { - hints.push(PERSISTENT_MEMORY_HINT); - } - if (!message.includes(MISSING_RUN_SSE_HINT)) { - hints.push(MISSING_RUN_SSE_HINT); - } - if (hints.length === 0) { + if (!SESSION_NOT_FOUND_PATTERN.test(message) || message.includes(PERSISTENT_MEMORY_HINT)) { return message; } - return `${message}\n\n${hints.join("\n\n")}`; + return `${message}\n\n${PERSISTENT_MEMORY_HINT}`; } diff --git a/frontend/tests/runSseError.test.mjs b/frontend/tests/runSseError.test.mjs index 7611c837..253bbc5b 100644 --- a/frontend/tests/runSseError.test.mjs +++ b/frontend/tests/runSseError.test.mjs @@ -28,14 +28,6 @@ test("adds persistent-memory guidance to run_sse 404 variants", () => { } }); -test("adds upgrade guidance when run_sse may be missing", () => { - const formatted = formatRunSseError("run_sse failed: 404"); - - assert.match(formatted, /接口不存在/); - assert.match(formatted, /升级/); - assert.match(formatted, /Skill 沙箱/); -}); - test("leaves unrelated errors unchanged", () => { for (const error of ["run_sse failed: 500", "create_session failed: 404"]) { assert.equal(formatRunSseError(error), error); diff --git a/tests/tools/builtin_tools/test_agentkit.py b/tests/tools/builtin_tools/test_agentkit.py index 64ef0614..51276c11 100644 --- a/tests/tools/builtin_tools/test_agentkit.py +++ b/tests/tools/builtin_tools/test_agentkit.py @@ -297,6 +297,63 @@ def get_session(self, _request): }, ) + def test_uses_create_session_endpoint_without_getting_session(self): + class FakeCreateSessionRequest: + def __init__(self, **_kwargs): + pass + + class FakeClient: + def __init__(self, **_kwargs): + pass + + def create_session(self, _request): + return types.SimpleNamespace( + session_id="session-1", + endpoint="https://public.example", + internal_endpoint="http://internal.example", + ) + + def get_session(self, _request): + raise AssertionError( + "get_session should not be called when create returns endpoint" + ) + + fake_tools_types = types.ModuleType("agentkit.sdk.tools.types") + fake_tools_types.CreateSessionRequest = FakeCreateSessionRequest + fake_tools_client = types.ModuleType("agentkit.sdk.tools.client") + fake_tools_client.AgentkitToolsClient = FakeClient + fake_tools_package = types.ModuleType("agentkit.sdk.tools") + fake_tools_package.types = fake_tools_types + + with patch.dict( + sys.modules, + { + "agentkit": types.ModuleType("agentkit"), + "agentkit.sdk": types.ModuleType("agentkit.sdk"), + "agentkit.sdk.tools": fake_tools_package, + "agentkit.sdk.tools.types": fake_tools_types, + "agentkit.sdk.tools.client": fake_tools_client, + }, + ): + with ( + patch.object( + self.agentkit_module, + "get_agentkit_endpoint_config", + return_value=("agentkit", "cn-beijing", "host", "https"), + ), + patch.object( + self.agentkit_module, + "get_agentkit_credentials", + return_value=("ak", "sk", {}), + ), + ): + endpoint = self.agentkit_module.ensure_agentkit_session_endpoint( + tool_id="tool-1", + tool_user_session_id="user-session-1", + ) + + self.assertEqual(endpoint, "https://public.example") + if __name__ == "__main__": unittest.main() diff --git a/tests/tools/builtin_tools/test_run_sandbox_agent.py b/tests/tools/builtin_tools/test_run_sandbox_agent.py index 4905bc82..f0dbc4d9 100644 --- a/tests/tools/builtin_tools/test_run_sandbox_agent.py +++ b/tests/tools/builtin_tools/test_run_sandbox_agent.py @@ -77,8 +77,9 @@ def _load_run_sandbox_agent_module(): def _load_execute_skills_module( - run_sandbox_agent, + *, ensure_agentkit_session_endpoint=lambda **_kwargs: "", + run_sandbox_agent=lambda **_kwargs: "", ): module_path = ( Path(__file__).resolve().parents[3] @@ -196,30 +197,6 @@ def test_runner_code_overrides_the_sandbox_process_environment(self): self.assertIn('srv_pythonpath = env.get("SRV_PYTHONPATH")', code) -class TestExecuteSkillsEnvVars(unittest.TestCase): - def test_passes_custom_env_vars_to_each_sandbox_execution(self): - captured_kwargs = {} - - def fake_run_sandbox_agent(**kwargs): - captured_kwargs.update(kwargs) - return "done" - - module = _load_execute_skills_module(fake_run_sandbox_agent) - tool_context = types.SimpleNamespace(state={}) - - result = module.execute_skills( - "do work", - tool_context=tool_context, - env_vars={"CUSTOM_VALUE": "custom", "TOS_SKILLS_DIR": ""}, - ) - - self.assertEqual(result, "done") - self.assertEqual( - captured_kwargs["extra_env_vars"], - {"CUSTOM_VALUE": "custom", "TOS_SKILLS_DIR": ""}, - ) - - class TestExecuteSkillsSkillApi(unittest.TestCase): def _tool_context(self): invocation_context = types.SimpleNamespace( @@ -249,14 +226,7 @@ def fake_urlopen(request, timeout=None): captured_requests.append((request, timeout)) return FakeResponse() - fallback_calls = [] - - def fake_run_sandbox_agent(**kwargs): - fallback_calls.append(kwargs) - return "fallback" - module = _load_execute_skills_module( - fake_run_sandbox_agent, ensure_agentkit_session_endpoint=lambda **_kwargs: "https://sandbox.test", ) @@ -264,7 +234,6 @@ def fake_run_sandbox_agent(**kwargs): result = module.execute_skills("do work", tool_context=self._tool_context()) self.assertEqual(result, "api result") - self.assertEqual([], fallback_calls) self.assertEqual(1, len(captured_requests)) request_obj, timeout = captured_requests[0] self.assertEqual("https://sandbox.test/v1/skills/execute", request_obj.full_url) @@ -273,9 +242,34 @@ def fake_run_sandbox_agent(**kwargs): self.assertEqual("tip-from-state", request_obj.headers["X-tip-token-key"]) self.assertIn(b'"prompt": "do work"', request_obj.data) + def test_env_vars_use_legacy_runcode_execution(self): + captured_kwargs = {} + + def fake_run_sandbox_agent(**kwargs): + captured_kwargs.update(kwargs) + return "legacy result" + + module = _load_execute_skills_module( + ensure_agentkit_session_endpoint=lambda **_kwargs: self.fail( + "Skill API must not be used when env_vars are provided" + ), + run_sandbox_agent=fake_run_sandbox_agent, + ) + + result = module.execute_skills( + "do work", + tool_context=self._tool_context(), + env_vars={"CUSTOM_VALUE": "custom", "TOS_SKILLS_DIR": ""}, + ) + + self.assertEqual(result, "legacy result") + self.assertEqual( + {"CUSTOM_VALUE": "custom", "TOS_SKILLS_DIR": ""}, + captured_kwargs["extra_env_vars"], + ) + def test_skill_api_url_preserves_agentkit_endpoint_query_auth(self): module = _load_execute_skills_module( - lambda **_kwargs: "fallback", ensure_agentkit_session_endpoint=lambda **_kwargs: "", ) @@ -304,26 +298,17 @@ def fake_urlopen(_request, **_kwargs): fp=NotFoundResponse(), ) - fallback_calls = [] - - def fake_run_sandbox_agent(**kwargs): - fallback_calls.append(kwargs) - return "legacy result" - module = _load_execute_skills_module( - fake_run_sandbox_agent, ensure_agentkit_session_endpoint=lambda **_kwargs: "https://sandbox.test", ) with patch.object(module.request, "urlopen", fake_urlopen): with self.assertRaisesRegex( RuntimeError, - r"HTTP 404.*(?:升级|upgrade).*Skill HTTP API", + r"HTTP 404.*(?:升级|upgrade).*Skill", ): module.execute_skills("do work", tool_context=self._tool_context()) - self.assertEqual([], fallback_calls) - def test_requires_sandbox_upgrade_when_skill_api_returns_405(self): class MethodNotAllowedResponse: def read(self): @@ -341,119 +326,39 @@ def fake_urlopen(_request, **_kwargs): fp=MethodNotAllowedResponse(), ) - fallback_calls = [] - - def fake_run_sandbox_agent(**kwargs): - fallback_calls.append(kwargs) - return "legacy result" - module = _load_execute_skills_module( - fake_run_sandbox_agent, ensure_agentkit_session_endpoint=lambda **_kwargs: "https://sandbox.test", ) with patch.object(module.request, "urlopen", fake_urlopen): with self.assertRaisesRegex( RuntimeError, - r"HTTP 405.*(?:升级|upgrade).*Skill HTTP API", + r"HTTP 405.*(?:升级|upgrade).*Skill", ): module.execute_skills("do work", tool_context=self._tool_context()) - self.assertEqual([], fallback_calls) - - def test_falls_back_to_legacy_runcode_when_session_endpoint_is_unavailable(self): - fallback_calls = [] - - def fake_run_sandbox_agent(**kwargs): - fallback_calls.append(kwargs) - return "legacy result" - + def test_raises_runtime_error_when_session_endpoint_is_unavailable(self): def raise_endpoint_error(**_kwargs): raise RuntimeError("session unsupported") module = _load_execute_skills_module( - fake_run_sandbox_agent, ensure_agentkit_session_endpoint=raise_endpoint_error, ) - result = module.execute_skills("do work", tool_context=self._tool_context()) - - self.assertEqual(result, "legacy result") - self.assertEqual(1, len(fallback_calls)) - self.assertEqual("do work", fallback_calls[0]["workflow_prompt"]) - - def test_env_vars_disable_skill_api_and_preserve_legacy_per_request_env(self): - fallback_calls = [] - - def fake_run_sandbox_agent(**kwargs): - fallback_calls.append(kwargs) - return "legacy result" - - module = _load_execute_skills_module( - fake_run_sandbox_agent, - ensure_agentkit_session_endpoint=lambda **_kwargs: self.fail( - "Skill API should be disabled when env_vars are provided" - ), - ) - - result = module.execute_skills( - "do work", - tool_context=self._tool_context(), - env_vars={"CUSTOM_VALUE": "custom"}, - ) - - self.assertEqual(result, "legacy result") - self.assertEqual( - { - "TOS_SKILLS_DIR": "tos://agentkit-platform-test-account/skills/", - "CUSTOM_VALUE": "custom", - }, - fallback_calls[0]["extra_env_vars"], - ) - - def test_legacy_protocol_env_disables_skill_api(self): - fallback_calls = [] - - def fake_run_sandbox_agent(**kwargs): - fallback_calls.append(kwargs) - return "legacy result" - - module = _load_execute_skills_module( - fake_run_sandbox_agent, - ensure_agentkit_session_endpoint=lambda **_kwargs: self.fail( - "Skill API should be disabled by VEADK_EXECUTE_SKILLS_PROTOCOL" - ), - ) - - with patch.dict( - module.os.environ, - {"VEADK_EXECUTE_SKILLS_PROTOCOL": "legacy"}, - clear=False, + with self.assertRaisesRegex( + RuntimeError, r"AgentKit session endpoint is not available" ): - result = module.execute_skills("do work", tool_context=self._tool_context()) - - self.assertEqual(result, "legacy result") - self.assertEqual(1, len(fallback_calls)) - - def test_missing_tool_context_keeps_legacy_execution_path(self): - fallback_calls = [] - - def fake_run_sandbox_agent(**kwargs): - fallback_calls.append(kwargs) - return "legacy result" + module.execute_skills("do work", tool_context=self._tool_context()) + def test_missing_tool_context_raises_value_error(self): module = _load_execute_skills_module( - fake_run_sandbox_agent, ensure_agentkit_session_endpoint=lambda **_kwargs: self.fail( "Skill API requires a tool_context" ), ) - result = module.execute_skills("do work", tool_context=None) - - self.assertEqual(result, "legacy result") - self.assertEqual(1, len(fallback_calls)) - self.assertIsNone(fallback_calls[0]["tool_context"]) + with self.assertRaisesRegex(ValueError, r"tool_context is required"): + module.execute_skills("do work", tool_context=None) def test_non_compatibility_skill_api_http_error_is_not_swallowed(self): class ServerErrorResponse: @@ -472,9 +377,7 @@ def fake_urlopen(_request, **_kwargs): fp=ServerErrorResponse(), ) - fallback_calls = [] module = _load_execute_skills_module( - lambda **kwargs: fallback_calls.append(kwargs) or "legacy result", ensure_agentkit_session_endpoint=lambda **_kwargs: "https://sandbox.test", ) @@ -482,8 +385,6 @@ def fake_urlopen(_request, **_kwargs): with self.assertRaisesRegex(RuntimeError, "HTTP 500: internal error"): module.execute_skills("do work", tool_context=self._tool_context()) - self.assertEqual([], fallback_calls) - def test_stream_mode_aggregates_text_chunks_from_skill_api_sse(self): sse_body = ( "event: chunk\n" @@ -504,7 +405,10 @@ def __exit__(self, *_args): return None def read(self): - return sse_body + raise AssertionError("stream response must not be buffered with read()") + + def __iter__(self): + return iter(sse_body.splitlines(keepends=True)) captured_urls = [] @@ -513,7 +417,6 @@ def fake_urlopen(request, **_kwargs): return FakeResponse() module = _load_execute_skills_module( - lambda **_kwargs: "fallback", ensure_agentkit_session_endpoint=lambda **_kwargs: "https://sandbox.test/", ) @@ -541,10 +444,12 @@ def __exit__(self, *_args): return None def read(self): - return sse_body + raise AssertionError("stream response must not be buffered with read()") + + def __iter__(self): + return iter(sse_body.splitlines(keepends=True)) module = _load_execute_skills_module( - lambda **_kwargs: "fallback", ensure_agentkit_session_endpoint=lambda **_kwargs: "https://sandbox.test", ) diff --git a/veadk/tools/builtin_tools/_agentkit.py b/veadk/tools/builtin_tools/_agentkit.py index 90f99b23..3e4eef43 100644 --- a/veadk/tools/builtin_tools/_agentkit.py +++ b/veadk/tools/builtin_tools/_agentkit.py @@ -236,9 +236,19 @@ def ensure_agentkit_session_endpoint( Ttl=ttl, ) ) + public_endpoint = getattr(session, "endpoint", None) + internal_endpoint = getattr(session, "internal_endpoint", None) + endpoint = ( + internal_endpoint or public_endpoint + if prefer_internal_endpoint + else public_endpoint or internal_endpoint + ) + if endpoint: + return endpoint + session_id = session.session_id if not session_id: - return session.internal_endpoint or session.endpoint or "" + return "" current_session = client.get_session( tools_types.GetSessionRequest( diff --git a/veadk/tools/builtin_tools/execute_skills.py b/veadk/tools/builtin_tools/execute_skills.py index 1947a4d9..266c38c9 100644 --- a/veadk/tools/builtin_tools/execute_skills.py +++ b/veadk/tools/builtin_tools/execute_skills.py @@ -16,6 +16,7 @@ import json import os +from collections.abc import Iterable from typing import Optional from urllib import error, request from urllib.parse import urlsplit, urlunsplit @@ -28,32 +29,22 @@ resolve_agentkit_tool_id, ) from veadk.tools.builtin_tools.run_sandbox_agent import run_sandbox_agent -from veadk.utils.logger import get_logger -logger = get_logger(__name__) - -_SKILL_API_UPGRADE_STATUS_CODES = frozenset({404}) +_SKILL_API_UPGRADE_STATUS_CODES = frozenset({404, 405}) _SKILL_API_TIMEOUT = 900 -_SKILL_API_UPGRADE_HINT = ( - "提示:当前 Skill 沙箱镜像未实现 Skill HTTP API " - "(/v1/skills/execute|stream),可能是旧版沙箱镜像。" - "请升级 Skill 沙箱镜像或切换到支持 Skill HTTP API 的新版沙箱。" -) - -_SKILL_STREAM_MISSING_HINT = ( - "提示:当前 Skill 沙箱镜像未实现 /v1/skills/stream 接口(HTTP 404)。" - "请升级 Skill 沙箱镜像到支持 /v1/skills/stream 的新版沙箱。" -) - - -class _SkillApiCompatibilityMiss(Exception): - """Raised when the sandbox endpoint is unreachable so legacy fallback should apply.""" - -class _SkillApiUpgradeRequired(Exception): - """Raised when the sandbox is reachable but does not expose the Skill HTTP API.""" +def _skill_api_upgrade_hint(path: str) -> str: + api_path = ( + "/v1/skills/stream" + if path.rstrip("/").endswith("/stream") + else "/v1/skills/execute" + ) + return ( + f"提示:当前 Skill 沙箱镜像未实现 {api_path} 接口,可能是旧版沙箱镜像。" + "请升级 Skill 沙箱镜像或切换到支持 Skill HTTP API 的新版沙箱。" + ) def _tool_user_session_id(tool_context: ToolContext) -> str: @@ -64,9 +55,7 @@ def _tool_user_session_id(tool_context: ToolContext) -> str: return agent_name + "_" + user_id + "_" + session_id -def _tip_token_key(tool_context: ToolContext | None) -> str | None: - if tool_context is None: - return os.getenv("TIP_TOKEN_KEY") or None +def _tip_token_key(tool_context: ToolContext) -> str | None: state = tool_context.state or {} return ( state.get("TIP_TOKEN_KEY") @@ -76,18 +65,9 @@ def _tip_token_key(tool_context: ToolContext | None) -> str | None: ) -def _skill_api_enabled(env_vars: Optional[dict[str, str]]) -> bool: - protocol = os.getenv("VEADK_EXECUTE_SKILLS_PROTOCOL", "auto").strip().lower() - if protocol in {"legacy", "runcode", "run_code"}: - return False - # Per-execution env vars are guaranteed by the legacy RunCode path. The new - # Skill HTTP API does not accept arbitrary per-request env overrides. - return not env_vars - - def _skill_api_url(endpoint: str, path: str) -> str: if not endpoint: - raise _SkillApiCompatibilityMiss("AgentKit session endpoint is empty") + raise RuntimeError("AgentKit session endpoint is empty") parts = urlsplit(endpoint) endpoint_path = parts.path.rstrip("/") skill_path = path.lstrip("/") @@ -104,7 +84,8 @@ def _post_skill_api_json( payload: dict[str, object], tip_token_key: str | None, timeout: int, -) -> bytes: + stream: bool, +) -> str: headers = { "Content-Type": "application/json", "Accept": "application/json, text/event-stream", @@ -120,23 +101,21 @@ def _post_skill_api_json( ) try: with request.urlopen(req, timeout=timeout) as response: - return response.read() + if stream: + return _parse_skill_stream_response(response) + return _parse_skill_execute_response(response.read()) except error.HTTPError as exc: if exc.code in _SKILL_API_UPGRADE_STATUS_CODES: - hint = ( - _SKILL_STREAM_MISSING_HINT - if path.rstrip("/").endswith("/v1/skills/stream") - else _SKILL_API_UPGRADE_HINT - ) - raise _SkillApiUpgradeRequired( - f"Skill HTTP API returned HTTP {exc.code}. {hint}" + raise RuntimeError( + f"Skill HTTP API returned HTTP {exc.code}. " + f"{_skill_api_upgrade_hint(path)}" ) from exc detail = exc.read().decode("utf-8", errors="replace") raise RuntimeError( f"Skill HTTP API request failed with HTTP {exc.code}: {detail}" ) from exc except error.URLError as exc: - raise _SkillApiCompatibilityMiss( + raise RuntimeError( f"Skill HTTP API endpoint is not reachable: {exc.reason}" ) from exc @@ -156,7 +135,7 @@ def _parse_skill_execute_response(raw: bytes) -> str: return json.dumps(payload, ensure_ascii=False) -def _parse_skill_stream_response(raw: bytes) -> str: +def _parse_skill_stream_response(raw: bytes | Iterable[bytes]) -> str: chunks: list[str] = [] event_name = "message" data_lines: list[str] = [] @@ -185,7 +164,9 @@ def flush_event() -> None: event_name = "message" data_lines = [] - for line in raw.decode("utf-8", errors="replace").splitlines(): + raw_lines = raw.splitlines() if isinstance(raw, bytes) else raw + for raw_line in raw_lines: + line = raw_line.decode("utf-8", errors="replace").rstrip("\r\n") if not line: flush_event() continue @@ -212,24 +193,22 @@ def _execute_skills_via_skill_api( endpoint = ensure_agentkit_session_endpoint( tool_id=tool_id, tool_user_session_id=_tool_user_session_id(tool_context), - tool_state=tool_context.state if tool_context else None, + tool_state=tool_context.state, ttl=max(timeout, 1800), ) except Exception as exc: - raise _SkillApiCompatibilityMiss( + raise RuntimeError( f"AgentKit session endpoint is not available: {exc}" ) from exc path = "/v1/skills/stream" if prefer_stream else "/v1/skills/execute" - raw = _post_skill_api_json( + return _post_skill_api_json( endpoint=endpoint, path=path, payload={"prompt": workflow_prompt}, tip_token_key=_tip_token_key(tool_context), timeout=timeout, + stream=prefer_stream, ) - if prefer_stream: - return _parse_skill_stream_response(raw) - return _parse_skill_execute_response(raw) def execute_skills( @@ -245,43 +224,36 @@ def execute_skills( Args: workflow_prompt (str): instruction of workflow env_vars (Optional[dict[str, str]]): Environment variables passed to the - skill agent process for this execution only. + skill agent process for this execution only. Requests with custom + environment variables use the legacy RunCode execution path. Returns: str: The output of the code execution. """ - timeout = _SKILL_API_TIMEOUT - tool_id = resolve_agentkit_tool_id("AGENTKIT_TOOL_ID_SKILLS") + if tool_context is None: + raise ValueError("tool_context is required for execute_skills") - if tool_context is not None and _skill_api_enabled(env_vars): - try: - return _execute_skills_via_skill_api( - workflow_prompt=workflow_prompt, - tool_id=tool_id, - tool_context=tool_context, - prefer_stream=prefer_stream, - timeout=timeout, - ) - except _SkillApiUpgradeRequired as exc: - raise RuntimeError(str(exc)) from exc - except _SkillApiCompatibilityMiss as exc: - logger.warning( - f"Skill HTTP API endpoint unreachable, falling back to legacy RunCode: {exc}" + tool_id = resolve_agentkit_tool_id("AGENTKIT_TOOL_ID_SKILLS") + if env_vars: + account_id = get_agentkit_account_id(tool_context.state) + extra_env_vars = dict(env_vars) + if account_id: + extra_env_vars.setdefault( + "TOS_SKILLS_DIR", + f"tos://agentkit-platform-{account_id}/skills/", ) - - account_id = get_agentkit_account_id(tool_context.state if tool_context else None) - extra_env_vars = {} - if account_id: - extra_env_vars["TOS_SKILLS_DIR"] = ( - f"tos://agentkit-platform-{account_id}/skills/" + return run_sandbox_agent( + workflow_prompt=workflow_prompt, + tool_id=tool_id, + tool_context=tool_context, + timeout=_SKILL_API_TIMEOUT, + extra_env_vars=extra_env_vars, ) - if env_vars: - extra_env_vars.update(env_vars) - return run_sandbox_agent( + return _execute_skills_via_skill_api( workflow_prompt=workflow_prompt, tool_id=tool_id, tool_context=tool_context, - timeout=timeout, - extra_env_vars=extra_env_vars, + prefer_stream=prefer_stream, + timeout=_SKILL_API_TIMEOUT, ) From 3c30111f597acf6086714d6620a349f1a658952e Mon Sep 17 00:00:00 2001 From: "lixuefei.nice" Date: Wed, 29 Jul 2026 11:31:52 +0800 Subject: [PATCH 7/8] feat(skil sandbox): fix 502 error --- tests/tools/builtin_tools/test_agentkit.py | 130 +++++++++++++++++- .../builtin_tools/test_run_sandbox_agent.py | 76 ++++++++++ veadk/tools/builtin_tools/_agentkit.py | 78 ++++++++--- veadk/tools/builtin_tools/execute_skills.py | 49 +++++++ 4 files changed, 309 insertions(+), 24 deletions(-) diff --git a/tests/tools/builtin_tools/test_agentkit.py b/tests/tools/builtin_tools/test_agentkit.py index 51276c11..71e1f7b6 100644 --- a/tests/tools/builtin_tools/test_agentkit.py +++ b/tests/tools/builtin_tools/test_agentkit.py @@ -230,6 +230,7 @@ def get_session(self, _request): return types.SimpleNamespace( endpoint="https://public.example", internal_endpoint="http://internal.example", + status="Ready", ) fake_tools_types = types.ModuleType("agentkit.sdk.tools.types") @@ -297,11 +298,17 @@ def get_session(self, _request): }, ) - def test_uses_create_session_endpoint_without_getting_session(self): + def test_waits_for_ready_even_when_create_session_returns_endpoint(self): + captured = {"get_calls": 0} + class FakeCreateSessionRequest: def __init__(self, **_kwargs): pass + class FakeGetSessionRequest: + def __init__(self, **_kwargs): + pass + class FakeClient: def __init__(self, **_kwargs): pass @@ -314,12 +321,76 @@ def create_session(self, _request): ) def get_session(self, _request): - raise AssertionError( - "get_session should not be called when create returns endpoint" + captured["get_calls"] += 1 + return types.SimpleNamespace( + session_id="session-1", + status="Ready", + endpoint="https://public.example", + internal_endpoint="http://internal.example", ) fake_tools_types = types.ModuleType("agentkit.sdk.tools.types") fake_tools_types.CreateSessionRequest = FakeCreateSessionRequest + fake_tools_types.GetSessionRequest = FakeGetSessionRequest + fake_tools_client = types.ModuleType("agentkit.sdk.tools.client") + fake_tools_client.AgentkitToolsClient = FakeClient + fake_tools_package = types.ModuleType("agentkit.sdk.tools") + fake_tools_package.types = fake_tools_types + + with patch.dict( + sys.modules, + { + "agentkit": types.ModuleType("agentkit"), + "agentkit.sdk": types.ModuleType("agentkit.sdk"), + "agentkit.sdk.tools": fake_tools_package, + "agentkit.sdk.tools.types": fake_tools_types, + "agentkit.sdk.tools.client": fake_tools_client, + }, + ): + with ( + patch.object( + self.agentkit_module, + "get_agentkit_endpoint_config", + return_value=("agentkit", "cn-beijing", "host", "https"), + ), + patch.object( + self.agentkit_module, + "get_agentkit_credentials", + return_value=("ak", "sk", {}), + ), + ): + endpoint = self.agentkit_module.ensure_agentkit_session_endpoint( + tool_id="tool-1", + tool_user_session_id="user-session-1", + ) + + self.assertEqual(endpoint, "https://public.example") + self.assertEqual(captured["get_calls"], 1) + + def test_polls_until_session_is_ready(self): + statuses = iter(["Starting", "Ready"]) + + class FakeRequest: + def __init__(self, **_kwargs): + pass + + class FakeClient: + def __init__(self, **_kwargs): + pass + + def create_session(self, _request): + return types.SimpleNamespace(session_id="session-1") + + def get_session(self, _request): + return types.SimpleNamespace( + status=next(statuses), + endpoint="https://public.example", + internal_endpoint=None, + ) + + fake_tools_types = types.ModuleType("agentkit.sdk.tools.types") + fake_tools_types.CreateSessionRequest = FakeRequest + fake_tools_types.GetSessionRequest = FakeRequest fake_tools_client = types.ModuleType("agentkit.sdk.tools.client") fake_tools_client.AgentkitToolsClient = FakeClient fake_tools_package = types.ModuleType("agentkit.sdk.tools") @@ -346,6 +417,7 @@ def get_session(self, _request): "get_agentkit_credentials", return_value=("ak", "sk", {}), ), + patch.object(self.agentkit_module.time, "sleep") as sleep, ): endpoint = self.agentkit_module.ensure_agentkit_session_endpoint( tool_id="tool-1", @@ -353,6 +425,58 @@ def get_session(self, _request): ) self.assertEqual(endpoint, "https://public.example") + sleep.assert_called_once_with(1.0) + + def test_raises_when_session_enters_failed_status(self): + class FakeRequest: + def __init__(self, **_kwargs): + pass + + class FakeClient: + def __init__(self, **_kwargs): + pass + + def create_session(self, _request): + return types.SimpleNamespace(session_id="session-1") + + def get_session(self, _request): + return types.SimpleNamespace(status="Failed") + + fake_tools_types = types.ModuleType("agentkit.sdk.tools.types") + fake_tools_types.CreateSessionRequest = FakeRequest + fake_tools_types.GetSessionRequest = FakeRequest + fake_tools_client = types.ModuleType("agentkit.sdk.tools.client") + fake_tools_client.AgentkitToolsClient = FakeClient + fake_tools_package = types.ModuleType("agentkit.sdk.tools") + fake_tools_package.types = fake_tools_types + + with patch.dict( + sys.modules, + { + "agentkit": types.ModuleType("agentkit"), + "agentkit.sdk": types.ModuleType("agentkit.sdk"), + "agentkit.sdk.tools": fake_tools_package, + "agentkit.sdk.tools.types": fake_tools_types, + "agentkit.sdk.tools.client": fake_tools_client, + }, + ): + with ( + patch.object( + self.agentkit_module, + "get_agentkit_endpoint_config", + return_value=("agentkit", "cn-beijing", "host", "https"), + ), + patch.object( + self.agentkit_module, + "get_agentkit_credentials", + return_value=("ak", "sk", {}), + ), + ): + with self.assertRaisesRegex(RuntimeError, "terminal status Failed"): + self.agentkit_module.ensure_agentkit_session_endpoint( + tool_id="tool-1", + tool_user_session_id="user-session-1", + ) if __name__ == "__main__": diff --git a/tests/tools/builtin_tools/test_run_sandbox_agent.py b/tests/tools/builtin_tools/test_run_sandbox_agent.py index f0dbc4d9..2b76fe24 100644 --- a/tests/tools/builtin_tools/test_run_sandbox_agent.py +++ b/tests/tools/builtin_tools/test_run_sandbox_agent.py @@ -80,6 +80,7 @@ def _load_execute_skills_module( *, ensure_agentkit_session_endpoint=lambda **_kwargs: "", run_sandbox_agent=lambda **_kwargs: "", + wait_for_skill_api_health=lambda **_kwargs: None, ): module_path = ( Path(__file__).resolve().parents[3] @@ -138,6 +139,8 @@ def _load_execute_skills_module( assert spec.loader is not None module = importlib.util.module_from_spec(spec) spec.loader.exec_module(module) + if wait_for_skill_api_health is not None: + module._wait_for_skill_api_health = wait_for_skill_api_health return module @@ -211,6 +214,7 @@ def _tool_context(self): def test_prefers_new_skill_execute_api_when_endpoint_is_available(self): captured_requests = [] + health_endpoints = [] class FakeResponse: def __enter__(self): @@ -228,12 +232,16 @@ def fake_urlopen(request, timeout=None): module = _load_execute_skills_module( ensure_agentkit_session_endpoint=lambda **_kwargs: "https://sandbox.test", + wait_for_skill_api_health=lambda **kwargs: health_endpoints.append( + kwargs["endpoint"] + ), ) with patch.object(module.request, "urlopen", fake_urlopen): result = module.execute_skills("do work", tool_context=self._tool_context()) self.assertEqual(result, "api result") + self.assertEqual(["https://sandbox.test"], health_endpoints) self.assertEqual(1, len(captured_requests)) request_obj, timeout = captured_requests[0] self.assertEqual("https://sandbox.test/v1/skills/execute", request_obj.full_url) @@ -242,6 +250,74 @@ def fake_urlopen(request, timeout=None): self.assertEqual("tip-from-state", request_obj.headers["X-tip-token-key"]) self.assertIn(b'"prompt": "do work"', request_obj.data) + def test_health_check_retries_502_until_upstream_is_ready(self): + attempts = [] + + class HealthyResponse: + def __enter__(self): + return self + + def __exit__(self, *_args): + return None + + class ErrorResponse: + def read(self): + return b"bad gateway" + + def close(self): + return None + + module = _load_execute_skills_module(wait_for_skill_api_health=None) + + def fake_urlopen(req, **_kwargs): + attempts.append((req.full_url, req.get_method())) + if len(attempts) == 1: + raise module.error.HTTPError( + url=req.full_url, + code=502, + msg="Bad Gateway", + hdrs={}, + fp=ErrorResponse(), + ) + return HealthyResponse() + + with ( + patch.object(module.request, "urlopen", fake_urlopen), + patch.object(module.time, "sleep") as sleep, + ): + module._wait_for_skill_api_health(endpoint="https://sandbox.test") + + self.assertEqual( + [ + ("https://sandbox.test/v1/skills/healthz", "GET"), + ("https://sandbox.test/v1/skills/healthz", "GET"), + ], + attempts, + ) + sleep.assert_called_once_with(1.0) + + def test_health_check_allows_images_without_health_endpoint(self): + class NotFoundResponse: + def read(self): + return b"not found" + + def close(self): + return None + + module = _load_execute_skills_module(wait_for_skill_api_health=None) + + def fake_urlopen(req, **_kwargs): + raise module.error.HTTPError( + url=req.full_url, + code=404, + msg="Not Found", + hdrs={}, + fp=NotFoundResponse(), + ) + + with patch.object(module.request, "urlopen", fake_urlopen): + module._wait_for_skill_api_health(endpoint="https://sandbox.test") + def test_env_vars_use_legacy_runcode_execution(self): captured_kwargs = {} diff --git a/veadk/tools/builtin_tools/_agentkit.py b/veadk/tools/builtin_tools/_agentkit.py index 3e4eef43..9a0ab0d8 100644 --- a/veadk/tools/builtin_tools/_agentkit.py +++ b/veadk/tools/builtin_tools/_agentkit.py @@ -14,6 +14,7 @@ import json import os +import time from typing import Any, Optional from veadk.auth.veauth.utils import get_credential_from_vefaas_iam @@ -24,6 +25,11 @@ logger = get_logger(__name__) +_SESSION_READY_TIMEOUT = 120.0 +_SESSION_POLL_INTERVAL = 1.0 +_SESSION_TERMINAL_STATUSES = frozenset({"failed", "terminating", "terminated"}) + + def resolve_agentkit_tool_id(*preferred_env_names: str) -> str: """Resolve the first configured AgentKit tool id with AGENTKIT_TOOL_ID fallback.""" for env_name in [*preferred_env_names, "AGENTKIT_TOOL_ID"]: @@ -215,11 +221,18 @@ def ensure_agentkit_session_endpoint( tool_state: Optional[dict[str, Any]] = None, ttl: int = 1800, prefer_internal_endpoint: bool = False, + ready_timeout: float = _SESSION_READY_TIMEOUT, + poll_interval: float = _SESSION_POLL_INTERVAL, ) -> str: - """Create or reuse an AgentKit tool session and return its HTTP endpoint.""" + """Create or reuse a Ready AgentKit tool session and return its endpoint.""" from agentkit.sdk.tools import types as tools_types from agentkit.sdk.tools.client import AgentkitToolsClient + if ready_timeout < 0: + raise ValueError("ready_timeout must be greater than or equal to 0") + if poll_interval <= 0: + raise ValueError("poll_interval must be greater than 0") + _, region, _, _ = get_agentkit_endpoint_config() ak, sk, header = get_agentkit_credentials(tool_state) session_token = header.get("X-Security-Token", "") @@ -236,26 +249,49 @@ def ensure_agentkit_session_endpoint( Ttl=ttl, ) ) - public_endpoint = getattr(session, "endpoint", None) - internal_endpoint = getattr(session, "internal_endpoint", None) - endpoint = ( - internal_endpoint or public_endpoint - if prefer_internal_endpoint - else public_endpoint or internal_endpoint - ) - if endpoint: - return endpoint - session_id = session.session_id if not session_id: - return "" - - current_session = client.get_session( - tools_types.GetSessionRequest( - ToolId=tool_id, - SessionId=session_id, + raise RuntimeError("AgentKit CreateSession response is missing SessionId") + + deadline = time.monotonic() + ready_timeout + last_status = "Unknown" + while True: + current_session = client.get_session( + tools_types.GetSessionRequest( + ToolId=tool_id, + SessionId=session_id, + ) ) - ) - if prefer_internal_endpoint: - return current_session.internal_endpoint or current_session.endpoint or "" - return current_session.endpoint or current_session.internal_endpoint or "" + status = (getattr(current_session, "status", None) or "").strip() + last_status = status or "Unknown" + logger.debug(f"AgentKit session {session_id} status: {last_status}") + normalized_status = status.lower() + if normalized_status == "ready": + public_endpoint = getattr(current_session, "endpoint", None) or getattr( + session, "endpoint", None + ) + internal_endpoint = getattr( + current_session, "internal_endpoint", None + ) or getattr(session, "internal_endpoint", None) + endpoint = ( + internal_endpoint or public_endpoint + if prefer_internal_endpoint + else public_endpoint or internal_endpoint + ) + if endpoint: + return endpoint + raise RuntimeError( + f"AgentKit session {session_id} is Ready but has no endpoint" + ) + if normalized_status in _SESSION_TERMINAL_STATUSES: + raise RuntimeError( + f"AgentKit session {session_id} entered terminal status {last_status}" + ) + + remaining = deadline - time.monotonic() + if remaining <= 0: + raise TimeoutError( + f"Timed out waiting for AgentKit session {session_id} to become " + f"Ready; last status: {last_status}" + ) + time.sleep(min(poll_interval, remaining)) diff --git a/veadk/tools/builtin_tools/execute_skills.py b/veadk/tools/builtin_tools/execute_skills.py index 266c38c9..79dddfaf 100644 --- a/veadk/tools/builtin_tools/execute_skills.py +++ b/veadk/tools/builtin_tools/execute_skills.py @@ -16,6 +16,7 @@ import json import os +import time from collections.abc import Iterable from typing import Optional from urllib import error, request @@ -32,7 +33,11 @@ _SKILL_API_UPGRADE_STATUS_CODES = frozenset({404, 405}) +_SKILL_API_TRANSIENT_STATUS_CODES = frozenset({502, 503, 504}) _SKILL_API_TIMEOUT = 900 +_SKILL_API_HEALTH_TIMEOUT = 30.0 +_SKILL_API_HEALTH_POLL_INTERVAL = 1.0 +_SKILL_API_HEALTH_REQUEST_TIMEOUT = 5.0 def _skill_api_upgrade_hint(path: str) -> str: @@ -120,6 +125,49 @@ def _post_skill_api_json( ) from exc +def _wait_for_skill_api_health( + *, + endpoint: str, + timeout: float = _SKILL_API_HEALTH_TIMEOUT, + poll_interval: float = _SKILL_API_HEALTH_POLL_INTERVAL, +) -> None: + """Wait until the Skill API upstream is reachable through the session endpoint.""" + deadline = time.monotonic() + timeout + last_error = "unknown error" + while True: + req = request.Request( + _skill_api_url(endpoint, "/v1/skills/healthz"), + headers={"Accept": "application/json"}, + method="GET", + ) + try: + remaining = max(0.001, deadline - time.monotonic()) + with request.urlopen( + req, + timeout=min(_SKILL_API_HEALTH_REQUEST_TIMEOUT, remaining), + ): + return + except error.HTTPError as exc: + if exc.code in _SKILL_API_UPGRADE_STATUS_CODES: + # Some compatible images predate the dedicated health endpoint. + return + if exc.code not in _SKILL_API_TRANSIENT_STATUS_CODES: + detail = exc.read().decode("utf-8", errors="replace") + raise RuntimeError( + f"Skill HTTP API health check failed with HTTP {exc.code}: {detail}" + ) from exc + last_error = f"HTTP {exc.code}" + except error.URLError as exc: + last_error = str(exc.reason) + + remaining = deadline - time.monotonic() + if remaining <= 0: + raise RuntimeError( + f"Timed out waiting for Skill HTTP API health check: {last_error}" + ) + time.sleep(min(poll_interval, remaining)) + + def _parse_skill_execute_response(raw: bytes) -> str: try: payload = json.loads(raw.decode("utf-8")) @@ -200,6 +248,7 @@ def _execute_skills_via_skill_api( raise RuntimeError( f"AgentKit session endpoint is not available: {exc}" ) from exc + _wait_for_skill_api_health(endpoint=endpoint) path = "/v1/skills/stream" if prefer_stream else "/v1/skills/execute" return _post_skill_api_json( endpoint=endpoint, From 6687149f677c3b73afb54a7700c441af73b72421 Mon Sep 17 00:00:00 2001 From: "lixuefei.nice" Date: Wed, 29 Jul 2026 11:40:55 +0800 Subject: [PATCH 8/8] feat(skil sandbox): only change agentkiy.py about skill --- tests/tools/builtin_tools/test_agentkit.py | 13 ++++--- .../builtin_tools/test_run_sandbox_agent.py | 6 +++- veadk/tools/builtin_tools/_agentkit.py | 36 ++++++++++++++++--- veadk/tools/builtin_tools/execute_skills.py | 1 + 4 files changed, 43 insertions(+), 13 deletions(-) diff --git a/tests/tools/builtin_tools/test_agentkit.py b/tests/tools/builtin_tools/test_agentkit.py index 71e1f7b6..c5a5aa69 100644 --- a/tests/tools/builtin_tools/test_agentkit.py +++ b/tests/tools/builtin_tools/test_agentkit.py @@ -298,7 +298,7 @@ def get_session(self, _request): }, ) - def test_waits_for_ready_even_when_create_session_returns_endpoint(self): + def test_uses_create_session_endpoint_without_waiting_by_default(self): captured = {"get_calls": 0} class FakeCreateSessionRequest: @@ -322,11 +322,8 @@ def create_session(self, _request): def get_session(self, _request): captured["get_calls"] += 1 - return types.SimpleNamespace( - session_id="session-1", - status="Ready", - endpoint="https://public.example", - internal_endpoint="http://internal.example", + raise AssertionError( + "get_session should not be called when waiting is disabled" ) fake_tools_types = types.ModuleType("agentkit.sdk.tools.types") @@ -365,7 +362,7 @@ def get_session(self, _request): ) self.assertEqual(endpoint, "https://public.example") - self.assertEqual(captured["get_calls"], 1) + self.assertEqual(captured["get_calls"], 0) def test_polls_until_session_is_ready(self): statuses = iter(["Starting", "Ready"]) @@ -422,6 +419,7 @@ def get_session(self, _request): endpoint = self.agentkit_module.ensure_agentkit_session_endpoint( tool_id="tool-1", tool_user_session_id="user-session-1", + wait_until_ready=True, ) self.assertEqual(endpoint, "https://public.example") @@ -476,6 +474,7 @@ def get_session(self, _request): self.agentkit_module.ensure_agentkit_session_endpoint( tool_id="tool-1", tool_user_session_id="user-session-1", + wait_until_ready=True, ) diff --git a/tests/tools/builtin_tools/test_run_sandbox_agent.py b/tests/tools/builtin_tools/test_run_sandbox_agent.py index 2b76fe24..59391a02 100644 --- a/tests/tools/builtin_tools/test_run_sandbox_agent.py +++ b/tests/tools/builtin_tools/test_run_sandbox_agent.py @@ -215,6 +215,7 @@ def _tool_context(self): def test_prefers_new_skill_execute_api_when_endpoint_is_available(self): captured_requests = [] health_endpoints = [] + session_kwargs = [] class FakeResponse: def __enter__(self): @@ -231,7 +232,9 @@ def fake_urlopen(request, timeout=None): return FakeResponse() module = _load_execute_skills_module( - ensure_agentkit_session_endpoint=lambda **_kwargs: "https://sandbox.test", + ensure_agentkit_session_endpoint=lambda **kwargs: ( + session_kwargs.append(kwargs) or "https://sandbox.test" + ), wait_for_skill_api_health=lambda **kwargs: health_endpoints.append( kwargs["endpoint"] ), @@ -241,6 +244,7 @@ def fake_urlopen(request, timeout=None): result = module.execute_skills("do work", tool_context=self._tool_context()) self.assertEqual(result, "api result") + self.assertTrue(session_kwargs[0]["wait_until_ready"]) self.assertEqual(["https://sandbox.test"], health_endpoints) self.assertEqual(1, len(captured_requests)) request_obj, timeout = captured_requests[0] diff --git a/veadk/tools/builtin_tools/_agentkit.py b/veadk/tools/builtin_tools/_agentkit.py index 9a0ab0d8..10abdeec 100644 --- a/veadk/tools/builtin_tools/_agentkit.py +++ b/veadk/tools/builtin_tools/_agentkit.py @@ -221,17 +221,19 @@ def ensure_agentkit_session_endpoint( tool_state: Optional[dict[str, Any]] = None, ttl: int = 1800, prefer_internal_endpoint: bool = False, + wait_until_ready: bool = False, ready_timeout: float = _SESSION_READY_TIMEOUT, poll_interval: float = _SESSION_POLL_INTERVAL, ) -> str: - """Create or reuse a Ready AgentKit tool session and return its endpoint.""" + """Create or reuse an AgentKit tool session and return its endpoint.""" from agentkit.sdk.tools import types as tools_types from agentkit.sdk.tools.client import AgentkitToolsClient - if ready_timeout < 0: - raise ValueError("ready_timeout must be greater than or equal to 0") - if poll_interval <= 0: - raise ValueError("poll_interval must be greater than 0") + if wait_until_ready: + if ready_timeout < 0: + raise ValueError("ready_timeout must be greater than or equal to 0") + if poll_interval <= 0: + raise ValueError("poll_interval must be greater than 0") _, region, _, _ = get_agentkit_endpoint_config() ak, sk, header = get_agentkit_credentials(tool_state) @@ -249,6 +251,30 @@ def ensure_agentkit_session_endpoint( Ttl=ttl, ) ) + if not wait_until_ready: + public_endpoint = getattr(session, "endpoint", None) + internal_endpoint = getattr(session, "internal_endpoint", None) + endpoint = ( + internal_endpoint or public_endpoint + if prefer_internal_endpoint + else public_endpoint or internal_endpoint + ) + if endpoint: + return endpoint + + session_id = session.session_id + if not session_id: + return "" + current_session = client.get_session( + tools_types.GetSessionRequest( + ToolId=tool_id, + SessionId=session_id, + ) + ) + if prefer_internal_endpoint: + return current_session.internal_endpoint or current_session.endpoint or "" + return current_session.endpoint or current_session.internal_endpoint or "" + session_id = session.session_id if not session_id: raise RuntimeError("AgentKit CreateSession response is missing SessionId") diff --git a/veadk/tools/builtin_tools/execute_skills.py b/veadk/tools/builtin_tools/execute_skills.py index 79dddfaf..75059ab5 100644 --- a/veadk/tools/builtin_tools/execute_skills.py +++ b/veadk/tools/builtin_tools/execute_skills.py @@ -243,6 +243,7 @@ def _execute_skills_via_skill_api( tool_user_session_id=_tool_user_session_id(tool_context), tool_state=tool_context.state, ttl=max(timeout, 1800), + wait_until_ready=True, ) except Exception as exc: raise RuntimeError(