diff --git a/gateway/kanban_watchers.py b/gateway/kanban_watchers.py index 2f5137a8a84f..216d342b9d92 100644 --- a/gateway/kanban_watchers.py +++ b/gateway/kanban_watchers.py @@ -493,8 +493,8 @@ def _collect(): ) if sub.get("thread_id") and not metadata.get("thread_id"): metadata["thread_id"] = sub["thread_id"] - # Adapters with no push channel (the API server — - # ``supports_async_delivery = False``) can NEVER + # Adapters with no push channel (the API server sets + # ``supports_push_delivery = False``) can NEVER # satisfy a text-send: ``send()`` always reports # SendResult(success=False) by design (see # ApiServerAdapter.send()). Treating that as a diff --git a/gateway/platforms/api_server.py b/gateway/platforms/api_server.py index 08fb7d215293..e4d0ee786892 100644 --- a/gateway/platforms/api_server.py +++ b/gateway/platforms/api_server.py @@ -1246,15 +1246,11 @@ class APIServerAdapter(BasePlatformAdapter): and routes them through hermes-agent's AIAgent. """ - # Stateless request/response: every route (the OpenAI-spec - # /v1/chat/completions and /v1/responses, and the proprietary /v1/runs SSE - # stream) tears down its channel when the turn ends. There is no persistent - # outbound channel to push a background completion to a client that already - # received its response, and ``send()`` is a no-op stub. So async-delivery - # tools (terminal notify_on_complete / watch_patterns, delegate_task - # background=True) must NOT promise delivery on this path — see - # ``async_delivery_supported()``. - supports_async_delivery: bool = False + # The HTTP response channel is stateless, but the gateway can wake the raw + # session by self-POSTing through /v1/chat/completions. Async jobs are + # therefore supported even though direct push via handle_message is not. + supports_async_delivery: bool = True + supports_push_delivery: bool = False # Same statelessness applies to the startup auto-resume prompt: no client # is waiting to answer "session restored — what next?", so a resumed turn @@ -6270,12 +6266,9 @@ def _bind_api_server_session( """Bind session contextvars for an API-server agent run. This is the SINGLE structural chokepoint every API-server agent-entry - path must use to seed session context — it hardwires - ``platform="api_server"`` and ``async_delivery=False`` so a new route - physically cannot reintroduce the silent-no-op bug (#10760) by - forgetting to mark the channel as non-delivering. There is no - ``async_delivery`` parameter to get wrong; the stateless HTTP path can - never wake the agent after the turn ends, on ANY route. + path must use to seed session context. API sessions support asynchronous + completion via the gateway's authenticated self-post wake path, while + remaining non-push adapters. Returns reset tokens; pass them to ``clear_session_vars`` in a ``finally`` block (the binding is request-scoped and must not outlive @@ -6289,7 +6282,7 @@ def _bind_api_server_session( chat_id=chat_id, session_key=session_key, session_id=session_id, - async_delivery=False, + async_delivery=True, ) async def _run_agent( diff --git a/gateway/run.py b/gateway/run.py index cf8f86e36800..0808bb192797 100644 --- a/gateway/run.py +++ b/gateway/run.py @@ -17660,7 +17660,8 @@ def _set_session_env(self, context: SessionContext) -> list: # (terminal notify_on_complete / watch_patterns, delegate_task # background=True) know whether this channel can wake a later turn. # Default True keeps CLI / unknown paths working; stateless adapters - # (api_server) declare supports_async_delivery=False. Use getattr so + # that cannot arrange any later wake declare supports_async_delivery=False. + # The API server supports self-post wakes while remaining non-push. Use getattr so # bare runners built via object.__new__ (tests) without self.adapters # don't blow up — they simply default to supported. _adapters = getattr(self, "adapters", None) or {} diff --git a/gateway/session_context.py b/gateway/session_context.py index cf227fac1211..71d822658701 100644 --- a/gateway/session_context.py +++ b/gateway/session_context.py @@ -100,8 +100,10 @@ def session_context_engaged() -> bool: # True — long-lived CLI sessions (in-process completion_queue drain) and the # real gateway platforms (Telegram/Discord/Slack/...), which hold a # persistent outbound channel and run the watcher/drain loops. -# False — finite runtimes that can end before a detached completion returns: -# stateless API-server requests and dispatcher-spawned Kanban workers. +# False — finite runtimes that cannot arrange any later wake, such as +# dispatcher-spawned one-shot Kanban workers. The API server is +# stateless but supports wake delivery through an authenticated +# self-post to the raw session, so it binds True. # # Tools that promise async delivery (terminal notify_on_complete / # watch_patterns, delegate_task background=True) read this via @@ -109,7 +111,7 @@ def session_context_engaged() -> bool: # can't keep — turning a silent no-op into an explicit contract. # # Default _UNSET => treated as supported, so CLI (which never sets a platform) -# and any contextvar-unaware path keep working. Stateless adapters opt OUT by +# and any contextvar-unaware path keep working. Non-delivering adapters opt OUT by # setting ``supports_async_delivery = False`` on the adapter class; the gateway # propagates that into this contextvar at session-bind time. _SESSION_ASYNC_DELIVERY: ContextVar = ContextVar("HERMES_SESSION_ASYNC_DELIVERY", default=_UNSET) diff --git a/gateway/wake.py b/gateway/wake.py index cee8d45e50c3..2518b79601ce 100644 --- a/gateway/wake.py +++ b/gateway/wake.py @@ -1,14 +1,14 @@ """Wake an existing agent session from a background completion event. Two delivery strategies, selected by the target adapter's -``supports_async_delivery`` capability flag: +``supports_push_delivery`` capability flag: * Push-capable adapters (telegram, discord, plugin platforms, ...): inject a synthetic ``MessageEvent(internal=True)`` through ``adapter.handle_message`` — the pre-existing wake path, preserved exactly. * Stateless request/response adapters (the API server, - ``supports_async_delivery = False``): ``handle_message`` would run the wake + ``supports_push_delivery = False``): ``handle_message`` would run the wake turn under a ``build_session_key()``-derived key (``agent:main:api_server:group:``) that NEVER matches the raw ``X-Hermes-Session-Id`` key real gateway/HQ turns run under @@ -45,11 +45,12 @@ def adapter_supports_push(adapter: Any) -> bool: """Whether this adapter can push a message to the user after a turn ends. - Mirrors ``gateway.session_context.async_delivery_supported`` but reads the - capability off the adapter class (``supports_async_delivery``) instead of - the request-scoped contextvar — background watchers run outside any bound - session context. Adapters that don't declare the flag are push-capable. + ``supports_async_delivery`` answers whether a later wake is possible; + ``supports_push_delivery`` answers how it is delivered. Older adapters + that only declare the former retain their previous behavior. """ + if hasattr(adapter, "supports_push_delivery"): + return bool(getattr(adapter, "supports_push_delivery")) return bool(getattr(adapter, "supports_async_delivery", True)) @@ -88,7 +89,7 @@ async def deliver_wake( if not session_id: raise ValueError( - "deliver_wake: non-push adapter (supports_async_delivery=False) " + "deliver_wake: non-push adapter " "requires the raw session id to self-post the wake turn" ) await _self_post_chat_completion(adapter, text=text, session_id=session_id) diff --git a/tests/gateway/test_async_delivery_capability.py b/tests/gateway/test_async_delivery_capability.py index 65b9691e3ab4..188bcdd12c2f 100644 --- a/tests/gateway/test_async_delivery_capability.py +++ b/tests/gateway/test_async_delivery_capability.py @@ -1,15 +1,13 @@ """Tests for the async-delivery capability gate (issue #10760). -Stateless request/response adapters (the API server / WebUI path) cannot route -a background completion back to the agent after a turn ends — there is no -persistent channel and ``APIServerAdapter.send()`` is a no-op stub. So tools -that promise async delivery (``terminal`` notify_on_complete / watch_patterns, -``delegate_task`` background=True) must refuse the promise on that path instead -of silently registering a watcher that never fires. +The API server has no persistent push channel, but it can route a background +completion back to the raw session through its authenticated self-post wake +path. Async capability and push capability are therefore separate contracts. This is wired through: - ``BasePlatformAdapter.supports_async_delivery`` (default True) - - ``APIServerAdapter.supports_async_delivery = False`` + - ``APIServerAdapter.supports_async_delivery = True`` + - ``APIServerAdapter.supports_push_delivery = False`` - ``gateway.session_context._SESSION_ASYNC_DELIVERY`` contextvar + ``async_delivery_supported()`` helper, bound per-session. @@ -208,15 +206,16 @@ def test_base_default_true(self): assert BasePlatformAdapter.supports_async_delivery is True - def test_api_server_false(self): + def test_api_server_supports_async_self_post_but_not_push(self): from gateway.platforms.api_server import APIServerAdapter - assert APIServerAdapter.supports_async_delivery is False + assert APIServerAdapter.supports_async_delivery is True + assert APIServerAdapter.supports_push_delivery is False - def test_api_server_bind_chokepoint_hardwires_no_delivery(self): + def test_api_server_bind_chokepoint_enables_self_post_delivery(self): """Every API-server agent-entry path binds through - _bind_api_server_session, which hardwires async_delivery=False — a new - route physically cannot reintroduce the silent no-op (#10760).""" + _bind_api_server_session, which hardwires async_delivery=True so a new + route cannot accidentally disable the self-post wake contract.""" from gateway.platforms.api_server import APIServerAdapter from gateway.session_context import clear_session_vars, get_session_env @@ -224,21 +223,19 @@ def test_api_server_bind_chokepoint_hardwires_no_delivery(self): chat_id="c1", session_key="sk1", session_id="sid1" ) try: - assert async_delivery_supported() is False + assert async_delivery_supported() is True assert get_session_env("HERMES_SESSION_PLATFORM") == "api_server" finally: clear_session_vars(tokens) def test_api_server_binding_does_not_outlive_turn(self): - """The no-delivery decision is request-scoped, NOT stuck to the session. - After clear, a session resumed on a delivering interface re-binds fresh - and is NOT blocked.""" + """The delivery decision remains request-scoped and clears cleanly.""" from gateway.platforms.api_server import APIServerAdapter from gateway.session_context import clear_session_vars - # Turn 1: same session over the API server -> blocked. + # Turn 1: same session over the API server -> self-post capable. tokens = APIServerAdapter._bind_api_server_session(session_key="shared-key") - assert async_delivery_supported() is False + assert async_delivery_supported() is True clear_session_vars(tokens) # Turn 2: SAME session_key resumed on a delivering interface (CLI/gateway) @@ -311,6 +308,26 @@ def test_gateway_registers_watcher(self): assert len(process_registry.pending_watchers) == 1 assert process_registry.pending_watchers[0]["platform"] == "telegram" + def test_api_server_registers_self_post_watcher(self): + from tools.process_registry import process_registry + + tokens = set_session_vars( + platform="api_server", + chat_id="raw-session", + session_key="raw-session", + session_id="raw-session", + async_delivery=True, + ) + try: + d = self._run_bg("sleep 30 && echo DONE") + finally: + clear_session_vars(tokens) + + assert d.get("notify_on_complete") is True + assert not d.get("notify_unsupported") + assert len(process_registry.pending_watchers) == 1 + assert process_registry.pending_watchers[0]["platform"] == "api_server" + def test_cli_stays_supported(self): """CLI delivers via the in-process completion_queue: notify stays on, no false 'unsupported' note, and no pending_watcher (empty platform).""" diff --git a/tests/gateway/test_wake_delivery.py b/tests/gateway/test_wake_delivery.py index 3c09ef87c18f..9f248a5fd395 100644 --- a/tests/gateway/test_wake_delivery.py +++ b/tests/gateway/test_wake_delivery.py @@ -2,7 +2,7 @@ Two strategies: * push-capable adapters keep the synthetic MessageEvent / handle_message path; -* the stateless API server (supports_async_delivery=False) self-POSTs +* the stateless API server (supports_push_delivery=False) self-POSTs /v1/chat/completions with the RAW session id in X-Hermes-Session-Id, so the wake turn resumes the REAL session instead of a parallel invisible one keyed by build_session_key(). @@ -53,6 +53,14 @@ def test_adapter_supports_push_default_true(): assert adapter_supports_push(ApiServerLikeAdapter()) is False +def test_explicit_push_capability_is_independent_from_async_capability(): + adapter = ApiServerLikeAdapter() + adapter.supports_async_delivery = True + adapter.supports_push_delivery = False + + assert adapter_supports_push(adapter) is False + + def test_deliver_wake_push_adapter_uses_handle_message(): adapter = PushAdapter() asyncio.run(deliver_wake(adapter, text="wake up", source=_source())) diff --git a/tests/tools/test_browser_cdp_override.py b/tests/tools/test_browser_cdp_override.py index fdf9e311ae0a..485a91a94a67 100644 --- a/tests/tools/test_browser_cdp_override.py +++ b/tests/tools/test_browser_cdp_override.py @@ -125,7 +125,7 @@ def test_normalizes_provider_returned_http_cdp_url_when_creating_session(self, m monkeypatch.setattr(browser_tool, "_session_last_activity", {}) monkeypatch.setattr(browser_tool, "_start_browser_cleanup_thread", lambda: None) monkeypatch.setattr(browser_tool, "_update_session_activity", lambda task_id: None) - monkeypatch.setattr(browser_tool, "_get_cdp_override", lambda: "") + monkeypatch.setattr(browser_tool, "_get_cdp_override", lambda *_: "") monkeypatch.setattr(browser_tool, "_get_cloud_provider", lambda: provider) with patch("tools.browser_tool.requests.get", return_value=response) as mock_get: @@ -140,6 +140,52 @@ def test_normalizes_provider_returned_http_cdp_url_when_creating_session(self, m class TestGetCdpOverride: + def test_template_routes_each_conversation_to_its_own_endpoint(self, monkeypatch): + import tools.browser_tool as browser_tool + + monkeypatch.setenv( + "BROWSER_CDP_URL_TEMPLATE", + "ws://toolbox.test/sessions/{session_id}", + ) + monkeypatch.delenv("BROWSER_CDP_URL", raising=False) + monkeypatch.setattr( + "gateway.session_context.get_session_env", + lambda _name: "", + ) + + assert browser_tool._get_cdp_override_raw("conversation-a") == ( + "ws://toolbox.test/sessions/conversation-a" + ) + assert browser_tool._get_cdp_override_raw("conversation-b") == ( + "ws://toolbox.test/sessions/conversation-b" + ) + + def test_template_url_encodes_session_id(self, monkeypatch): + import tools.browser_tool as browser_tool + + monkeypatch.setenv( + "BROWSER_CDP_URL_TEMPLATE", + "ws://toolbox.test/sessions/{session_id}", + ) + monkeypatch.setattr( + "gateway.session_context.get_session_env", + lambda _name: "", + ) + + assert browser_tool._get_cdp_override_raw("child/task 1") == ( + "ws://toolbox.test/sessions/child%2Ftask%201" + ) + + def test_template_without_session_placeholder_is_rejected(self, monkeypatch): + import tools.browser_tool as browser_tool + + monkeypatch.setenv( + "BROWSER_CDP_URL_TEMPLATE", + "ws://toolbox.test/shared-browser", + ) + + assert browser_tool._get_cdp_override_raw("conversation-a") == "" + def test_prefers_env_var_over_config(self, monkeypatch): import tools.browser_tool as browser_tool @@ -202,6 +248,15 @@ def test_camofox_yields_to_config_cdp_override(self, monkeypatch): with patch("hermes_cli.config.read_raw_config", return_value={}): assert bc.is_camofox_mode() is False + # A conversation-scoped gateway also suppresses camofox. + monkeypatch.delenv("BROWSER_CDP_URL", raising=False) + monkeypatch.setenv( + "BROWSER_CDP_URL_TEMPLATE", + "ws://toolbox.test/sessions/{session_id}", + ) + with patch("hermes_cli.config.read_raw_config", return_value={}): + assert bc.is_camofox_mode() is False + class TestCreateCdpSession: """_create_cdp_session() must sanitize the CDP URL before logging. diff --git a/tests/tools/test_browser_cdp_tool.py b/tests/tools/test_browser_cdp_tool.py index 512c68319ab6..387cd788f133 100644 --- a/tests/tools/test_browser_cdp_tool.py +++ b/tests/tools/test_browser_cdp_tool.py @@ -135,7 +135,7 @@ def cdp_server(monkeypatch): server = _CDPServer() ws_url = server.start() monkeypatch.setattr( - browser_cdp_tool, "_resolve_cdp_endpoint", lambda: ws_url + browser_cdp_tool, "_resolve_cdp_endpoint", lambda *_: ws_url ) try: yield server @@ -163,7 +163,7 @@ def test_non_string_method_returns_error(): def test_non_dict_params_returns_error(monkeypatch): monkeypatch.setattr( - browser_cdp_tool, "_resolve_cdp_endpoint", lambda: "ws://localhost:9999" + browser_cdp_tool, "_resolve_cdp_endpoint", lambda *_: "ws://localhost:9999" ) result = json.loads( browser_cdp_tool.browser_cdp(method="Target.getTargets", params="not-a-dict") # type: ignore[arg-type] @@ -178,7 +178,7 @@ def test_non_dict_params_returns_error(monkeypatch): def test_no_endpoint_returns_helpful_error(monkeypatch): - monkeypatch.setattr(browser_cdp_tool, "_resolve_cdp_endpoint", lambda: "") + monkeypatch.setattr(browser_cdp_tool, "_resolve_cdp_endpoint", lambda *_: "") result = json.loads(browser_cdp_tool.browser_cdp(method="Target.getTargets")) assert "error" in result assert "/browser connect" in result["error"] @@ -187,7 +187,7 @@ def test_no_endpoint_returns_helpful_error(monkeypatch): def test_non_ws_endpoint_returns_error(monkeypatch): monkeypatch.setattr( - browser_cdp_tool, "_resolve_cdp_endpoint", lambda: "http://localhost:9222" + browser_cdp_tool, "_resolve_cdp_endpoint", lambda *_: "http://localhost:9222" ) result = json.loads(browser_cdp_tool.browser_cdp(method="Target.getTargets")) assert "error" in result @@ -230,6 +230,14 @@ def test_browser_level_success(cdp_server): def test_browser_level_redacts_secret_result(cdp_server): fake_key = "sk-" + "CDPSECRETRESULT1234567890" + cdp_server.on( + "Target.getTargets", + lambda params, sid: {"targetInfos": [{"targetId": "tab-A", "type": "page"}]}, + ) + cdp_server.on( + "Target.attachToTarget", + lambda params, sid: {"sessionId": "sess-tab-A"}, + ) cdp_server.on( "Runtime.evaluate", lambda params, sid: {"result": {"type": "string", "value": fake_key}}, @@ -285,6 +293,50 @@ def test_target_attach_then_call(cdp_server): assert calls[1]["sessionId"] == "sess-tab-A" +def test_page_method_without_target_auto_selects_page(cdp_server): + cdp_server.on( + "Target.getTargets", + lambda params, sid: { + "targetInfos": [ + {"targetId": "worker-A", "type": "service_worker"}, + {"targetId": "tab-A", "type": "page"}, + ] + }, + ) + cdp_server.on( + "Target.attachToTarget", + lambda params, sid: {"sessionId": f"sess-{params['targetId']}"}, + ) + cdp_server.on("Page.reload", lambda params, sid: {"reloaded": sid}) + + result = json.loads(browser_cdp_tool.browser_cdp(method="Page.reload")) + + assert result["success"] is True + assert result["result"] == {"reloaded": "sess-tab-A"} + assert [call["method"] for call in cdp_server.received()] == [ + "Target.getTargets", + "Target.attachToTarget", + "Page.reload", + ] + + +def test_page_method_requires_target_when_multiple_pages_are_open(cdp_server): + cdp_server.on( + "Target.getTargets", + lambda params, sid: { + "targetInfos": [ + {"targetId": "tab-A", "type": "page"}, + {"targetId": "tab-B", "type": "page"}, + ] + }, + ) + + result = json.loads(browser_cdp_tool.browser_cdp(method="Page.reload")) + + assert "Multiple page targets" in result["error"] + assert "target_id" in result["error"] + + # --------------------------------------------------------------------------- # CDP error responses # --------------------------------------------------------------------------- @@ -324,6 +376,14 @@ def slow(params, sid): time.sleep(10) return {} + cdp_server.on( + "Target.getTargets", + lambda params, sid: {"targetInfos": [{"targetId": "tab-A", "type": "page"}]}, + ) + cdp_server.on( + "Target.attachToTarget", + lambda params, sid: {"sessionId": "sess-tab-A"}, + ) cdp_server.on("Page.slowMethod", slow) result = json.loads( browser_cdp_tool.browser_cdp( @@ -402,7 +462,7 @@ def test_runtime_evaluate_blocked_when_current_page_is_private(monkeypatch): monkeypatch.setattr( browser_cdp_tool, "_resolve_cdp_endpoint", - lambda: "ws://127.0.0.1:9222/devtools/browser/mock", + lambda *_: "ws://127.0.0.1:9222/devtools/browser/mock", ) import tools.browser_tool as bt @@ -500,7 +560,7 @@ def test_page_navigate_to_private_url_blocked_before_cdp(monkeypatch): monkeypatch.setattr( browser_cdp_tool, "_resolve_cdp_endpoint", - lambda: "ws://127.0.0.1:9222/devtools/browser/mock", + lambda *_: "ws://127.0.0.1:9222/devtools/browser/mock", ) import tools.browser_tool as bt @@ -527,6 +587,14 @@ async def fake_call(*args, **kwargs): def test_private_guard_inactive_does_not_probe(monkeypatch, cdp_server): + cdp_server.on( + "Target.getTargets", + lambda params, sid: {"targetInfos": [{"targetId": "tab-A", "type": "page"}]}, + ) + cdp_server.on( + "Target.attachToTarget", + lambda params, sid: {"sessionId": "sess-tab-A"}, + ) cdp_server.on("Runtime.evaluate", lambda params, sid: {"result": {"value": "ok"}}) import tools.browser_tool as bt diff --git a/tests/tools/test_browser_cleanup.py b/tests/tools/test_browser_cleanup.py index 047703b2df57..4d148a2c737e 100644 --- a/tests/tools/test_browser_cleanup.py +++ b/tests/tools/test_browser_cleanup.py @@ -67,6 +67,23 @@ def test_cleanup_browser_clears_tracking_state(self): mock_stop.assert_called_once_with("task-1") mock_run.assert_called_once_with("task-1", "close", [], timeout=10) + def test_cleanup_remote_cdp_never_closes_provider_browser(self): + browser_tool = self.browser_tool + browser_tool._active_sessions["task-remote"] = { + "session_name": "cdp-remote", + "bb_session_id": None, + "cdp_url": "ws://toolbox/sessions/conversation-a", + } + + with ( + patch("tools.browser_tool._maybe_stop_recording"), + patch("tools.browser_tool._run_browser_command") as mock_run, + patch("tools.browser_tool.os.path.exists", return_value=False), + ): + browser_tool.cleanup_browser("task-remote") + + mock_run.assert_not_called() + def test_cleanup_camofox_managed_persistence_skips_close(self): """When camofox mode + managed persistence, soft_cleanup fires instead of close.""" browser_tool = self.browser_tool diff --git a/tests/tools/test_browser_resilience.py b/tests/tools/test_browser_resilience.py index 88f86cb5c9c1..93eef9d6cc16 100644 --- a/tests/tools/test_browser_resilience.py +++ b/tests/tools/test_browser_resilience.py @@ -215,7 +215,34 @@ def test_snapshot_renderer_failure_recovers_without_retrying_blank_page( ["-c"], timeout=90, ) - mock_recover.assert_called_once_with() + mock_recover.assert_called_once_with("brand-a") + + +def test_vision_renderer_failure_uses_same_scoped_recovery_contract( + monkeypatch, + tmp_path, +) -> None: + monkeypatch.setattr(browser_tool, "_is_camofox_mode", lambda: False) + monkeypatch.setattr(browser_tool, "_last_session_key", lambda _task: "conversation-a") + monkeypatch.setattr(browser_tool, "_is_local_backend", lambda: True) + monkeypatch.setattr(browser_tool, "_get_browser_engine", lambda: "auto") + monkeypatch.setattr(browser_tool, "_session_is_cdp_backed", lambda _task: True) + monkeypatch.setattr("hermes_constants.get_hermes_dir", lambda *_args: tmp_path) + + with ( + patch("tools.browser_tool._run_browser_command", return_value=_cdp_failure()), + patch( + "tools.browser_tool._recover_omnio_browser", + return_value=True, + ) as mock_recover, + ): + result = json.loads( + browser_tool.browser_vision("what is visible?", task_id="conversation-a") + ) + + assert result["code"] == "browser_session_reset" + assert result["next_action"] == "browser_navigate" + mock_recover.assert_called_once_with("conversation-a") def test_snapshot_bare_timeout_does_not_recover(monkeypatch) -> None: @@ -358,7 +385,7 @@ def test_navigate_retryable_failure_recovers_and_retries_once( ), call("nav-recover", "snapshot", ["-c"], timeout=90), ] - mock_recover.assert_called_once_with() + mock_recover.assert_called_once_with("nav-recover") browser_tool._last_active_session_key.pop("nav-recover", None) @@ -393,7 +420,7 @@ def test_snapshot_preserves_failure_when_recover_is_unsupported( ["-c"], timeout=90, ) - mock_recover.assert_called_once_with() + mock_recover.assert_called_once_with("brand-a") def test_toolbox_v5_does_not_receive_recover_operation(monkeypatch) -> None: @@ -413,7 +440,7 @@ def test_toolbox_v5_does_not_receive_recover_operation(monkeypatch) -> None: ) as mock_get, patch("tools.browser_tool.requests.post") as mock_post, ): - recovered = browser_tool._recover_omnio_browser() + recovered = browser_tool._recover_omnio_browser("conversation-a") assert recovered is False mock_get.assert_called_once_with( @@ -422,6 +449,7 @@ def test_toolbox_v5_does_not_receive_recover_operation(monkeypatch) -> None: "Authorization": "Bearer pair-bearer", "Content-Type": "application/json", "X-Omnio-Brand": "brand-a", + "X-Hermes-Session-Id": "conversation-a", }, timeout=3, ) @@ -453,7 +481,7 @@ def test_toolbox_recover_404_degrades_gracefully(monkeypatch) -> None: return_value=unsupported_response, ) as mock_post, ): - recovered = browser_tool._recover_omnio_browser() + recovered = browser_tool._recover_omnio_browser("conversation-a") assert recovered is False mock_post.assert_called_once_with( @@ -462,6 +490,7 @@ def test_toolbox_recover_404_degrades_gracefully(monkeypatch) -> None: "Authorization": "Bearer pair-bearer", "Content-Type": "application/json", "X-Omnio-Brand": "brand-a", + "X-Hermes-Session-Id": "conversation-a", }, json={"operation": "recover"}, timeout=20, diff --git a/tests/tools/test_file_tools_tilde_profile.py b/tests/tools/test_file_tools_tilde_profile.py index 003e95b797fd..60826712b2c2 100644 --- a/tests/tools/test_file_tools_tilde_profile.py +++ b/tests/tools/test_file_tools_tilde_profile.py @@ -108,3 +108,15 @@ def test_absolute_tilde_in_workspace_root(self, tmp_path, monkeypatch): assert str(profile_home) in str(resolved) assert str(process_home) not in str(resolved) + + def test_container_tilde_is_preserved_for_remote_home_expansion(self, monkeypatch): + monkeypatch.setattr(ft, "_uses_container_paths", lambda _task: True) + + with patch( + "hermes_constants.get_subprocess_home", + return_value="/Users/host-user", + ): + resolved = ft._resolve_path_for_task("~/brand/brief.md", task_id="sprite") + + assert resolved == ft.PurePosixPath("~/brand/brief.md") + assert "/Users/host-user" not in str(resolved) diff --git a/tests/tools/test_process_registry.py b/tests/tools/test_process_registry.py index a9ea279fce8c..c86c07fc955f 100644 --- a/tests/tools/test_process_registry.py +++ b/tests/tools/test_process_registry.py @@ -781,9 +781,11 @@ class FakeEnv: def __init__(self): self.commands = [] self._responses = iter([ + {"output": ""}, {"output": "hello\n"}, {"output": "1\n"}, {"output": "0\n"}, + {"output": "hello\n"}, ]) def execute(self, command, **kwargs): @@ -802,9 +804,63 @@ def execute(self, command, **kwargs): "/path with spaces/hermes_bg.exit", ) - assert env.commands[0][0] == "cat '/path with spaces/hermes_bg.log' 2>/dev/null" - assert env.commands[1][0] == "kill -0 \"$(cat '/path with spaces/hermes_bg.pid' 2>/dev/null)\" 2>/dev/null; echo $?" - assert env.commands[2][0] == "cat '/path with spaces/hermes_bg.exit' 2>/dev/null" + assert env.commands[0][0] == "cat '/path with spaces/hermes_bg.exit' 2>/dev/null" + assert env.commands[1][0] == "cat '/path with spaces/hermes_bg.log' 2>/dev/null" + assert env.commands[2][0] == "kill -0 \"$(cat '/path with spaces/hermes_bg.pid' 2>/dev/null)\" 2>/dev/null; echo $?" + assert env.commands[3][0] == "cat '/path with spaces/hermes_bg.exit' 2>/dev/null" + assert env.commands[4][0] == "cat '/path with spaces/hermes_bg.log' 2>/dev/null" + + def test_env_poller_prefers_durable_exit_marker(self, registry): + session = _make_session(sid="proc_exit_marker") + session.exited = False + + class FakeEnv: + def execute(self, command, **kwargs): + if "hermes_bg.exit" in command: + return {"output": "7\n"} + if "job.log" in command: + return {"output": "finished output\n"} + raise AssertionError(f"unexpected command: {command}") + + with patch("tools.process_registry.time.sleep", return_value=None), \ + patch.object(registry, "_move_to_finished") as move: + registry._env_poller_loop( + session, FakeEnv(), "/tmp/job.log", "/tmp/job.pid", "/tmp/hermes_bg.exit" + ) + + assert session.exited is True + assert session.exit_code == 7 + assert session.completion_reason == "exited" + assert session.output_buffer == "finished output\n" + move.assert_called_once_with(session) + + def test_env_poller_retries_transient_backend_failure(self, registry): + session = _make_session(sid="proc_transient") + session.exited = False + calls = {"count": 0} + + class FakeEnv: + def execute(self, command, **kwargs): + calls["count"] += 1 + if calls["count"] == 1: + raise ConnectionError("temporary toolbox disconnect") + if "hermes_bg.exit" in command: + return {"output": "0\n"} + if "job.log" in command: + return {"output": "done after reconnect\n"} + raise AssertionError(f"unexpected command: {command}") + + with patch("tools.process_registry.time.sleep", return_value=None), \ + patch.object(registry, "_move_to_finished") as move: + registry._env_poller_loop( + session, FakeEnv(), "/tmp/job.log", "/tmp/job.pid", "/tmp/hermes_bg.exit" + ) + + assert calls["count"] == 3 + assert session.exit_code == 0 + assert session.completion_reason == "exited" + assert session.output_buffer == "done after reconnect\n" + move.assert_called_once_with(session) # ========================================================================= diff --git a/tests/tools/test_sprites_environment.py b/tests/tools/test_sprites_environment.py index 481cfb78bf7a..4ac7a73b8dd9 100644 --- a/tests/tools/test_sprites_environment.py +++ b/tests/tools/test_sprites_environment.py @@ -1083,6 +1083,40 @@ def file_request(self, payload): assert 'HTTP 400: {"detail":"bad path"}' in caplog.text +def test_sprites_file_operations_expand_tilde_with_sprite_home(): + from tools.environments.sprites import SpritesEnvironment, SpritesFileOperations + + class FakeEnv: + cwd = "/brand" + config = None + + def __init__(self): + self.requests = [] + + def execute(self, command, cwd=None, **kwargs): + assert command == "echo $HOME" + return {"output": "/home/oai/share\n", "returncode": 0} + + def file_request(self, payload): + self.requests.append(payload) + return {"content": "brief", "totalLines": 1, "fileSize": 5} + + env = FakeEnv() + ops = SpritesFileOperations(cast(SpritesEnvironment, env)) + + result = ops.read_file("~/brand/brief.md") + + assert result.error is None + assert env.requests == [ + { + "operation": "read", + "path": "/home/oai/share/brand/brief.md", + "offset": 1, + "limit": 500, + } + ] + + def test_sprites_environment_should_write_content_through_files_endpoint(): from tools.environments.sprites import SpritesEnvironment diff --git a/tools/browser_camofox.py b/tools/browser_camofox.py index fc60d1cb5b6a..acda3d04f734 100644 --- a/tools/browser_camofox.py +++ b/tools/browser_camofox.py @@ -98,7 +98,7 @@ def _config_cdp_url() -> str: Read here (instead of importing ``browser_tool._get_cdp_override`` to avoid a circular import) so Camofox can yield to a config-based CDP override the - same way it already yields to the ``BROWSER_CDP_URL`` env override. + same way it already yields to CDP environment overrides. """ try: from hermes_cli.config import read_raw_config @@ -117,12 +117,15 @@ def is_camofox_mode() -> bool: A CDP override takes priority over Camofox so the browser tools operate on the real CDP browser (and a CDP backend is treated as non-local for SSRF checks) instead of being silently routed to Camofox. The override may come - from the ``BROWSER_CDP_URL`` env var (set by ``/browser connect``) OR a - persistent ``browser.cdp_url`` in config.yaml — both are honored, matching + from ``BROWSER_CDP_URL_TEMPLATE`` (conversation-scoped gateways), the + ``BROWSER_CDP_URL`` env var (set by ``/browser connect``), OR a persistent + ``browser.cdp_url`` in config.yaml — all are honored, matching ``browser_tool._get_cdp_override()``'s precedence. (Previously only the env var suppressed Camofox, so ``CAMOFOX_URL`` + a config CDP override still routed navigation through Camofox.) """ + if os.getenv("BROWSER_CDP_URL_TEMPLATE", "").strip(): + return False if os.getenv("BROWSER_CDP_URL", "").strip(): return False if _config_cdp_url(): @@ -947,4 +950,3 @@ def camofox_console(clear: bool = False, task_id: Optional[str] = None) -> str: "note": "Console log capture is not available with the Camofox backend. " "Use browser_snapshot or browser_vision to inspect page state.", }) - diff --git a/tools/browser_cdp_tool.py b/tools/browser_cdp_tool.py index f5c90c0de71a..8b5e51f5467b 100644 --- a/tools/browser_cdp_tool.py +++ b/tools/browser_cdp_tool.py @@ -4,9 +4,9 @@ Exposes a single tool, ``browser_cdp``, that sends arbitrary CDP commands to the browser's DevTools WebSocket endpoint. Works when a CDP URL is -configured — either via ``/browser connect`` (sets ``BROWSER_CDP_URL``) or -``browser.cdp_url`` in ``config.yaml`` — or when a CDP-backed cloud provider -session is active. +configured — via a conversation-scoped ``BROWSER_CDP_URL_TEMPLATE``, +``/browser connect`` (sets ``BROWSER_CDP_URL``), or ``browser.cdp_url`` in +``config.yaml`` — or when a CDP-backed cloud provider session is active. This is the escape hatch for browser operations not covered by the main browser tool surface (``browser_navigate``, ``browser_click``, @@ -96,19 +96,20 @@ def _run_async(coro): # --------------------------------------------------------------------------- -def _resolve_cdp_endpoint() -> str: +def _resolve_cdp_endpoint(task_id: Optional[str] = None) -> str: """Return the normalized CDP WebSocket URL, or empty string if unavailable. Delegates to ``tools.browser_tool._get_cdp_override`` so precedence stays consistent with the rest of the browser tool surface: - 1. ``BROWSER_CDP_URL`` env var (live override from ``/browser connect``) - 2. ``browser.cdp_url`` in ``config.yaml`` + 1. ``BROWSER_CDP_URL_TEMPLATE`` (conversation-scoped gateway) + 2. ``BROWSER_CDP_URL`` env var (live override from ``/browser connect``) + 3. ``browser.cdp_url`` in ``config.yaml`` """ try: from tools.browser_tool import _get_cdp_override # type: ignore[import-not-found] - return (_get_cdp_override() or "").strip() + return (_get_cdp_override(task_id) or "").strip() except Exception as exc: # pragma: no cover — defensive logger.debug("browser_cdp: failed to resolve CDP endpoint: %s", exc) return "" @@ -192,14 +193,15 @@ async def _cdp_call( target_id: Optional[str], timeout: float, ) -> Dict[str, Any]: - """Make a single CDP call, optionally attaching to a target first. + """Make a single CDP call, attaching page-scoped methods to a page target. When ``target_id`` is provided, we call ``Target.attachToTarget`` with ``flatten=True`` to multiplex a page-level session over the same browser-level WebSocket, then send ``method`` with that ``sessionId``. - When ``target_id`` is None, ``method`` is sent at browser level — which - works for ``Target.*``, ``Browser.*``, ``Storage.*`` and a few other - globally-scoped domains. + When ``target_id`` is omitted for a page-scoped domain, the sole page + target is selected automatically. Multiple pages require an explicit + target so commands cannot land in the wrong conversation tab. + Browser-scoped methods are sent at the browser level. """ assert websockets is not None # guarded by _WS_AVAILABLE at call-site @@ -213,7 +215,53 @@ async def _cdp_call( next_id = 1 session_id: Optional[str] = None - # --- Step 1: attach to target if requested --- + page_scoped_prefixes = ( + "Page.", + "Runtime.", + "DOM.", + "Emulation.", + "Network.", + "CSS.", + "Input.", + "Overlay.", + "Performance.", + ) + if target_id is None and method.startswith(page_scoped_prefixes): + discover_id = next_id + next_id += 1 + await ws.send(json.dumps({ + "id": discover_id, + "method": "Target.getTargets", + "params": {}, + })) + deadline = asyncio.get_running_loop().time() + timeout + while True: + remaining = deadline - asyncio.get_running_loop().time() + if remaining <= 0: + raise TimeoutError("Timed out discovering a page target") + raw = await asyncio.wait_for(ws.recv(), timeout=remaining) + msg = json.loads(raw) + if msg.get("id") != discover_id: + continue + if "error" in msg: + raise RuntimeError(f"Target.getTargets failed: {msg['error']}") + targets = msg.get("result", {}).get("targetInfos", []) + page_target_ids = [ + str(item.get("targetId")) + for item in targets + if item.get("type") == "page" and item.get("targetId") + ] + if not page_target_ids: + raise RuntimeError("No page target is available for this CDP method") + if len(page_target_ids) > 1: + raise RuntimeError( + "Multiple page targets are open; call Target.getTargets " + "and provide target_id explicitly" + ) + target_id = page_target_ids[0] + break + + # --- Step 1: attach to target if requested or auto-selected --- if target_id: attach_id = next_id next_id += 1 @@ -404,10 +452,9 @@ def browser_cdp( Args: method: CDP method name, e.g. ``"Target.getTargets"``. params: Method-specific parameters; defaults to ``{}``. - target_id: Optional target/tab ID for page-level methods. When set, - we first attach to the target (``flatten=True``) and send - ``method`` with the resulting ``sessionId``. Uses a fresh - stateless CDP connection. + target_id: Optional target/tab ID for page-level methods. When omitted, + Hermes selects the page automatically only when exactly one exists. + Uses a fresh stateless CDP connection. frame_id: Optional cross-origin (OOPIF) iframe ``frame_id`` from ``browser_snapshot.frame_tree.children[]``. When set (and the frame is an OOPIF with a live session tracked by the CDP @@ -457,7 +504,7 @@ def browser_cdp( "Install it with: pip install websockets" ) - endpoint = _resolve_cdp_endpoint() + endpoint = _resolve_cdp_endpoint(effective_task_id) if not endpoint: return tool_error( "No CDP endpoint is available. Run '/browser connect' to attach " @@ -565,7 +612,8 @@ def browser_cdp( "'mobile': false}, target_id=\n\n" "**Usage rules:**\n" "- Browser-level methods (Target.*, Browser.*, Storage.*): omit " - "target_id and frame_id.\n" + "target_id and frame_id. If omitted for a page-scoped method, Hermes " + "selects it automatically only when exactly one page exists.\n" "- Page-level methods (Page.*, Runtime.*, DOM.*, Emulation.*, " "Network.* scoped to a tab): pass target_id from Target.getTargets.\n" "- **Cross-origin iframe scope** (Runtime.evaluate inside an OOPIF, " @@ -638,8 +686,9 @@ def _browser_cdp_check() -> bool: """Availability check for browser_cdp. The tool is only offered when a static CDP URL is configured — via - ``/browser connect`` (``BROWSER_CDP_URL``) or ``browser.cdp_url`` in - ``config.yaml``. Configuration, not liveness: the endpoint's health at + ``BROWSER_CDP_URL_TEMPLATE``, ``/browser connect`` (``BROWSER_CDP_URL``), + or ``browser.cdp_url`` in ``config.yaml``. Configuration, not liveness: + the endpoint's health at registration says nothing about its health when the model calls the tool, and probing it here stalls registration whenever the browser host has not come up yet. diff --git a/tools/browser_tool.py b/tools/browser_tool.py index c064d88c70f9..8ced0c38f061 100644 --- a/tools/browser_tool.py +++ b/tools/browser_tool.py @@ -511,12 +511,31 @@ def _resolve_cdp_override(cdp_url: str) -> str: return ws_url -def _get_cdp_override_raw() -> str: +def _browser_session_id(task_id: Optional[str] = None) -> str: + """Return the conversation identity used by an Omnio browser process. + + API-server tool calls carry the raw Hermes session id in a ContextVar. It + is the stable conversation identity, whereas a delegated tool call may use + a child task id. Non-gateway callers fall back to their task id. + """ + try: + from gateway.session_context import get_session_env + + session_id = str(get_session_env("HERMES_SESSION_ID") or "").strip() + if session_id: + return session_id + except Exception: + pass + return str(task_id or "").strip() + + +def _get_cdp_override_raw(task_id: Optional[str] = None) -> str: """Return the CONFIGURED CDP URL override verbatim, or empty string. Precedence is: - 1. ``BROWSER_CDP_URL`` env var (live override from ``/browser connect``) - 2. ``browser.cdp_url`` in config.yaml (persistent config) + 1. ``BROWSER_CDP_URL_TEMPLATE`` env var (conversation-scoped Omnio relay) + 2. ``BROWSER_CDP_URL`` env var (live override from ``/browser connect``) + 3. ``browser.cdp_url`` in config.yaml (persistent config) When either is set, we skip both Browserbase and the local headless launcher and connect directly to the supplied Chrome DevTools Protocol @@ -528,6 +547,21 @@ def _get_cdp_override_raw() -> str: (is the endpoint live, and what is its concrete websocket) at the cost of seconds when the endpoint is not up yet. """ + template = os.environ.get("BROWSER_CDP_URL_TEMPLATE", "").strip() + if template: + if "{session_id}" not in template: + logger.warning( + "Ignoring BROWSER_CDP_URL_TEMPLATE because it does not contain " + "the required {session_id} placeholder" + ) + return "" + session_id = _browser_session_id(task_id) + if session_id: + from urllib.parse import quote + + return template.replace("{session_id}", quote(session_id, safe="")) + return template + env_override = os.environ.get("BROWSER_CDP_URL", "").strip() if env_override: return env_override @@ -545,7 +579,7 @@ def _get_cdp_override_raw() -> str: return "" -def _get_cdp_override() -> str: +def _get_cdp_override(task_id: Optional[str] = None) -> str: """Return the configured CDP override RESOLVED to a connectable websocket. Discovery is a network probe, so this belongs at the point where a browser @@ -560,13 +594,33 @@ def _get_cdp_override() -> str: would outlive a Chrome restart behind an unchanged discovery URL and hand out a dead endpoint. Discovery is cheap once it is only reached at connect time. """ - raw_override = _get_cdp_override_raw() + raw_override = _get_cdp_override_raw(task_id) if not raw_override: return "" + if "{session_id}" in raw_override: + logger.warning( + "BROWSER_CDP_URL_TEMPLATE is configured but no Hermes session id " + "is available; refusing to connect to an unresolved endpoint" + ) + return "" return _resolve_cdp_override(raw_override) -def _toolbox_browser_binding() -> Optional[Tuple[str, Dict[str, str]]]: +def _get_task_cdp_override(task_id: str) -> str: + """Resolve a task only when the configured endpoint is session-templated. + + Legacy static overrides have no task dimension. Keeping their zero-argument + path also preserves compatibility for downstream browser integrations that + wrap or monkeypatch ``_get_cdp_override``. + """ + if os.environ.get("BROWSER_CDP_URL_TEMPLATE", "").strip(): + return _get_cdp_override(task_id) + return _get_cdp_override() + + +def _toolbox_browser_binding( + task_id: Optional[str] = None, +) -> Optional[Tuple[str, Dict[str, str]]]: """Return the paired Toolbox browser endpoint and authenticated headers.""" base_url = os.environ.get("OMNIO_TOOLBOX_URL", "").strip().rstrip("/") bearer = os.environ.get("OMNIO_TOOLBOX_BEARER", "").strip() @@ -577,19 +631,23 @@ def _toolbox_browser_binding() -> Optional[Tuple[str, Dict[str, str]]]: or os.environ.get("OMNIO_BRAND_ID", "").strip() or "__default__" ) + headers = { + "Authorization": f"Bearer {bearer}", + "Content-Type": "application/json", + "X-Omnio-Brand": brand, + } + session_id = _browser_session_id(task_id) + if session_id: + headers["X-Hermes-Session-Id"] = session_id return ( f"{base_url}/browser", - { - "Authorization": f"Bearer {bearer}", - "Content-Type": "application/json", - "X-Omnio-Brand": brand, - }, + headers, ) -def _recover_omnio_browser() -> bool: - """Ask a recovery-capable paired Toolbox to replace the Brand context.""" - binding = _toolbox_browser_binding() +def _recover_omnio_browser(task_id: Optional[str] = None) -> bool: + """Ask a recovery-capable paired Toolbox to replace this conversation's browser.""" + binding = _toolbox_browser_binding(task_id) if binding is None: return False browser_url, headers = binding @@ -732,13 +790,14 @@ def _recover_browser_after_failure( result: Dict[str, Any], *, cdp_backed: bool, + task_id: Optional[str] = None, ) -> bool: """Recover once after a structured retryable renderer failure.""" if result.get("success") or not cdp_backed: return False if not _should_recover_browser_failure(result): return False - return _recover_omnio_browser() + return _recover_omnio_browser(task_id) def _retry_navigation_after_recovery( @@ -750,7 +809,11 @@ def _retry_navigation_after_recovery( cdp_backed: bool, ) -> Dict[str, Any]: """Recover an Omnio CDP browser and retry one navigation once.""" - if not _recover_browser_after_failure(result, cdp_backed=cdp_backed): + if not _recover_browser_after_failure( + result, + cdp_backed=cdp_backed, + task_id=task_id, + ): return result logger.warning( "Recovered paired browser after navigation failure; retrying once " @@ -815,8 +878,9 @@ def _ensure_cdp_supervisor(task_id: str) -> None: double-attach. Resolves the CDP URL in this order: - 1. ``BROWSER_CDP_URL`` / ``browser.cdp_url`` — covers ``/browser connect`` - and config-set overrides. + 1. ``BROWSER_CDP_URL_TEMPLATE`` / ``BROWSER_CDP_URL`` / + ``browser.cdp_url`` — covers conversation-scoped gateways, + ``/browser connect``, and config-set overrides. 2. ``_active_sessions[task_id]["cdp_url"]`` — covers Browserbase + any other cloud provider whose ``create_session`` returns a raw CDP URL. @@ -824,7 +888,7 @@ def _ensure_cdp_supervisor(task_id: str) -> None: the browser session itself. The agent simply won't see ``pending_dialogs`` / ``frame_tree`` fields in snapshots. """ - cdp_url = _get_cdp_override() + cdp_url = _get_task_cdp_override(task_id) if not cdp_url: # Fallback: active session may carry a per-session CDP URL from a # cloud provider (Browserbase sets this). @@ -2496,7 +2560,7 @@ def _get_session_info(task_id: Optional[str] = None) -> Dict[str, Any]: force_local = _is_local_sidecar_key(task_id) # Create session outside the lock (network call in cloud mode) - cdp_override = _get_cdp_override() + cdp_override = _get_task_cdp_override(task_id) if cdp_override and not force_local: session_info = _create_cdp_session(task_id, cdp_override) elif force_local: @@ -3280,7 +3344,7 @@ def _session_is_cdp_backed( return bool(active_session.get("cdp_url")) if _is_local_sidecar_key(task_id): return False - return bool(os.environ.get("BROWSER_CDP_URL", "").strip()) + return bool(_get_cdp_override_raw(task_id)) def _decode_navigation_probe(result: Dict[str, Any]) -> Dict[str, str]: @@ -3764,6 +3828,7 @@ def browser_snapshot( session_reset = _recover_browser_after_failure( result, cdp_backed=cdp_backed, + task_id=effective_task_id, ) if session_reset: return json.dumps( @@ -4909,6 +4974,27 @@ def browser_vision( ) if not result.get("success"): + session_reset = _recover_browser_after_failure( + result, + cdp_backed=_session_is_cdp_backed(effective_task_id), + task_id=effective_task_id, + ) + if session_reset: + return json.dumps( + { + "success": False, + "error": ( + "The browser renderer became unresponsive and the " + "session was reset; re-open the page with " + "browser_navigate to continue." + ), + "code": "browser_session_reset", + "retryable": True, + "session_reset": True, + "next_action": "browser_navigate", + }, + ensure_ascii=False, + ) error_detail = result.get("error", "Unknown error") _cp = _get_cloud_provider() mode = "local" if _cp is None else f"cloud ({_cp.provider_name()})" @@ -5194,12 +5280,15 @@ def _cleanup_single_browser_session(task_id: str) -> None: # Stop auto-recording before closing (saves the file) _maybe_stop_recording(task_id) - # Try to close via agent-browser first (needs session in _active_sessions) - try: - _run_browser_command(task_id, "close", [], timeout=10) - logger.debug("agent-browser close command completed for task %s", task_id) - except Exception as e: - logger.warning("agent-browser close failed for task %s: %s", task_id, e) + # A CDP-backed session is owned by its provider (Omnio Toolbox or a + # cloud browser). ``agent-browser close`` may issue Browser.close and + # kill that remote browser, so only use it for a locally-owned Chrome. + if not session_info.get("cdp_url"): + try: + _run_browser_command(task_id, "close", [], timeout=10) + logger.debug("agent-browser close command completed for task %s", task_id) + except Exception as e: + logger.warning("agent-browser close failed for task %s: %s", task_id, e) # Now remove from tracking under lock with _cleanup_lock: @@ -5221,17 +5310,7 @@ def _cleanup_single_browser_session(task_id: str) -> None: if session_name: socket_dir = os.path.join(_socket_safe_tmpdir(), f"agent-browser-{session_name}") if os.path.exists(socket_dir): - # agent-browser writes {session}.pid in the socket dir - pid_file = os.path.join(socket_dir, f"{session_name}.pid") - if os.path.isfile(pid_file): - try: - from tools.process_registry import ProcessRegistry - daemon_pid = int(Path(pid_file).read_text(encoding="utf-8").strip()) - ProcessRegistry._terminate_host_pid(daemon_pid) - logger.debug("Killed daemon pid %s for %s", daemon_pid, session_name) - except (ProcessLookupError, ValueError, PermissionError, OSError): - logger.debug("Could not kill daemon pid for %s (already dead or inaccessible)", session_name) - shutil.rmtree(socket_dir, ignore_errors=True) + _terminate_timed_out_browser_daemon(session_info, socket_dir) logger.debug("Removed task %s from active sessions", task_id) else: diff --git a/tools/environments/sprites.py b/tools/environments/sprites.py index af9f12ef91a7..a9412ff56716 100644 --- a/tools/environments/sprites.py +++ b/tools/environments/sprites.py @@ -519,6 +519,14 @@ def __init__(self, terminal_env: SpritesEnvironment): self.env: SpritesEnvironment = terminal_env def _files(self, payload: dict[str, Any]) -> dict[str, Any]: + # The Toolbox file API does not perform shell expansion. Resolve tilde + # paths through the Sprite terminal so file tools and `cd ~` agree on + # the sandbox user's home instead of leaking the gateway host HOME. + payload = dict(payload) + for key in ("path", "src", "dst"): + value = payload.get(key) + if isinstance(value, str) and value.startswith("~"): + payload[key] = self._expand_path(value) try: return self.env.file_request(payload) except SpritesToolboxError as error: diff --git a/tools/file_tools.py b/tools/file_tools.py index 093523fc2fb0..77d4d49bd043 100644 --- a/tools/file_tools.py +++ b/tools/file_tools.py @@ -373,6 +373,11 @@ def _resolve_path_for_task(filepath: str, task_id: str = "default") -> Path | Pu """ container_paths = _uses_container_paths(task_id) if container_paths: + # Preserve the tilde until ShellFileOperations can expand it against + # the sandbox's HOME. Expanding here would use the gateway host HOME + # and turn ~/foo into an invalid host path inside Sprites/Docker/SSH. + if filepath == "~" or filepath.startswith("~/"): + return PurePosixPath(filepath) expanded = _expand_tilde(filepath) if posixpath.isabs(expanded): return _normalize_without_host_deref(expanded) diff --git a/tools/process_registry.py b/tools/process_registry.py index e37e6417c770..c37fc76c8cc1 100644 --- a/tools/process_registry.py +++ b/tools/process_registry.py @@ -1055,9 +1055,25 @@ def _env_poller_loop( quoted_pid_path = shlex.quote(pid_path) quoted_exit_path = shlex.quote(exit_path) prev_output_len = 0 # track delta for watch pattern scanning + consecutive_backend_failures = 0 + dead_without_exit_marker = 0 while not session.exited: time.sleep(2) # Poll every 2 seconds + recorded_exit = None try: + # The exit marker is the source of truth. Read it before a + # liveness probe so PID reuse, a delayed kill -0 result, or a + # transient backend response cannot turn a completed job into + # an invented exit_code=-1. + exit_result = env.execute( + f"cat {quoted_exit_path} 2>/dev/null", + timeout=5, + ) + exit_str = exit_result.get("output", "").strip() + try: + recorded_exit = int(exit_str.splitlines()[-1].strip()) + except (ValueError, IndexError): + recorded_exit = None # Read new output from the log file result = env.execute(f"cat {quoted_log_path} 2>/dev/null", timeout=10) new_output = result.get("output", "") @@ -1073,6 +1089,14 @@ def _env_poller_loop( self._check_watch_patterns(session, delta) self._emit_output(session, delta) + if recorded_exit is not None: + session.exit_code = recorded_exit + session.exited = True + if session.completion_reason != "killed": + session.completion_reason = "exited" + self._move_to_finished(session) + return + # Check if process is still running check = env.execute( f"kill -0 \"$(cat {quoted_pid_path} 2>/dev/null)\" 2>/dev/null; echo $?", @@ -1080,30 +1104,46 @@ def _env_poller_loop( ) check_output = check.get("output", "").strip() if check_output and check_output.splitlines()[-1].strip() != "0": - # Process has exited -- get exit code captured by the wrapper shell. - exit_result = env.execute( - f"cat {quoted_exit_path} 2>/dev/null", - timeout=5, - ) - exit_str = exit_result.get("output", "").strip() - try: - session.exit_code = int(exit_str.splitlines()[-1].strip()) - except (ValueError, IndexError): + # The wrapper writes the marker before it exits. Allow a + # few polls for storage visibility instead of fabricating + # an exit code when the marker briefly lags the PID state. + dead_without_exit_marker += 1 + if dead_without_exit_marker >= 3: + session.exited = True session.exit_code = -1 + session.completion_reason = "lost" + session.termination_source = "exit_marker_missing" + self._move_to_finished(session) + return + else: + dead_without_exit_marker = 0 + consecutive_backend_failures = 0 + + except Exception as exc: + if recorded_exit is not None: + session.exit_code = recorded_exit session.exited = True if session.completion_reason != "killed": session.completion_reason = "exited" self._move_to_finished(session) return - - except Exception: - # Environment might be gone (sandbox reaped, etc.) - session.exited = True - session.exit_code = -1 - session.completion_reason = "lost" - session.termination_source = "backend_lost" - self._move_to_finished(session) - return + # One failed Toolbox/backend poll is not evidence that the + # sandbox or process disappeared. Require repeated failures; + # the next successful poll can still observe the exit marker. + consecutive_backend_failures += 1 + logger.warning( + "Background process poll failed (%d/3, session=%s): %s", + consecutive_backend_failures, + session.id, + exc, + ) + if consecutive_backend_failures >= 3: + session.exited = True + session.exit_code = -1 + session.completion_reason = "lost" + session.termination_source = "backend_lost" + self._move_to_finished(session) + return def _pty_reader_loop(self, session: ProcessSession): """Background thread: read output from a PTY process.""" diff --git a/website/docs/reference/environment-variables.md b/website/docs/reference/environment-variables.md index 362652974035..67cd067907cb 100644 --- a/website/docs/reference/environment-variables.md +++ b/website/docs/reference/environment-variables.md @@ -143,6 +143,7 @@ For native Anthropic auth, Hermes prefers Claude Code's own credential files whe | `BROWSER_USE_API_KEY` | Browser Use cloud browser API key ([browser-use.com](https://browser-use.com/)) | | `FIRECRAWL_BROWSER_TTL` | Firecrawl browser session TTL in seconds (default: 300) | | `BROWSER_CDP_URL` | Chrome DevTools Protocol URL for local browser (set via `/browser connect`, e.g. `ws://localhost:9222`) | +| `BROWSER_CDP_URL_TEMPLATE` | Conversation-scoped CDP gateway URL containing `{session_id}`; takes precedence over `BROWSER_CDP_URL` | | `CAMOFOX_URL` | Camofox local anti-detection browser URL (default: `http://localhost:9377`) | | `CAMOFOX_USER_ID` | Optional externally managed Camofox user ID for shared visible sessions | | `CAMOFOX_SESSION_KEY` | Optional Camofox session key used when creating tabs for `CAMOFOX_USER_ID` |