From 4bb023ff9023007029b7c14ed182cd0d6ff76092 Mon Sep 17 00:00:00 2001 From: CMGS Date: Thu, 3 Sep 2026 14:00:03 +0800 Subject: [PATCH 01/10] fix(python): typed error for a refused connection --- sdk/python/cocoonsandbox/conn.py | 5 ++++- sdk/python/tests/test_hardening.py | 5 +++++ 2 files changed, 9 insertions(+), 1 deletion(-) diff --git a/sdk/python/cocoonsandbox/conn.py b/sdk/python/cocoonsandbox/conn.py index 7daa9baf..c9fe2e67 100644 --- a/sdk/python/cocoonsandbox/conn.py +++ b/sdk/python/cocoonsandbox/conn.py @@ -84,7 +84,10 @@ def dial_agent(addr: str, sandbox_id: str, token: str, timeout: float) -> Conn: if any(c in value for c in "\r\n\0"): raise APIError("agent upgrade", 0, f"{name} contains a control character") host, port = addr.rsplit(":", 1) - sock = socket.create_connection((host, int(port)), timeout=timeout) + try: + sock = socket.create_connection((host, int(port)), timeout=timeout) + except OSError as exc: + raise ProtocolError(f"dial {addr}: {exc}") from exc # Nagle off: exec/write send small back-to-back frames before the first read. sock.setsockopt(socket.IPPROTO_TCP, socket.TCP_NODELAY, 1) reader = None diff --git a/sdk/python/tests/test_hardening.py b/sdk/python/tests/test_hardening.py index 78311274..a9161831 100644 --- a/sdk/python/tests/test_hardening.py +++ b/sdk/python/tests/test_hardening.py @@ -17,6 +17,11 @@ def test_dial_agent_rejects_control_chars_in_identity(): dial_agent("127.0.0.1:1", "sb_1", "tok\r\nX-Evil: 1", 0.5) +def test_dial_agent_wraps_refused_connection(dead_addr): + with pytest.raises(ProtocolError): + dial_agent(dead_addr, "sb_1", "tok", 0.5) + + def test_watcher_propagates_silkd_error(): client_sock, guest_sock = socket.socketpair() guest_sock.sendall(json.dumps({"type": "error", "kind": "not_found", "message": "gone"}).encode() + b"\n") From 4e8f43fe8e3dcc519b886ad6ba5fa3f160abb8ac Mon Sep 17 00:00:00 2001 From: CMGS Date: Thu, 3 Sep 2026 14:01:28 +0800 Subject: [PATCH 02/10] review(python): narrow the background-thread catches; one-line comments --- sdk/python/cocoonsandbox/checkpoint.py | 3 +-- sdk/python/cocoonsandbox/client.py | 15 ++++----------- sdk/python/cocoonsandbox/frames.py | 8 ++------ sdk/python/cocoonsandbox/sandbox.py | 18 +++++++----------- sdk/python/cocoonsandbox/template.py | 3 +-- 5 files changed, 15 insertions(+), 32 deletions(-) diff --git a/sdk/python/cocoonsandbox/checkpoint.py b/sdk/python/cocoonsandbox/checkpoint.py index 0a52b866..5dce1d2d 100644 --- a/sdk/python/cocoonsandbox/checkpoint.py +++ b/sdk/python/cocoonsandbox/checkpoint.py @@ -27,8 +27,7 @@ def new(self, ttl_seconds: int = 0) -> Sandbox: redirect to the node that actually holds it; if every candidate fails transiently, the claim falls back to the origin once so it heals (pulls the checkpoint) locally.""" - # Local import: a top-level one would close the client -> sandbox -> - # checkpoint cycle. + # local import: a top-level one closes the client -> sandbox -> checkpoint cycle. from .client import _redirect_fallback claim = {"ttl_seconds": ttl_seconds} if ttl_seconds else {} diff --git a/sdk/python/cocoonsandbox/client.py b/sdk/python/cocoonsandbox/client.py index 4f3d2b65..7ee377bd 100644 --- a/sdk/python/cocoonsandbox/client.py +++ b/sdk/python/cocoonsandbox/client.py @@ -49,8 +49,7 @@ def delete_template(self, template: str, net: str = "", size: str = "") -> None: candidates = (reply or {}).get("redirect") or [] if not candidates: return - # The retry carries no_redirect, mirroring the claim protocol: the - # owner answers for itself, never a second hop. + # the owner answers for itself under no_redirect, never a second hop. query["no_redirect"] = "1" path = "/v1/templates?" + urllib.parse.urlencode(query) _try_each(candidates, lambda peer: self._request(peer, "DELETE", path, None, "delete template")) @@ -60,8 +59,7 @@ def lookup(self, id: str, token: str) -> Sandbox: 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 client timeout after the winner has answered. + # 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, @@ -128,11 +126,7 @@ def post(peer): return self._handle_from(owner, reply) def _peers(self) -> list: - # /v1/peers is tenant-accessible (cluster topology); /v1/info is - # operator-only, so a tenant lookup cannot read peers from it. A - # single node has none — degrade to just the entry node. Bounded - # tighter than the client default (mirrors the Go SDK's peersTimeout) - # so one slow entry node cannot stall the scatter it feeds. + # /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 [] @@ -297,8 +291,7 @@ def attempt(addr): try: return attempt(origin) except APIError as origin_exc: - # Both halves matter to whoever reads this: the peers' failure says - # why the claim left the origin, the origin's why returning did not help. + # both halves matter: why the claim left the origin, and why returning did not help. combined = f"{origin_exc.message} (after redirect targets failed: {exc.message})" raise APIError(verb, origin_exc.status, combined) from origin_exc diff --git a/sdk/python/cocoonsandbox/frames.py b/sdk/python/cocoonsandbox/frames.py index 2d52d48d..68760364 100644 --- a/sdk/python/cocoonsandbox/frames.py +++ b/sdk/python/cocoonsandbox/frames.py @@ -11,8 +11,7 @@ PROTO_VERSION = 1 MAX_FRAME = 8 * 1024 * 1024 FS_CHUNK = 256 * 1024 # silkd's per-frame chunk size, distinct from BULK_CHUNK below -# Bulk streams (push tars, port bytes) chunk larger — fewer frames for the -# same bytes, still far under MAX_FRAME after base64; mirrors the Go SDK. +# bulk streams chunk larger: fewer frames per byte, still under MAX_FRAME after base64. BULK_CHUNK = 1 << 20 @@ -32,10 +31,7 @@ def encode_request(op: str, **fields) -> bytes: def decode_response(line: bytes) -> dict: """Parses one response frame; the returned dict carries its tag under "type" and any binary payload decoded under "data".""" - # base64 is JSON-escape-free, so a frame shaped exactly {"type":..., - # "data":...} can be sliced directly, skipping json.loads; any other - # shape (extra fields, trailing bytes, non-alphabet bytes) falls through - # to the full parse instead of risking a silently wrong slice. + # base64 is JSON-escape-free, so an exactly-shaped data frame slices without json.loads. if line.startswith(b'{"type":"'): te = line.find(b'"', 9) if te > 0 and line[9:te] in (b"stdout", b"stderr", b"data") and line.startswith(b'","data":"', te): diff --git a/sdk/python/cocoonsandbox/sandbox.py b/sdk/python/cocoonsandbox/sandbox.py index 99f7a472..fe04a921 100644 --- a/sdk/python/cocoonsandbox/sandbox.py +++ b/sdk/python/cocoonsandbox/sandbox.py @@ -13,7 +13,7 @@ from .checkpoint import Checkpoint from .conn import Conn, _Closeable, dial_agent -from .errors import APIError, ExitError, ProtocolError, StreamTimeout +from .errors import APIError, ExitError, ProtocolError, SandboxError, StreamTimeout from .frames import BULK_CHUNK, FS_CHUNK from .template import Template @@ -66,8 +66,7 @@ def run(self, argv: list[str], cwd: str = "", env: dict | None = None, 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) - # The guest stops draining stdin while blocked writing stdout, so - # feeding it to completion before reading deadlocks. + # 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() code = _pump_stdio(conn, on_stdout, on_stderr) @@ -335,12 +334,12 @@ def _proxy_accept_loop(self, listener: socket.socket, port: int) -> None: def _proxy_conn(self, local: socket.socket, port: int) -> None: try: guest = self.dial_port(port) - except Exception: + except (SandboxError, OSError): local.close() return def pump_out(): - with contextlib.suppress(Exception): + with contextlib.suppress(SandboxError, OSError): while True: chunk = guest.recv() if not chunk: @@ -357,8 +356,7 @@ def pump_out(): if not chunk: break guest.send(chunk) - # Half-close, then let the guest finish answering: closing here - # would cut the reply the local client is still waiting for. + # half-close, not close: the guest's reply is still in flight. guest.close_write() pump.join() finally: @@ -401,9 +399,7 @@ def __init__(self, conn: Conn): self.error: Exception | None = None def __iter__(self) -> Iterator[dict]: - # Connection-bound: a close, drop, or undecodable frame ends iteration; - # a real server error frame (SilkdError) propagates. error tells a - # clean close (None) from a relay that dropped mid-stream. + # a transport failure ends iteration and sets error; a SilkdError frame propagates. while True: try: frame = self._conn.recv() @@ -503,7 +499,7 @@ def _send_chunks(conn: Conn, data: bytes, op: str = "data", chunk: int = FS_CHUN def _feed_stdin(conn: Conn, stdin: bytes) -> None: - with contextlib.suppress(Exception): # the reader reports the real failure + with contextlib.suppress(SandboxError, OSError): # the reader reports the real failure if stdin: _send_chunks(conn, stdin, op="stdin") conn.send("stdin_close") diff --git a/sdk/python/cocoonsandbox/template.py b/sdk/python/cocoonsandbox/template.py index ad9c9cd3..c9863555 100644 --- a/sdk/python/cocoonsandbox/template.py +++ b/sdk/python/cocoonsandbox/template.py @@ -26,8 +26,7 @@ def new(self, ttl_seconds: int = 0, volumes: list[str | Mapping[str, str]] | Non 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 would close the client → sandbox → - # template cycle. + # local import: a top-level one closes the client -> sandbox -> template cycle. from .client import _claim_body claim = _claim_body(self.name, self.net, self.size, ttl_seconds, volumes, mount) From af69c1f23a041e1af39b1314a09846bd2aa80186 Mon Sep 17 00:00:00 2001 From: CMGS Date: Thu, 3 Sep 2026 14:09:46 +0800 Subject: [PATCH 03/10] review(go): layout, one func type, WalkDir, logger name, test comments --- e2e/e2e_test.go | 8 ----- e2e/fakeengine_test.go | 6 ---- protocol/wire/frame.go | 50 ++++++++++++++++---------------- protocol/wire/frame_test.go | 3 -- sandboxd/ca.go | 1 - sandboxd/main.go | 1 - sandboxd/mesh/mesh.go | 6 ++-- sandboxd/pool/archive.go | 4 +-- sandboxd/store/peer/transport.go | 44 +++++++++++++++------------- sandboxd/types/types.go | 20 ++++++------- sdk/go/client.go | 6 +++- 11 files changed, 70 insertions(+), 79 deletions(-) diff --git a/e2e/e2e_test.go b/e2e/e2e_test.go index b4ae082a..fc49d0c2 100644 --- a/e2e/e2e_test.go +++ b/e2e/e2e_test.go @@ -379,8 +379,6 @@ func TestWritableVolumeEndToEnd(t *testing.T) { } } -// The read-only leg runs second: it is admitted only once the writer's -// release has cleared the dirty marker. func TestVolumeModeWireShape(t *testing.T) { scratch := writeVolumeImage(t, "scratch.img", "scratch-bytes") stack := startTenantStack(t, "node-token", nil, @@ -523,8 +521,6 @@ func TestAttachOnlyVolumeWireShape(t *testing.T) { } } -// TestDirtyVolumeRefusesReader: the marker a crashed writer leaves behind -// (pre-created here) turns read-only claims into 409s over the wire. func TestDirtyVolumeRefusesReader(t *testing.T) { scratch := writeVolumeImage(t, "scratch.img", "scratch-bytes") if err := os.WriteFile(scratch+".dirty", nil, 0o600); err != nil { @@ -601,8 +597,6 @@ func writeVolumeImage(t *testing.T, name, content string) string { return image } -// assertNoDirtyMarker fails if the image carries the write-ahead marker: an -// attach-only claim makes no consistency promise, so it must never write one. func assertNoDirtyMarker(t *testing.T, image, when string) { t.Helper() if _, err := os.Stat(image + ".dirty"); !errors.Is(err, os.ErrNotExist) { @@ -610,8 +604,6 @@ func assertNoDirtyMarker(t *testing.T, image, when string) { } } -// rawClaimResponse decodes the volume entries generically, so the assertion -// is the server's own JSON rather than the SDK's mirror of it. type rawClaimResponse struct { ID string `json:"id"` Token string `json:"token"` diff --git a/e2e/fakeengine_test.go b/e2e/fakeengine_test.go index 4ebb5cc8..f5f05914 100644 --- a/e2e/fakeengine_test.go +++ b/e2e/fakeengine_test.go @@ -16,8 +16,6 @@ import ( "github.com/cocoonstack/sandbox/sdk/go/silkd/silkdtest" ) -// fakeEngine replaces only the cocoon CLI: every "VM" is a silkdtest daemon -// behind a real hybrid-vsock UDS, so the data plane runs production code. type fakeEngine struct { real *engine.Engine dir string @@ -81,8 +79,6 @@ func (f *fakeEngine) SnapshotRemove(_ context.Context, _ string) error { return func (f *fakeEngine) SnapshotList(_ context.Context) ([]string, error) { return nil, nil } -// Hibernate closes the VM's silkd listener and Restore brings a fresh one up, -// mirroring the stop/resume the control plane observes. func (f *fakeEngine) Hibernate(ctx context.Context, name, _ string) error { return f.Remove(ctx, name) } @@ -143,8 +139,6 @@ func (f *fakeEngine) SyncGuest(_ context.Context, _ string) error { return nil } -// Removals join the trace only for VMs that carried a volume, so warm-pool -// churn cannot perturb the order. func (f *fakeEngine) volumeOpsLog() []string { f.mu.Lock() defer f.mu.Unlock() diff --git a/protocol/wire/frame.go b/protocol/wire/frame.go index 499809f2..a7e0cc9b 100644 --- a/protocol/wire/frame.go +++ b/protocol/wire/frame.go @@ -518,6 +518,16 @@ type InfoResp struct { func (InfoResp) RespType() string { return "info" } +// ProcInfo is one entry of Procs; ExitCode is absent while running. +type ProcInfo struct { + PID uint32 `json:"pid"` + Argv []string `json:"argv"` + Detached bool `json:"detached"` + State string `json:"state"` + ExitCode *int32 `json:"exit_code,omitempty"` + StartedAtEpochSecs uint64 `json:"started_at_epoch_secs"` +} + // Procs answers Ps. type Procs struct { Procs []ProcInfo `json:"procs"` @@ -532,6 +542,13 @@ type DataResp struct { func (DataResp) RespType() string { return "data" } +// DirEntry is one entry of Entries; Kind is one of the FileKind* consts. +type DirEntry struct { + Name string `json:"name"` + Kind string `json:"kind"` + Size uint64 `json:"size"` +} + // Entries answers FsList. type Entries struct { Entries []DirEntry `json:"entries"` @@ -539,6 +556,14 @@ type Entries struct { func (Entries) RespType() string { return "entries" } +// FileInfo is the Stat payload; Mode carries permission bits only. +type FileInfo struct { + Kind string `json:"kind"` + Size uint64 `json:"size"` + Mode uint32 `json:"mode"` + MtimeEpochSecs uint64 `json:"mtime_epoch_secs"` +} + // Stat answers FsStat. type Stat struct { Info FileInfo `json:"info"` @@ -625,31 +650,6 @@ type GitBranches struct { func (GitBranches) RespType() string { return "git_branches" } -// ProcInfo is one entry of Procs; ExitCode is absent while running. -type ProcInfo struct { - PID uint32 `json:"pid"` - Argv []string `json:"argv"` - Detached bool `json:"detached"` - State string `json:"state"` - ExitCode *int32 `json:"exit_code,omitempty"` - StartedAtEpochSecs uint64 `json:"started_at_epoch_secs"` -} - -// DirEntry is one entry of Entries; Kind is one of the FileKind* consts. -type DirEntry struct { - Name string `json:"name"` - Kind string `json:"kind"` - Size uint64 `json:"size"` -} - -// FileInfo is the Stat payload; Mode carries permission bits only. -type FileInfo struct { - Kind string `json:"kind"` - Size uint64 `json:"size"` - Mode uint32 `json:"mode"` - MtimeEpochSecs uint64 `json:"mtime_epoch_secs"` -} - // EncodeRequest renders {"v":1,"op":...,fields} without a trailing newline. func EncodeRequest(r Request) ([]byte, error) { return encodeTagged(requestHead+r.Op()+`"`, r) diff --git a/protocol/wire/frame_test.go b/protocol/wire/frame_test.go index 4efef372..19a7e47c 100644 --- a/protocol/wire/frame_test.go +++ b/protocol/wire/frame_test.go @@ -158,7 +158,6 @@ func TestEveryVerbHasAFixture(t *testing.T) { } } -// Twin of silkd's enum_value_sets_match_fixture; the lists are order-sensitive. func TestEnumValueSetsMatchFixture(t *testing.T) { raw, err := os.ReadFile(filepath.Join(fixtureDir, "enums.json")) if err != nil { @@ -244,8 +243,6 @@ func TestBulkDecodeStrictShape(t *testing.T) { } } -// TestTagAfterOtherKeys pins the tokenizer fallback: producers emit the tag -// first, but the protocol never promised order. func TestTagAfterOtherKeys(t *testing.T) { resp, err := DecodeResponse([]byte(`{"data":"aGk=","type":"stdout"}`)) if err != nil { diff --git a/sandboxd/ca.go b/sandboxd/ca.go index 840ee6cf..7547b32c 100644 --- a/sandboxd/ca.go +++ b/sandboxd/ca.go @@ -102,7 +102,6 @@ func writeCAFiles(dir, name string, certPEM, keyPEM []byte, force bool) error { return nil } -// writeKeyMaterial writes path with O_EXCL unless force. func writeKeyMaterial(path string, data []byte, perm os.FileMode, force bool) error { flags := os.O_WRONLY | os.O_CREATE | os.O_TRUNC if !force { diff --git a/sandboxd/main.go b/sandboxd/main.go index a866bc9b..44ddff2d 100644 --- a/sandboxd/main.go +++ b/sandboxd/main.go @@ -207,7 +207,6 @@ func startMesh(ctx context.Context, cfg *config.Config, mgr *pool.Manager) (*mes return msh, nil } -// gossipNodeState republishes this node's counts, templates, and volumes every tick. func gossipNodeState(ctx context.Context, msh *mesh.Mesh, mgr *pool.Manager) { t := time.NewTicker(gossipInterval) defer t.Stop() diff --git a/sandboxd/mesh/mesh.go b/sandboxd/mesh/mesh.go index d915d02c..31138dc3 100644 --- a/sandboxd/mesh/mesh.go +++ b/sandboxd/mesh/mesh.go @@ -39,6 +39,8 @@ type NodeState struct { Digest string `json:"digest,omitempty"` // cluster-invariant config digest } +type nodeMatch func(NodeState) bool + // Mesh is the node's view of the cluster and its own gossiped state. type Mesh struct { ml *memberlist.Memberlist @@ -223,7 +225,7 @@ func (m *Mesh) Shutdown() error { return m.ml.Shutdown() } -func (m *Mesh) warmCandidates(keyHash string, match func(NodeState) bool) []string { +func (m *Mesh) warmCandidates(keyHash string, match nodeMatch) []string { m.mu.Lock() type cand struct { addr string @@ -262,7 +264,7 @@ func (m *Mesh) persistEpoch(epoch uint64) error { return storeEpoch(m.epochPath, epoch) } -func (m *Mesh) owners(match func(NodeState) bool) []string { +func (m *Mesh) owners(match nodeMatch) []string { m.mu.Lock() var owners []string for id, st := range m.view { diff --git a/sandboxd/pool/archive.go b/sandboxd/pool/archive.go index 39a24bc5..4768d74e 100644 --- a/sandboxd/pool/archive.go +++ b/sandboxd/pool/archive.go @@ -348,12 +348,12 @@ func (m *Manager) retryArchiveDelete(ctx context.Context, ckID string) { l.Unlock() if err != nil { m.recDone(ckID) - log.WithFunc("pool.retryArchiveDeletes").Warnf(ctx, "delete %s: %v", ckID, err) + log.WithFunc("pool.retryArchiveDelete").Warnf(ctx, "delete %s: %v", ckID, err) return } m.recDoneEvict(ckID) if err := m.clearArchiveCk(ckID); err != nil { - log.WithFunc("pool.retryArchiveDeletes").Warnf(ctx, "clear %s: %v", ckID, err) + log.WithFunc("pool.retryArchiveDelete").Warnf(ctx, "clear %s: %v", ckID, err) } } diff --git a/sandboxd/store/peer/transport.go b/sandboxd/store/peer/transport.go index 9ec9c035..50dcb97d 100644 --- a/sandboxd/store/peer/transport.go +++ b/sandboxd/store/peer/transport.go @@ -10,6 +10,7 @@ import ( "errors" "fmt" "io" + "io/fs" "net/http" "net/url" "os" @@ -196,7 +197,7 @@ func safeJoin(root, name string) (string, error) { // tarInto emits only regular files and directories: a symlink would steer the reader outside dst. func tarInto(src, prefix string, tw *tar.Writer) error { - err := filepath.Walk(src, func(path string, fi os.FileInfo, err error) error { + err := filepath.WalkDir(src, func(path string, entry fs.DirEntry, err error) error { if err != nil { return err } @@ -207,33 +208,36 @@ func tarInto(src, prefix string, tw *tar.Writer) error { if rel == "." { return nil } + if !entry.IsDir() && !entry.Type().IsRegular() { + return nil + } + fi, err := entry.Info() + if err != nil { + return err + } name := prefix + rel - switch { - case fi.IsDir(): + if entry.IsDir() { return tw.WriteHeader(&tar.Header{ Name: name + "/", Mode: int64(fi.Mode().Perm()), Typeflag: tar.TypeDir, }) - case fi.Mode().IsRegular(): - if err := tw.WriteHeader(&tar.Header{ - Name: name, - Mode: int64(fi.Mode().Perm()), - Size: fi.Size(), - Typeflag: tar.TypeReg, - }); err != nil { - return err - } - f, err := os.Open(path) //nolint:gosec // path comes from Walk over src - if err != nil { - return err - } - defer func() { _ = f.Close() }() - _, err = io.Copy(tw, f) + } + if err = tw.WriteHeader(&tar.Header{ + Name: name, + Mode: int64(fi.Mode().Perm()), + Size: fi.Size(), + Typeflag: tar.TypeReg, + }); err != nil { return err - default: - return nil } + f, err := os.Open(path) //nolint:gosec // path comes from WalkDir over src + if err != nil { + return err + } + defer func() { _ = f.Close() }() + _, err = io.Copy(tw, f) + return err }) if err != nil { return fmt.Errorf("tar %s: %w", src, err) diff --git a/sandboxd/types/types.go b/sandboxd/types/types.go index 5d3fb180..74fc4747 100644 --- a/sandboxd/types/types.go +++ b/sandboxd/types/types.go @@ -196,6 +196,16 @@ type Checkpoint struct { Archive bool `json:"archive,omitempty"` } +// VMNetConfig is the per-NIC host tap the egress-lane nft lock binds. +type VMNetConfig struct { + TAP string `json:"tap"` +} + +// VMConfig is the config subset of VMRecord. +type VMConfig struct { + Name string `json:"name"` +} + // VMRecord is the subset of cocoon's VM records the control plane reads. type VMRecord struct { State string `json:"state"` @@ -213,16 +223,6 @@ func (r VMRecord) TapDevice() string { return r.NetworkConfigs[0].TAP } -// VMNetConfig is the per-NIC host tap the egress-lane nft lock binds. -type VMNetConfig struct { - TAP string `json:"tap"` -} - -// VMConfig is the config subset of VMRecord. -type VMConfig struct { - Name string `json:"name"` -} - // Volume is one dataset mount; Mode is normalized to "" (read-only) or VolumeModeRW. type Volume struct { Name string `json:"name"` diff --git a/sdk/go/client.go b/sdk/go/client.go index aef31e58..7ca8920c 100644 --- a/sdk/go/client.go +++ b/sdk/go/client.go @@ -306,10 +306,14 @@ func retryTransient(err error) bool { } } +type claimEncoder func(noRedirect, requirePromoted bool) ([]byte, error) + +type claimPoster func(addr string, body []byte) (claimResponse, error) + // claimFollow runs the claim protocol from origin: claim there, and on a // redirect re-encode with no_redirect and follow via redirectFallback. Only // the fallback error carries the verb — first-contact errors return raw. -func claimFollow(origin, verb string, encode func(noRedirect, requirePromoted bool) ([]byte, error), claimAt func(addr string, body []byte) (claimResponse, error)) (string, claimResponse, error) { +func claimFollow(origin, verb string, encode claimEncoder, claimAt claimPoster) (string, claimResponse, error) { body, err := encode(false, false) if err != nil { return "", claimResponse{}, err From 5ee7bf6fbc270c046b52a84a6c560d6abdbbd9f1 Mon Sep 17 00:00:00 2001 From: CMGS Date: Thu, 3 Sep 2026 14:16:58 +0800 Subject: [PATCH 04/10] review(rust): ownership on the silkd request path --- silkd/src/fs.rs | 28 ++++++++++++++++++---------- silkd/src/git.rs | 25 +++++++++++-------------- silkd/src/lsp.rs | 5 +++-- silkd/src/proc.rs | 5 +++-- silkd/src/proto.rs | 9 +++++---- silkd/src/server.rs | 3 ++- silkd/src/tree.rs | 10 +++++----- 7 files changed, 47 insertions(+), 38 deletions(-) diff --git a/silkd/src/fs.rs b/silkd/src/fs.rs index a804cb46..5f66369a 100644 --- a/silkd/src/fs.rs +++ b/silkd/src/fs.rs @@ -5,6 +5,7 @@ use std::io; use std::os::unix::fs::PermissionsExt; +use std::path::{Path, PathBuf}; use std::time::UNIX_EPOCH; use tokio::fs; @@ -32,7 +33,8 @@ where R: AsyncBufRead + Unpin, W: AsyncWrite + Unpin, { - let tmp = tmp_name(&path); + let path = Path::new(&path); + let tmp = tmp_name(path); let mut file = match fs::File::create(&tmp).await { Ok(f) => f, Err(e) => return err_frame(w, &e, "create").await, @@ -48,7 +50,7 @@ where let _ = fs::remove_file(&tmp).await; return proto::write_feed_error(w, fail).await; } - if let Err(e) = commit_tmp(&tmp, &path, mode).await { + if let Err(e) = commit_tmp(&tmp, path, mode).await { return err_frame(w, &e, "commit").await; } proto::write_frame(w, &Response::Done).await @@ -92,14 +94,13 @@ pub async fn list(w: &mut W, path: String) -> io::Result< /// Writes `bytes` to `path` via a sibling temp file committed into place, so /// a crash never leaves a truncated file. -pub async fn write_atomic(path: &std::path::Path, bytes: &[u8]) -> io::Result<()> { - let path = path.display().to_string(); - let tmp = tmp_name(&path); +pub async fn write_atomic(path: &Path, bytes: &[u8]) -> io::Result<()> { + let tmp = tmp_name(path); if let Err(e) = fs::write(&tmp, bytes).await { let _ = fs::remove_file(&tmp).await; return Err(e); } - commit_tmp(&tmp, &path, None).await + commit_tmp(&tmp, path, None).await } /// Reports metadata for `path` (following symlinks). @@ -177,7 +178,10 @@ fn scan_dir(path: &str, tx: &mpsc::Sender>) -> io::Result<()> { Err(_) => (FileKind::Other, 0), }; entries.push(DirEntry { - name: ent.file_name().to_string_lossy().into_owned(), + name: ent + .file_name() + .into_string() + .unwrap_or_else(|os| os.to_string_lossy().into_owned()), kind, size, }); @@ -191,8 +195,12 @@ fn scan_dir(path: &str, tx: &mpsc::Sender>) -> io::Result<()> { Ok(()) } -fn tmp_name(path: &str) -> String { - format!("{path}.silkd-{}.tmp", crate::sysutil::tmp_suffix()) +fn tmp_name(path: &Path) -> PathBuf { + let mut tmp = path.as_os_str().to_os_string(); + tmp.push(".silkd-"); + tmp.push(crate::sysutil::tmp_suffix()); + tmp.push(".tmp"); + PathBuf::from(tmp) } /// Commits a fully-written temp file over `path`. An explicit mode wins; @@ -200,7 +208,7 @@ fn tmp_name(path: &str) -> String { /// replacing an executable script doesn't silently strip its exec bit /// (rename alone would leave the temp's create default). The temp is /// removed on any failure. -async fn commit_tmp(tmp: &str, path: &str, mode: Option) -> io::Result<()> { +async fn commit_tmp(tmp: &Path, path: &Path, mode: Option) -> io::Result<()> { let outcome = async { let effective = match mode { Some(m) => Some(m), diff --git a/silkd/src/git.rs b/silkd/src/git.rs index a0e32250..5a498994 100644 --- a/silkd/src/git.rs +++ b/silkd/src/git.rs @@ -276,8 +276,8 @@ fn parse_file_line(line: &str) -> Option { let mut fields = line.split(' '); let xy = fields.nth(1)?; // field 2 let mut chars = xy.chars(); - let staged = chars.next()?.to_string(); - let unstaged = chars.next()?.to_string(); + let staged = chars.next()?; + let unstaged = chars.next()?; let skip = match kind { "2" => 7, "u" => 8, @@ -286,21 +286,24 @@ fn parse_file_line(line: &str) -> Option { let path = fields.nth(skip)?; // Rejoin any spaces the split consumed, then drop a rename's \t. let rest: Vec<&str> = fields.collect(); - let full = if rest.is_empty() { + let mut full = if rest.is_empty() { path.to_string() } else { format!("{path} {}", rest.join(" ")) }; + if let Some(tab) = full.find('\t') { + full.truncate(tab); + } Some(GitFileStatus { - path: full.split('\t').next()?.to_string(), + path: full, staged, unstaged, }) } "?" => Some(GitFileStatus { path: line.get(2..)?.to_string(), - staged: "?".to_string(), - unstaged: "?".to_string(), + staged: '?', + unstaged: '?', }), _ => None, } @@ -344,10 +347,7 @@ mod tests { fn parses_ordinary_rename_untracked_and_conflict() { let ordinary = parse_file_line("1 .M N... 100644 100644 100644 h1 h2 src/main.rs").unwrap(); assert_eq!(ordinary.path, "src/main.rs"); - assert_eq!( - (ordinary.staged.as_str(), ordinary.unstaged.as_str()), - (".", "M") - ); + assert_eq!((ordinary.staged, ordinary.unstaged), ('.', 'M')); let rename = parse_file_line("2 R. N... 100644 100644 100644 h1 h2 R100 new.rs\told.rs").unwrap(); @@ -359,9 +359,6 @@ mod tests { let conflict = parse_file_line("u UU N... 100644 100644 100644 100644 h1 h2 h3 conflict.rs").unwrap(); assert_eq!(conflict.path, "conflict.rs"); - assert_eq!( - (conflict.staged.as_str(), conflict.unstaged.as_str()), - ("U", "U") - ); + assert_eq!((conflict.staged, conflict.unstaged), ('U', 'U')); } } diff --git a/silkd/src/lsp.rs b/silkd/src/lsp.rs index 3b94b1f5..ce3b3803 100644 --- a/silkd/src/lsp.rs +++ b/silkd/src/lsp.rs @@ -16,6 +16,7 @@ //! image ships no manifests, so `lsp_start` for any language there answers a //! typed `not_found` naming the flavor that provides one. +use std::borrow::Cow; use std::collections::HashMap; use std::process::Stdio; use std::sync::atomic::{AtomicU64, Ordering}; @@ -233,8 +234,8 @@ async fn pump_stdout( /// manifest_dir allows tests (and an operator) to relocate the manifest dir /// via SILKD_LSP_DIR, mirroring silkd's other env overrides. -fn manifest_dir() -> String { - std::env::var("SILKD_LSP_DIR").unwrap_or_else(|_| MANIFEST_DIR.to_string()) +fn manifest_dir() -> Cow<'static, str> { + std::env::var("SILKD_LSP_DIR").map_or(Cow::Borrowed(MANIFEST_DIR), Cow::Owned) } async fn read_manifest(language: &str) -> Option> { diff --git a/silkd/src/proc.rs b/silkd/src/proc.rs index b6ca1388..9911167c 100644 --- a/silkd/src/proc.rs +++ b/silkd/src/proc.rs @@ -3,6 +3,7 @@ //! Detached processes keep a bounded output ring so a later `logs`/`attach` //! can replay what already streamed. +use std::borrow::Cow; use std::collections::{HashMap, VecDeque}; use std::os::fd::{AsRawFd, OwnedFd}; use std::sync::atomic::{AtomicU64, Ordering}; @@ -209,8 +210,8 @@ impl Proc { fn info(&self) -> ProcInfo { let (state, exit_code) = match *sysutil::lock(&self.state) { - State::Running => ("running".to_string(), None), - State::Exited(c) => ("exited".to_string(), Some(c)), + State::Running => (Cow::Borrowed("running"), None), + State::Exited(c) => (Cow::Borrowed("exited"), Some(c)), }; ProcInfo { pid: self.pid, diff --git a/silkd/src/proto.rs b/silkd/src/proto.rs index d462235d..b9fa6fbe 100644 --- a/silkd/src/proto.rs +++ b/silkd/src/proto.rs @@ -4,6 +4,7 @@ //! (exit / done / error). Binary payloads ride as base64 (`data` fields), //! matching Go's default []byte JSON encoding for the SDK side. +use std::borrow::Cow; use std::collections::HashMap; use std::io; use std::sync::Arc; @@ -233,7 +234,7 @@ pub enum Response { message: String, }, Info { - version: String, + version: Cow<'static, str>, proto: u32, uptime_secs: u64, procs: usize, @@ -323,8 +324,8 @@ pub enum GitBranchOp { #[derive(Debug, Deserialize, Serialize)] pub struct GitFileStatus { pub path: String, - pub staged: String, - pub unstaged: String, + pub staged: char, + pub unstaged: char, } #[derive(Clone, Copy, Debug, Deserialize, PartialEq, Serialize)] @@ -365,7 +366,7 @@ pub struct ProcInfo { pub pid: u32, pub argv: Vec, pub detached: bool, - pub state: String, + pub state: Cow<'static, str>, #[serde(default, skip_serializing_if = "Option::is_none")] pub exit_code: Option, pub started_at_epoch_secs: u64, diff --git a/silkd/src/server.rs b/silkd/src/server.rs index 563160b0..70921341 100644 --- a/silkd/src/server.rs +++ b/silkd/src/server.rs @@ -1,6 +1,7 @@ //! Connection handling: read the leading request frame, dispatch to a //! handler, and for exec forward subsequent client frames over a channel. +use std::borrow::Cow; use std::time::{Instant, SystemTime, UNIX_EPOCH}; use tokio::io::{AsyncBufRead, AsyncWrite}; @@ -199,7 +200,7 @@ impl State { proto::write_frame( w, &Response::Info { - version: env!("CARGO_PKG_VERSION").to_string(), + version: Cow::Borrowed(env!("CARGO_PKG_VERSION")), proto: proto::PROTO_VERSION, uptime_secs: self.started.elapsed().as_secs(), procs: self.table.len(), diff --git a/silkd/src/tree.rs b/silkd/src/tree.rs index 5a5c5e33..573afb5f 100644 --- a/silkd/src/tree.rs +++ b/silkd/src/tree.rs @@ -49,7 +49,7 @@ where pub async fn pull(w: &mut W, path: String) -> io::Result<()> { let p = Path::new(&path); let (parent, name) = match (p.parent(), p.file_name()) { - (Some(par), Some(n)) if !n.is_empty() => (par.to_path_buf(), n.to_os_string()), + (Some(par), Some(n)) if !n.is_empty() => (par, n), _ => return proto::error_frame(w, ErrorKind::BadRequest, "invalid path").await, }; // symlink_metadata, not metadata: a dangling symlink is a valid tar source @@ -58,7 +58,7 @@ pub async fn pull(w: &mut W, path: String) -> io::Result< return err_frame(w, &e, "stat source").await; } let cwd = if parent.as_os_str().is_empty() { - Path::new(".").to_path_buf() + Path::new(".") } else { parent }; @@ -66,9 +66,9 @@ pub async fn pull(w: &mut W, path: String) -> io::Result< // `--` so a basename starting with `-` is a path, not a tar option. cmd.arg("-c") .arg("-C") - .arg(&cwd) + .arg(cwd) .arg("--") - .arg(&name) + .arg(name) .stdout(Stdio::piped()); let (mut child, err_task) = match spawn_tar(cmd) { Ok(pair) => pair, @@ -203,5 +203,5 @@ async fn drain(mut stderr: tokio::process::ChildStderr) -> String { out.extend_from_slice(&buf[..(n.min(CAP - out.len()))]); } } - String::from_utf8_lossy(&out).into_owned() + String::from_utf8(out).unwrap_or_else(|e| String::from_utf8_lossy(e.as_bytes()).into_owned()) } From 0d89d94af1d2852487b279e8a8d02b5a4b9768d5 Mon Sep 17 00:00:00 2001 From: CMGS Date: Thu, 3 Sep 2026 14:21:16 +0800 Subject: [PATCH 05/10] review(rust): layout, shared cmdline tokenizer, one-line docs --- boot/init/src/boot.rs | 6 +- boot/init/src/cfg.rs | 25 +++-- silkd/src/exec.rs | 14 +-- silkd/src/find.rs | 8 +- silkd/src/git.rs | 14 +-- silkd/src/lsp.rs | 3 +- silkd/src/proto.rs | 10 +- silkd/src/session.rs | 16 +-- silkd/src/sysutil.rs | 198 ++++++++++++++++++------------------- silkd/tests/common/mod.rs | 16 +-- silkd/tests/exec_e2e.rs | 4 +- silkd/tests/forward_e2e.rs | 5 +- silkd/tests/git_e2e.rs | 3 +- silkd/tests/lsp_e2e.rs | 7 +- silkd/tests/pty_e2e.rs | 2 - silkd/tests/session_e2e.rs | 4 +- silkd/tests/tree_e2e.rs | 3 +- 17 files changed, 131 insertions(+), 207 deletions(-) diff --git a/boot/init/src/boot.rs b/boot/init/src/boot.rs index 22eb62b0..62d81806 100644 --- a/boot/init/src/boot.rs +++ b/boot/init/src/boot.rs @@ -145,11 +145,7 @@ fn assemble(cfg: &BootCfg, marks: &mut Marks) -> Result<(), String> { Ok(()) } -/// Materializes kernel ip= params (cocoon CNI static flow) as MAC-matched -/// networkd units in the new root — persistence only, nothing is configured -/// in the initramfs; networkd applies them once the real init is up. A NIC -/// that never shows up degrades that interface to the DHCP fallback instead -/// of failing the boot, matching the old init-bottom hook. +/// Persists kernel ip= params as MAC-matched networkd units in the new root; a missing NIC degrades to the DHCP fallback. fn persist_network(cfg: &BootCfg) { if cfg.ips.is_empty() { return; diff --git a/boot/init/src/cfg.rs b/boot/init/src/cfg.rs index aa24d1ef..3303b2cd 100644 --- a/boot/init/src/cfg.rs +++ b/boot/init/src/cfg.rs @@ -54,8 +54,7 @@ pub fn parse(cmdline: &str) -> Result { debug: false, trace: false, }; - for tok in cmdline.split_ascii_whitespace() { - let (key, val) = tok.split_once('=').unwrap_or((tok, "")); + for (key, val) in params(cmdline) { match key { "cocoon.layers" => { cfg.layers = val @@ -92,18 +91,11 @@ pub fn parse(cmdline: &str) -> Result { Ok(cfg) } -/// Debug check for the path where parse() itself failed and cfg.debug is -/// unavailable. Token handling mirrors parse() exactly: same value -/// predicate, last occurrence wins. +/// Debug check for the path where parse() itself failed and cfg.debug is unavailable. pub fn debug_requested(cmdline: &str) -> bool { - let mut debug = false; - for tok in cmdline.split_ascii_whitespace() { - let (key, val) = tok.split_once('=').unwrap_or((tok, "")); - if key == "sandbox.debug" { - debug = debug_token(val); - } - } - debug + params(cmdline) + .rfind(|&(key, _)| key == "sandbox.debug") + .is_some_and(|(_, val)| debug_token(val)) } /// Overlay mount data. Layer mountpoints are index-based (/l/0, /l/1, …) so @@ -132,6 +124,13 @@ pub fn network_unit(ip: &IpParam, mac: &str) -> String { unit } +/// Kernel cmdline tokens as (key, value); a bare token carries an empty value. +fn params(cmdline: &str) -> impl DoubleEndedIterator { + cmdline + .split_ascii_whitespace() + .map(|tok| tok.split_once('=').unwrap_or((tok, ""))) +} + /// sandbox.debug value semantics, shared by parse() and debug_requested(). fn debug_token(val: &str) -> bool { val.is_empty() || val == "1" diff --git a/silkd/src/exec.rs b/silkd/src/exec.rs index d764249b..31c1401c 100644 --- a/silkd/src/exec.rs +++ b/silkd/src/exec.rs @@ -82,12 +82,7 @@ where let stdout = child.stdout.take(); let stderr = child.stderr.take(); - // Foreground consumes a backpressured mpsc: when the client falls behind, - // the sends block, the pipe fills, and the child slows down — nothing - // drops while the child lives (POST_EXIT_DRAIN bounds only the post-exit - // tail). Attachers still ride the best-effort broadcast (a secondary - // observer may drop under extreme lag). Detached has no client to pace - // it, so it gets no foreground sender. + // a backpressured mpsc paces the child to the foreground client; attachers ride the lossy broadcast. let (fg_tx, fg_rx) = mpsc::channel::(FG_CAP); let pump_fg = (!req.detach).then(|| fg_tx.clone()); let sup_fg = (!req.detach).then(|| fg_tx.clone()); @@ -97,12 +92,7 @@ where let pump_abort = pump.abort_handle(); tokio::spawn(pump_stdin(stdin, client, req.detach)); - // One supervisor per exec: reap the child, drain output within a grace - // window (so a daemonizer holding the pipe can't wedge us), then publish - // the terminal Exit — to the broadcast (attachers) and the foreground mpsc. - // The reaped code is recorded before the drain: a disconnect can abort the - // supervisor mid-drain, and its fallback must publish the real code, not - // fabricate -1 for a child that already exited cleanly. + // record the reaped code before the drain: an abort mid-drain must publish the real code, not -1. let reaped: Arc> = Arc::new(OnceLock::new()); let sup_reaped = Arc::clone(&reaped); let sup_proc = Arc::clone(&proc); diff --git a/silkd/src/find.rs b/silkd/src/find.rs index d8cc679b..3a40e833 100644 --- a/silkd/src/find.rs +++ b/silkd/src/find.rs @@ -151,13 +151,7 @@ where find_bounded(reader, w, path, pattern, glob, MATCH_QUEUE_BYTES).await } -/// Rewrites every `pattern` match to `replacement` in each of `files`, -/// streaming one `replaced` frame per file (with its match count) and a -/// terminal `done`. A file over the find size bound is skipped with a zero -/// count rather than read into memory. A read/write failure on one file ends -/// the stream with an error — files whose `replaced` frame already went out -/// are committed (each file is atomic; the list is not). An invalid pattern is -/// rejected before any file is touched. +/// Rewrites `pattern` to `replacement` in each of `files`, one `replaced` frame each; per-file atomic, not per-list. pub async fn replace( w: &mut W, files: Vec, diff --git a/silkd/src/git.rs b/silkd/src/git.rs index 5a498994..8a3d1337 100644 --- a/silkd/src/git.rs +++ b/silkd/src/git.rs @@ -165,11 +165,7 @@ async fn net_verb( terminal(w, verb, &out).await } -/// Builds a `git -C dir` command with config and stdio policy applied. Config -/// (an auth token, and quotePath=false so paths come back raw) rides in -/// `GIT_CONFIG_*` env vars, not `-c` args: the process environ is root-only, -/// whereas argv is world-readable via /proc//cmdline — a de-escalated -/// exec could otherwise scrape the token. +/// Builds a `git -C dir` command; config rides in `GIT_CONFIG_*` env vars because argv is world-readable via /proc. fn git_cmd(dir: &str, auth: Option<&str>) -> Command { let mut cmd = Command::new("git"); sysutil::align_proxy_env(&mut cmd); @@ -262,13 +258,7 @@ fn parse_ahead_behind(rest: &str) -> (u32, u32) { (ahead, behind) } -/// Parses one porcelain-v2 entry. XY is field 2; the path is the last -/// space-field (kept intact even with spaces — core.quotePath=false keeps it -/// raw). Ordinary changes ("1") have the path at field 9; renames/copies ("2") -/// add an Xscore field, so the path is field 10 and carries "\t", -/// of which we keep the new path; unmerged ("u") entries carry four modes and -/// three hashes, putting the bare path at field 11. Untracked ("?") is a -/// bare path. +/// Parses one porcelain-v2 entry: XY is field 2, the path the last space-field at 9 ("1"), 10 ("2") or 11 ("u"). fn parse_file_line(line: &str) -> Option { let kind = line.split(' ').next()?; match kind { diff --git a/silkd/src/lsp.rs b/silkd/src/lsp.rs index ce3b3803..9d9e1f56 100644 --- a/silkd/src/lsp.rs +++ b/silkd/src/lsp.rs @@ -232,8 +232,7 @@ async fn pump_stdout( proto::write_frame(w, &Response::Done).await } -/// manifest_dir allows tests (and an operator) to relocate the manifest dir -/// via SILKD_LSP_DIR, mirroring silkd's other env overrides. +/// SILKD_LSP_DIR relocates the manifest dir for tests and operators. fn manifest_dir() -> Cow<'static, str> { std::env::var("SILKD_LSP_DIR").map_or(Cow::Borrowed(MANIFEST_DIR), Cow::Owned) } diff --git a/silkd/src/proto.rs b/silkd/src/proto.rs index b9fa6fbe..0bc43882 100644 --- a/silkd/src/proto.rs +++ b/silkd/src/proto.rs @@ -387,10 +387,7 @@ pub async fn read_frame(r: &mut R) -> io::Result