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
3 changes: 2 additions & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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/
Expand Down
6 changes: 4 additions & 2 deletions docs/sandboxd-api.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
23 changes: 18 additions & 5 deletions protocol/wire/frame.go
Original file line number Diff line number Diff line change
@@ -1,9 +1,9 @@
// 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
// 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 (
Expand Down Expand Up @@ -707,6 +707,19 @@ func AppendBulkRequest(buf []byte, op string, data []byte) []byte {
return append(buf, '"', '}', '\n')
}

// IsContinuation reports whether a request line feeds the RPC in flight instead of opening one.
func IsContinuation(line []byte) bool {
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
// json.Unmarshal that dominates downloads; any other shape falls back to slow.
func fastBulk(tag string, slow respDecoder, mk func([]byte) Response) respDecoder {
Expand Down
29 changes: 29 additions & 0 deletions protocol/wire/frame_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -311,6 +311,35 @@ 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},
{`{"op":"data","v":1,"data":"aGk="}`, 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
Expand Down
2 changes: 1 addition & 1 deletion sandboxd/engine/silkd.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
2 changes: 1 addition & 1 deletion sandboxd/pool/telemetry.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
}
Expand Down
60 changes: 40 additions & 20 deletions sandboxd/server/relay.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
)
Expand Down Expand Up @@ -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)
}}
Expand All @@ -134,18 +134,19 @@ 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
}

// 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, 4096)
for !t.done {
buf := make([]byte, wire.BulkChunk)
for {
n, err := t.Read(buf)
if n > 0 {
wn, werr := w.Write(buf[:n])
Expand All @@ -161,28 +162,47 @@ func (t *auditTee) WriteTo(w io.Writer) (int64, error) {
return total, err
}
}
n, err := io.Copy(w, t.r)
return total + n, err
}

func (t *auditTee) Read(p []byte) (int, error) {
n, err := t.r.Read(p)
if t.done || n == 0 {
return n, err
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:]
}
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
}
chunk := p[:n]
if i := bytes.IndexByte(chunk, '\n'); i >= 0 {
t.buf = append(t.buf, chunk[:i+1]...)
t.done = true
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
}
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
}
80 changes: 66 additions & 14 deletions sandboxd/server/relay_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -8,19 +8,29 @@ import (
"net"
"net/http"
"net/http/httptest"
"slices"
"strings"
"testing"
"testing/iotest"
"time"
)

const (
execFrame = `{"v":1,"op":"exec","argv":["echo","42"],"detach":false}` + "\n"
"github.com/cocoonstack/sandbox/protocol/wire"
"github.com/cocoonstack/sandbox/sandboxd/pool"
)

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) {
Expand Down Expand Up @@ -127,12 +137,46 @@ 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 := 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, []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
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)
}
})
}
}
}

Expand Down Expand Up @@ -187,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)
Expand Down