diff --git a/packages/sandbox/daemon-go/internal/content/content.go b/packages/sandbox/daemon-go/internal/content/content.go new file mode 100644 index 0000000000..9b14328e46 --- /dev/null +++ b/packages/sandbox/daemon-go/internal/content/content.go @@ -0,0 +1,250 @@ +package content + +// Turns a storage snapshot into the entry map, applying the file-name rule: +// decode once, resolve spellings, report the rest. Ported from +// server/content.ts and server/bodyCache.ts. + +import ( + "container/list" + "sort" + "strconv" + "sync" +) + +const ( + maxBlockBytes = 1 * 1024 * 1024 + maxSnapshotReads = 3 + defaultCacheBytes = 32 * 1024 * 1024 +) + +type parsedBody struct { + ok bool + value *Object + kind string // invalid-json | not-an-object | too-large + message string + bytes int +} + +// parseBody parses one stored entry, enforcing the per-entry byte limit. +func parseBody(text string) parsedBody { + bytes := len(text) + if bytes > maxBlockBytes { + return parsedBody{kind: "too-large", bytes: bytes, message: "the file is over " + strconv.Itoa(maxBlockBytes) + " bytes"} + } + v, err := ParseJSON(text) + if err != nil { + // OPEN: the message is this parser's, not V8's. + return parsedBody{kind: "invalid-json", bytes: bytes, message: err.Error()} + } + o, ok := asObject(v) + if !ok { + return parsedBody{kind: "not-an-object", bytes: bytes, message: "the file doesn't hold a JSON object"} + } + return parsedBody{ok: true, value: o, bytes: bytes} +} + +// bodyCache is a bounded LRU of parsed bodies keyed by file and version. +type bodyCache struct { + mu sync.Mutex + maxBytes int + total int + order *list.List // front: oldest + entries map[string]*list.Element +} + +type cacheItem struct { + key string + body parsedBody +} + +func newBodyCache(maxBytes int) *bodyCache { + return &bodyCache{maxBytes: maxBytes, order: list.New(), entries: map[string]*list.Element{}} +} + +func (c *bodyCache) get(file, version string) (parsedBody, bool) { + c.mu.Lock() + defer c.mu.Unlock() + e, ok := c.entries[file+"\x00"+version] + if !ok { + return parsedBody{}, false + } + c.order.MoveToBack(e) + return e.Value.(*cacheItem).body, true +} + +func (c *bodyCache) set(file, version string, body parsedBody) { + if body.bytes > c.maxBytes { + return + } + c.mu.Lock() + defer c.mu.Unlock() + key := file + "\x00" + version + if e, ok := c.entries[key]; ok { + c.total -= e.Value.(*cacheItem).body.bytes + c.order.Remove(e) + } + c.entries[key] = c.order.PushBack(&cacheItem{key: key, body: body}) + c.total += body.bytes + for c.total > c.maxBytes { + oldest := c.order.Front() + item := oldest.Value.(*cacheItem) + c.order.Remove(oldest) + delete(c.entries, item.key) + c.total -= item.body.bytes + } +} + +type loadedEntry struct { + File string + Version string + Value *Object // nil when the body wasn't needed +} + +type namedEntry struct { + Name string + loadedEntry +} + +type diagnostic struct { + File string + Kind string + Message string + Name string // shadowed only + Winner string // shadowed only +} + +func (d diagnostic) json() *Object { + o := NewObject("file", d.File, "kind", d.Kind) + if d.Kind == "shadowed" { + o.Set("name", d.Name) + o.Set("winner", d.Winner) + } + o.Set("message", d.Message) + return o +} + +type loadedContent struct { + snapshot storageSnapshot + // entries in file-name order of the winning files + entries []namedEntry + // every saved-block file by spelling key, keys in first-seen order + groupOrder []string + groups map[string][]storageFile + diagnostics []diagnostic + moved bool +} + +func (h *Handler) loadContent(snap storageSnapshot, readAll bool) (*loadedContent, error) { + var files []storageFile + for _, f := range snap.Files { + if IsBlockFileName(f.File) { + files = append(files, f) + } + } + c := &loadedContent{snapshot: snap, groups: map[string][]storageFile{}} + for _, f := range files { + key, _ := FullyDecodeFileName(f.File) + if _, ok := c.groups[key]; !ok { + c.groupOrder = append(c.groupOrder, key) + } + c.groups[key] = append(c.groups[key], f) + } + + parsed := map[string]parsedBody{} + var toRead []storageFile + for _, key := range c.groupOrder { + group := c.groups[key] + if !readAll && len(group) < 2 { + continue + } + for _, f := range group { + if f.Size > maxBlockBytes { + parsed[f.File] = parsedBody{kind: "too-large", bytes: int(f.Size), message: "the file is over " + strconv.Itoa(maxBlockBytes) + " bytes"} + continue + } + if hit, ok := h.cache.get(f.File, f.Version); ok { + parsed[f.File] = hit + } else { + toRead = append(toRead, f) + } + } + } + if len(toRead) > 0 { + names := make([]string, len(toRead)) + for i, f := range toRead { + names[i] = f.File + } + bodies, err := h.store.readFiles(names) + if err != nil { + return nil, err + } + for _, f := range toRead { + read, ok := bodies[f.File] + if !ok { + c.moved = true // vanished since the snapshot + continue + } + body := parseBody(read.Text) + // Cached under the version of the bytes actually read. + h.cache.set(f.File, read.Version, body) + if read.Version != f.Version { + c.moved = true + continue + } + parsed[f.File] = body + } + } + + var candidates []*spellingCandidate + for _, f := range files { + body, ok := parsed[f.File] + if !ok { + if !readAll { + candidates = append(candidates, &spellingCandidate{File: f.File, Version: f.Version}) + } + continue + } + if !body.ok { + c.diagnostics = append(c.diagnostics, diagnostic{File: f.File, Kind: body.kind, Message: body.message}) + continue + } + candidates = append(candidates, &spellingCandidate{File: f.File, Version: f.Version, HasPath: entryHasPath(body.value), Value: body.value}) + } + + for _, r := range resolveSpellings(candidates) { + c.entries = append(c.entries, namedEntry{Name: r.Name, loadedEntry: loadedEntry{File: r.Winner.File, Version: r.Winner.Version, Value: r.Winner.Value}}) + for _, loser := range r.Shadowed { + c.diagnostics = append(c.diagnostics, diagnostic{ + File: loser.File, Kind: "shadowed", Name: r.Name, Winner: r.Winner.File, + Message: `another spelling of "` + r.Name + `" wins: ` + r.Winner.File, + }) + } + } + sort.SliceStable(c.diagnostics, func(i, j int) bool { + return compareJS(c.diagnostics[i].File, c.diagnostics[j].File) < 0 + }) + return c, nil +} + +// loadCurrentContent takes a snapshot and loads it, taking a new one when a +// file changed while its body was read. shortCircuit ends the read right +// after the snapshot (a conditional read that's "not modified"). +func (h *Handler) loadCurrentContent(readAll bool, shortCircuit func(storageSnapshot) bool) (storageSnapshot, *loadedContent, error) { + for read := 1; read <= maxSnapshotReads; read++ { + snap, err := h.store.snapshot() + if err != nil { + return snap, nil, err + } + if shortCircuit != nil && shortCircuit(snap) { + return snap, nil, nil + } + c, err := h.loadContent(snap, readAll) + if err != nil { + return snap, nil, err + } + if !c.moved { + return snap, c, nil + } + } + return storageSnapshot{}, nil, errUnavailable("saved blocks kept changing while being read; retry shortly", 250) +} diff --git a/packages/sandbox/daemon-go/internal/content/handler.go b/packages/sandbox/daemon-go/internal/content/handler.go new file mode 100644 index 0000000000..22cd73cfe9 --- /dev/null +++ b/packages/sandbox/daemon-go/internal/content/handler.go @@ -0,0 +1,264 @@ +// Package content serves the Deco content protocol (`deco-content` v1: +// JSON-RPC 2.0; here the read methods describe, schema.get and blocks.list) +// over the sandbox's working tree. +// +// It is a port of the TypeScript reference in `@decocms/blocks/protocol` +// (server/, storage/fs, keys, secrets), behaviour-identical by design: same +// files and bytes, same error codes and shapes, same secret guard. The +// conformance suite `@decocms/blocks/protocol/conformance` runs against it in +// daemon-e2e, so drift fails CI. Port changes from the TS source; don't invent. +package content + +import ( + "bytes" + "compress/gzip" + "errors" + "fmt" + "io" + "log/slog" + "net/http" + "strconv" + "strings" + "sync" +) + +const maxRequestBytes = 8 * 1024 * 1024 + +// Options configure a Handler. +type Options struct { + Store *FSStore + ServerName string + ServerVersion string + // Logf receives errors that become Internal errors. + Logf func(format string, args ...any) +} + +// Handler serves the protocol for one app root. +type Handler struct { + store *FSStore + serverName string + serverVersion string + pollIntervalMs int + cache *bodyCache + logFn func(string, ...any) + schemaMu sync.Mutex + schemaVersion string + schemaMeta *Object +} + +// NewHandler builds a handler over a filesystem storage. +func NewHandler(o Options) *Handler { + h := &Handler{ + store: o.Store, + serverName: o.ServerName, + serverVersion: o.ServerVersion, + pollIntervalMs: defaultPollMs, + cache: newBodyCache(defaultCacheBytes), + logFn: o.Logf, + } + if h.serverName == "" { + h.serverName = "deco-blocks" + } + if h.serverVersion == "" { + h.serverVersion = "unknown" + } + return h +} + +func (h *Handler) logf(format string, args ...any) { + if h.logFn != nil { + h.logFn(format, args...) + return + } + slog.Warn("content protocol", "msg", fmt.Sprintf(format, args...)) +} + +// ---------------------------------------------------------------- HTTP plumbing + +var errBodyTooLarge = errors.New("body too large") + +type bodyEncodingError struct{ msg string } + +func (e *bodyEncodingError) Error() string { return e.msg } + +func readLimited(r io.Reader, limit int64) ([]byte, error) { + b, err := io.ReadAll(io.LimitReader(r, limit+1)) + if err != nil { + return nil, err + } + if int64(len(b)) > limit { + return nil, errBodyTooLarge + } + return b, nil +} + +// readBody reads at most limit bytes, decompressing a gzip body (the limit +// applies after decompression, so gzip can't bypass it). +func readBody(r *http.Request, limit int64) ([]byte, error) { + encoding := strings.ToLower(strings.TrimSpace(r.Header.Get("Content-Encoding"))) + if encoding == "" { + encoding = "identity" + } + if encoding == "identity" && r.ContentLength > limit { + return nil, errBodyTooLarge + } + if r.Body == nil || r.Body == http.NoBody { + return []byte{}, nil + } + if encoding == "identity" { + return readLimited(r.Body, limit) + } + if encoding != "gzip" { + return nil, &bodyEncodingError{msg: `unsupported Content-Encoding "` + encoding + `"`} + } + zr, err := gzip.NewReader(r.Body) + if err != nil { + return nil, &bodyEncodingError{msg: "the gzip body is corrupt"} + } + b, err := readLimited(zr, limit) + if err != nil { + if errors.Is(err, errBodyTooLarge) { + return nil, err + } + return nil, &bodyEncodingError{msg: "the gzip body is corrupt"} + } + return b, nil +} + +// jsNumber is JS's Number(string), for the q-values of Accept-Encoding. +func jsNumber(s string) float64 { + s = jsTrim(s) + if s == "" { + return 0 + } + switch s { + case "Infinity", "+Infinity": + return 1 + case "-Infinity": + return -1 + } + if len(s) > 2 && s[0] == '0' && strings.ContainsRune("xXoObB", rune(s[1])) { + base := map[byte]int{'x': 16, 'X': 16, 'o': 8, 'O': 8, 'b': 2, 'B': 2}[s[1]] + n, err := strconv.ParseUint(s[2:], base, 64) + if err != nil { + return -1 + } + return float64(n) + } + for _, c := range s { + if !(c >= '0' && c <= '9' || c == '.' || c == 'e' || c == 'E' || c == '+' || c == '-') { + return -1 // NaN: never > 0 + } + } + f, err := strconv.ParseFloat(s, 64) + if err != nil { + var ne *strconv.NumError + if errors.As(err, &ne) && errors.Is(ne.Err, strconv.ErrRange) { + return f + } + return -1 + } + return f +} + +func acceptsGzip(r *http.Request) bool { + header := strings.Join(r.Header.Values("Accept-Encoding"), ", ") + if header == "" { + return false + } + for _, part := range strings.Split(header, ",") { + pieces := strings.Split(strings.ToLower(strings.TrimSpace(part)), ";") + coding := strings.TrimSpace(pieces[0]) + if coding != "gzip" && coding != "*" { + continue + } + q, hasQ := "", false + for _, p := range pieces[1:] { + if p = strings.TrimSpace(p); strings.HasPrefix(p, "q=") { + q, hasQ = p[2:], true + break + } + } + if !hasQ || jsNumber(q) > 0 { + return true + } + } + return false +} + +const gzipThresholdBytes = 1024 + +// writeJSON is the TS jsonResponse: JSON, no-store, gzip when accepted. +func writeJSON(w http.ResponseWriter, r *http.Request, body string, status int, extra map[string]string) { + h := w.Header() + h.Set("Content-Type", "application/json; charset=utf-8") + h.Set("Cache-Control", "no-store") + h.Set("Vary", "Accept-Encoding") + for k, v := range extra { + h.Set(k, v) + } + if len(body) < gzipThresholdBytes || !acceptsGzip(r) { + h.Set("Content-Length", strconv.Itoa(len(body))) + w.WriteHeader(status) + io.WriteString(w, body) + return + } + var buf bytes.Buffer + zw := gzip.NewWriter(&buf) + io.WriteString(zw, body) + zw.Close() + h.Set("Content-Encoding", "gzip") + h.Set("Content-Length", strconv.Itoa(buf.Len())) + w.WriteHeader(status) + w.Write(buf.Bytes()) +} + +func errorBody(e *ProtocolError) string { + return Stringify(NewObject("jsonrpc", "2.0", "id", nil, "error", e.JSON()), 0) +} + +func isJSONContentType(r *http.Request) bool { + values := r.Header.Values("Content-Type") + if len(values) == 0 { + return false + } + t := strings.Join(values, ", ") + return strings.ToLower(strings.TrimSpace(strings.SplitN(t, ";", 2)[0])) == "application/json" +} + +// ServeRPC serves the protocol endpoint. +func (h *Handler) ServeRPC(w http.ResponseWriter, r *http.Request) { + if r.Method != http.MethodPost { + writeJSON(w, r, errorBody(errInvalidRequest("use POST")), 405, map[string]string{"Allow": "POST"}) + return + } + if !isJSONContentType(r) { + writeJSON(w, r, errorBody(errInvalidRequest("Content-Type must be application/json")), 415, nil) + return + } + raw, err := readBody(r, maxRequestBytes) + if err != nil { + var enc *bodyEncodingError + switch { + case errors.Is(err, errBodyTooLarge): + writeJSON(w, r, errorBody(errLimitExceeded("the request body is over "+strconv.Itoa(maxRequestBytes)+" bytes", NewObject("limit", "maxRequestBytes"))), 413, nil) + case errors.As(err, &enc): + writeJSON(w, r, errorBody(errInvalidRequest(enc.msg)), 415, nil) + default: + // The client went away mid-body. + writeJSON(w, r, errorBody(errInternal()), 500, nil) + } + return + } + text, ok := decodeUTF8(raw, true, true) + if !ok { + writeJSON(w, r, errorBody(errParse()), 200, nil) + return + } + body, err := ParseJSON(text) + if err != nil { + writeJSON(w, r, errorBody(errParse()), 200, nil) + return + } + writeJSON(w, r, h.dispatch(body), 200, nil) +} diff --git a/packages/sandbox/daemon-go/internal/content/methods.go b/packages/sandbox/daemon-go/internal/content/methods.go new file mode 100644 index 0000000000..ad028f10b7 --- /dev/null +++ b/packages/sandbox/daemon-go/internal/content/methods.go @@ -0,0 +1,149 @@ +package content + +// The read methods (describe, schema.get, blocks.list), ported from +// server/methods/read.ts. + +import ( + "strings" +) + +const ( + maxPublicKeyBytes = 16 * 1024 + assetsURLPrefix = "/assets/" + defaultPollMs = 2000 +) + +// servablePublicKey is the key describe may serve: a single PUBLIC KEY PEM +// block of at most 16 KiB; anything else is never broadcast. +func servablePublicKey(text *string) *string { + if text == nil || len(*text) > maxPublicKeyBytes || strings.Contains(*text, "PRIVATE KEY") { + return nil + } + if publicKeyDerFromPem(*text) == nil { + return nil + } + return text +} + +func (h *Handler) describe() (any, error) { + desc := h.store.describe() + stored, err := h.store.readSecretsPublicKey() + if err != nil { + return nil, err + } + publicKey := servablePublicKey(stored) + if stored != nil && publicKey == nil { + h.logf(".deco/secrets.pub isn't a single PUBLIC KEY PEM block; describe reports no key") + } + var secrets any + if publicKey != nil { + secrets = NewObject("publicKey", *publicKey) + } + return NewObject( + "protocol", "deco-content", + "version", NewObject("major", 1.0, "minor", 0.0), + "server", NewObject("name", h.serverName, "version", h.serverVersion), + "kind", "working-tree", + "readOnly", false, + "root", desc.Root, + "schemaFormat", "deco-meta@1", + "pollIntervalMs", float64(h.pollIntervalMs), + "preview", nil, + "assets", NewObject("dir", desc.AssetsDir, "urlPrefix", assetsURLPrefix, "maxBytes", float64(desc.AssetsMaxBytes)), + "secrets", secrets, + ), nil +} + +// parseSchema parses schema text; a file caught mid-write is Unavailable, +// never served torn. +func parseSchema(text string) (*Object, error) { + v, err := ParseJSON(text) + if err != nil { + return nil, errUnavailable("the schema file is being written; retry shortly", 500) + } + o, ok := asObject(v) + if !ok { + return nil, errUnavailable("the schema file doesn't hold a JSON object", 500) + } + return o, nil +} + +func ifNoneMatchOf(params any) (string, bool) { + v, ok := prop(params, "ifNoneMatch").(string) + return v, ok +} + +func (h *Handler) schemaGet(params any) (any, error) { + stored, err := h.store.readSchema() + if err != nil { + return nil, err + } + if stored == nil { + // No schema yet is a state, not an error; the snapshot still refuses a + // site without a .deco folder. + if _, err := h.store.snapshot(); err != nil { + return nil, err + } + return NewObject("notModified", false, "version", nil, "resolvedRef", nil, "schema", nil), nil + } + if inm, ok := ifNoneMatchOf(params); ok && inm == stored.Version { + return NewObject("notModified", true, "version", stored.Version), nil + } + meta, err := h.parsedSchemaFor(stored) + if err != nil { + return nil, err + } + return NewObject("notModified", false, "version", stored.Version, "resolvedRef", nil, "schema", meta), nil +} + +// parsedSchemaFor parses a stored schema, reusing the last parse of the same version. +func (h *Handler) parsedSchemaFor(stored *storedSchema) (*Object, error) { + h.schemaMu.Lock() + if h.schemaVersion == stored.Version && h.schemaMeta != nil { + meta := h.schemaMeta + h.schemaMu.Unlock() + return meta, nil + } + h.schemaMu.Unlock() + meta, err := parseSchema(stored.Text) + if err != nil { + return nil, err + } + h.schemaMu.Lock() + h.schemaVersion, h.schemaMeta = stored.Version, meta + h.schemaMu.Unlock() + return meta, nil +} + +func (h *Handler) blocksList(params any) (any, error) { + inm, hasINM := ifNoneMatchOf(params) + snap, content, err := h.loadCurrentContent(true, func(s storageSnapshot) bool { + return hasINM && inm == s.Revision + }) + if err != nil { + return nil, err + } + if content == nil { + return NewObject("notModified", true, "revision", snap.Revision, "resolvedRef", nil), nil + } + blocks, versions := NewObject(), NewObject() + for _, e := range content.entries { + if e.Value == nil { + continue + } + blocks.Set(e.Name, e.Value) + versions.Set(e.Name, e.Version) + } + diagnostics := make([]any, len(content.diagnostics)) + for i, d := range content.diagnostics { + diagnostics[i] = d.json() + } + return NewObject( + "notModified", false, + "revision", snap.Revision, + "resolvedRef", nil, + "blocks", blocks, + "versions", versions, + "diagnostics", diagnostics, + ), nil +} diff --git a/packages/sandbox/daemon-go/internal/content/params.go b/packages/sandbox/daemon-go/internal/content/params.go new file mode 100644 index 0000000000..74fc7a0622 --- /dev/null +++ b/packages/sandbox/daemon-go/internal/content/params.go @@ -0,0 +1,114 @@ +package content + +// Parameter validation for the read methods, ported from params.ts (a zod 4 +// strictObject per method). Unknown parameters are refused, so a guard the +// server doesn't understand never turns into an unguarded write. Messages +// follow zod 4's wording; like zod, a `__proto__` key is never validated. + +import ( + "strings" +) + +const maxOpaqueLength = 1024 + +func jsTypeName(v any) string { + switch v.(type) { + case nil: + return "null" + case bool: + return "boolean" + case float64: + return "number" + case string: + return "string" + case []any: + return "array" + } + return "object" +} + +type issues []string + +func (is *issues) add(path, message string) { + *is = append(*is, path+": "+message) +} + +// opaque is z.string().min(1).max(1024) (nullable: also null). +func (is *issues) opaque(path string, v any, nullable bool) { + if v == nil && nullable { + return + } + if s, ok := v.(string); ok { + n := jsLength(s) + if n < 1 { + is.add(path, "Too small: expected string to have >=1 characters") + } else if n > maxOpaqueLength { + is.add(path, "Too big: expected string to have <=1024 characters") + } + return + } + is.add(path, "Invalid input: expected string, received "+jsTypeName(v)) + // zod 4 still runs the length checks on anything with a length. + if list, ok := v.([]any); ok { + if len(list) < 1 { + is.add(path, "Too small: expected array to have >=1 items") + } else if len(list) > maxOpaqueLength { + is.add(path, "Too big: expected array to have <=1024 items") + } + } +} + +var paramShapes = map[string][]string{ + "describe": {}, + "schema.get": {"ifNoneMatch"}, + "blocks.list": {"ifNoneMatch"}, +} + +// validateParams checks params for method (present is false when the request +// had no `params`, which reads as {}). +func validateParams(method string, params any, present bool) *ProtocolError { + if !present { + return nil + } + obj, ok := asObject(params) + if !ok { + return errInvalidParams("params must be an object") + } + shape := paramShapes[method] + var is issues + for _, key := range shape { + v, has := obj.Get(key) + if !has { + continue + } + switch key { + case "ifNoneMatch": + is.opaque(key, v, false) + } + } + var unknown []string + for _, k := range obj.Keys() { + if k == "__proto__" { + continue + } + known := false + for _, s := range shape { + if s == k { + known = true + break + } + } + if !known { + unknown = append(unknown, `"`+k+`"`) + } + } + if len(unknown) == 1 { + is = append(is, "unknown parameter "+unknown[0]) + } else if len(unknown) > 1 { + is = append(is, "unknown parameters "+strings.Join(unknown, ", ")) + } + if len(is) > 0 { + return errInvalidParams(strings.Join(is, "; ")) + } + return nil +} diff --git a/packages/sandbox/daemon-go/internal/content/params_test.go b/packages/sandbox/daemon-go/internal/content/params_test.go new file mode 100644 index 0000000000..37cd6e78ea --- /dev/null +++ b/packages/sandbox/daemon-go/internal/content/params_test.go @@ -0,0 +1,37 @@ +package content + +import "testing" + +// Ported from params.test.ts; the messages are zod 4's, captured from the +// published package. +func TestValidateParams(t *testing.T) { + ok := func(method, params string) { + t.Helper() + if err := validateParams(method, mustParse(t, params), true); err != nil { + t.Errorf("%s %s: %v", method, params, err.Message) + } + } + fail := func(method, params, message string) { + t.Helper() + err := validateParams(method, mustParse(t, params), true) + if err == nil || err.Code != CodeInvalidParams || err.Message != message { + t.Errorf("%s %s:\n got %+v\nwant %q", method, params, err, message) + } + } + if validateParams("describe", nil, false) != nil { + t.Error("absent params read as {}") + } + ok("describe", `{}`) + fail("describe", `null`, "params must be an object") + fail("describe", `[]`, "params must be an object") + fail("describe", `"x"`, "params must be an object") + fail("describe", `{"a":1}`, `unknown parameter "a"`) + fail("describe", `{"a":1,"b":2}`, `unknown parameters "a", "b"`) + ok("describe", `{"__proto__":5}`) + fail("blocks.list", `{"ifNoneMatch":5}`, "ifNoneMatch: Invalid input: expected string, received number") + fail("blocks.list", `{"ifNoneMatch":""}`, "ifNoneMatch: Too small: expected string to have >=1 characters") + fail("blocks.list", `{"ifNoneMatch":[]}`, "ifNoneMatch: Invalid input: expected string, received array; ifNoneMatch: Too small: expected array to have >=1 items") + fail("blocks.list", `{"zz":1,"ifNoneMatch":5}`, `ifNoneMatch: Invalid input: expected string, received number; unknown parameter "zz"`) + fail("schema.get", `{"ifNoneMatch":null}`, "ifNoneMatch: Invalid input: expected string, received null") + ok("schema.get", `{"ifNoneMatch":"v"}`) +} diff --git a/packages/sandbox/daemon-go/internal/content/rpc.go b/packages/sandbox/daemon-go/internal/content/rpc.go new file mode 100644 index 0000000000..50418984bf --- /dev/null +++ b/packages/sandbox/daemon-go/internal/content/rpc.go @@ -0,0 +1,145 @@ +package content + +// The JSON-RPC 2.0 layer, ported from server/rpc.ts: every request needs an +// id; a batch runs in order, holds at most 10 calls and isn't atomic. + +import ( + "errors" + "strings" +) + +const maxBatchCalls = 10 + +var methodNames = map[string]bool{"describe": true, "schema.get": true, "blocks.list": true} + +var envelopeMembers = map[string]bool{"jsonrpc": true, "id": true, "method": true, "params": true} + +type envelope struct { + id any // string or float64 + method string + params any + hasParams bool +} + +func parseEnvelope(value any) (*envelope, *ProtocolError, any) { + raw, isObject := asObject(value) + var id any + if raw != nil { + switch v := prop(raw, "id").(type) { + case string, float64: + id = v + } + } + if !isObject { + return nil, errInvalidRequest("a request must be an object"), nil + } + if jsonrpc, ok := prop(raw, "jsonrpc").(string); !ok || jsonrpc != "2.0" { + return nil, errInvalidRequest(`jsonrpc must be "2.0"`), id + } + if id == nil { + return nil, errInvalidRequest("every request needs a string or number id"), nil + } + method, ok := prop(raw, "method").(string) + if !ok { + return nil, errInvalidRequest("method must be a string"), id + } + for _, key := range raw.Keys() { + if !envelopeMembers[key] { + return nil, errInvalidRequest(`unknown request member "` + key + `"`), id + } + } + if !methodNames[method] { + return nil, errMethodNotFound(method), id + } + params, has := raw.Get("params") + return &envelope{id: id, method: method, params: params, hasParams: has}, nil, id +} + +func errorResponse(id any, e *ProtocolError) string { + return Stringify(NewObject("jsonrpc", "2.0", "id", id, "error", e.JSON()), 0) +} + +// toProtocolError maps storage errors to protocol errors; nil for unknown ones. +func toProtocolError(err error) *ProtocolError { + var pe *ProtocolError + var nf *storageNotFound + var inv *storageInvalidFile + var un *storageUnavailable + switch { + case errors.As(err, &pe): + return pe + case errors.As(err, &nf): + return errNotFound(nf.msg) + case errors.As(err, &inv): + name := BlockNameFromFile(inv.file) + return errInvalidBlock([]Violation{{Name: name, Rule: "unsupported-name", Message: `this storage can't hold the entry "` + name + `"`}}) + case errors.As(err, &un): + return errUnavailable(un.msg, un.retryAfterMs) + } + return nil +} + +func (h *Handler) run(e *envelope) (any, error) { + if perr := validateParams(e.method, e.params, e.hasParams); perr != nil { + return nil, perr + } + if e.method != "describe" { + if err := h.store.checkContained(); err != nil { + return nil, err + } + } + switch e.method { + case "describe": + return h.describe() + case "schema.get": + return h.schemaGet(e.params) + default: + return h.blocksList(e.params) + } +} + +func (h *Handler) call(value any) (out string) { + e, perr, id := parseEnvelope(value) + if perr != nil { + return errorResponse(id, perr) + } + defer func() { + // An unexpected throw (a URIError from a lone surrogate, a bug) is an + // Internal error for this call only, as in the TS server. + if r := recover(); r != nil { + if _, isURI := r.(uriError); !isURI { + h.logf("content protocol: %s panicked: %v", e.method, r) + } + out = errorResponse(e.id, errInternal()) + } + }() + result, err := h.run(e) + if err != nil { + if known := toProtocolError(err); known != nil { + return errorResponse(e.id, known) + } + h.logf("content protocol: %s failed: %v", e.method, err) + return errorResponse(e.id, errInternal()) + } + return Stringify(NewObject("jsonrpc", "2.0", "id", e.id, "result", result), 0) +} + +// dispatch runs a parsed body (one request or a batch) and returns the +// serialized response body. +func (h *Handler) dispatch(body any) string { + list, isBatch := body.([]any) + if !isBatch { + return h.call(body) + } + if len(list) == 0 { + return errorResponse(nil, errInvalidRequest("an empty batch")) + } + if len(list) > maxBatchCalls { + return errorResponse(nil, errLimitExceeded("a batch holds at most 10 calls", NewObject("limit", "maxBatchCalls"))) + } + parts := make([]string, len(list)) + for i, item := range list { + parts[i] = h.call(item) + } + return "[" + strings.Join(parts, ",") + "]" +} diff --git a/packages/sandbox/daemon-go/internal/content/secrets_test.go b/packages/sandbox/daemon-go/internal/content/secrets_test.go new file mode 100644 index 0000000000..9dd20f8197 --- /dev/null +++ b/packages/sandbox/daemon-go/internal/content/secrets_test.go @@ -0,0 +1,227 @@ +package content + +import ( + "bytes" + "crypto/rand" + "crypto/rsa" + "crypto/x509" + "encoding/base64" + "encoding/pem" + "reflect" + "strings" + "testing" +) + +// Ported from secrets.test.ts, ciphertext.test.ts and __tests__/fixtures.ts. + +const schemaFixtureJSON = `{ + "manifest": {"blocks": { + "sections": {"hero": {"$ref": "#/definitions/aGVybw=="}, "newsletter": {"$ref": "#/definitions/bmV3c2xldHRlcg=="}}, + "loaders": {"multivariate": {"$ref": "#/definitions/bXY="}, "lazy": {"$ref": "#/definitions/bGF6eQ=="}, + "website/loaders/secret.ts": {"$ref": "#/definitions/djc="}}, + "content": {"settings": {"$ref": "#/definitions/c2V0dGluZ3M="}} + }}, + "schema": {"definitions": { + "aGVybw==": {"type": "object", "properties": {"title": {"type": "string"}, "padding": {"type": "string"}}}, + "bmV3c2xldHRlcg==": {"type": "object", "properties": {"listId": {"type": "string"}, "apiKey": {"type": "string", "format": "secret"}}}, + "c2V0dGluZ3M=": {"type": "object", "properties": { + "integrations": {"type": "array", "items": {"type": "object", "properties": {"token": {"$ref": "#/definitions/U2VjcmV0"}, "label": {"type": "string"}}}}, + "nested": {"anyOf": [{"type": "object", "properties": {"key": {"$ref": "#/definitions/U2VjcmV0"}}}]}, + "extra": {"type": "object", "additionalProperties": {"$ref": "#/definitions/U2VjcmV0"}} + }}, + "U2VjcmV0": {"type": "string", "format": "secret"}, + "bXY=": {"type": "object", "properties": {"variants": {"type": "array"}}}, + "bGF6eQ==": {"type": "object", "properties": {"value": {}}}, + "djc=": {"type": "object", "properties": {"name": {"type": "string"}, "encrypted": {"type": "string", "format": "secret"}}} + }} +}` + +func mustParse(t *testing.T, text string) any { + t.Helper() + v, err := ParseJSON(text) + if err != nil { + t.Fatalf("ParseJSON: %v", err) + } + return v +} + +func b64(n int, b byte) string { + return base64.RawURLEncoding.EncodeToString(bytes.Repeat([]byte{b}, n)) +} + +func ciphertextWithLengths(key, iv, ct int) string { + return "v1." + b64(key, 3) + "." + b64(iv, 4) + "." + b64(ct, 5) +} + +var validCiphertext = ciphertextWithLengths(384, 12, 32) + +func secretJSON(c string) string { return `{"__resolveType":"secret","ciphertext":"` + c + `"}` } + +func ruleList(t *testing.T, entry string) []string { + t.Helper() + meta, _ := asObject(mustParse(t, schemaFixtureJSON)) + var out []string + for _, v := range checkSecrets("entry", mustParse(t, entry), meta) { + out = append(out, v.Rule+"@"+*v.Pointer) + } + return out +} + +func TestCiphertextFormat(t *testing.T) { + for _, c := range []string{ciphertextWithLengths(256, 12, 16), ciphertextWithLengths(384, 12, 32), ciphertextWithLengths(512, 12, 1024)} { + if !isWellFormedCiphertext(c) { + t.Errorf("refused %s…", c[:20]) + } + } + refused := []any{ + "", "v1.", "v1.hunter2", "v1.my-api-key_123", "v1.QUJD.ZGVm", + "v2" + validCiphertext[2:], ciphertextWithLengths(255, 12, 16), ciphertextWithLengths(384, 16, 32), + ciphertextWithLengths(384, 12, 15), validCiphertext + ".QUJD", validCiphertext + "==", + validCiphertext[:len(validCiphertext)-1] + "+", strings.Replace(validCiphertext, ".", ". ", 1), + "hunter2", 42.0, nil, + } + for _, c := range refused { + if isWellFormedCiphertext(c) { + t.Errorf("accepted %v", c) + } + } + // Non-canonical base64url (stray low bits) is refused: one value, one spelling. + if _, ok := decodeBase64URL("QUJE"); !ok { + t.Error("canonical") + } + if _, ok := decodeBase64URL("QUI"); !ok { + t.Error("unpadded") + } + if _, ok := decodeBase64URL("QUJ"); ok { + t.Error("stray low bits") + } +} + +func publicKeyPEM(t *testing.T) string { + t.Helper() + key, err := rsa.GenerateKey(rand.Reader, 2048) + if err != nil { + t.Fatal(err) + } + der, _ := x509.MarshalPKIXPublicKey(&key.PublicKey) + return string(pem.EncodeToMemory(&pem.Block{Type: "PUBLIC KEY", Bytes: der})) +} + +func TestPublicKeyDerFromPem(t *testing.T) { + pub := publicKeyPEM(t) + if publicKeyDerFromPem(pub) == nil || publicKeyDerFromPem("\n "+pub+"\n\n") == nil { + t.Error("a single PUBLIC KEY block is accepted, surrounding whitespace too") + } + for _, bad := range []string{ + "", "hello", pub + pub, pub + "trailing text", + strings.ReplaceAll(pub, "PUBLIC KEY", "RSA PRIVATE KEY"), + strings.Replace(pub, "-----END PUBLIC KEY-----", "-----END PRIVATE KEY-----", 1), + strings.Replace(pub, "M", "*", 1), + } { + if publicKeyDerFromPem(bad) != nil { + t.Errorf("accepted %q", bad[:min(len(bad), 40)]) + } + } + if s := servablePublicKey(&pub); s == nil { + t.Error("servable") + } + withPrivate := pub + "\n# PRIVATE KEY" + if servablePublicKey(&withPrivate) != nil { + t.Error("a PRIVATE KEY mention is never served") + } + huge := pub + strings.Repeat(" ", 16*1024) + if servablePublicKey(&huge) != nil { + t.Error("over 16 KiB") + } +} + +func TestSecretGuard(t *testing.T) { + newsletter := func(apiKey string) string { + return `{"__resolveType":"newsletter","listId":"l1","apiKey":` + apiKey + `}` + } + expect := func(entry string, want ...string) { + t.Helper() + got := ruleList(t, entry) + if len(want) == 0 { + want = nil + } + if !reflect.DeepEqual(got, want) { + t.Errorf("%s:\n got %v\nwant %v", entry, got, want) + } + } + expect(newsletter(`"hunter2"`), "secret-field@/apiKey") + for _, v := range []string{`42`, `null`, `{"value":"hunter2"}`, `{"__resolveType":"MyKey"}`} { + expect(newsletter(v), "secret-field@/apiKey") + } + expect(newsletter(secretJSON(validCiphertext))) + expect(newsletter(secretJSON("hunter2")), "secret-ciphertext@/apiKey") + expect(newsletter(`{"__resolveType":"secret"}`), "secret-ciphertext@/apiKey") + expect(`{"__resolveType":"newsletter","listId":"l1"}`) + expect(`{"__resolveType":"settings", + "integrations":[{"token":`+secretJSON(validCiphertext)+`,"label":"ok"},{"token":"plain","label":"leak"}], + "nested":{"key":"plain"},"extra":{"a":`+secretJSON(validCiphertext)+`,"b":"plain"}}`, + "secret-field@/integrations/1/token", "secret-field@/nested/key", "secret-field@/extra/b") + expect(`{"__resolveType":"page","sections":[{"__resolveType":"newsletter","apiKey":"plain"},{"__resolveType":"hero","title":"x"}]}`, + "secret-field@/sections/0/apiKey") + variants := func(values ...string) string { + var parts []string + for _, v := range values { + parts = append(parts, `{"rule":{"__resolveType":"always"},"value":{"__resolveType":"lazy","value":`+v+`}}`) + } + return newsletter(`{"__resolveType":"multivariate","variants":[` + strings.Join(parts, ",") + `]}`) + } + expect(variants(secretJSON(validCiphertext), secretJSON(validCiphertext))) + expect(variants(secretJSON(validCiphertext), `"plain"`), "secret-field@/apiKey/variants/1/value/value") + expect(newsletter(`{"__resolveType":"multivariate"}`), "secret-field@/apiKey") + expect(newsletter(`{"__resolveType":"website/flags/multivariate.ts","variants":[{"rule":{"__resolveType":"always"},"value":"plain"}]}`), + "secret-field@/apiKey/variants/0/value") + expect(`{"__resolveType":"hero","title":`+secretJSON("bad")+`,"list":[`+secretJSON(validCiphertext)+`]}`, "secret-ciphertext@/title") + expect(`{"__resolveType":"hero","title":"plain text is fine here"}`) + expect(`{"__resolveType":"unknown-type","apiKey":"not a known Secret field"}`) + expect(`{"__resolveType":"settings","extra":{"constructor":"plain","x/~y":"plain"}}`, + "secret-field@/extra/constructor", "secret-field@/extra/x~1~0y") + + // Where a schema has `properties`, an Object.prototype name is found there + // (`key in properties` in JS) and never falls back to additionalProperties. + box, _ := asObject(mustParse(t, `{"manifest":{"blocks":{"s":{"box":{"type":"object","properties":{"title":{"type":"string"}},"additionalProperties":{"type":"string","format":"secret"}}}}}}`)) + inherited := checkSecrets("e", mustParse(t, `{"__resolveType":"box","constructor":"plain","toString":"plain","other":"plain","title":"t"}`), box) + if len(inherited) != 1 || *inherited[0].Pointer != "/other" { + t.Errorf("inherited names: %+v", inherited) + } + + // Without a schema only the ciphertexts are checked. + noSchema := checkSecrets("e", mustParse(t, `{"__resolveType":"hero","title":`+secretJSON("bad")+`}`), nil) + if len(noSchema) != 1 || noSchema[0].Rule != "secret-ciphertext" { + t.Errorf("no schema: %+v", noSchema) + } + top := checkSecrets("e", mustParse(t, secretJSON("bad")), nil) + if len(top) != 1 || *top[0].Pointer != "" { + t.Errorf("top-level secret block: %+v", top) + } +} + +func TestLegacySecretLoaderIsExempt(t *testing.T) { + // The v7 loader's `encrypted` is marked format: secret but holds the + // site's own ciphertext: it is walked without its definition. + const loader = `"website/loaders/secret.ts"` + got := ruleList(t, `{"__resolveType":"settings","apps":{"__resolveType":`+loader+`,"name":"API_KEY","encrypted":"0a1b2c"}}`) + if got != nil { + t.Errorf("legacy loader: %v", got) + } + got = ruleList(t, `{"__resolveType":`+loader+`,"encrypted":`+secretJSON("bad")+`}`) + if !reflect.DeepEqual(got, []string{"secret-ciphertext@/encrypted"}) { + t.Errorf("a secret block inside the loader is still checked: %v", got) + } + got = ruleList(t, `{"__resolveType":"site/loaders/secret.ts","encrypted":"0a1b2c"}`) + if got != nil { + t.Errorf("any app's loaders/secret.ts: %v", got) + } +} + +func TestSecretGuardNamesTheEntry(t *testing.T) { + meta, _ := asObject(mustParse(t, schemaFixtureJSON)) + v := checkSecrets("Newsletter", mustParse(t, `{"__resolveType":"newsletter","apiKey":"x"}`), meta) + if len(v) != 1 || v[0].Name != "Newsletter" || v[0].Rule != "secret-field" || *v[0].Pointer != "/apiKey" { + t.Errorf("%+v", v) + } +}