From 2b505654245212ddc8403380f264f6ba0deaab97 Mon Sep 17 00:00:00 2001 From: CMGS Date: Wed, 16 Sep 2026 23:18:10 +0800 Subject: [PATCH 1/2] sdk/python: keep a handle's relay connection across calls Every data-plane call dialed the owner node, upgraded and, behind an edge, shook hands with TLS, then paid a vsock connect on the node (#195). With silkd serving RPCs back to back (#196) a handle now parks its connection after a call and the next call sends on it: a Sandbox carries a ConnPool (8 connections, 30 s idle; Client(keep_alive=...) tunes or disables it) and every RPC helper runs under _park_after, which parks the connection after a terminal frame, error frames included, and drops it on anything else. Streams (watch, pty, port_forward, lsp) still own their connection; run joins its stdin thread before the connection goes back, so the pump's frames never land on the next RPC. close and hibernate drain the pool first. Before a parked connection is reused, a zero-timeout select tells a peer that hung up (hibernate, release, the idle sweep) or spoke unprompted from a quiet one, so the next call redials and wakes the guest as before. A handle's first dial asks the daemon's proto with one info round trip: a daemon before proto 2 closes after answering, so the handle redials and dials per call from then on. logs' trailing done is consumed so the connection stays in frame. Hot path: a pooled call costs one lock pair and one select(2) instead of a TCP dial, an HTTP upgrade, a TLS handshake and a vsock connect; the first call on a handle pays one info round trip. --- docs/sdk-python.md | 10 ++ sdk/python/cocoonsandbox/client.py | 10 +- sdk/python/cocoonsandbox/conn.py | 70 +++++++++++- sdk/python/cocoonsandbox/frames.py | 2 + sdk/python/cocoonsandbox/sandbox.py | 82 +++++++++++--- sdk/python/tests/test_keepalive.py | 157 ++++++++++++++++++++++++++ sdk/python/tests/test_proc.py | 2 + sdk/python/tests/test_stream.py | 17 ++- sdk/python/tests/test_wire_binding.py | 7 +- 9 files changed, 331 insertions(+), 26 deletions(-) create mode 100644 sdk/python/tests/test_keepalive.py diff --git a/docs/sdk-python.md b/docs/sdk-python.md index 7ab7dae9..3b8f1f5f 100644 --- a/docs/sdk-python.md +++ b/docs/sdk-python.md @@ -170,6 +170,16 @@ automatically with the same transparent wake. A claim with a connection live when the sweep checks it (a relay stream, a buffered exec, a preview dial, an egress request) is not swept; the idle clock restarts when that connection ends. +Data-plane calls share a handle's relay connection: after a call the SDK +keeps the connection for 30 seconds (`Client(..., keep_alive=...)` tunes the +window; 0 dials per call) and the next call on that handle sends its request +on it, so a busy handle pays the dial, upgrade and TLS handshake once. A kept +connection counts as live for `idle_hibernate_seconds` until it closes, so +keep the window below that setting; `close` and `hibernate` drop it at once. +Streams (`watch`, `open_pty`, `dial_port`, an LSP session) take a connection +of their own, and a guest whose silkd predates the back-to-back protocol +gets one connection per call as before. + If that deployment also enables `archive_after_seconds`, archiving replaces the original claim deadline with the archive-retention deadline (or no deadline when archives are kept forever). Waking an archive starts a fresh diff --git a/sdk/python/cocoonsandbox/client.py b/sdk/python/cocoonsandbox/client.py index 0febb4d1..c3287a0e 100644 --- a/sdk/python/cocoonsandbox/client.py +++ b/sdk/python/cocoonsandbox/client.py @@ -28,8 +28,15 @@ class Client: """Talks to one sandboxd node (and, transparently, its cluster).""" def __init__( - self, addr: str, api_token: str = "", timeout: float = 120.0, *, ssl_context: ssl.SSLContext | None = None + self, + addr: str, + api_token: str = "", + timeout: float = 120.0, + *, + ssl_context: ssl.SSLContext | None = None, + keep_alive: float = 30.0, ) -> None: + """keep_alive bounds how long a handle keeps an idle relay connection for its next call; 0 dials per call.""" endpoint = _endpoint_url(addr.split(",")[0].strip()) self.addr = endpoint.geturl().removeprefix("http://") self._scheme = endpoint.scheme @@ -39,6 +46,7 @@ def __init__( self._opener: urllib.request.OpenerDirector | None = None self.api_token = api_token self.timeout = timeout + self.keep_alive = keep_alive def new( self, diff --git a/sdk/python/cocoonsandbox/conn.py b/sdk/python/cocoonsandbox/conn.py index 5c87ffab..2e000120 100644 --- a/sdk/python/cocoonsandbox/conn.py +++ b/sdk/python/cocoonsandbox/conn.py @@ -3,15 +3,19 @@ from __future__ import annotations import contextlib +import select import socket import ssl +import threading import time import urllib.parse -from collections.abc import Iterator +from collections.abc import Callable, Iterator from typing import Any, BinaryIO, Protocol, TypeVar from .errors import APIError, ProtocolError, SilkdError -from .frames import MAX_FRAME, decode_response, encode_request +from .frames import KEEP_ALIVE_PROTO, MAX_FRAME, decode_response, encode_request + +KEEP_ALIVE_CONNS = 8 _CloseableT = TypeVar("_CloseableT", bound="_Closeable") @@ -69,6 +73,14 @@ def recv_until(self, *terminal: str) -> Iterator[dict[str, Any]]: if frame["type"] in terminal: return + def quiet(self) -> bool: + """Reports whether the peer has neither hung up nor spoken since the last frame.""" + try: + readable, _, _ = select.select([self._sock], [], [], 0) + except (OSError, ValueError): + return False + return not readable + def close(self) -> None: self.abort() try: @@ -77,6 +89,60 @@ def close(self) -> None: self._sock.close() +class _Parked: + def __init__(self, conn: Conn, idle: float, evict: Callable[[_Parked], None]) -> None: + self.conn = conn + self.timer = threading.Timer(idle, evict, (self,)) + self.timer.daemon = True + + +class ConnPool: + """Parks a handle's idle relay connections between calls; silkd serves RPCs back to back from proto 2.""" + + def __init__(self, idle: float) -> None: + self.proto = 0 + self._idle = idle + self._lock = threading.Lock() + self._parked: list[_Parked] = [] + + def take(self) -> Conn | None: + """Returns a parked connection whose peer is still there, or None.""" + while True: + with self._lock: + if not self._parked: + return None + entry = self._parked.pop() + entry.timer.cancel() + if entry.conn.quiet(): + return entry.conn + entry.conn.close() + + def park(self, conn: Conn) -> None: + """Keeps conn for the next call until idle passes; a daemon before proto 2 or a zero window closes it.""" + with self._lock: + keep = self._idle > 0 and self.proto >= KEEP_ALIVE_PROTO and len(self._parked) < KEEP_ALIVE_CONNS + if keep: + entry = _Parked(conn, self._idle, self._evict) + self._parked.append(entry) + entry.timer.start() + if not keep: + conn.close() + + def drain(self) -> None: + with self._lock: + parked, self._parked = self._parked, [] + for entry in parked: + entry.timer.cancel() + entry.conn.close() + + def _evict(self, entry: _Parked) -> None: + with self._lock: + if entry not in self._parked: + return + self._parked.remove(entry) + entry.conn.close() + + def dial_agent( addr: str, sandbox_id: str, diff --git a/sdk/python/cocoonsandbox/frames.py b/sdk/python/cocoonsandbox/frames.py index a1cea6bf..2bfcefce 100644 --- a/sdk/python/cocoonsandbox/frames.py +++ b/sdk/python/cocoonsandbox/frames.py @@ -8,6 +8,8 @@ from typing import Any PROTO_VERSION = 1 +# the info proto from which silkd serves RPCs back to back on one connection +KEEP_ALIVE_PROTO = 2 MAX_FRAME = 8 * 1024 * 1024 FS_CHUNK = 256 * 1024 # bulk streams chunk larger than silkd's FS_CHUNK: fewer frames per byte, still under MAX_FRAME after base64. diff --git a/sdk/python/cocoonsandbox/sandbox.py b/sdk/python/cocoonsandbox/sandbox.py index e8ab84cb..6c412d0c 100644 --- a/sdk/python/cocoonsandbox/sandbox.py +++ b/sdk/python/cocoonsandbox/sandbox.py @@ -7,12 +7,13 @@ import threading import time from collections.abc import Callable, Iterator +from contextlib import AbstractContextManager from typing import TYPE_CHECKING, Any, cast from .checkpoint import Checkpoint -from .conn import Conn, _Closeable -from .errors import APIError, ExitError, ProtocolError, SandboxError -from .frames import BULK_CHUNK, FS_CHUNK +from .conn import Conn, ConnPool, _Closeable +from .errors import APIError, ExitError, ProtocolError, SandboxError, SilkdError +from .frames import BULK_CHUNK, FS_CHUNK, KEEP_ALIVE_PROTO from .template import Template if TYPE_CHECKING: @@ -41,6 +42,7 @@ def __init__( self.from_checkpoint = from_checkpoint self.template_digest = template_digest self.volumes = [dict(volume) for volume in volumes or []] + self._pool = ConnPool(client.keep_alive) def __enter__(self) -> Sandbox: return self @@ -97,13 +99,14 @@ def run( deadline = None if timeout is None else time.monotonic() + timeout expired = threading.Event() try: - conn = self._dial(deadline) + conn = self._connect(deadline) except (ProtocolError, TimeoutError): if deadline is not None and time.monotonic() >= deadline: raise TimeoutError(f"command did not finish within {timeout}s") from None raise - with conn: + with self._park_after(conn): watchdog = _arm_watchdog(conn, deadline, expired) + pump = threading.Thread(target=_feed_stdin, args=(conn, stdin), daemon=True) try: conn.send( "exec", @@ -114,17 +117,19 @@ def run( detach=False, session=session or None, ) - pump = threading.Thread(target=_feed_stdin, args=(conn, stdin), daemon=True) pump.start() code = _pump_stdio(conn, on_stdout, on_stderr) except (ProtocolError, OSError): if expired.is_set(): raise TimeoutError(f"command did not finish within {timeout}s") from None raise + except SilkdError: + pump.join() + raise finally: if watchdog is not None: watchdog.cancel() - pump.join() # the closed conn fails a stalled send, so this cannot hang + pump.join() # the pump's frames must not land on the next RPC; silkd drains what the command left unread if code is None: raise ProtocolError("exec stream ended without an exit frame") return code @@ -164,19 +169,19 @@ def attach( def write_file(self, path: str, data: bytes, mode: int | None = None) -> None: """Writes data to path atomically (temp + rename on the guest).""" - with self._dial() as conn: + with self._lease() as conn: conn.send("fs_write", path=path, mode=mode) _send_chunks(conn, data) conn.send("data_end") _expect(conn, "done") def read_file(self, path: str) -> bytes: - with self._dial() as conn: + with self._lease() as conn: conn.send("fs_read", path=path) return _drain_data(conn) def list_dir(self, path: str) -> list[dict[str, Any]]: - with self._dial() as conn: + with self._lease() as conn: conn.send("fs_list", path=path) entries: list[dict[str, Any]] = [] for frame in conn.recv_until("done"): @@ -198,7 +203,7 @@ def rename(self, src: str, dst: str) -> None: def push(self, dest: str, tar_stream: bytes) -> None: """Extracts a tar stream into dest; a truncated stream leaves dest untouched.""" - with self._dial() as conn: + with self._lease() as conn: conn.send("fs_push", dest=dest) _send_chunks(conn, tar_stream, chunk=BULK_CHUNK) conn.send("data_end") @@ -206,7 +211,7 @@ def push(self, dest: str, tar_stream: bytes) -> None: def pull(self, path: str) -> bytes: """Returns path (file or tree) as a tar archive.""" - with self._dial() as conn: + with self._lease() as conn: conn.send("fs_pull", path=path) return _drain_data(conn) @@ -215,14 +220,14 @@ def find(self, path: str, pattern: str, glob: str = "") -> list[dict[str, Any]]: def find_iter(self, path: str, pattern: str, glob: str = "") -> Iterator[dict[str, Any]]: """Yields matches as they stream.""" - with self._dial() as conn: + with self._lease() as conn: conn.send("fs_find", path=path, pattern=pattern, glob=glob or None) for f in conn.recv_until("done"): if f["type"] == "match": yield f def replace(self, files: list[str], pattern: str, replacement: str) -> list[dict[str, Any]]: - with self._dial() as conn: + with self._lease() as conn: conn.send("fs_replace", files=files, pattern=pattern, replacement=replacement) return [f for f in conn.recv_until("done") if f["type"] == "replaced"] @@ -284,6 +289,7 @@ def fork(self, count: int, ttl_seconds: int = 0) -> list[Sandbox]: def hibernate(self) -> None: """Snapshots and stops the VM; the next guest call restores its state.""" + self._pool.drain() self._client._request( self.owner, "POST", f"/v1/sandboxes/{self.id}/hibernate", None, "hibernate", bearer=self.token ) @@ -347,6 +353,7 @@ def dial_port(self, port: int) -> PortConn: def close(self) -> None: """Releases the sandbox; its VM is destroyed.""" + self._pool.drain() try: self._client._request( self.owner, "POST", f"/v1/sandboxes/{self.id}/release", None, "release", bearer=self.token @@ -358,8 +365,44 @@ def close(self) -> None: def _dial(self, deadline: float | None = None) -> Conn: return self._client._dial(self.owner, self.id, self.token, deadline) + def _connect(self, deadline: float | None = None) -> Conn: + """Takes the parked connection or dials; the first dial asks the daemon's proto and redials one that closed.""" + conn = self._pool.take() + if conn is not None: + return conn + conn = self._dial(deadline) + if self._pool.proto: + return conn + try: + conn.send("info") + proto = int(_expect(conn, "info").get("proto") or 1) + except BaseException: + conn.close() + raise + self._pool.proto = proto + if proto >= KEEP_ALIVE_PROTO: + return conn + conn.close() + return self._dial(deadline) + + @contextlib.contextmanager + def _park_after(self, conn: Conn) -> Iterator[Conn]: + """Runs one RPC on conn: a terminal frame, error frames included, parks it; anything else drops it.""" + try: + yield conn + except SilkdError: + self._pool.park(conn) + raise + except BaseException: + conn.close() + raise + self._pool.park(conn) + + def _lease(self) -> AbstractContextManager[Conn]: + return self._park_after(self._connect()) + def _open_stream(self, op: str, expect: str = "ready", **fields: object) -> tuple[Conn, dict[str, Any]]: - conn = self._dial() + conn = self._connect() try: conn.send(op, **fields) frame = _expect(conn, expect) @@ -409,7 +452,7 @@ def pump_out() -> None: local.close() def _call(self, op: str, expect: str, **fields: object) -> dict[str, Any]: - with self._dial() as conn: + with self._lease() as conn: conn.send(op, **fields) return _expect(conn, expect) @@ -423,9 +466,12 @@ def _drain_proc( on_stdout: Callable[[bytes], object] | None, on_stderr: Callable[[bytes], object] | None, ) -> int | None: - with self._dial() as conn: + with self._lease() as conn: conn.send(op, pid=pid) - return _pump_stdio(conn, on_stdout, on_stderr) + code = _pump_stdio(conn, on_stdout, on_stderr) + if op == "logs" and code is not None: + _expect(conn, "done") # logs closes with done after the exit frame + return code class Session(_Closeable): diff --git a/sdk/python/tests/test_keepalive.py b/sdk/python/tests/test_keepalive.py new file mode 100644 index 00000000..181f0fae --- /dev/null +++ b/sdk/python/tests/test_keepalive.py @@ -0,0 +1,157 @@ +"""One relay connection serves a handle's calls back to back; an old daemon, +a zero window, an idle window and a peer that hung up each fall back to a +fresh dial.""" + +import base64 +import contextlib +import json +import socket +import threading +import time + +from cocoonsandbox import Client, Sandbox + +INPUT_OPS = ("stdin", "stdin_close", "data", "data_end") + + +class FakeAgent: + """A relay-side silkd stand-in answering info, fs_stat and exec echo; proto 1 hangs up after one RPC.""" + + def __init__(self, proto: int = 2, hang_up_after: int = 0) -> None: + self.proto = proto + self.hang_up_after = hang_up_after + self.upgrades = 0 + self.hangups = 0 + self.closed = 0 + self._server = socket.create_server(("127.0.0.1", 0)) + self.addr = f"127.0.0.1:{self._server.getsockname()[1]}" + threading.Thread(target=self._accept, daemon=True).start() + + def stop(self) -> None: + self._server.close() + + def _accept(self) -> None: + with contextlib.suppress(OSError): + while True: + conn, _ = self._server.accept() + self.upgrades += 1 + threading.Thread(target=self._serve, args=(conn,), daemon=True).start() + + def _serve(self, conn: socket.socket) -> None: + reader = conn.makefile("rb") + with conn, contextlib.suppress(OSError): + while reader.readline() not in (b"\r\n", b""): + pass + conn.sendall(b"HTTP/1.1 101 Switching Protocols\r\n\r\n") + replies = 0 + while True: + line = reader.readline() + if not line: + break + op = json.loads(line)["op"] + if op in INPUT_OPS: + continue + for frame in self._answer(op): + conn.sendall(json.dumps(frame).encode() + b"\n") + replies += 1 + if self.proto < 2 or replies == self.hang_up_after: + conn.shutdown(socket.SHUT_WR) + self.hangups += 1 + while reader.readline(): + pass + break + self.closed += 1 + + def _answer(self, op: str) -> list: + if op == "info": + return [{"type": "info", "version": "fake", "proto": self.proto, "uptime_secs": 0, "procs": 0}] + if op == "fs_stat": + return [{"type": "stat", "info": {"kind": "dir", "size": 0, "mode": 0o755, "mtime_epoch_secs": 0}}] + if op == "exec": + out = base64.b64encode(b"42\n").decode() + return [{"type": "started", "pid": 1}, {"type": "stdout", "data": out}, {"type": "exit", "code": 0}] + return [{"type": "error", "kind": "unimplemented", "message": op}] + + +def handle(agent: FakeAgent, **client_kwargs) -> Sandbox: + return Sandbox(client=Client(agent.addr, **client_kwargs), id="sb_1", token="tok", owner=agent.addr) + + +def wait_until(cond, message: str) -> None: + deadline = time.monotonic() + 3 + while time.monotonic() < deadline: + if cond(): + return + time.sleep(0.01) + raise AssertionError(message) + + +def test_calls_share_one_connection(): + agent = FakeAgent() + sb = handle(agent) + try: + for _ in range(3): + assert sb.stat("/")["kind"] == "dir" + assert sb.exec("echo", "42") == "42\n" + finally: + sb._pool.drain() + agent.stop() + assert agent.upgrades == 1 + + +def test_old_daemon_dials_per_call(): + agent = FakeAgent(proto=1) + sb = handle(agent) + try: + for _ in range(3): + assert sb.exec("echo", "42") == "42\n" + finally: + agent.stop() + assert agent.upgrades == 4, "the proto probe plus one dial per call" + + +def test_keep_alive_off_dials_per_call(): + agent = FakeAgent() + sb = handle(agent, keep_alive=0) + try: + for _ in range(3): + sb.stat("/") + finally: + agent.stop() + assert agent.upgrades == 3 + + +def test_idle_connection_closes(): + agent = FakeAgent() + sb = handle(agent, keep_alive=0.05) + try: + sb.stat("/") + wait_until(lambda: agent.closed == 1, "idle connection still open after the keep-alive window") + finally: + sb._pool.drain() + agent.stop() + + +def test_peer_hang_up_is_noticed_before_reuse(): + agent = FakeAgent(hang_up_after=2) + sb = handle(agent) + try: + sb.stat("/") + wait_until(lambda: agent.hangups == 1, "the fake never hung up") + assert sb.stat("/")["kind"] == "dir" + finally: + sb._pool.drain() + agent.stop() + assert agent.upgrades == 2 + + +def test_close_drains_the_parked_connection(monkeypatch): + agent = FakeAgent() + sb = handle(agent) + monkeypatch.setattr(sb._client, "_request", lambda *args, **kwargs: {}) + try: + sb.stat("/") + sb.close() + wait_until(lambda: agent.closed == 1, "parked connection survived close") + finally: + agent.stop() diff --git a/sdk/python/tests/test_proc.py b/sdk/python/tests/test_proc.py index d03e22dd..b60f1aa5 100644 --- a/sdk/python/tests/test_proc.py +++ b/sdk/python/tests/test_proc.py @@ -44,6 +44,7 @@ def test_attach_returns_exit_code(monkeypatch): def test_run_pumps_stdin_while_reading_output(monkeypatch): sb = Sandbox(client=Client("127.0.0.1:1"), id="sb_1", token="tok", owner="127.0.0.1:1") + sb._pool.proto = 1 blocking = BlockingStdinConn([{"type": "exit", "code": 0}], buffer_frames=1) monkeypatch.setattr(sb, "_dial", lambda deadline=None: blocking) @@ -104,6 +105,7 @@ def recv(self): def fake_sandbox(monkeypatch, frames): sb = Sandbox(client=Client("127.0.0.1:1"), id="sb_1", token="tok", owner="127.0.0.1:1") + sb._pool.proto = 1 conn = FakeConn(frames) monkeypatch.setattr(sb, "_dial", lambda deadline=None: conn) return sb, conn diff --git a/sdk/python/tests/test_stream.py b/sdk/python/tests/test_stream.py index 15a3858f..2f33e5c0 100644 --- a/sdk/python/tests/test_stream.py +++ b/sdk/python/tests/test_stream.py @@ -30,6 +30,15 @@ def send(self, op: str, **fields) -> None: def abort(self) -> None: self.aborted.set() + def close(self) -> None: + pass + + +def legacy_sandbox(addr: str) -> Sandbox: + sb = Sandbox(client=Client(addr, timeout=TIMEOUT), id="sb_1", token="tok", owner=addr) + sb._pool.proto = 1 + return sb + def serve_port_forward(server: socket.socket, quiet: float, ops: list[str]) -> None: conn, _ = server.accept() @@ -68,7 +77,7 @@ def test_port_stream_outlives_the_client_timeout(): addr = f"127.0.0.1:{server.getsockname()[1]}" ops: list[str] = [] threading.Thread(target=serve_port_forward, args=(server, 3 * TIMEOUT, ops), daemon=True).start() - sb = Sandbox(client=Client(addr, timeout=TIMEOUT), id="sb_1", token="tok", owner=addr) + sb = legacy_sandbox(addr) try: with sb.dial_port(5000) as port: assert port.recv() == b"late" @@ -82,7 +91,7 @@ def test_run_timeout_cuts_a_silent_command(): server = socket.create_server(("127.0.0.1", 0)) addr = f"127.0.0.1:{server.getsockname()[1]}" threading.Thread(target=serve_started_then_hang, args=(server, 5 * TIMEOUT), daemon=True).start() - sb = Sandbox(client=Client(addr, timeout=TIMEOUT), id="sb_1", token="tok", owner=addr) + sb = legacy_sandbox(addr) started = time.monotonic() try: with pytest.raises(TimeoutError): @@ -93,7 +102,7 @@ def test_run_timeout_cuts_a_silent_command(): def test_run_timeout_cuts_a_blocked_exec_send(monkeypatch): - sb = Sandbox(client=Client("127.0.0.1:1"), id="sb_1", token="tok", owner="127.0.0.1:1") + sb = legacy_sandbox("127.0.0.1:1") conn = BlockedSendConn() monkeypatch.setattr(sb, "_dial", lambda deadline=None: conn) with pytest.raises(TimeoutError): @@ -111,7 +120,7 @@ def test_dial_is_still_bounded_by_the_client_timeout(): server = socket.create_server(("127.0.0.1", 0)) addr = f"127.0.0.1:{server.getsockname()[1]}" threading.Thread(target=serve_silence, args=(server,), daemon=True).start() - sb = Sandbox(client=Client(addr, timeout=TIMEOUT), id="sb_1", token="tok", owner=addr) + sb = legacy_sandbox(addr) started = time.monotonic() try: with pytest.raises(OSError): diff --git a/sdk/python/tests/test_wire_binding.py b/sdk/python/tests/test_wire_binding.py index 4327ec2a..5b128f0d 100644 --- a/sdk/python/tests/test_wire_binding.py +++ b/sdk/python/tests/test_wire_binding.py @@ -64,7 +64,7 @@ ("req_exec_detach", [{"type": "started", "pid": 7}], lambda sb, f: sb.spawn(*f["argv"])), ("req_ps", [{"type": "procs", "procs": []}], lambda sb, f: sb.ps()), ("req_kill", [{"type": "done"}], lambda sb, f: sb.kill(f["pid"], signal=f["signal"])), - ("req_logs", [{"type": "exit", "code": 0}], lambda sb, f: sb.logs(f["pid"])), + ("req_logs", [{"type": "exit", "code": 0}, {"type": "done"}], lambda sb, f: sb.logs(f["pid"])), ("req_attach", [{"type": "exit", "code": 0}], lambda sb, f: sb.attach(f["pid"])), ("req_fs_watch", [{"type": "ready"}], lambda sb, f: sb.watch(f["path"], recursive=f["recursive"]).close()), ("req_git_branch", [{"type": "done"}], lambda sb, f: sb.git_create_branch(f["path"], f["name"])), @@ -120,6 +120,7 @@ def test_enum_value_sets_match_corpus(): def test_git_branch_actions_come_from_the_corpus(monkeypatch): enums = json.loads((FIXTURES / "enums.json").read_text()) sb = Sandbox(client=Client("127.0.0.1:1"), id="sb_1", token="tok", owner="127.0.0.1:1") + sb._pool.proto = 1 sent = [] monkeypatch.setattr(sb, "_dial", lambda deadline=None: BranchActionConn(sent)) @@ -155,6 +156,9 @@ def recv(self): def recv_until(self, *terminal): yield self.recv() + def close(self): + pass + def fake_sandbox(monkeypatch, replies): """A Sandbox whose _dial yields a real Conn over a socketpair; a guest @@ -182,5 +186,6 @@ def guest(): thread.start() sb = Sandbox(client=Client("127.0.0.1:1"), id="sb_1", token="tok", owner="127.0.0.1:1") + sb._pool.proto = 1 monkeypatch.setattr(sb, "_dial", lambda deadline=None: Conn(client_sock, client_sock.makefile("rb"))) return sb, sent, thread From f7804a5313125a4f8e654afc622124dbd44d251e Mon Sep 17 00:00:00 2001 From: CMGS Date: Thu, 17 Sep 2026 00:22:10 +0800 Subject: [PATCH 2/2] sdk/python: sweep parked connections with one timer MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Every park armed a threading.Timer, a real OS thread, which measured at 42.6 µs, over half the cost of a kept-connection call, and left eight sleeping threads per handle; the pool now records a deadline per parked connection and one self-rearming sweeper closes them as they expire, so park and take touch no thread at all. run() sends stdin_close inline when there is no stdin instead of spawning the pump thread for one frame. The daemon's proto lives on the Sandbox, a zero keep_alive window skips the probe, _park_after and _lease fold into one context manager, _Parked goes, and logs' trailing done is a parameter rather than a string compare. The tests share sandbox_at, accept_upgrade and wait_until through conftest and pin the old one-shot behaviour with keep_alive=0 instead of a private field. --- sdk/python/cocoonsandbox/client.py | 2 +- sdk/python/cocoonsandbox/conn.py | 68 ++++++------ sdk/python/cocoonsandbox/sandbox.py | 50 +++++---- sdk/python/tests/conftest.py | 29 ++++- sdk/python/tests/test_keepalive.py | 149 ++++++++++++-------------- sdk/python/tests/test_proc.py | 9 +- sdk/python/tests/test_stream.py | 19 ++-- sdk/python/tests/test_wire_binding.py | 9 +- 8 files changed, 173 insertions(+), 162 deletions(-) diff --git a/sdk/python/cocoonsandbox/client.py b/sdk/python/cocoonsandbox/client.py index c3287a0e..85ad3be7 100644 --- a/sdk/python/cocoonsandbox/client.py +++ b/sdk/python/cocoonsandbox/client.py @@ -36,7 +36,7 @@ def __init__( ssl_context: ssl.SSLContext | None = None, keep_alive: float = 30.0, ) -> None: - """keep_alive bounds how long a handle keeps an idle relay connection for its next call; 0 dials per call.""" + """keep_alive keeps a handle's idle relay connection, which holds the sandbox's idle clock; 0 dials per call.""" endpoint = _endpoint_url(addr.split(",")[0].strip()) self.addr = endpoint.geturl().removeprefix("http://") self._scheme = endpoint.scheme diff --git a/sdk/python/cocoonsandbox/conn.py b/sdk/python/cocoonsandbox/conn.py index 2e000120..d1f21596 100644 --- a/sdk/python/cocoonsandbox/conn.py +++ b/sdk/python/cocoonsandbox/conn.py @@ -9,11 +9,11 @@ import threading import time import urllib.parse -from collections.abc import Callable, Iterator +from collections.abc import Iterator from typing import Any, BinaryIO, Protocol, TypeVar from .errors import APIError, ProtocolError, SilkdError -from .frames import KEEP_ALIVE_PROTO, MAX_FRAME, decode_response, encode_request +from .frames import MAX_FRAME, decode_response, encode_request KEEP_ALIVE_CONNS = 8 @@ -89,58 +89,58 @@ def close(self) -> None: self._sock.close() -class _Parked: - def __init__(self, conn: Conn, idle: float, evict: Callable[[_Parked], None]) -> None: - self.conn = conn - self.timer = threading.Timer(idle, evict, (self,)) - self.timer.daemon = True - - class ConnPool: - """Parks a handle's idle relay connections between calls; silkd serves RPCs back to back from proto 2.""" + """Parks a handle's idle relay connections between calls; one sweeper timer closes them as they expire.""" def __init__(self, idle: float) -> None: - self.proto = 0 self._idle = idle self._lock = threading.Lock() - self._parked: list[_Parked] = [] + self._parked: list[tuple[float, Conn]] = [] + self._sweep: threading.Timer | None = None def take(self) -> Conn | None: - """Returns a parked connection whose peer is still there, or None.""" while True: with self._lock: if not self._parked: return None - entry = self._parked.pop() - entry.timer.cancel() - if entry.conn.quiet(): - return entry.conn - entry.conn.close() + expires, conn = self._parked.pop() + if time.monotonic() < expires and conn.quiet(): + return conn + conn.close() def park(self, conn: Conn) -> None: - """Keeps conn for the next call until idle passes; a daemon before proto 2 or a zero window closes it.""" with self._lock: - keep = self._idle > 0 and self.proto >= KEEP_ALIVE_PROTO and len(self._parked) < KEEP_ALIVE_CONNS - if keep: - entry = _Parked(conn, self._idle, self._evict) - self._parked.append(entry) - entry.timer.start() - if not keep: - conn.close() + if self._idle > 0 and len(self._parked) < KEEP_ALIVE_CONNS: + self._parked.append((time.monotonic() + self._idle, conn)) + if self._sweep is None: + self._arm(self._idle) + return + conn.close() def drain(self) -> None: with self._lock: parked, self._parked = self._parked, [] - for entry in parked: - entry.timer.cancel() - entry.conn.close() + if self._sweep is not None: + self._sweep.cancel() + self._sweep = None + for _, conn in parked: + conn.close() + + def _arm(self, delay: float) -> None: + self._sweep = threading.Timer(delay, self._evict) + self._sweep.daemon = True + self._sweep.start() - def _evict(self, entry: _Parked) -> None: + def _evict(self) -> None: + now = time.monotonic() with self._lock: - if entry not in self._parked: - return - self._parked.remove(entry) - entry.conn.close() + expired = [conn for expires, conn in self._parked if expires <= now] + self._parked = [entry for entry in self._parked if entry[0] > now] + self._sweep = None + if self._parked: + self._arm(self._parked[0][0] - now) + for conn in expired: + conn.close() def dial_agent( diff --git a/sdk/python/cocoonsandbox/sandbox.py b/sdk/python/cocoonsandbox/sandbox.py index 6c412d0c..0d6e988c 100644 --- a/sdk/python/cocoonsandbox/sandbox.py +++ b/sdk/python/cocoonsandbox/sandbox.py @@ -7,7 +7,6 @@ import threading import time from collections.abc import Callable, Iterator -from contextlib import AbstractContextManager from typing import TYPE_CHECKING, Any, cast from .checkpoint import Checkpoint @@ -43,6 +42,7 @@ def __init__( self.template_digest = template_digest self.volumes = [dict(volume) for volume in volumes or []] self._pool = ConnPool(client.keep_alive) + self._proto = 0 def __enter__(self) -> Sandbox: return self @@ -104,9 +104,9 @@ def run( if deadline is not None and time.monotonic() >= deadline: raise TimeoutError(f"command did not finish within {timeout}s") from None raise - with self._park_after(conn): + with self._lease(conn): watchdog = _arm_watchdog(conn, deadline, expired) - pump = threading.Thread(target=_feed_stdin, args=(conn, stdin), daemon=True) + pump = threading.Thread(target=_feed_stdin, args=(conn, stdin), daemon=True) if stdin else None try: conn.send( "exec", @@ -117,19 +117,25 @@ def run( detach=False, session=session or None, ) - pump.start() + if pump is not None: + pump.start() + else: + with contextlib.suppress(OSError): + conn.send("stdin_close") code = _pump_stdio(conn, on_stdout, on_stderr) except (ProtocolError, OSError): if expired.is_set(): raise TimeoutError(f"command did not finish within {timeout}s") from None raise except SilkdError: - pump.join() + if pump is not None: + pump.join() raise finally: if watchdog is not None: watchdog.cancel() - pump.join() # the pump's frames must not land on the next RPC; silkd drains what the command left unread + if pump is not None: + pump.join() # the pump's frames must not land on the next RPC if code is None: raise ProtocolError("exec stream ended without an exit frame") return code @@ -156,7 +162,7 @@ def logs( on_stderr: Callable[[bytes], object] | None = None, ) -> int | None: """Replays buffered output and returns the exit code, or None if the process still runs.""" - return self._drain_proc("logs", pid, on_stdout, on_stderr) + return self._drain_proc("logs", pid, on_stdout, on_stderr, trailing_done=True) def attach( self, @@ -366,40 +372,45 @@ def _dial(self, deadline: float | None = None) -> Conn: return self._client._dial(self.owner, self.id, self.token, deadline) def _connect(self, deadline: float | None = None) -> Conn: - """Takes the parked connection or dials; the first dial asks the daemon's proto and redials one that closed.""" + """Takes a parked connection or dials one, asking the daemon's proto on a handle's first kept dial.""" conn = self._pool.take() if conn is not None: return conn conn = self._dial(deadline) - if self._pool.proto: + if self._proto or self._client.keep_alive <= 0: return conn try: conn.send("info") proto = int(_expect(conn, "info").get("proto") or 1) - except BaseException: + except Exception: conn.close() raise - self._pool.proto = proto + self._proto = proto if proto >= KEEP_ALIVE_PROTO: return conn conn.close() return self._dial(deadline) @contextlib.contextmanager - def _park_after(self, conn: Conn) -> Iterator[Conn]: - """Runs one RPC on conn: a terminal frame, error frames included, parks it; anything else drops it.""" + def _lease(self, conn: Conn | None = None) -> Iterator[Conn]: + """Runs one RPC on conn, dialed when absent; a terminal frame parks it and anything else drops it.""" + if conn is None: + conn = self._connect() try: yield conn except SilkdError: - self._pool.park(conn) + self._park(conn) raise except BaseException: conn.close() raise - self._pool.park(conn) + self._park(conn) - def _lease(self) -> AbstractContextManager[Conn]: - return self._park_after(self._connect()) + def _park(self, conn: Conn) -> None: + if self._proto >= KEEP_ALIVE_PROTO: + self._pool.park(conn) + else: + conn.close() def _open_stream(self, op: str, expect: str = "ready", **fields: object) -> tuple[Conn, dict[str, Any]]: conn = self._connect() @@ -465,12 +476,13 @@ def _drain_proc( pid: int, on_stdout: Callable[[bytes], object] | None, on_stderr: Callable[[bytes], object] | None, + trailing_done: bool = False, ) -> int | None: with self._lease() as conn: conn.send(op, pid=pid) code = _pump_stdio(conn, on_stdout, on_stderr) - if op == "logs" and code is not None: - _expect(conn, "done") # logs closes with done after the exit frame + if trailing_done and code is not None: + _expect(conn, "done") return code diff --git a/sdk/python/tests/conftest.py b/sdk/python/tests/conftest.py index 8155508a..6101223b 100644 --- a/sdk/python/tests/conftest.py +++ b/sdk/python/tests/conftest.py @@ -1,13 +1,38 @@ -"""Fixtures shared by the cluster-behavior suites: in-process fake nodes and -addresses that refuse a connection.""" +"""Fixtures and helpers shared by the suites: in-process fake nodes, addresses +that refuse a connection, and the relay-side upgrade handshake.""" import socket import threading +import time from http.server import HTTPServer +from typing import BinaryIO, Callable import pytest from test_client import FakeNode +from cocoonsandbox import Client, Sandbox + + +def sandbox_at(addr: str, **client_kwargs) -> Sandbox: + return Sandbox(client=Client(addr, **client_kwargs), id="sb_1", token="tok", owner=addr) + + +def accept_upgrade(conn: socket.socket) -> BinaryIO: + reader = conn.makefile("rb") + while reader.readline() not in (b"\r\n", b""): + pass + conn.sendall(b"HTTP/1.1 101 Switching Protocols\r\n\r\n") + return reader + + +def wait_until(cond: Callable[[], bool], message: str) -> None: + deadline = time.monotonic() + 3 + while time.monotonic() < deadline: + if cond(): + return + time.sleep(0.01) + raise AssertionError(message) + @pytest.fixture def spawn_node(): diff --git a/sdk/python/tests/test_keepalive.py b/sdk/python/tests/test_keepalive.py index 181f0fae..8710ae4f 100644 --- a/sdk/python/tests/test_keepalive.py +++ b/sdk/python/tests/test_keepalive.py @@ -1,94 +1,21 @@ -"""One relay connection serves a handle's calls back to back; an old daemon, -a zero window, an idle window and a peer that hung up each fall back to a -fresh dial.""" +"""One relay connection serves a handle's calls back to back; the fallbacks dial afresh.""" import base64 import contextlib import json import socket import threading -import time -from cocoonsandbox import Client, Sandbox +from conftest import accept_upgrade, sandbox_at, wait_until -INPUT_OPS = ("stdin", "stdin_close", "data", "data_end") - - -class FakeAgent: - """A relay-side silkd stand-in answering info, fs_stat and exec echo; proto 1 hangs up after one RPC.""" - - def __init__(self, proto: int = 2, hang_up_after: int = 0) -> None: - self.proto = proto - self.hang_up_after = hang_up_after - self.upgrades = 0 - self.hangups = 0 - self.closed = 0 - self._server = socket.create_server(("127.0.0.1", 0)) - self.addr = f"127.0.0.1:{self._server.getsockname()[1]}" - threading.Thread(target=self._accept, daemon=True).start() - - def stop(self) -> None: - self._server.close() - - def _accept(self) -> None: - with contextlib.suppress(OSError): - while True: - conn, _ = self._server.accept() - self.upgrades += 1 - threading.Thread(target=self._serve, args=(conn,), daemon=True).start() - - def _serve(self, conn: socket.socket) -> None: - reader = conn.makefile("rb") - with conn, contextlib.suppress(OSError): - while reader.readline() not in (b"\r\n", b""): - pass - conn.sendall(b"HTTP/1.1 101 Switching Protocols\r\n\r\n") - replies = 0 - while True: - line = reader.readline() - if not line: - break - op = json.loads(line)["op"] - if op in INPUT_OPS: - continue - for frame in self._answer(op): - conn.sendall(json.dumps(frame).encode() + b"\n") - replies += 1 - if self.proto < 2 or replies == self.hang_up_after: - conn.shutdown(socket.SHUT_WR) - self.hangups += 1 - while reader.readline(): - pass - break - self.closed += 1 - - def _answer(self, op: str) -> list: - if op == "info": - return [{"type": "info", "version": "fake", "proto": self.proto, "uptime_secs": 0, "procs": 0}] - if op == "fs_stat": - return [{"type": "stat", "info": {"kind": "dir", "size": 0, "mode": 0o755, "mtime_epoch_secs": 0}}] - if op == "exec": - out = base64.b64encode(b"42\n").decode() - return [{"type": "started", "pid": 1}, {"type": "stdout", "data": out}, {"type": "exit", "code": 0}] - return [{"type": "error", "kind": "unimplemented", "message": op}] - - -def handle(agent: FakeAgent, **client_kwargs) -> Sandbox: - return Sandbox(client=Client(agent.addr, **client_kwargs), id="sb_1", token="tok", owner=agent.addr) +from cocoonsandbox.frames import KEEP_ALIVE_PROTO - -def wait_until(cond, message: str) -> None: - deadline = time.monotonic() + 3 - while time.monotonic() < deadline: - if cond(): - return - time.sleep(0.01) - raise AssertionError(message) +INPUT_OPS = ("stdin", "stdin_close", "data", "data_end") def test_calls_share_one_connection(): agent = FakeAgent() - sb = handle(agent) + sb = sandbox_at(agent.addr) try: for _ in range(3): assert sb.stat("/")["kind"] == "dir" @@ -101,7 +28,7 @@ def test_calls_share_one_connection(): def test_old_daemon_dials_per_call(): agent = FakeAgent(proto=1) - sb = handle(agent) + sb = sandbox_at(agent.addr) try: for _ in range(3): assert sb.exec("echo", "42") == "42\n" @@ -112,7 +39,7 @@ def test_old_daemon_dials_per_call(): def test_keep_alive_off_dials_per_call(): agent = FakeAgent() - sb = handle(agent, keep_alive=0) + sb = sandbox_at(agent.addr, keep_alive=0) try: for _ in range(3): sb.stat("/") @@ -123,7 +50,7 @@ def test_keep_alive_off_dials_per_call(): def test_idle_connection_closes(): agent = FakeAgent() - sb = handle(agent, keep_alive=0.05) + sb = sandbox_at(agent.addr, keep_alive=0.05) try: sb.stat("/") wait_until(lambda: agent.closed == 1, "idle connection still open after the keep-alive window") @@ -134,7 +61,7 @@ def test_idle_connection_closes(): def test_peer_hang_up_is_noticed_before_reuse(): agent = FakeAgent(hang_up_after=2) - sb = handle(agent) + sb = sandbox_at(agent.addr) try: sb.stat("/") wait_until(lambda: agent.hangups == 1, "the fake never hung up") @@ -147,7 +74,7 @@ def test_peer_hang_up_is_noticed_before_reuse(): def test_close_drains_the_parked_connection(monkeypatch): agent = FakeAgent() - sb = handle(agent) + sb = sandbox_at(agent.addr) monkeypatch.setattr(sb._client, "_request", lambda *args, **kwargs: {}) try: sb.stat("/") @@ -155,3 +82,59 @@ def test_close_drains_the_parked_connection(monkeypatch): wait_until(lambda: agent.closed == 1, "parked connection survived close") finally: agent.stop() + + +class FakeAgent: + """A relay-side silkd stand-in; proto 1 half-closes after one RPC and drains like the relay does.""" + + def __init__(self, proto: int = KEEP_ALIVE_PROTO, hang_up_after: int = 0) -> None: + self.proto = proto + self.hang_up_after = hang_up_after + self.upgrades = 0 + self.hangups = 0 + self.closed = 0 + self._server = socket.create_server(("127.0.0.1", 0)) + self.addr = f"127.0.0.1:{self._server.getsockname()[1]}" + threading.Thread(target=self._accept, daemon=True).start() + + def stop(self) -> None: + self._server.close() + + def _accept(self) -> None: + with contextlib.suppress(OSError): + while True: + conn, _ = self._server.accept() + self.upgrades += 1 + threading.Thread(target=self._serve, args=(conn,), daemon=True).start() + + def _serve(self, conn: socket.socket) -> None: + with conn, contextlib.suppress(OSError): + reader = accept_upgrade(conn) + replies = 0 + while True: + line = reader.readline() + if not line: + break + op = json.loads(line)["op"] + if op in INPUT_OPS: + continue + for frame in self._answer(op): + conn.sendall(json.dumps(frame).encode() + b"\n") + replies += 1 + if self.proto < KEEP_ALIVE_PROTO or replies == self.hang_up_after: + conn.shutdown(socket.SHUT_WR) + self.hangups += 1 + while reader.readline(): + pass + break + self.closed += 1 + + def _answer(self, op: str) -> list: + if op == "info": + return [{"type": "info", "version": "fake", "proto": self.proto, "uptime_secs": 0, "procs": 0}] + if op == "fs_stat": + return [{"type": "stat", "info": {"kind": "dir", "size": 0, "mode": 0o755, "mtime_epoch_secs": 0}}] + if op == "exec": + out = base64.b64encode(b"42\n").decode() + return [{"type": "started", "pid": 1}, {"type": "stdout", "data": out}, {"type": "exit", "code": 0}] + return [{"type": "error", "kind": "unimplemented", "message": op}] diff --git a/sdk/python/tests/test_proc.py b/sdk/python/tests/test_proc.py index b60f1aa5..440d99b9 100644 --- a/sdk/python/tests/test_proc.py +++ b/sdk/python/tests/test_proc.py @@ -3,7 +3,8 @@ import threading -from cocoonsandbox import Client, Sandbox +from conftest import sandbox_at + from cocoonsandbox.frames import FS_CHUNK @@ -43,8 +44,7 @@ def test_attach_returns_exit_code(monkeypatch): def test_run_pumps_stdin_while_reading_output(monkeypatch): - sb = Sandbox(client=Client("127.0.0.1:1"), id="sb_1", token="tok", owner="127.0.0.1:1") - sb._pool.proto = 1 + sb = sandbox_at("127.0.0.1:1", keep_alive=0) blocking = BlockingStdinConn([{"type": "exit", "code": 0}], buffer_frames=1) monkeypatch.setattr(sb, "_dial", lambda deadline=None: blocking) @@ -104,8 +104,7 @@ def recv(self): def fake_sandbox(monkeypatch, frames): - sb = Sandbox(client=Client("127.0.0.1:1"), id="sb_1", token="tok", owner="127.0.0.1:1") - sb._pool.proto = 1 + sb = sandbox_at("127.0.0.1:1", keep_alive=0) conn = FakeConn(frames) monkeypatch.setattr(sb, "_dial", lambda deadline=None: conn) return sb, conn diff --git a/sdk/python/tests/test_stream.py b/sdk/python/tests/test_stream.py index 2f33e5c0..fb9e966c 100644 --- a/sdk/python/tests/test_stream.py +++ b/sdk/python/tests/test_stream.py @@ -7,8 +7,9 @@ import time import pytest +from conftest import accept_upgrade, sandbox_at -from cocoonsandbox import Client, Sandbox +from cocoonsandbox import Sandbox TIMEOUT = 0.2 @@ -35,17 +36,12 @@ def close(self) -> None: def legacy_sandbox(addr: str) -> Sandbox: - sb = Sandbox(client=Client(addr, timeout=TIMEOUT), id="sb_1", token="tok", owner=addr) - sb._pool.proto = 1 - return sb + return sandbox_at(addr, timeout=TIMEOUT, keep_alive=0) def serve_port_forward(server: socket.socket, quiet: float, ops: list[str]) -> None: conn, _ = server.accept() - reader = conn.makefile("rb") - while reader.readline() not in (b"\r\n", b""): - pass - conn.sendall(b"HTTP/1.1 101 Switching Protocols\r\n\r\n") + reader = accept_upgrade(conn) ops.append(json.loads(reader.readline())["op"]) conn.sendall(b'{"type":"ready"}\n') time.sleep(quiet) @@ -55,10 +51,7 @@ def serve_port_forward(server: socket.socket, quiet: float, ops: list[str]) -> N def serve_started_then_hang(server: socket.socket, quiet: float) -> None: conn, _ = server.accept() - reader = conn.makefile("rb") - while reader.readline() not in (b"\r\n", b""): - pass - conn.sendall(b"HTTP/1.1 101 Switching Protocols\r\n\r\n") + reader = accept_upgrade(conn) reader.readline() conn.sendall(b'{"type":"started","pid":7}\n') time.sleep(quiet) @@ -111,7 +104,7 @@ def test_run_timeout_cuts_a_blocked_exec_send(monkeypatch): def test_run_rejects_a_non_positive_timeout(): - sb = Sandbox(client=Client("127.0.0.1:1", timeout=TIMEOUT), id="sb_1", token="tok", owner="127.0.0.1:1") + sb = legacy_sandbox("127.0.0.1:1") with pytest.raises(ValueError): sb.run(["true"], timeout=0) diff --git a/sdk/python/tests/test_wire_binding.py b/sdk/python/tests/test_wire_binding.py index 5b128f0d..f6cecac4 100644 --- a/sdk/python/tests/test_wire_binding.py +++ b/sdk/python/tests/test_wire_binding.py @@ -16,8 +16,9 @@ import threading import pytest +from conftest import sandbox_at -from cocoonsandbox import Client, Lsp, Pty, Sandbox, Session +from cocoonsandbox import Lsp, Pty, Session from cocoonsandbox.conn import Conn from cocoonsandbox.frames import PROTO_VERSION @@ -119,8 +120,7 @@ def test_enum_value_sets_match_corpus(): def test_git_branch_actions_come_from_the_corpus(monkeypatch): enums = json.loads((FIXTURES / "enums.json").read_text()) - sb = Sandbox(client=Client("127.0.0.1:1"), id="sb_1", token="tok", owner="127.0.0.1:1") - sb._pool.proto = 1 + sb = sandbox_at("127.0.0.1:1", keep_alive=0) sent = [] monkeypatch.setattr(sb, "_dial", lambda deadline=None: BranchActionConn(sent)) @@ -185,7 +185,6 @@ def guest(): thread = threading.Thread(target=guest, daemon=True) thread.start() - sb = Sandbox(client=Client("127.0.0.1:1"), id="sb_1", token="tok", owner="127.0.0.1:1") - sb._pool.proto = 1 + sb = sandbox_at("127.0.0.1:1", keep_alive=0) monkeypatch.setattr(sb, "_dial", lambda deadline=None: Conn(client_sock, client_sock.makefile("rb"))) return sb, sent, thread