-
Notifications
You must be signed in to change notification settings - Fork 462
fix: clear 0.4.3 release blockers (discovery path/XDG + WebSocket close) #1170
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
56368e0
541c023
7321d73
10d3647
4d2fc5b
4b4a5c2
ebd608a
287c715
0d68e5c
b35fcea
d6d61f4
99cfaa0
5e3f9f9
440faf5
39e30dc
94b95e3
2831ed5
d140395
7dc8031
210b678
f336c19
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -23,6 +23,7 @@ | |
| import os | ||
| import re | ||
| import stat | ||
| import tempfile | ||
| from dataclasses import asdict, dataclass | ||
| from pathlib import Path | ||
| from typing import Any, Type | ||
|
|
@@ -345,9 +346,28 @@ def _default_cache_file() -> Path: | |
| shared, world-writable temporary directory. A fixed path under the shared temp dir | ||
| lets another local user pre-create the cache file and redirect discovery to | ||
| attacker-controlled import paths (`import_module` on a cached `client_module_path`). | ||
| Per the XDG Base Directory specification, relative `XDG_CACHE_HOME` values are | ||
| ignored so an untrusted working tree cannot supply a victim-owned cache file. | ||
| Relative `HOME` values are likewise rejected: `Path.home()` must not be | ||
| resolved against the current working directory. | ||
| """ | ||
| base = os.environ.get("XDG_CACHE_HOME") | ||
| root = Path(base) if base else Path.home() / ".cache" | ||
| if base and Path(base).is_absolute(): | ||
| root = Path(base) | ||
| else: | ||
| home = Path.home() | ||
| if home.is_absolute(): | ||
| root = home / ".cache" | ||
| else: | ||
| # Keep the fallback absolute and uid-scoped so a shared temp root | ||
| # cannot be turned into a fixed, cross-user planting target. | ||
| uid = os.getuid() if hasattr(os, "getuid") else os.getpid() | ||
| root = Path(tempfile.gettempdir()) / f"openenv-{uid}-cache" | ||
| if not root.is_absolute(): | ||
|
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Temp fallback can use checkout cacheHigh Severity The new relative- Reviewed by Cursor Bugbot for commit f336c19. Configure here. |
||
| raise RuntimeError( | ||
| "Cannot resolve an absolute discovery cache directory when " | ||
| "XDG_CACHE_HOME and HOME are both missing or relative" | ||
| ) | ||
| return root / "openenv" / "discovery_cache.json" | ||
|
|
||
|
|
||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -255,6 +255,15 @@ async def _best_effort_close(ws: ClientConnection) -> None: | |
| pass # Best effort | ||
|
|
||
|
|
||
| async def _best_effort_disconnect(ws: ClientConnection) -> None: | ||
| """Notify the server, then close the socket without propagating failures.""" | ||
| try: | ||
| await ws.send(json.dumps({"type": "close"})) | ||
| except (Exception, asyncio.CancelledError): | ||
| pass # Best effort | ||
| await _best_effort_close(ws) | ||
|
|
||
|
|
||
| class EnvClient(ABC, Generic[ActT, ObsT, StateT]): | ||
| """ | ||
| Async environment client for persistent sessions. | ||
|
|
@@ -538,6 +547,13 @@ async def _connect_async(self) -> "EnvClient": | |
| self._ws = None | ||
| self._ws_loop = None | ||
|
|
||
| # A timed-out request drops its socket immediately but closes it in the | ||
| # background so the timeout itself remains prompt. Wait for that close | ||
| # before opening a replacement: the old server-side session continues | ||
| # occupying a capacity slot until the close handshake finishes, and | ||
| # many environments allow only one session. | ||
| await self._drain_pending_close_tasks() | ||
|
|
||
| try: | ||
| self._start_provider_if_needed() | ||
| except Exception: | ||
|
|
@@ -573,24 +589,49 @@ async def _connect_async(self) -> "EnvClient": | |
| def disconnect(self) -> Any: | ||
| return self._dispatch(self._disconnect_async) | ||
|
|
||
| def _schedule_socket_close( | ||
| self, ws: ClientConnection, *, notify_server: bool = False | ||
| ) -> asyncio.Task[None]: | ||
| """Schedule and track a socket close on the current event loop.""" | ||
| close = _best_effort_disconnect(ws) if notify_server else _best_effort_close(ws) | ||
| close_task = asyncio.create_task(close) | ||
| self._pending_close_tasks.add(close_task) | ||
| close_task.add_done_callback(self._pending_close_tasks.discard) | ||
| return close_task | ||
|
|
||
| async def _disconnect_async(self) -> None: | ||
| """Close the WebSocket connection.""" | ||
| if self._ws is not None: | ||
| ws = self._ws | ||
| ws_loop = self._ws_loop | ||
| same_loop = ws_loop is asyncio.get_running_loop() | ||
| try: | ||
| if same_loop: | ||
| await ws.send(json.dumps({"type": "close"})) | ||
| except Exception: | ||
| pass # Best effort | ||
| try: | ||
| if same_loop: | ||
| await ws.close() | ||
| except Exception: | ||
| pass | ||
| # Detach first so cancellation during the close handshake cannot | ||
| # leave a stale socket cached for a later operation. | ||
| self._ws = None | ||
| self._ws_loop = None | ||
|
cursor[bot] marked this conversation as resolved.
|
||
| same_loop = ws_loop is asyncio.get_running_loop() | ||
| if same_loop: | ||
| # Track the detached socket before awaiting anything. If this | ||
| # caller is cancelled, the shielded task keeps closing and a | ||
| # later reconnect drains it before opening a replacement. | ||
| close_task = self._schedule_socket_close(ws, notify_server=True) | ||
| await asyncio.shield(close_task) | ||
|
|
||
| async def _drain_pending_close_tasks(self) -> None: | ||
| """Wait for background socket closes owned by the current event loop. | ||
|
|
||
| Shielding keeps cancellation of the caller from cancelling the close | ||
| tasks themselves. This matters both before reconnecting, when the old | ||
| server session must release its capacity slot, and during explicit | ||
| client shutdown. | ||
| """ | ||
| loop = asyncio.get_running_loop() | ||
| tasks = [ | ||
| task | ||
| for task in tuple(self._pending_close_tasks) | ||
| if not task.done() and task.get_loop() is loop | ||
| ] | ||
| if tasks: | ||
| await asyncio.shield(asyncio.gather(*tasks, return_exceptions=True)) | ||
|
|
||
| async def _ensure_connected(self) -> None: | ||
| """Ensure WebSocket connection is established on the current loop. | ||
|
|
@@ -641,9 +682,7 @@ async def _receive(self) -> Dict[str, Any]: | |
| # would actually block for up to 10s before its deadline was | ||
| # honored. Scheduling it lets the exception propagate | ||
| # immediately while the close still happens in the background. | ||
| close_task = asyncio.ensure_future(_best_effort_close(ws)) | ||
| self._pending_close_tasks.add(close_task) | ||
| close_task.add_done_callback(self._pending_close_tasks.discard) | ||
| self._schedule_socket_close(ws) | ||
| raise | ||
| return json.loads(raw) | ||
|
|
||
|
|
@@ -957,34 +996,36 @@ async def _close_async(self) -> None: | |
| If this client was created via from_docker_image() or from_env(), | ||
| this will also stop and remove the associated container/process. | ||
| """ | ||
| for child in list(self._child_clients): | ||
| with suppress(Exception): | ||
| await child.close() | ||
| self._child_clients.clear() | ||
|
|
||
| try: | ||
| # Wait out any backgrounded closes from a dropped socket (see | ||
| # `_receive()` / `_best_effort_close`) so a real close() call still | ||
| # sees the handshake through. SyncEnvClient.close() waits for | ||
| # `_close_async()` before stopping its loop, so the relevant risk | ||
| # is async-context cancellation of close itself — not `_stop_loop()`. | ||
| # Keep this gather inside the provider-teardown try/finally so a | ||
| # cancelled close cannot skip container/process cleanup. | ||
| if self._pending_close_tasks: | ||
| await asyncio.gather(*self._pending_close_tasks, return_exceptions=True) | ||
| await self._disconnect_async() | ||
| for child in list(self._child_clients): | ||
| with suppress(Exception): | ||
| await child.close() | ||
| finally: | ||
| # Parent teardown must run even when a child close is cancelled. | ||
| self._child_clients.clear() | ||
|
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Cancel drops remaining child sessionsMedium Severity Cancelling Reviewed by Cursor Bugbot for commit f336c19. Configure here. |
||
| try: | ||
| if self._provider is not None: | ||
| # Handle both ContainerProvider and RuntimeProvider | ||
| if hasattr(self._provider, "stop_container"): | ||
| self._provider.stop_container() | ||
| elif hasattr(self._provider, "stop"): | ||
| self._provider.stop() | ||
| try: | ||
| # A real close waits out backgrounded closes, but shield them | ||
| # from cancellation so their socket handshakes aren't | ||
| # abandoned midway. | ||
| await self._drain_pending_close_tasks() | ||
| finally: | ||
| # Run even when pending-close draining is cancelled. A client | ||
| # may already have reconnected, and that current socket must | ||
| # not remain cached or open during teardown. | ||
| await self._disconnect_async() | ||
| finally: | ||
| if self._start_provider_on_connect: | ||
| self._base_url = None | ||
| self._ws_url = None | ||
| try: | ||
| if self._provider is not None: | ||
| # Handle both ContainerProvider and RuntimeProvider | ||
| if hasattr(self._provider, "stop_container"): | ||
| self._provider.stop_container() | ||
| elif hasattr(self._provider, "stop"): | ||
| self._provider.stop() | ||
| finally: | ||
| if self._start_provider_on_connect: | ||
| self._base_url = None | ||
| self._ws_url = None | ||
|
|
||
| def _stop_provider_best_effort(self) -> None: | ||
| """Stop the underlying provider directly, ignoring any errors. | ||
|
|
||


There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Minor nit (non-blocking): this paragraph describes what the packaged schemas enforce (relative-path safety), but it's inserted immediately after "Other rules require semantic validation:" — whose trailing colon introduces the table of semantic-validation rules right below it. As written, the schema-enforced note interrupts that lead-in → table flow.
Consider either folding it into the preceding sentence ("The packaged schemas enforce object shape, required fields, relative-path safety, ...") or separating it from the "semantic validation:" clause (e.g. its own line before that sentence), so the colon flows directly into the table.