From 496d59131faf676282002006a80aba0db469764c Mon Sep 17 00:00:00 2001 From: CMGS Date: Wed, 16 Sep 2026 22:43:11 +0800 Subject: [PATCH 1/2] relay: audit every request line on a kept connection The audit tee recorded the first client line and then passed bytes through, on the premise that a connection carried one RPC. silkd now serves RPCs back to back on one connection (#196), so the tee scans every line: request frames are recorded, the input frames of the RPC in flight (stdin, stdin_close, data, data_end) are skipped by their canonical head, and a request frame past the 4 KB cap is recorded once as oversized. wire.IsContinuation names that head check for the relay and the SDKs alike. The Go-side docs drop the one-connection-per-RPC wording, and the agent endpoint's doc says that an open relay holds the sandbox's idle clock, which is what a client keeping a connection warm must know. Hot path: nothing changes without audit; with audit, one newline scan per relayed chunk and at most 4 KB copied per line, where only the first line was scanned before. --- README.md | 3 +- docs/sandboxd-api.md | 6 ++- protocol/wire/frame.go | 18 ++++++++- protocol/wire/frame_test.go | 28 ++++++++++++++ sandboxd/engine/silkd.go | 2 +- sandboxd/pool/telemetry.go | 2 +- sandboxd/server/relay.go | 69 ++++++++++++++++------------------- sandboxd/server/relay_test.go | 56 ++++++++++++++++++++++++---- 8 files changed, 133 insertions(+), 51 deletions(-) diff --git a/README.md b/README.md index b6004ea0..67cc4e79 100644 --- a/README.md +++ b/README.md @@ -29,7 +29,8 @@ performance) — source in persistent shell sessions, streaming fs, tar-stream tree push/pull, find/replace, watch (ready-acked), pty, structured git, guest port relay (`port_forward`), and an LSP broker for flavor-shipped language - servers — newline-JSON frames over vsock 2048, one connection per RPC; + servers — newline-JSON frames over vsock 2048, RPCs back to back on one + connection; baked into the base image - `sandboxd/` — per-node control plane (Go): warm pools refilled from golden snapshot exports (online-retunable), claim/release/hibernate/fork/promote/ diff --git a/docs/sandboxd-api.md b/docs/sandboxd-api.md index 533271e9..686a4e0c 100644 --- a/docs/sandboxd-api.md +++ b/docs/sandboxd-api.md @@ -520,8 +520,10 @@ cocoon's machine-level metering ledger for audit cross-checks. Auth: the sandbox's own token. Requires `Upgrade: silkd` + `Connection: Upgrade`; answers `101 Switching Protocols` and from then on -the connection is a byte-for-byte relay to the guest's silkd (one silkd RPC -per connection — see [silkd](silkd.md)). 426 without the upgrade header, 404 +the connection is a byte-for-byte relay to the guest's silkd, carrying RPCs +back to back (see [silkd](silkd.md)). An open relay holds the sandbox's idle +clock, so a client that keeps a connection warm must close it when idle for +`idle_hibernate_seconds` to apply. 426 without the upgrade header, 404 unknown sandbox, 502 guest unreachable. ## POST /v1/sandboxes/{id}/exec diff --git a/protocol/wire/frame.go b/protocol/wire/frame.go index 9a7a2ecc..50dcb237 100644 --- a/protocol/wire/frame.go +++ b/protocol/wire/frame.go @@ -1,7 +1,8 @@ // Package wire is the Go binding of the silkd wire protocol, shared by the // SDK and sandboxd: -// newline-delimited JSON frames over one connection per RPC, requests tagged -// by "op", responses by "type", binary payloads base64 in data fields. The +// newline-delimited JSON frames, RPCs back to back on one connection, +// requests tagged by "op", responses by "type", binary payloads base64 in +// data fields. The // authoritative contract is the shared corpus in protocol/wire/fixtures/v1 — // silkd's Rust tests and this package's tests round-trip the same files. package wire @@ -13,6 +14,7 @@ import ( "encoding/json" "fmt" "io" + "slices" "strconv" ) @@ -61,6 +63,13 @@ var ( reqTagHead = []byte(requestHead) respTagHead = []byte(`{"type":"`) + continuationHeads = [][]byte{ + []byte(requestHead + `data"`), + []byte(requestHead + `data_end"`), + []byte(requestHead + `stdin"`), + []byte(requestHead + `stdin_close"`), + } + requestDecoders = map[string]func([]byte) (Request, error){ "exec": decodeReq[Exec], "info": decodeReq[Info], //nolint:goconst // wire tag shared with the response type by design @@ -707,6 +716,11 @@ func AppendBulkRequest(buf []byte, op string, data []byte) []byte { return append(buf, '"', '}', '\n') } +// IsContinuation reports, from the canonical head every encoder here emits, whether a request line carries input for the RPC in flight rather than opening one. +func IsContinuation(line []byte) bool { + return slices.ContainsFunc(continuationHeads, func(head []byte) bool { return bytes.HasPrefix(line, head) }) +} + // fastBulk slices the base64 data out of a canonical bulk frame, skipping the // json.Unmarshal that dominates downloads; any other shape falls back to slow. func fastBulk(tag string, slow respDecoder, mk func([]byte) Response) respDecoder { diff --git a/protocol/wire/frame_test.go b/protocol/wire/frame_test.go index 1a6076a5..ceb53e46 100644 --- a/protocol/wire/frame_test.go +++ b/protocol/wire/frame_test.go @@ -311,6 +311,34 @@ func TestAppendBulkRequestMatchesEncodeRequest(t *testing.T) { } } +func TestIsContinuation(t *testing.T) { + encoded := func(r Request) string { + line, err := EncodeRequest(r) + if err != nil { + t.Fatalf("EncodeRequest(%s): %v", r.Op(), err) + } + return string(line) + } + for _, tt := range []struct { + line string + want bool + }{ + {encoded(Data{Data: []byte("hi")}), true}, + {encoded(DataEnd{}), true}, + {encoded(Stdin{Data: []byte("hi")}), true}, + {encoded(StdinClose{}), true}, + {string(AppendBulkRequest(nil, "stdin", []byte("hi"))), true}, + {encoded(Exec{Argv: []string{"true"}}), false}, + {encoded(FsStat{Path: "/"}), false}, + {encoded(Info{}), false}, + {"", false}, + } { + if got := IsContinuation([]byte(tt.line)); got != tt.want { + t.Errorf("IsContinuation(%q) = %v, want %v", tt.line, got, tt.want) + } + } +} + func jsonEqual(t *testing.T, a, b []byte) bool { t.Helper() var av, bv any diff --git a/sandboxd/engine/silkd.go b/sandboxd/engine/silkd.go index cbd06f77..e32bac80 100644 --- a/sandboxd/engine/silkd.go +++ b/sandboxd/engine/silkd.go @@ -55,7 +55,7 @@ func (e *Engine) dialSilkdSession(ctx context.Context, vsockSocket string) (*sil }, nil } -// silkdStream serves one request per dial: silkd handles one request per connection. +// silkdStream dials per request; the node's own calls are too rare to keep a connection. func (e *Engine) silkdStream(ctx context.Context, vsockSocket string, req wire.Request, onFrame func(wire.Response) error) error { s, err := e.dialSilkdSession(ctx, vsockSocket) if err != nil { diff --git a/sandboxd/pool/telemetry.go b/sandboxd/pool/telemetry.go index 64caa7e8..157ac5b8 100644 --- a/sandboxd/pool/telemetry.go +++ b/sandboxd/pool/telemetry.go @@ -105,7 +105,7 @@ func (m *Manager) Audit(ctx context.Context, id string, line []byte) { } var frame auditFrame if err := json.Unmarshal(line, &frame); err != nil || frame.Op == "" { - return // torn cap boundary or a non-frame first line: nothing to record + return // torn cap boundary or a non-frame line: nothing to record } m.recordAudit(ctx, id, frame) } diff --git a/sandboxd/server/relay.go b/sandboxd/server/relay.go index 1e9bff80..bc428384 100644 --- a/sandboxd/server/relay.go +++ b/sandboxd/server/relay.go @@ -12,6 +12,7 @@ import ( "github.com/projecteru2/core/log" + "github.com/cocoonstack/sandbox/protocol/wire" "github.com/cocoonstack/sandbox/sandboxd/pool" "github.com/cocoonstack/sandbox/sandboxd/utils" ) @@ -113,7 +114,6 @@ func (s *Server) relay(ctx context.Context, id string, client net.Conn, clientBu clientR = io.MultiReader(io.LimitReader(clientBuf, int64(n)), client) } if s.mgr.AuditEnabled() { - // one connection is one RPC, so the first line client→guest is the request frame clientR = &auditTee{r: clientR, record: func(line []byte) { s.mgr.Audit(ctx, id, line) }} @@ -134,55 +134,50 @@ func (s *Server) relay(ctx context.Context, id string, client net.Conn, clientBu <-done } -// auditTee captures the first line flowing through it, then degrades to a pass-through. +// auditTee records each request line it relays; input frames and the bytes past the cap pass unrecorded. type auditTee struct { r io.Reader record func([]byte) buf []byte - done bool + skip bool } -func (t *auditTee) WriteTo(w io.Writer) (int64, error) { - var total int64 - buf := make([]byte, 4096) - for !t.done { - n, err := t.Read(buf) - if n > 0 { - wn, werr := w.Write(buf[:n]) - total += int64(wn) - if werr != nil { - return total, werr - } - } - if err == io.EOF { - return total, nil - } - if err != nil { - return total, err +func (t *auditTee) Read(p []byte) (int, error) { + n, err := t.r.Read(p) + rest := p[:n] + for len(rest) > 0 { + i := bytes.IndexByte(rest, '\n') + if i < 0 { + t.take(rest) + break } + t.take(rest[:i]) + t.endLine() + rest = rest[i+1:] } - n, err := io.Copy(w, t.r) - return total + n, err + return n, err } -func (t *auditTee) Read(p []byte) (int, error) { - n, err := t.r.Read(p) - if t.done || n == 0 { - return n, err +// take keeps a line's head up to the cap; a request frame that outgrows it is recorded once, oversized. +func (t *auditTee) take(b []byte) { + if t.skip { + return + } + t.buf = append(t.buf, b...) + if len(t.buf) <= pool.AuditLineCap { + return } - chunk := p[:n] - if i := bytes.IndexByte(chunk, '\n'); i >= 0 { - t.buf = append(t.buf, chunk[:i+1]...) - t.done = true + t.skip = true + if !wire.IsContinuation(t.buf) { t.record(t.buf) - t.buf = nil - return n, err } - t.buf = append(t.buf, chunk...) - if len(t.buf) > pool.AuditLineCap { - t.done = true // a frame this large is payload, not addressing + t.buf = t.buf[:0] +} + +func (t *auditTee) endLine() { + if !t.skip && len(t.buf) > 0 && !wire.IsContinuation(t.buf) { t.record(t.buf) - t.buf = nil } - return n, err + t.buf = t.buf[:0] + t.skip = false } diff --git a/sandboxd/server/relay_test.go b/sandboxd/server/relay_test.go index 20135de2..6150a235 100644 --- a/sandboxd/server/relay_test.go +++ b/sandboxd/server/relay_test.go @@ -8,13 +8,22 @@ import ( "net" "net/http" "net/http/httptest" + "slices" "strings" "testing" + "testing/iotest" "time" + + "github.com/cocoonstack/sandbox/sandboxd/pool" ) const ( - execFrame = `{"v":1,"op":"exec","argv":["echo","42"],"detach":false}` + "\n" + execFrame = `{"v":1,"op":"exec","argv":["echo","42"],"detach":false}` + "\n" + statFrame = `{"v":1,"op":"fs_stat","path":"/"}` + dataFrame = `{"v":1,"op":"data","data":"aGk="}` + "\n" + dataEndFrame = `{"v":1,"op":"data_end"}` + "\n" + stdinFrame = `{"v":1,"op":"stdin","data":"aGk="}` + "\n" + stdinCloseFrame = `{"v":1,"op":"stdin_close"}` + "\n" agentRequest = "GET /v1/sandboxes/sb_1/agent HTTP/1.1\r\n" + "Host: sandboxd\r\n" + @@ -127,12 +136,45 @@ func TestRelayClientDisconnectClosesGuest(t *testing.T) { } } -func TestAuditTeeWriteToEndsCleanlyBeforeALine(t *testing.T) { - tee := &auditTee{r: strings.NewReader("partial"), record: func([]byte) { t.Error("recorded a line that never ended") }} - var out bytes.Buffer - n, err := tee.WriteTo(&out) - if err != nil || n != 7 || out.String() != "partial" { - t.Errorf("WriteTo = %d, %v, %q; want 7, nil, the bytes", n, err, out.String()) +func TestAuditTeeRecordsRequestLines(t *testing.T) { + big := `{"v":1,"op":"exec","argv":["` + strings.Repeat("x", pool.AuditLineCap) + `"],"detach":false}` + chunk := `{"v":1,"op":"data","data":"` + strings.Repeat("A", 3*pool.AuditLineCap) + `"}` + for _, tt := range []struct { + name string + in string + want []string + }{ + {"requests back to back", execFrame + statFrame + "\n", []string{strings.TrimSuffix(execFrame, "\n"), statFrame}}, + {"upload input skipped", statFrame + "\n" + dataFrame + chunk + "\n" + dataEndFrame + execFrame, []string{statFrame, strings.TrimSuffix(execFrame, "\n")}}, + {"stdin skipped", execFrame + stdinFrame + stdinCloseFrame, []string{strings.TrimSuffix(execFrame, "\n")}}, + {"oversized request recorded once", big + "\n" + statFrame + "\n", []string{"oversized", statFrame}}, + {"partial tail never recorded", execFrame + "partial", []string{strings.TrimSuffix(execFrame, "\n")}}, + } { + for _, rd := range []struct { + name string + wrap func(io.Reader) io.Reader + }{ + {"whole", func(r io.Reader) io.Reader { return r }}, + {"byte at a time", iotest.OneByteReader}, + } { + t.Run(tt.name+"/"+rd.name, func(t *testing.T) { + var got []string + tee := &auditTee{r: rd.wrap(strings.NewReader(tt.in)), record: func(line []byte) { + if len(line) > pool.AuditLineCap { + got = append(got, "oversized") + return + } + got = append(got, string(line)) + }} + var out bytes.Buffer + if _, err := io.Copy(&out, tee); err != nil || out.String() != tt.in { + t.Fatalf("copy: %v, relayed %d bytes of %d", err, out.Len(), len(tt.in)) + } + if !slices.Equal(got, tt.want) { + t.Errorf("recorded %q, want %q", got, tt.want) + } + }) + } } } From 67f281065ea85f51eeb3d3d73802d85508a73fd3 Mon Sep 17 00:00:00 2001 From: CMGS Date: Thu, 17 Sep 2026 00:06:58 +0800 Subject: [PATCH 2/2] relay: classify audit lines by op and keep bulk-sized reads IsContinuation matched four canonical heads by prefix, which misread a producer that orders its keys differently and restated the op strings the encoders own; it now takes the op through frameTag, the canonical fast path with the token-walk fallback, and switches on the request types' own Op(). The audit tee regains a WriteTo with a bulk-sized buffer, since without it io.Copy's 32 KB buffer split every upload frame into eight reads, and it copies at most the cap plus one byte of an oversized frame instead of the whole chunk. The relay test builds its frames with wire.EncodeRequest so the predicate is tested against the encoders, not against literals. --- protocol/wire/frame.go | 31 +++++++++++----------- protocol/wire/frame_test.go | 1 + sandboxd/server/relay.go | 25 ++++++++++++++++++ sandboxd/server/relay_test.go | 50 +++++++++++++++++++++-------------- 4 files changed, 71 insertions(+), 36 deletions(-) diff --git a/protocol/wire/frame.go b/protocol/wire/frame.go index 50dcb237..838e29e1 100644 --- a/protocol/wire/frame.go +++ b/protocol/wire/frame.go @@ -1,10 +1,9 @@ // Package wire is the Go binding of the silkd wire protocol, shared by the -// SDK and sandboxd: -// newline-delimited JSON frames, RPCs back to back on one connection, -// requests tagged by "op", responses by "type", binary payloads base64 in -// data fields. The -// authoritative contract is the shared corpus in protocol/wire/fixtures/v1 — -// silkd's Rust tests and this package's tests round-trip the same files. +// SDK and sandboxd: newline-delimited JSON frames, RPCs back to back on one +// connection, requests tagged by "op", responses by "type", binary payloads +// base64 in data fields. The authoritative contract is the shared corpus in +// protocol/wire/fixtures/v1 — silkd's Rust tests and this package's tests +// round-trip the same files. package wire import ( @@ -14,7 +13,6 @@ import ( "encoding/json" "fmt" "io" - "slices" "strconv" ) @@ -63,13 +61,6 @@ var ( reqTagHead = []byte(requestHead) respTagHead = []byte(`{"type":"`) - continuationHeads = [][]byte{ - []byte(requestHead + `data"`), - []byte(requestHead + `data_end"`), - []byte(requestHead + `stdin"`), - []byte(requestHead + `stdin_close"`), - } - requestDecoders = map[string]func([]byte) (Request, error){ "exec": decodeReq[Exec], "info": decodeReq[Info], //nolint:goconst // wire tag shared with the response type by design @@ -716,9 +707,17 @@ func AppendBulkRequest(buf []byte, op string, data []byte) []byte { return append(buf, '"', '}', '\n') } -// IsContinuation reports, from the canonical head every encoder here emits, whether a request line carries input for the RPC in flight rather than opening one. +// IsContinuation reports whether a request line feeds the RPC in flight instead of opening one. func IsContinuation(line []byte) bool { - return slices.ContainsFunc(continuationHeads, func(head []byte) bool { return bytes.HasPrefix(line, head) }) + op, err := frameTag(line, reqTagHead, "op") + if err != nil { + return false + } + switch string(op) { + case Data{}.Op(), DataEnd{}.Op(), Stdin{}.Op(), StdinClose{}.Op(): + return true + } + return false } // fastBulk slices the base64 data out of a canonical bulk frame, skipping the diff --git a/protocol/wire/frame_test.go b/protocol/wire/frame_test.go index ceb53e46..eb9ee14a 100644 --- a/protocol/wire/frame_test.go +++ b/protocol/wire/frame_test.go @@ -328,6 +328,7 @@ func TestIsContinuation(t *testing.T) { {encoded(Stdin{Data: []byte("hi")}), true}, {encoded(StdinClose{}), true}, {string(AppendBulkRequest(nil, "stdin", []byte("hi"))), true}, + {`{"op":"data","v":1,"data":"aGk="}`, true}, {encoded(Exec{Argv: []string{"true"}}), false}, {encoded(FsStat{Path: "/"}), false}, {encoded(Info{}), false}, diff --git a/sandboxd/server/relay.go b/sandboxd/server/relay.go index bc428384..87bcc370 100644 --- a/sandboxd/server/relay.go +++ b/sandboxd/server/relay.go @@ -142,6 +142,28 @@ type auditTee struct { skip bool } +// WriteTo keeps an audited upload at bulk-sized reads; io.Copy's own buffer would split each frame into eight. +func (t *auditTee) WriteTo(w io.Writer) (int64, error) { + var total int64 + buf := make([]byte, wire.BulkChunk) + for { + n, err := t.Read(buf) + if n > 0 { + wn, werr := w.Write(buf[:n]) + total += int64(wn) + if werr != nil { + return total, werr + } + } + if err == io.EOF { + return total, nil + } + if err != nil { + return total, err + } + } +} + func (t *auditTee) Read(p []byte) (int, error) { n, err := t.r.Read(p) rest := p[:n] @@ -163,6 +185,9 @@ func (t *auditTee) take(b []byte) { if t.skip { return } + if room := pool.AuditLineCap + 1 - len(t.buf); len(b) > room { + b = b[:room] + } t.buf = append(t.buf, b...) if len(t.buf) <= pool.AuditLineCap { return diff --git a/sandboxd/server/relay_test.go b/sandboxd/server/relay_test.go index 6150a235..df9af569 100644 --- a/sandboxd/server/relay_test.go +++ b/sandboxd/server/relay_test.go @@ -14,22 +14,23 @@ import ( "testing/iotest" "time" + "github.com/cocoonstack/sandbox/protocol/wire" "github.com/cocoonstack/sandbox/sandboxd/pool" ) -const ( - execFrame = `{"v":1,"op":"exec","argv":["echo","42"],"detach":false}` + "\n" - statFrame = `{"v":1,"op":"fs_stat","path":"/"}` - dataFrame = `{"v":1,"op":"data","data":"aGk="}` + "\n" - dataEndFrame = `{"v":1,"op":"data_end"}` + "\n" - stdinFrame = `{"v":1,"op":"stdin","data":"aGk="}` + "\n" - stdinCloseFrame = `{"v":1,"op":"stdin_close"}` + "\n" - - agentRequest = "GET /v1/sandboxes/sb_1/agent HTTP/1.1\r\n" + - "Host: sandboxd\r\n" + - "Authorization: Bearer tok\r\n" + - "Upgrade: silkd\r\n" + - "Connection: Upgrade\r\n\r\n" +const agentRequest = "GET /v1/sandboxes/sb_1/agent HTTP/1.1\r\n" + + "Host: sandboxd\r\n" + + "Authorization: Bearer tok\r\n" + + "Upgrade: silkd\r\n" + + "Connection: Upgrade\r\n\r\n" + +var ( + execFrame = line(wire.Exec{Argv: []string{"echo", "42"}}) + statFrame = line(wire.FsStat{Path: "/"}) + dataFrame = line(&wire.Data{Data: []byte("hi")}) + dataEndFrame = line(wire.DataEnd{}) + stdinFrame = line(&wire.Stdin{Data: []byte("hi")}) + stdinCloseFrame = line(wire.StdinClose{}) ) func TestCloseRelaysRefusesLateRelays(t *testing.T) { @@ -137,18 +138,19 @@ func TestRelayClientDisconnectClosesGuest(t *testing.T) { } func TestAuditTeeRecordsRequestLines(t *testing.T) { - big := `{"v":1,"op":"exec","argv":["` + strings.Repeat("x", pool.AuditLineCap) + `"],"detach":false}` - chunk := `{"v":1,"op":"data","data":"` + strings.Repeat("A", 3*pool.AuditLineCap) + `"}` + big := line(wire.Exec{Argv: []string{strings.Repeat("x", pool.AuditLineCap)}}) + chunk := line(&wire.Data{Data: make([]byte, 3*pool.AuditLineCap)}) + exec, stat := strings.TrimSuffix(execFrame, "\n"), strings.TrimSuffix(statFrame, "\n") for _, tt := range []struct { name string in string want []string }{ - {"requests back to back", execFrame + statFrame + "\n", []string{strings.TrimSuffix(execFrame, "\n"), statFrame}}, - {"upload input skipped", statFrame + "\n" + dataFrame + chunk + "\n" + dataEndFrame + execFrame, []string{statFrame, strings.TrimSuffix(execFrame, "\n")}}, - {"stdin skipped", execFrame + stdinFrame + stdinCloseFrame, []string{strings.TrimSuffix(execFrame, "\n")}}, - {"oversized request recorded once", big + "\n" + statFrame + "\n", []string{"oversized", statFrame}}, - {"partial tail never recorded", execFrame + "partial", []string{strings.TrimSuffix(execFrame, "\n")}}, + {"requests back to back", execFrame + statFrame, []string{exec, stat}}, + {"upload input skipped", statFrame + dataFrame + chunk + dataEndFrame + execFrame, []string{stat, exec}}, + {"stdin skipped", execFrame + stdinFrame + stdinCloseFrame, []string{exec}}, + {"oversized request recorded once", big + statFrame, []string{"oversized", stat}}, + {"partial tail never recorded", execFrame + "partial", []string{exec}}, } { for _, rd := range []struct { name string @@ -229,6 +231,14 @@ func upgradeConn(t *testing.T, ts *httptest.Server) (net.Conn, *bufio.Reader) { return conn, r } +func line(r wire.Request) string { + frame, err := wire.EncodeRequest(r) + if err != nil { + panic(err) + } + return string(frame) + "\n" +} + func readStatus101(t *testing.T, r *bufio.Reader) { t.Helper() resp, err := http.ReadResponse(r, nil)