Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
43 changes: 43 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -295,6 +295,49 @@ sandbox.files.write(
)
```

### Run commands and collect output

`sandbox.exec()` and `sandbox.processes.start()` stream command output from the
receiver as soon as execution starts. The SDK collects stdout and stderr in
memory, so a completed result is not limited to the receiver's replay buffer.
This requires a receiver supporting streaming `POST /sandbox/processes`; roll
out the receiver before upgrading the SDK.

```python
result = sandbox.exec("make test", max_output_bytes=128 * 1024 * 1024)
print(result.stdout, result.stderr, result.exit_code)

process = sandbox.processes.start("make test")
try:
for event in process.stream():
if event.type == "stdout":
print(event.data, end="")
result = process.wait()
finally:
process.disconnect()
```

The combined output limit defaults to 64 MiB per command and can be adjusted
with `max_output_bytes`. Exceeding it raises `output_limit_exceeded`. A broken
stream, missing output, or receiver truncation raises `incomplete_output`.
These errors include the process ID and do not automatically rerun the command.

Process streams use a separate 60-second read-idle timeout once response headers
arrive. Output and the receiver's 15-second heartbeats reset this timeout, so quiet
commands can run longer than the client's ordinary HTTP timeout. That ordinary
timeout still applies to connection setup and waiting for response headers.

`start()` returns after the process starts and collects in the background.
`wait(timeout_sec=...)` limits the local wait; collection continues after a wait
timeout. The timeout passed to `start()` or `exec()` limits command execution.
A local wait timeout raises `TimeoutError` (`asyncio.TimeoutError` in the async API).
`disconnect()` stops collection and leaves the command running. Use `kill()` to
stop it. Reattaching with `get()` can retrieve only retained receiver output;
`wait()` raises if that output has been truncated.

The async API has the same behavior: await `exec()`, `start()`, `wait()`, and
`disconnect()`, and use `async for` with `stream()`.

### Resume terminal output after reconnect

```python
Expand Down
2 changes: 2 additions & 0 deletions hyperbrowser/client/managers/async_manager/sandbox.py
Original file line number Diff line number Diff line change
Expand Up @@ -307,6 +307,7 @@ async def exec(
timeout_ms: Optional[int] = None,
timeout_sec: Optional[int] = None,
run_as: Optional[str] = None,
max_output_bytes: int = 64 * 1024 * 1024,
):
return await self.processes.exec(
input,
Expand All @@ -315,6 +316,7 @@ async def exec(
timeout_ms=timeout_ms,
timeout_sec=timeout_sec,
run_as=run_as,
max_output_bytes=max_output_bytes,
)

async def get_process(self, process_id: str) -> SandboxProcessHandle:
Expand Down
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import asyncio
import base64
from typing import AsyncIterator, Dict, Optional, Union

Expand All @@ -16,6 +17,11 @@
SandboxProcessStdinParams as SandboxProcessStdinParamsDict,
)
from ...sandboxes.shared import _normalize_exec_params
from ...sandboxes.process_output import (
DEFAULT_MAX_PROCESS_OUTPUT_BYTES,
ProcessOutput,
validate_output_limit,
)
from .sandbox_transport import RuntimeTransport

DEFAULT_PROCESS_KILL_WAIT_SECONDS = 5.0
Expand All @@ -25,6 +31,10 @@ class SandboxProcessHandle:
def __init__(self, transport: RuntimeTransport, summary: SandboxProcessSummary):
self._transport = transport
self._summary = summary
self._output = None
self._collector = None
self._events = None
self._changed = asyncio.Event()

