Skip to content

Commit 9088723

Browse files
committed
Route browser fs and logs endpoints directly to the VM
Add fs and logs to the default direct-to-VM subresource prefixes so filesystem operations and log streaming use the cached browser base_url and JWT instead of the control plane. Serialize multipart array entries with indexed names so a file part stays associated with the sibling fields of its array entry, which repeated `files[][file]` names cannot express. Only retry a stale direct-to-VM auth failure on the control plane when the request body can be rebuilt byte for byte; a streamed body is consumed by the direct attempt, so retrying would send a truncated body. The stale route is evicted either way.
1 parent 4d26fb9 commit 9088723

7 files changed

Lines changed: 504 additions & 35 deletions

File tree

‎src/kernel/_base_client.py‎

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -587,7 +587,10 @@ def _serialize_multipartform(self, data: Mapping[object, object]) -> dict[str, o
587587
# TODO: type ignore is required as stringify_items is well typed but we can't be
588588
# well typed without heavy validation.
589589
data, # type: ignore
590-
array_format="brackets",
590+
# Indexed names (`files[0][dest_path]`) keep each array entry's fields
591+
# grouped together; repeated `files[][dest_path]` parts cannot be
592+
# matched back to their file part. `extract_files` uses the same format.
593+
array_format="indices",
591594
)
592595
serialized: dict[str, object] = {}
593596
for key, value in items:

‎src/kernel/_client.py‎

Lines changed: 11 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -41,6 +41,7 @@
4141
strip_direct_vm_auth,
4242
rewrite_direct_vm_options,
4343
browser_routing_config_from_env,
44+
is_stale_direct_vm_auth_response,
4445
should_retry_stale_direct_vm_auth,
4546
maybe_evict_browser_route_from_response,
4647
maybe_populate_browser_route_cache_from_response,
@@ -365,9 +366,12 @@ def _prepare_request(self, request: httpx.Request) -> None:
365366

366367
@override
367368
def _should_retry(self, response: httpx.Response) -> bool:
368-
if should_retry_stale_direct_vm_auth(response):
369+
if is_stale_direct_vm_auth_response(response):
369370
maybe_evict_browser_route_from_response(response, cache=self.browser_route_cache)
370-
return True
371+
# The route is evicted either way; only retry when the body can be
372+
# rebuilt, otherwise the caller sees the original auth failure and a
373+
# later call goes to the control plane.
374+
return should_retry_stale_direct_vm_auth(response)
371375
return super()._should_retry(response)
372376

373377
@override
@@ -748,9 +752,12 @@ async def _prepare_request(self, request: httpx.Request) -> None:
748752

749753
@override
750754
def _should_retry(self, response: httpx.Response) -> bool:
751-
if should_retry_stale_direct_vm_auth(response):
755+
if is_stale_direct_vm_auth_response(response):
752756
maybe_evict_browser_route_from_response(response, cache=self.browser_route_cache)
753-
return True
757+
# The route is evicted either way; only retry when the body can be
758+
# rebuilt, otherwise the caller sees the original auth failure and a
759+
# later call goes to the control plane.
760+
return should_retry_stale_direct_vm_auth(response)
754761
return super()._should_retry(response)
755762

756763
@override

‎src/kernel/_utils/_utils.py‎

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -40,15 +40,17 @@ def extract_files(
4040
query: Mapping[str, object],
4141
*,
4242
paths: Sequence[Sequence[str]],
43-
array_format: ArrayFormat = "brackets",
43+
array_format: ArrayFormat = "indices",
4444
) -> list[tuple[str, FileTypes]]:
4545
"""Recursively extract files from the given dictionary based on specified paths.
4646
4747
A path may look like this ['foo', 'files', '<array>', 'data'].
4848
4949
``array_format`` controls how ``<array>`` segments contribute to the emitted
5050
field name. Supported values: ``"brackets"`` (``foo[]``), ``"repeat"`` and
51-
``"comma"`` (``foo``), ``"indices"`` (``foo[0]``, ``foo[1]``).
51+
``"comma"`` (``foo``), ``"indices"`` (``foo[0]``, ``foo[1]``). Indexed names are
52+
the default so that a file part stays associated with the sibling fields of the
53+
same array entry, which repeated ``foo[]`` names cannot express.
5254
5355
Note: this mutates the given dictionary.
5456
"""

‎src/kernel/lib/browser_routing/routing.py‎

Lines changed: 54 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -44,7 +44,17 @@ def browser_routing_config_from_env() -> BrowserRoutingConfig:
4444
# Path prefixes eligible for direct-to-VM routing. "telemetry/stream" is
4545
# the live SSE endpoint (VM); "telemetry/events" is a historical read
4646
# served by the control plane (S2) and must NOT be here.
47-
return BrowserRoutingConfig(subresources=("curl", "telemetry/stream", "computer", "playwright", "process"))
47+
return BrowserRoutingConfig(
48+
subresources=(
49+
"curl",
50+
"telemetry/stream",
51+
"computer",
52+
"playwright",
53+
"process",
54+
"fs",
55+
"logs",
56+
)
57+
)
4858
if raw.strip() == "":
4959
return BrowserRoutingConfig()
5060

@@ -189,7 +199,49 @@ def is_stale_direct_vm_auth_response(response: httpx.Response) -> bool:
189199

190200

191201
def should_retry_stale_direct_vm_auth(response: httpx.Response) -> bool:
192-
return is_stale_direct_vm_auth_response(response)
202+
"""Whether a stale direct-to-VM auth failure can be retried on the control plane.
203+
204+
A retry rebuilds the request from the original options, so it is only safe when
205+
the body can be serialized again byte for byte. Streamed bodies (e.g. a file
206+
object passed to fs.write_file) are consumed by the direct request, so retrying
207+
would send a truncated or empty body to the control plane.
208+
"""
209+
if not is_stale_direct_vm_auth_response(response):
210+
return False
211+
return direct_vm_request_body_is_replayable(response.request)
212+
213+
214+
def direct_vm_request_body_is_replayable(request: httpx.Request) -> bool:
215+
try:
216+
_ = request.content
217+
except httpx.RequestNotRead:
218+
pass
219+
else:
220+
# httpx already buffered the body, so rebuilding it yields the same bytes.
221+
return True
222+
223+
# httpx encodes multipart bodies as a stream of fields it re-renders per attempt.
224+
fields = getattr(request.stream, "fields", None)
225+
if fields is None:
226+
# A streamed body (file object, iterator or async iterator) cannot be replayed.
227+
return False
228+
return all(_multipart_field_is_replayable(field) for field in cast("list[Any]", fields))
229+
230+
231+
def _multipart_field_is_replayable(field: Any) -> bool:
232+
file = getattr(field, "file", None)
233+
if file is None:
234+
# A data field renders from an in-memory value.
235+
return True
236+
if isinstance(file, (bytes, str)):
237+
return True
238+
if getattr(file, "closed", False):
239+
return False
240+
if not callable(getattr(file, "seek", None)):
241+
return False
242+
seekable = getattr(file, "seekable", None)
243+
# httpx rewinds seekable file fields before rendering them again.
244+
return bool(seekable()) if callable(seekable) else True
193245

194246

195247
def _session_id_from_browser_delete_path(path: str) -> str | None:

0 commit comments

Comments
 (0)