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
14 changes: 7 additions & 7 deletions protocol/wire/frame.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"`
Expand All @@ -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"`
Expand Down
56 changes: 42 additions & 14 deletions sdk/python/cocoonsandbox/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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 []

Expand All @@ -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)."""
Expand Down Expand Up @@ -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
Expand All @@ -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)
Expand Down Expand Up @@ -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"):
Expand Down
118 changes: 77 additions & 41 deletions sdk/python/cocoonsandbox/sandbox.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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()
Expand All @@ -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]:
Expand All @@ -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)."""
Expand Down Expand Up @@ -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)
Expand All @@ -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)
Expand Down Expand Up @@ -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
Expand All @@ -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", ""),
)

Expand All @@ -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:
Expand All @@ -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:
Expand All @@ -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
Expand All @@ -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:
Expand Down Expand Up @@ -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:
Expand Down
10 changes: 6 additions & 4 deletions sdk/python/cocoonsandbox/template.py
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand All @@ -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"
)
Loading