@property
def id(self) -> str:
Expand All @@ -50,6 +60,15 @@ async def wait(
timeout_ms: Optional[int] = None,
timeout_sec: Optional[int] = None,
) -> SandboxProcessResult:
if self._collector is not None:
timeout = None
if timeout_sec is not None and timeout_sec > 0:
timeout = timeout_sec
elif timeout_ms is not None and timeout_ms > 0:
timeout = timeout_ms / 1000
if not self._collector.done():
await asyncio.wait_for(asyncio.shield(self._collector), timeout)
return self._collected_result()
payload = await self._transport.request_json(
f"/sandbox/processes/{self.id}/wait",
method="POST",
Expand All @@ -60,6 +79,10 @@ async def wait(
headers={"content-type": "application/json"},
)
result = SandboxProcessResult(**payload["result"])
if result.output_truncated:
raise ProcessOutput(self.id, 0).failure(
"Retained process output is incomplete; collect output from process start"
)
self._summary = SandboxProcessSummary(
id=result.id,
status=result.status,
Expand Down Expand Up @@ -133,6 +156,21 @@ async def write_stdin(
)

async def stream(self, from_seq: Optional[int] = None) -> AsyncIterator[object]:
if self._output is not None:
index = 0
while True:
self._changed.clear()
while index < len(self._output.events):
event = self._output.events[index]
index += 1
if from_seq is None or event.seq >= from_seq:
yield event
if self._collector.done():
yield SandboxProcessExitEvent(
type="exit", result=self._collected_result()
)
return
await self._changed.wait()
params = {"from_seq": from_seq} if from_seq and from_seq > 0 else None
async for event in self._transport.stream_sse(
f"/sandbox/processes/{self.id}/stream",
Expand All @@ -153,6 +191,70 @@ async def stream(self, from_seq: Optional[int] = None) -> AsyncIterator[object]:
result=SandboxProcessResult(**data),
)

def _collected_result(self) -> SandboxProcessResult:
if self._output.error is not None:
raise self._output.error
if self._output.result is None:
raise self._output.failure(
"Command stream ended before its completion event"
)
result = self._output.result
self._summary = self._summary.model_copy(
update={
"status": result.status,
"exit_code": result.exit_code,
"completed_at": result.completed_at,
}
)
return result

async def _collect(self) -> None:
try:
async for event in self._events:
self._output.consume(event)
self._changed.set()
if self._output.result is not None:
return
self._output.error = self._output.failure(
"Command stream ended before its completion event"
)
except asyncio.CancelledError:
if self._output.result is None:
self._output.error = self._output.failure(
"Command output collection disconnected"
)
except Exception as error:
self._output.error = (
self._output.failure(str(error))
if not hasattr(error, "code")
else error
)
finally:
try:
await self._events.aclose()
except Exception as error:
if self._output.result is None and self._output.error is None:
self._output.error = self._output.failure(str(error))
finally:
self._changed.set()

async def disconnect(self) -> None:
"""Stop collecting output; the detached command continues running."""
if self._collector is not None and not self._collector.done():
if self._output.result is None and self._output.error is None:
self._output.error = self._output.failure(
"Command output collection disconnected"
)
self._collector.cancel()
try:
await self._collector
except asyncio.CancelledError:
pass
finally:
# Cancellation may happen before the collector gets its first turn.
await self._events.aclose()
self._changed.set()

async def result(self) -> SandboxProcessResult:
return await self.wait()

Expand All @@ -170,22 +272,21 @@ async def exec(
timeout_ms: Optional[int] = None,
timeout_sec: Optional[int] = None,
run_as: Optional[str] = None,
max_output_bytes: int = DEFAULT_MAX_PROCESS_OUTPUT_BYTES,
) -> SandboxProcessResult:
params = _normalize_exec_params(
handle = await self.start(
input,
cwd=cwd,
env=env,
timeout_ms=timeout_ms,
timeout_sec=timeout_sec,
run_as=run_as,
max_output_bytes=max_output_bytes,
)
payload = await self._transport.request_json(
"/sandbox/exec",
method="POST",
json_body=dump_request(params, SandboxExecParams),
headers={"content-type": "application/json"},
)
return SandboxProcessResult(**payload["result"])
try:
return await handle.wait()
finally:
await handle.disconnect()

async def start(
self,
Expand All @@ -196,7 +297,9 @@ async def start(
timeout_ms: Optional[int] = None,
timeout_sec: Optional[int] = None,
run_as: Optional[str] = None,
max_output_bytes: int = DEFAULT_MAX_PROCESS_OUTPUT_BYTES,
) -> SandboxProcessHandle:
validate_output_limit(max_output_bytes)
params = _normalize_exec_params(
input,
cwd=cwd,
Expand All @@ -205,16 +308,25 @@ async def start(
timeout_sec=timeout_sec,
run_as=run_as,
)
payload = await self._transport.request_json(
events = self._transport.stream_sse(
"/sandbox/processes",
method="POST",
json_body=dump_request(params, SandboxExecParams),
headers={"content-type": "application/json"},
)
return SandboxProcessHandle(
self._transport,
SandboxProcessSummary(**payload["process"]),
)
try:
started = await events.__anext__()
if started["event"] != "started":
raise RuntimeError("Expected process start event")
handle = SandboxProcessHandle(
self._transport, SandboxProcessSummary(**started["data"])
)
except BaseException:
await events.aclose()
raise
handle._events = events
handle._output = ProcessOutput(handle.id, max_output_bytes)
handle._collector = asyncio.create_task(handle._collect())
return handle

async def get(self, process_id: str) -> SandboxProcessHandle:
payload = await self._transport.request_json(f"/sandbox/processes/{process_id}")
Expand Down
Loading