From 3b405f13ed697ede64f19d47459f47dc66e59005 Mon Sep 17 00:00:00 2001 From: CMGS Date: Thu, 3 Sep 2026 15:22:29 +0800 Subject: [PATCH 1/2] review: GitFileStatus ahead of the result that carries it Vocabulary types cluster ahead of the type that consumes them; pure move, sorted lines identical. --- protocol/wire/frame.go | 14 +++++++------- 1 file changed, 7 insertions(+), 7 deletions(-) diff --git a/protocol/wire/frame.go b/protocol/wire/frame.go index a7e0cc9b..54345a79 100644 --- a/protocol/wire/frame.go +++ b/protocol/wire/frame.go @@ -617,6 +617,13 @@ type Event struct { func (Event) RespType() string { return "event" } +// GitFileStatus is one porcelain-v2 entry; Staged/Unstaged are XY status codes. +type GitFileStatus struct { + Path string `json:"path"` + Staged string `json:"staged"` + Unstaged string `json:"unstaged"` +} + // GitStatusResult answers GitStatus. type GitStatusResult struct { Branch string `json:"branch"` @@ -628,13 +635,6 @@ type GitStatusResult struct { func (GitStatusResult) RespType() string { return "git_status_result" } -// GitFileStatus is one porcelain-v2 entry; Staged/Unstaged are XY status codes. -type GitFileStatus struct { - Path string `json:"path"` - Staged string `json:"staged"` - Unstaged string `json:"unstaged"` -} - // GitCommitResult answers GitCommit with the new commit hash. type GitCommitResult struct { Hash string `json:"hash"` From 9c3a15d0e5d8211252cce8be46d0863f6c5d4e0e Mon Sep 17 00:00:00 2001 From: CMGS Date: Thu, 3 Sep 2026 15:22:29 +0800 Subject: [PATCH 2/2] review(python): ruff format --- sdk/python/cocoonsandbox/client.py | 56 ++++++--- sdk/python/cocoonsandbox/sandbox.py | 118 +++++++++++------- sdk/python/cocoonsandbox/template.py | 10 +- sdk/python/e2e.py | 3 +- sdk/python/tests/test_client.py | 169 +++++++++++++++++--------- sdk/python/tests/test_fault_matrix.py | 30 +++-- sdk/python/tests/test_proc.py | 7 +- sdk/python/tests/test_wire_binding.py | 126 +++++++++---------- 8 files changed, 315 insertions(+), 204 deletions(-) diff --git a/sdk/python/cocoonsandbox/client.py b/sdk/python/cocoonsandbox/client.py index 7ee377bd..435dad58 100644 --- a/sdk/python/cocoonsandbox/client.py +++ b/sdk/python/cocoonsandbox/client.py @@ -28,8 +28,16 @@ def __init__(self, addr: str, api_token: str = "", timeout: float = 120.0): self.api_token = api_token self.timeout = timeout - def new(self, template: str, net: str = "", size: str = "", ttl_seconds: int = 0, claim_ref: str = "", - volumes: list[str | Mapping[str, str]] | None = None, mount: bool = True) -> Sandbox: + def new( + self, + template: str, + net: str = "", + size: str = "", + ttl_seconds: int = 0, + claim_ref: str = "", + volumes: list[str | Mapping[str, str]] | None = None, + mount: bool = True, + ) -> Sandbox: """Claims a sandbox; a warm hit is milliseconds. On a cluster a warm miss may redirect to a peer, followed transparently; if every candidate fails transiently, the claim falls back to the origin @@ -58,12 +66,19 @@ def lookup(self, id: str, token: str) -> Sandbox: """Relocates a handle from id + token: asks the entry node and every mesh peer concurrently, binding to whichever confirms ownership first — one dead peer must not cost its full timeout.""" + def probe(addr: str) -> Sandbox: # bounded like _peers: a scatter loser must not hold a socket for the full timeout. - reply = self._request(addr, "GET", f"/v1/sandboxes/{id}/owner", None, "owner", - bearer=token, timeout=min(_PEERS_TIMEOUT, self.timeout)) - return Sandbox(client=self, id=id, token=token, - owner=reply.get("owner_addr") or addr) + reply = self._request( + addr, + "GET", + f"/v1/sandboxes/{id}/owner", + None, + "owner", + bearer=token, + timeout=min(_PEERS_TIMEOUT, self.timeout), + ) + return Sandbox(client=self, id=id, token=token, owner=reply.get("owner_addr") or addr) addrs = [self.addr, *self._peers()] try: @@ -128,8 +143,12 @@ def post(peer): def _peers(self) -> list: # /v1/peers is tenant-accessible; /v1/info is operator-only, so a tenant cannot read peers from it. try: - return self._request(self.addr, "GET", "/v1/peers", None, "peers", - timeout=min(_PEERS_TIMEOUT, self.timeout)).get("peers") or [] + return ( + self._request( + self.addr, "GET", "/v1/peers", None, "peers", timeout=min(_PEERS_TIMEOUT, self.timeout) + ).get("peers") + or [] + ) except APIError: return [] @@ -148,8 +167,9 @@ def _handle_from(self, dialed: str, reply: dict) -> Sandbox: def _post_json(self, addr: str, path: str, body: dict, verb: str) -> dict: return self._request(addr, "POST", path, body, verb) - def _request(self, addr: str, method: str, path: str, body, verb: str, bearer: str = "", - timeout: float = 0.0) -> dict: + def _request( + self, addr: str, method: str, path: str, body, verb: str, bearer: str = "", timeout: float = 0.0 + ) -> dict: """Issues one control-plane request. bearer overrides the api token — sandbox-scoped verbs (release, hibernate) authenticate with the per-sandbox token instead; timeout overrides the client default (0 keeps it).""" @@ -182,9 +202,15 @@ def _request(self, addr: str, method: str, path: str, body, verb: str, bearer: s raise APIError(verb, 0, "malformed JSON in response") from exc -def _claim_body(template: str, net: str, size: str, ttl_seconds: int, - volumes: list[str | Mapping[str, str]] | None = None, mount: bool = True, - claim_ref: str = "") -> dict: +def _claim_body( + template: str, + net: str, + size: str, + ttl_seconds: int, + volumes: list[str | Mapping[str, str]] | None = None, + mount: bool = True, + claim_ref: str = "", +) -> dict: claim = {"template": template} if net: claim["net"] = net @@ -209,7 +235,8 @@ def _volume_body(volume: str | Mapping[str, str], mount: bool) -> dict: unknown = sorted(set(volume) - {"name", "mount", "mode"}) if unknown: raise TypeError( - f"volume mapping accepts only name, mount, and mode, got unexpected key(s): {', '.join(unknown)}") + f"volume mapping accepts only name, mount, and mode, got unexpected key(s): {', '.join(unknown)}" + ) if not mount and "mount" in volume: raise TypeError("volume mount is meaningless with mount=False, which leaves mounting to the caller") body = dict(volume) @@ -277,6 +304,7 @@ def _redirect_fallback(origin: str, candidates: list, post, verb: str): would fail the same way. A second-level redirect (a compliant server never sends one once no_redirect is set) fails the candidate rather than being followed. Returns (addr, reply).""" + def attempt(addr): reply = post(addr) if reply.get("redirect"): diff --git a/sdk/python/cocoonsandbox/sandbox.py b/sdk/python/cocoonsandbox/sandbox.py index fe04a921..997da7ad 100644 --- a/sdk/python/cocoonsandbox/sandbox.py +++ b/sdk/python/cocoonsandbox/sandbox.py @@ -24,9 +24,17 @@ class Sandbox: """One claimed microVM.""" - def __init__(self, client: Client, id: str, token: str, owner: str, - deadline: str = "", from_checkpoint: str = "", template_digest: str = "", - volumes: list[dict] | None = None): + def __init__( + self, + client: Client, + id: str, + token: str, + owner: str, + deadline: str = "", + from_checkpoint: str = "", + template_digest: str = "", + volumes: list[dict] | None = None, + ): self._client = client self.id = id self.token = token @@ -46,26 +54,43 @@ def __exit__(self, *exc) -> None: if exc[0] is None: # a clean block surfaces a real release failure raise - def exec(self, *argv: str, cwd: str = "", env: dict | None = None, - user: str = "", session: str = "", stdin: bytes = b"") -> str: + def exec( + self, *argv: str, cwd: str = "", env: dict | None = None, user: str = "", session: str = "", stdin: bytes = b"" + ) -> str: """Runs argv to completion and returns stdout; a non-zero exit raises ExitError carrying stderr.""" out, err = bytearray(), bytearray() - code = self.run(list(argv), cwd=cwd, env=env, user=user, session=session, - stdin=stdin, on_stdout=out.extend, on_stderr=err.extend) + code = self.run( + list(argv), + cwd=cwd, + env=env, + user=user, + session=session, + stdin=stdin, + on_stdout=out.extend, + on_stderr=err.extend, + ) if code != 0: raise ExitError(code, err.decode(errors="replace"), out.decode(errors="replace")) return out.decode(errors="replace") - def run(self, argv: list[str], cwd: str = "", env: dict | None = None, - user: str = "", session: str = "", stdin: bytes = b"", - on_stdout: Callable[[bytes], object] | None = None, - on_stderr: Callable[[bytes], object] | None = None) -> int: + def run( + self, + argv: list[str], + cwd: str = "", + env: dict | None = None, + user: str = "", + session: str = "", + stdin: bytes = b"", + on_stdout: Callable[[bytes], object] | None = None, + on_stderr: Callable[[bytes], object] | None = None, + ) -> int: """Runs argv streaming stdio through the callbacks (raw bytes — chunk boundaries may split multi-byte sequences); returns the exit code.""" with self._dial() as conn: - conn.send("exec", argv=argv, cwd=cwd or None, env=env, - user=user or None, detach=False, session=session or None) + conn.send( + "exec", argv=argv, cwd=cwd or None, env=env, user=user or None, detach=False, session=session or None + ) # the guest stops draining stdin while blocked on stdout, so feeding it fully first deadlocks. pump = threading.Thread(target=_feed_stdin, args=(conn, stdin), daemon=True) pump.start() @@ -75,12 +100,12 @@ def run(self, argv: list[str], cwd: str = "", env: dict | None = None, raise ProtocolError("exec stream ended without an exit frame") return code - def spawn(self, *argv: str, cwd: str = "", env: dict | None = None, - user: str = "") -> int: + def spawn(self, *argv: str, cwd: str = "", env: dict | None = None, user: str = "") -> int: """Starts argv detached, returning its pid immediately; the process keeps a bounded output ring readable later via logs()/attach().""" - started = self._call("exec", "started", argv=list(argv), cwd=cwd or None, - env=env, user=user or None, detach=True) + started = self._call( + "exec", "started", argv=list(argv), cwd=cwd or None, env=env, user=user or None, detach=True + ) return started["pid"] def ps(self) -> list[dict]: @@ -93,14 +118,22 @@ def kill(self, pid: int, signal: int | None = None) -> None: already exited is a no-op success.""" self._done_rpc("kill", pid=pid, signal=signal or None) - def logs(self, pid: int, on_stdout: Callable[[bytes], object] | None = None, - on_stderr: Callable[[bytes], object] | None = None) -> int | None: + def logs( + self, + pid: int, + on_stdout: Callable[[bytes], object] | None = None, + on_stderr: Callable[[bytes], object] | None = None, + ) -> int | None: """Replays a process's ring-buffered output through the callbacks; returns its exit code if it already exited, else None.""" return self._drain_proc("logs", pid, on_stdout, on_stderr) - def attach(self, pid: int, on_stdout: Callable[[bytes], object] | None = None, - on_stderr: Callable[[bytes], object] | None = None) -> int | None: + def attach( + self, + pid: int, + on_stdout: Callable[[bytes], object] | None = None, + on_stderr: Callable[[bytes], object] | None = None, + ) -> int | None: """Replays buffered output then follows live output until the process exits, returning its exit code (None only if the proc table dropped it mid-attach).""" @@ -176,8 +209,7 @@ def replace(self, files: list[str], pattern: str, replacement: str) -> list[dict def git_clone(self, url: str, path: str, branch: str = "", depth: int = 0, auth: str = "") -> None: """Clones into path (egress lane only; the none lane answers a typed unimplemented error pointing at push).""" - self._done_rpc("git_clone", url=url, path=path, branch=branch or None, - depth=depth or None, auth=auth or None) + self._done_rpc("git_clone", url=url, path=path, branch=branch or None, depth=depth or None, auth=auth or None) def git_status(self, path: str) -> dict: return self._call("git_status", "git_status_result", path=path) @@ -187,8 +219,7 @@ def git_add(self, path: str, files: list[str]) -> None: def git_commit(self, path: str, message: str, author: str) -> str: """Commits staged changes; returns the commit hash.""" - return self._call("git_commit", "git_commit_result", - path=path, message=message, author=author).get("hash", "") + return self._call("git_commit", "git_commit_result", path=path, message=message, author=author).get("hash", "") def git_push(self, path: str, auth: str = "") -> None: self._done_rpc("git_push", path=path, auth=auth or None) @@ -235,8 +266,9 @@ def fork(self, count: int, ttl_seconds: int = 0) -> list[Sandbox]: def hibernate(self) -> None: """Snapshots and stops the VM, freeing its memory; the next call that reaches the guest wakes it transparently, state intact.""" - self._client._request(self.owner, "POST", f"/v1/sandboxes/{self.id}/hibernate", - None, "hibernate", bearer=self.token) + self._client._request( + self.owner, "POST", f"/v1/sandboxes/{self.id}/hibernate", None, "hibernate", bearer=self.token + ) def checkpoint(self, name: str = "") -> Checkpoint: """Captures full state without stopping the sandbox; the returned @@ -250,11 +282,16 @@ def checkpoint(self, name: str = "") -> Checkpoint: def promote(self, template: str) -> Template: """Publishes this sandbox's state as a claimable template on its node; the returned handle is bound to that node.""" - reply = self._client._post_json(self.owner, f"/v1/sandboxes/{self.id}/promote", - {"token": self.token, "template": template}, "promote") + reply = self._client._post_json( + self.owner, f"/v1/sandboxes/{self.id}/promote", {"token": self.token, "template": template}, "promote" + ) key = reply["key"] return Template( - self._client, self.owner, key["template"], key.get("net", ""), key.get("size", ""), + self._client, + self.owner, + key["template"], + key.get("net", ""), + key.get("size", ""), reply.get("content_digest", ""), ) @@ -265,12 +302,12 @@ def start_lsp(self, language: str, root: str = "") -> Lsp: started = self._call("lsp_start", "lsp_started", language=language, root=root or None) return Lsp(self, started["server_id"]) - def open_pty(self, cols: int = 80, rows: int = 24, cwd: str = "", - env: dict | None = None, user: str = "") -> Pty: + def open_pty(self, cols: int = 80, rows: int = 24, cwd: str = "", env: dict | None = None, user: str = "") -> Pty: """Runs the guest shell under a pty; returns a byte-stream handle. A pty is a process guest-side: resize goes through its pid.""" - conn, started = self._open_stream("pty_open", expect="started", cols=cols, - rows=rows, cwd=cwd or None, env=env, user=user or None) + conn, started = self._open_stream( + "pty_open", expect="started", cols=cols, rows=rows, cwd=cwd or None, env=env, user=user or None + ) return Pty(self, conn, started["pid"]) def proxy_port(self, local_addr: str, port: int) -> socket.socket: @@ -279,8 +316,7 @@ def proxy_port(self, local_addr: str, port: int) -> socket.socket: is "host:port"; port 0 picks a free one.""" host, _, lport = local_addr.rpartition(":") listener = socket.create_server((host or "127.0.0.1", int(lport))) - threading.Thread(target=self._proxy_accept_loop, - args=(listener, port), daemon=True).start() + threading.Thread(target=self._proxy_accept_loop, args=(listener, port), daemon=True).start() return listener def preview_url(self, port: int, ttl_seconds: int = 0) -> str: @@ -303,8 +339,9 @@ def close(self) -> None: gone is not an error — double-release and reap races stay silent, matching the Go SDK.""" try: - self._client._request(self.owner, "POST", f"/v1/sandboxes/{self.id}/release", - None, "release", bearer=self.token) + self._client._request( + self.owner, "POST", f"/v1/sandboxes/{self.id}/release", None, "release", bearer=self.token + ) except APIError as exc: if exc.status != 404: raise @@ -328,8 +365,7 @@ def _proxy_accept_loop(self, listener: socket.socket, port: int) -> None: with contextlib.suppress(OSError): while True: local, _ = listener.accept() - threading.Thread(target=self._proxy_conn, - args=(local, port), daemon=True).start() + threading.Thread(target=self._proxy_conn, args=(local, port), daemon=True).start() def _proxy_conn(self, local: socket.socket, port: int) -> None: try: @@ -495,7 +531,7 @@ def close(self) -> None: def _send_chunks(conn: Conn, data: bytes, op: str = "data", chunk: int = FS_CHUNK) -> None: view = memoryview(data) for off in range(0, len(view), chunk): - conn.send(op, data=view[off:off + chunk]) + conn.send(op, data=view[off : off + chunk]) def _feed_stdin(conn: Conn, stdin: bytes) -> None: diff --git a/sdk/python/cocoonsandbox/template.py b/sdk/python/cocoonsandbox/template.py index c9863555..5c49aef4 100644 --- a/sdk/python/cocoonsandbox/template.py +++ b/sdk/python/cocoonsandbox/template.py @@ -22,8 +22,9 @@ def __init__(self, client: Client, addr: str, name: str, net: str, size: str, co self.size = size self.content_digest = content_digest - def new(self, ttl_seconds: int = 0, volumes: list[str | Mapping[str, str]] | None = None, - mount: bool = True) -> Sandbox: + def new( + self, ttl_seconds: int = 0, volumes: list[str | Mapping[str, str]] | None = None, mount: bool = True + ) -> Sandbox: """Claims the template, following placement when volumes require it. mount=False attaches the volumes without mounting them.""" # local import: a top-level one closes the client -> sandbox -> template cycle. @@ -42,5 +43,6 @@ def delete(self) -> None: query = _template_query(self.name, self.net, self.size) query["no_redirect"] = "1" - self._client._request(self._addr, "DELETE", "/v1/templates?" + urllib.parse.urlencode(query), - None, "delete template") + self._client._request( + self._addr, "DELETE", "/v1/templates?" + urllib.parse.urlencode(query), None, "delete template" + ) diff --git a/sdk/python/e2e.py b/sdk/python/e2e.py index 4ac3810d..af87f983 100644 --- a/sdk/python/e2e.py +++ b/sdk/python/e2e.py @@ -41,6 +41,7 @@ def step(name): def wrap(fn): steps.append((name, fn)) return fn + return wrap @step("exec") @@ -128,7 +129,7 @@ def _checkpoint(): branch = ckpt.new() try: assert branch.read_file("/root/ck.txt") == b"v1" # captured moment - assert sb.read_file("/root/ck.txt") == b"v2" # source unaffected + assert sb.read_file("/root/ck.txt") == b"v2" # source unaffected finally: branch.close() listed = [c.id for c in client.checkpoints()] diff --git a/sdk/python/tests/test_client.py b/sdk/python/tests/test_client.py index 4d7b1806..77c43e9d 100644 --- a/sdk/python/tests/test_client.py +++ b/sdk/python/tests/test_client.py @@ -55,7 +55,9 @@ def node(): def test_claim_happy_path(node): FakeNode.routes[("POST", "/v1/claim")] = lambda body, path: ( - 200, {"id": "sb_1", "token": "tok", "owner_addr": node, "template_digest": "sha256:task"}) + 200, + {"id": "sb_1", "token": "tok", "owner_addr": node, "template_digest": "sha256:task"}, + ) sb = Client(node).new("rt:24.04") assert sb.id == "sb_1" and sb.owner == node assert sb.template_digest == "sha256:task" @@ -66,29 +68,40 @@ def test_claim_sends_volumes(node): def claim(body, path): seen.append(body) - return 200, {"id": "sb_1", "token": "tok", "volumes": [ - {"name": "imagenet", "mount": "/volumes/imagenet"}, - {"name": "weights-llama", "mount": "/models"}, - ]} + return 200, { + "id": "sb_1", + "token": "tok", + "volumes": [ + {"name": "imagenet", "mount": "/volumes/imagenet"}, + {"name": "weights-llama", "mount": "/models"}, + ], + } FakeNode.routes[("POST", "/v1/claim")] = claim - sb = Client(node).new("rt:24.04", volumes=[ - "imagenet", {"name": "weights-llama", "mount": "/models"}]) - assert seen == [{"template": "rt:24.04", "volumes": [ - {"name": "imagenet"}, - {"name": "weights-llama", "mount": "/models"}, - ]}] + sb = Client(node).new("rt:24.04", volumes=["imagenet", {"name": "weights-llama", "mount": "/models"}]) + assert seen == [ + { + "template": "rt:24.04", + "volumes": [ + {"name": "imagenet"}, + {"name": "weights-llama", "mount": "/models"}, + ], + } + ] assert sb.volumes == [ {"name": "imagenet", "mount": "/volumes/imagenet"}, {"name": "weights-llama", "mount": "/models"}, ] -@pytest.mark.parametrize(("volumes", "match"), [ - ([("imagenet", "/datasets/imagenet")], "name string or mapping"), - ([{"name": "imagenet", "mode": "rwx"}], "volume 'imagenet': mode must be 'rw' or 'ro', got 'rwx'"), - ([{"name": "imagenet", "bogus": "x"}], r"unexpected key\(s\): bogus"), -]) +@pytest.mark.parametrize( + ("volumes", "match"), + [ + ([("imagenet", "/datasets/imagenet")], "name string or mapping"), + ([{"name": "imagenet", "mode": "rwx"}], "volume 'imagenet': mode must be 'rw' or 'ro', got 'rwx'"), + ([{"name": "imagenet", "bogus": "x"}], r"unexpected key\(s\): bogus"), + ], +) def test_claim_rejects_invalid_volumes(node, volumes, match): with pytest.raises(TypeError, match=match): Client(node).new("rt:24.04", volumes=volumes) @@ -99,15 +112,24 @@ def test_claim_sends_volume_mode_rw(node): def claim(body, path): seen.append(body) - return 200, {"id": "sb_1", "token": "tok", "volumes": [ - {"name": "scratch", "mount": "/data", "mode": "rw"}, - ]} + return 200, { + "id": "sb_1", + "token": "tok", + "volumes": [ + {"name": "scratch", "mount": "/data", "mode": "rw"}, + ], + } FakeNode.routes[("POST", "/v1/claim")] = claim sb = Client(node).new("rt:24.04", volumes=[{"name": "scratch", "mount": "/data", "mode": "rw"}]) - assert seen == [{"template": "rt:24.04", "volumes": [ - {"name": "scratch", "mount": "/data", "mode": "rw"}, - ]}] + assert seen == [ + { + "template": "rt:24.04", + "volumes": [ + {"name": "scratch", "mount": "/data", "mode": "rw"}, + ], + } + ] assert sb.volumes == [{"name": "scratch", "mount": "/data", "mode": "rw"}] @@ -129,17 +151,24 @@ def test_claim_attaches_volumes_without_mounting(node): def claim(body, path): seen.append(body) - return 200, {"id": "sb_1", "token": "tok", "volumes": [ - {"name": "imagenet"}, {"name": "scratch", "mode": "rw"}, - ]} + return 200, { + "id": "sb_1", + "token": "tok", + "volumes": [ + {"name": "imagenet"}, + {"name": "scratch", "mode": "rw"}, + ], + } FakeNode.routes[("POST", "/v1/claim")] = claim sb = Client(node).new("rt:24.04", volumes=["imagenet", {"name": "scratch", "mode": "rw"}], mount=False) - assert seen == [{ - "template": "rt:24.04", - "volumes": [{"name": "imagenet"}, {"name": "scratch", "mode": "rw"}], - "volumes_attach_only": True, - }] + assert seen == [ + { + "template": "rt:24.04", + "volumes": [{"name": "imagenet"}, {"name": "scratch", "mode": "rw"}], + "volumes_attach_only": True, + } + ] assert sb.volumes == [{"name": "imagenet"}, {"name": "scratch", "mode": "rw"}] @@ -151,8 +180,7 @@ def claim(body, path): return 200, {"id": "sb_2", "token": "tok", "volumes": [{"name": "imagenet"}]} FakeNode.routes[("POST", "/v1/claim")] = claim - sb = Template(Client(node), node, "task:v1", "none", "small").new( - volumes=["imagenet"], mount=False) + sb = Template(Client(node), node, "task:v1", "none", "small").new(volumes=["imagenet"], mount=False) assert seen[0]["volumes_attach_only"] is True assert seen[0]["volumes"] == [{"name": "imagenet"}] assert sb.volumes == [{"name": "imagenet"}] @@ -163,7 +191,8 @@ def test_claim_rejects_mount_without_mounting(node): Client(node).new("rt:24.04", volumes=[{"name": "imagenet", "mount": "/datasets"}], mount=False) with pytest.raises(TypeError, match="meaningless with mount=False"): Template(Client(node), node, "task:v1", "none", "small").new( - volumes=[{"name": "imagenet", "mount": "/datasets"}], mount=False) + volumes=[{"name": "imagenet", "mount": "/datasets"}], mount=False + ) def test_claim_keeps_mounting_by_default(node): @@ -183,19 +212,26 @@ def test_template_claim_sends_volumes(node): def claim(body, path): seen.append(body) - return 200, {"id": "sb_2", "token": "tok", "volumes": [ - {"name": "imagenet", "mount": "/datasets/imagenet"}, - ]} + return 200, { + "id": "sb_2", + "token": "tok", + "volumes": [ + {"name": "imagenet", "mount": "/datasets/imagenet"}, + ], + } FakeNode.routes[("POST", "/v1/claim")] = claim sb = Template(Client(node), node, "task:v1", "none", "small").new( - volumes=[{"name": "imagenet", "mount": "/datasets/imagenet"}]) - assert seen == [{ - "template": "task:v1", - "net": "none", - "size": "small", - "volumes": [{"name": "imagenet", "mount": "/datasets/imagenet"}], - }] + volumes=[{"name": "imagenet", "mount": "/datasets/imagenet"}] + ) + assert seen == [ + { + "template": "task:v1", + "net": "none", + "size": "small", + "volumes": [{"name": "imagenet", "mount": "/datasets/imagenet"}], + } + ] assert sb.volumes == [{"name": "imagenet", "mount": "/datasets/imagenet"}] @@ -210,7 +246,8 @@ def claim(body, path): FakeNode.routes[("POST", "/v1/claim")] = claim sb = Template(Client(node), node, "task:v1", "none", "small").new( - volumes=[{"name": "imagenet", "mount": "/datasets/imagenet"}]) + volumes=[{"name": "imagenet", "mount": "/datasets/imagenet"}] + ) assert "no_redirect" not in seen[0] assert seen[1]["no_redirect"] is True assert seen[1]["require_promoted"] is True @@ -219,13 +256,15 @@ def claim(body, path): def test_volume_catalog(node): - want = [{ - "name": "imagenet", - "default_mount": "/volumes/imagenet", - "size_bytes": 42, - "available": True, - "nodes": 3, - }] + want = [ + { + "name": "imagenet", + "default_mount": "/volumes/imagenet", + "size_bytes": 42, + "available": True, + "nodes": 3, + } + ] FakeNode.routes[("GET", "/v1/volumes")] = lambda body, path: (200, {"volumes": want}) assert Client(node).volumes() == want @@ -257,12 +296,16 @@ def test_volume_catalog_surfaces_writable(node): def test_promote_returns_content_digest(node): FakeNode.routes[("POST", "/v1/claim")] = lambda body, path: ( - 200, {"id": "sb_1", "token": "tok", "owner_addr": node}) + 200, + {"id": "sb_1", "token": "tok", "owner_addr": node}, + ) FakeNode.routes[("POST", "/v1/sandboxes/sb_1/promote")] = lambda body, path: ( - 200, { + 200, + { "key": {"template": "task:v1", "net": "none", "size": "small"}, "content_digest": "sha256:promoted", - }) + }, + ) tpl = Client(node).new("rt:24.04").promote("task:v1") assert tpl.name == "task:v1" @@ -284,9 +327,13 @@ def claim(body, path): assert sb.id == "sb_2" assert "no_redirect" not in seen[0] assert seen[1]["no_redirect"] is True - assert seen[0]["volumes"] == seen[1]["volumes"] == [ - {"name": "imagenet", "mount": "/datasets/imagenet"}, - ] + assert ( + seen[0]["volumes"] + == seen[1]["volumes"] + == [ + {"name": "imagenet", "mount": "/datasets/imagenet"}, + ] + ) def test_api_error_carries_server_message(node): @@ -298,11 +345,15 @@ def test_api_error_carries_server_message(node): def test_checkpoint_listing_binds_handles(node): FakeNode.routes[("GET", "/v1/checkpoints")] = lambda body, path: ( - 200, {"checkpoints": [{"id": "ck_0011223344556677", "name": "s1", "sandbox_id": "sb_1"}]}) + 200, + {"checkpoints": [{"id": "ck_0011223344556677", "name": "s1", "sandbox_id": "sb_1"}]}, + ) ckpts = Client(node).checkpoints() assert len(ckpts) == 1 and ckpts[0].id == "ck_0011223344556677" FakeNode.routes[("POST", "/v1/checkpoints/ck_0011223344556677/claim")] = lambda body, path: ( - 200, {"id": "sb_branch", "token": "t2"}) + 200, + {"id": "sb_branch", "token": "t2"}, + ) branch = ckpts[0].new() assert branch.id == "sb_branch" diff --git a/sdk/python/tests/test_fault_matrix.py b/sdk/python/tests/test_fault_matrix.py index 3bcaa512..b16b17e1 100644 --- a/sdk/python/tests/test_fault_matrix.py +++ b/sdk/python/tests/test_fault_matrix.py @@ -171,10 +171,12 @@ def ok(body, path): def test_lookup_dead_peer_resolves_fast(spawn_node, dead_addr): owner = spawn_node({("GET", "/v1/sandboxes/sb_1/owner"): lambda body, path: (200, {"owner_addr": "10.0.0.9:7777"})}) - entry = spawn_node({ - ("GET", "/v1/sandboxes/sb_1/owner"): lambda body, path: (404, {"error": "not here"}), - ("GET", "/v1/peers"): lambda body, path: (200, {"peers": [dead_addr, owner]}), - }) + entry = spawn_node( + { + ("GET", "/v1/sandboxes/sb_1/owner"): lambda body, path: (404, {"error": "not here"}), + ("GET", "/v1/peers"): lambda body, path: (200, {"peers": [dead_addr, owner]}), + } + ) start = time.monotonic() sb = Client(entry, timeout=10.0).lookup("sb_1", "tok") elapsed = time.monotonic() - start @@ -184,10 +186,12 @@ def test_lookup_dead_peer_resolves_fast(spawn_node, dead_addr): def test_lookup_hung_peer_resolves_fast(spawn_node, hung_addr): owner = spawn_node({("GET", "/v1/sandboxes/sb_1/owner"): lambda body, path: (200, {"owner_addr": "10.0.0.9:7777"})}) - entry = spawn_node({ - ("GET", "/v1/sandboxes/sb_1/owner"): lambda body, path: (404, {"error": "not here"}), - ("GET", "/v1/peers"): lambda body, path: (200, {"peers": [hung_addr, owner]}), - }) + entry = spawn_node( + { + ("GET", "/v1/sandboxes/sb_1/owner"): lambda body, path: (404, {"error": "not here"}), + ("GET", "/v1/peers"): lambda body, path: (200, {"peers": [hung_addr, owner]}), + } + ) start = time.monotonic() sb = Client(entry, timeout=10.0).lookup("sb_1", "tok") elapsed = time.monotonic() - start @@ -196,10 +200,12 @@ def test_lookup_hung_peer_resolves_fast(spawn_node, hung_addr): def test_lookup_all_miss(spawn_node, dead_addr): - entry = spawn_node({ - ("GET", "/v1/sandboxes/sb_1/owner"): lambda body, path: (404, {"error": "not here"}), - ("GET", "/v1/peers"): lambda body, path: (200, {"peers": [dead_addr]}), - }) + entry = spawn_node( + { + ("GET", "/v1/sandboxes/sb_1/owner"): lambda body, path: (404, {"error": "not here"}), + ("GET", "/v1/peers"): lambda body, path: (200, {"peers": [dead_addr]}), + } + ) with pytest.raises(APIError) as exc: Client(entry).lookup("sb_1", "tok") assert exc.value.status == 404 and "no owner found" in exc.value.message diff --git a/sdk/python/tests/test_proc.py b/sdk/python/tests/test_proc.py index 139f79a8..394a8d4a 100644 --- a/sdk/python/tests/test_proc.py +++ b/sdk/python/tests/test_proc.py @@ -52,8 +52,7 @@ def test_spawn_returns_pid(monkeypatch): def test_ps_lists_procs(monkeypatch): - procs = [{"pid": 41, "argv": ["sleep"], "detached": True, "state": "running", - "started_at_epoch_secs": 1}] + procs = [{"pid": 41, "argv": ["sleep"], "detached": True, "state": "running", "started_at_epoch_secs": 1}] sb, _ = fake_sandbox(monkeypatch, [{"type": "procs", "procs": procs}]) assert sb.ps() == procs @@ -65,9 +64,7 @@ def test_kill_sends_signal(monkeypatch): def test_logs_running_proc_has_no_code(monkeypatch): - frames = [{"type": "stdout", "data": b"hello"}, - {"type": "stderr", "data": b"oops"}, - {"type": "done"}] + frames = [{"type": "stdout", "data": b"hello"}, {"type": "stderr", "data": b"oops"}, {"type": "done"}] sb, _ = fake_sandbox(monkeypatch, frames) out, errs = [], [] assert sb.logs(41, on_stdout=out.append, on_stderr=errs.append) is None diff --git a/sdk/python/tests/test_wire_binding.py b/sdk/python/tests/test_wire_binding.py index 08d74016..ca03884c 100644 --- a/sdk/python/tests/test_wire_binding.py +++ b/sdk/python/tests/test_wire_binding.py @@ -24,74 +24,64 @@ FIXTURES = pathlib.Path(__file__).parents[3] / "protocol" / "wire" / "fixtures" / "v1" CASES = [ - ("req_exec", [{"type": "started", "pid": 7}, {"type": "exit", "code": 0}], - lambda sb, f: sb.run(f["argv"], cwd=f["cwd"], env=f["env"], user=f["user"], session=f["session"])), - ("req_fs_write", [{"type": "done"}], - lambda sb, f: sb.write_file(f["path"], b"x", mode=f.get("mode"))), - ("req_fs_read", [{"type": "done"}], - lambda sb, f: sb.read_file(f["path"])), - ("req_fs_list", [{"type": "done"}], - lambda sb, f: sb.list_dir(f["path"])), - ("req_fs_stat", [{"type": "stat", "info": {"kind": "file", "size": 1}}], - lambda sb, f: sb.stat(f["path"])), - ("req_fs_mkdir", [{"type": "done"}], - lambda sb, f: sb.mkdir(f["path"], parents=f.get("parents", False))), - ("req_fs_rm", [{"type": "done"}], - lambda sb, f: sb.remove(f["path"], recursive=f.get("recursive", False))), - ("req_fs_rename", [{"type": "done"}], - lambda sb, f: sb.rename(f["from"], f["to"])), - ("req_fs_push", [{"type": "done"}], - lambda sb, f: sb.push(f["dest"], b"tar")), - ("req_fs_pull", [{"type": "done"}], - lambda sb, f: sb.pull(f["path"])), - ("req_fs_find", [{"type": "done"}], - lambda sb, f: sb.find(f["path"], f["pattern"], glob=f.get("glob", ""))), - ("req_fs_replace", [{"type": "done"}], - lambda sb, f: sb.replace(f["files"], f["pattern"], f["replacement"])), - ("req_session_create", [{"type": "session_created", "id": "sess-1"}], - lambda sb, f: sb.session(cwd=f["cwd"], env=f["env"])), - ("req_session_list", [{"type": "sessions", "sessions": []}], - lambda sb, f: sb.sessions()), - ("req_git_clone", [{"type": "done"}], - lambda sb, f: sb.git_clone(f["url"], f["path"], branch=f["branch"], depth=f["depth"], auth=f["auth"])), - ("req_git_status", [{"type": "git_status_result", "branch": "main"}], - lambda sb, f: sb.git_status(f["path"])), - ("req_git_add", [{"type": "done"}], - lambda sb, f: sb.git_add(f["path"], f["files"])), - ("req_git_commit", [{"type": "git_commit_result", "hash": "h"}], - lambda sb, f: sb.git_commit(f["path"], f["message"], f["author"])), - ("req_git_push", [{"type": "done"}], - lambda sb, f: sb.git_push(f["path"], auth=f["auth"])), - ("req_git_pull", [{"type": "done"}], - lambda sb, f: sb.git_pull(f["path"], auth=f["auth"])), - ("req_pty_resize", [{"type": "done"}], - lambda sb, f: _pty_stub(sb, f["pid"]).resize(f["cols"], f["rows"])), - ("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_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"])), - ("req_session_rm", [{"type": "done"}], - lambda sb, f: _session_stub(sb, f["id"]).close()), - ("req_lsp_start", [{"type": "lsp_started", "server_id": "lsp-1"}], - lambda sb, f: sb.start_lsp(f["language"], root=f["root"])), - ("req_lsp_request", [{"type": "ready"}], - lambda sb, f: _lsp_stub(sb, f["server_id"]).request().close()), - ("req_lsp_stop", [{"type": "done"}], - lambda sb, f: _lsp_stub(sb, f["server_id"]).stop()), - ("req_port_forward", [{"type": "ready"}], - lambda sb, f: sb.dial_port(f["port"]).close()), - ("req_pty_open", [{"type": "started", "pid": 7}], - lambda sb, f: sb.open_pty(cols=f["cols"], rows=f["rows"], cwd=f["cwd"], env=f["env"], user=f["user"]).close()), + ( + "req_exec", + [{"type": "started", "pid": 7}, {"type": "exit", "code": 0}], + lambda sb, f: sb.run(f["argv"], cwd=f["cwd"], env=f["env"], user=f["user"], session=f["session"]), + ), + ("req_fs_write", [{"type": "done"}], lambda sb, f: sb.write_file(f["path"], b"x", mode=f.get("mode"))), + ("req_fs_read", [{"type": "done"}], lambda sb, f: sb.read_file(f["path"])), + ("req_fs_list", [{"type": "done"}], lambda sb, f: sb.list_dir(f["path"])), + ("req_fs_stat", [{"type": "stat", "info": {"kind": "file", "size": 1}}], lambda sb, f: sb.stat(f["path"])), + ("req_fs_mkdir", [{"type": "done"}], lambda sb, f: sb.mkdir(f["path"], parents=f.get("parents", False))), + ("req_fs_rm", [{"type": "done"}], lambda sb, f: sb.remove(f["path"], recursive=f.get("recursive", False))), + ("req_fs_rename", [{"type": "done"}], lambda sb, f: sb.rename(f["from"], f["to"])), + ("req_fs_push", [{"type": "done"}], lambda sb, f: sb.push(f["dest"], b"tar")), + ("req_fs_pull", [{"type": "done"}], lambda sb, f: sb.pull(f["path"])), + ("req_fs_find", [{"type": "done"}], lambda sb, f: sb.find(f["path"], f["pattern"], glob=f.get("glob", ""))), + ("req_fs_replace", [{"type": "done"}], lambda sb, f: sb.replace(f["files"], f["pattern"], f["replacement"])), + ( + "req_session_create", + [{"type": "session_created", "id": "sess-1"}], + lambda sb, f: sb.session(cwd=f["cwd"], env=f["env"]), + ), + ("req_session_list", [{"type": "sessions", "sessions": []}], lambda sb, f: sb.sessions()), + ( + "req_git_clone", + [{"type": "done"}], + lambda sb, f: sb.git_clone(f["url"], f["path"], branch=f["branch"], depth=f["depth"], auth=f["auth"]), + ), + ("req_git_status", [{"type": "git_status_result", "branch": "main"}], lambda sb, f: sb.git_status(f["path"])), + ("req_git_add", [{"type": "done"}], lambda sb, f: sb.git_add(f["path"], f["files"])), + ( + "req_git_commit", + [{"type": "git_commit_result", "hash": "h"}], + lambda sb, f: sb.git_commit(f["path"], f["message"], f["author"]), + ), + ("req_git_push", [{"type": "done"}], lambda sb, f: sb.git_push(f["path"], auth=f["auth"])), + ("req_git_pull", [{"type": "done"}], lambda sb, f: sb.git_pull(f["path"], auth=f["auth"])), + ("req_pty_resize", [{"type": "done"}], lambda sb, f: _pty_stub(sb, f["pid"]).resize(f["cols"], f["rows"])), + ("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_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"])), + ("req_session_rm", [{"type": "done"}], lambda sb, f: _session_stub(sb, f["id"]).close()), + ( + "req_lsp_start", + [{"type": "lsp_started", "server_id": "lsp-1"}], + lambda sb, f: sb.start_lsp(f["language"], root=f["root"]), + ), + ("req_lsp_request", [{"type": "ready"}], lambda sb, f: _lsp_stub(sb, f["server_id"]).request().close()), + ("req_lsp_stop", [{"type": "done"}], lambda sb, f: _lsp_stub(sb, f["server_id"]).stop()), + ("req_port_forward", [{"type": "ready"}], lambda sb, f: sb.dial_port(f["port"]).close()), + ( + "req_pty_open", + [{"type": "started", "pid": 7}], + lambda sb, f: sb.open_pty(cols=f["cols"], rows=f["rows"], cwd=f["cwd"], env=f["env"], user=f["user"]).close(), + ), ]