diff --git a/packages/sandbox/daemon-go/internal/content/handler.go b/packages/sandbox/daemon-go/internal/content/handler.go index 22cd73cfe9..079dd57059 100644 --- a/packages/sandbox/daemon-go/internal/content/handler.go +++ b/packages/sandbox/daemon-go/internal/content/handler.go @@ -1,6 +1,6 @@ // 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. +// JSON-RPC 2.0 with describe, schema.get, blocks.list and blocks.apply, plus +// `PUT …/assets/` uploads) 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 @@ -29,32 +29,46 @@ type Options struct { Store *FSStore ServerName string ServerVersion string + // OnCommit is called after blocks.apply lands, with the block files + // (inside .deco/blocks) it wrote or deleted. + OnCommit func(files []string) + // OnAsset is called after an upload, with the stored file name. + OnAsset func(name string) // Logf receives errors that become Internal errors. Logf func(format string, args ...any) } -// Handler serves the protocol for one app root. +// Handler serves the protocol and uploads 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 + store *FSStore + serverName string + serverVersion string + pollIntervalMs int + maxCommitAttempts int + retryMinMs, retryMaxMs int + cache *bodyCache + onCommit func([]string) + onAsset func(string) + 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, + store: o.Store, + serverName: o.ServerName, + serverVersion: o.ServerVersion, + pollIntervalMs: defaultPollMs, + maxCommitAttempts: defaultCommitTries, + retryMinMs: 50, + retryMaxMs: 200, + cache: newBodyCache(defaultCacheBytes), + onCommit: o.OnCommit, + onAsset: o.OnAsset, + logFn: o.Logf, } if h.serverName == "" { h.serverName = "deco-blocks" diff --git a/packages/sandbox/daemon-go/internal/content/methods.go b/packages/sandbox/daemon-go/internal/content/methods.go index ad028f10b7..55d0308960 100644 --- a/packages/sandbox/daemon-go/internal/content/methods.go +++ b/packages/sandbox/daemon-go/internal/content/methods.go @@ -1,16 +1,20 @@ package content -// The read methods (describe, schema.get, blocks.list), ported from -// server/methods/read.ts. +// The four methods, ported from server/methods/read.ts and apply.ts. import ( + "math/rand" + "strconv" "strings" + "time" ) const ( - maxPublicKeyBytes = 16 * 1024 - assetsURLPrefix = "/assets/" - defaultPollMs = 2000 + maxOpsPerApply = 500 + maxPublicKeyBytes = 16 * 1024 + assetsURLPrefix = "/assets/" + defaultPollMs = 2000 + defaultCommitTries = 3 ) // servablePublicKey is the key describe may serve: a single PUBLIC KEY PEM @@ -147,3 +151,361 @@ func (h *Handler) blocksList(params any) (any, error) { "diagnostics", diagnostics, ), nil } + +// ---------------------------------------------------------------- blocks.apply + +type setOp struct { + name string + value any +} + +type guardOp struct { + name string + version any // a string, or nil (must not exist); unvalidated under __proto__ +} + +type normalizedApply struct { + set []setOp + delete []string // names to delete that aren't also set + ifMatch []guardOp +} + +func normalize(params any) (*normalizedApply, error) { + a := &normalizedApply{} + setNames := map[string]bool{} + if set, ok := asObject(prop(params, "set")); ok { + for _, k := range set.Keys() { + v, _ := set.Get(k) + a.set = append(a.set, setOp{name: k, value: v}) + setNames[k] = true + } + } + if list, ok := prop(params, "delete").([]any); ok { + seen := map[string]bool{} + for _, item := range list { + name := item.(string) + if seen[name] { + continue + } + seen[name] = true + if !setNames[name] { + a.delete = append(a.delete, name) + } + } + } + ops := len(a.set) + len(a.delete) + if ops > maxOpsPerApply { + return nil, errLimitExceeded(strconv.Itoa(ops)+" names in one blocks.apply; the limit is "+strconv.Itoa(maxOpsPerApply), NewObject("limit", "maxOpsPerApply")) + } + if guards, ok := asObject(prop(params, "ifMatch")); ok { + for _, k := range guards.Keys() { + v, _ := guards.Get(k) + a.ifMatch = append(a.ifMatch, guardOp{name: k, version: v}) + } + } + return a, nil +} + +func nameViolations(name string, vs []NameViolation) []Violation { + out := make([]Violation, len(vs)) + for i, v := range vs { + out[i] = Violation{Name: name, Rule: v.Reason, Message: v.Message} + } + return out +} + +// validateStatic checks everything that doesn't depend on stored content. +func validateStatic(a *normalizedApply, meta *Object) ([]Violation, map[string]string) { + var violations []Violation + bodies := map[string]string{} + bySpelling := map[string]string{} + for _, op := range a.set { + violations = append(violations, nameViolations(op.name, checkBlockName(op.name, nil))...) + key := SpellingKey(op.name) + if twin, ok := bySpelling[key]; ok { + violations = append(violations, Violation{Name: op.name, Rule: "spelling-collision", Message: `"` + op.name + `" and "` + twin + `" are spellings of the same entry`}) + } else { + bySpelling[key] = op.name + } + if _, ok := asObject(op.value); !ok { + violations = append(violations, Violation{Name: op.name, Rule: "not-an-object", Message: "an entry must be a JSON object"}) + continue + } + body := SerializeBlock(op.value) + if len(body) > maxBlockBytes { + violations = append(violations, Violation{Name: op.name, Rule: "too-large", Message: "the entry is over " + strconv.Itoa(maxBlockBytes) + " bytes"}) + continue + } + bodies[op.name] = body + violations = append(violations, checkSecrets(op.name, op.value, meta)...) + } + for _, name := range a.delete { + violations = append(violations, nameViolations(name, checkDeletedName(name))...) + } + return violations, bodies +} + +type spelled struct { + name string + entry loadedEntry +} + +func entriesBySpelling(c *loadedContent) map[string]spelled { + out := map[string]spelled{} + for _, e := range c.entries { + key, _ := FullyDecodeFileName(e.File) + out[key] = spelled{name: e.Name, entry: e.loadedEntry} + } + return out +} + +// validateAgainst checks the rules that depend on the existing entries. +func validateAgainst(c *loadedContent, bySpelling map[string]spelled, a *normalizedApply) []Violation { + var violations []Violation + existing := make([]string, 0, len(c.entries)) + for _, e := range c.entries { + existing = append(existing, e.Name) + } + for _, op := range a.set { + if _, ok := bySpelling[SpellingKey(op.name)]; ok { + continue + } + names := append([]string(nil), existing...) + for _, other := range a.set { + if other.name != op.name { + names = append(names, other.name) + } + } + for _, v := range checkBlockName(op.name, names) { + if v.Reason == "case-collision" { + violations = append(violations, Violation{Name: op.name, Rule: v.Reason, Message: v.Message}) + } + } + } + return violations +} + +// orderedSet keeps insertion order, like a JS Set. +type orderedSet struct { + items []string + has map[string]bool +} + +func (s *orderedSet) add(v string) { + if s.has == nil { + s.has = map[string]bool{} + } + if !s.has[v] { + s.has[v] = true + s.items = append(s.items, v) + } +} + +func (s *orderedSet) remove(v string) { + if !s.has[v] { + return + } + delete(s.has, v) + for i, x := range s.items { + if x == v { + s.items = append(s.items[:i], s.items[i+1:]...) + return + } + } +} + +type applyPlan struct { + put []filePut + delete []string + expected map[string]*string +} + +// plan writes encode(name) and deletes every other spelling; only guarded +// entries (all their spellings) become commit expectations. +func plan(c *loadedContent, bySpelling map[string]spelled, a *normalizedApply, bodies map[string]string) applyPlan { + versionOf := map[string]string{} + for _, f := range c.snapshot.Files { + if IsBlockFileName(f.File) { + versionOf[f.File] = f.Version + } + } + groupFiles := func(name string) []string { + var out []string + for _, f := range c.groups[SpellingKey(name)] { + out = append(out, f.File) + } + return out + } + p := applyPlan{expected: map[string]*string{}} + putIndex := map[string]int{} + var deletes orderedSet + for _, op := range a.set { + file := BlockFileName(op.name) + if i, ok := putIndex[file]; ok { + p.put[i].Content = bodies[op.name] + } else { + putIndex[file] = len(p.put) + p.put = append(p.put, filePut{File: file, Content: bodies[op.name]}) + } + for _, other := range groupFiles(op.name) { + if other != file { + deletes.add(other) + } + } + } + for _, name := range a.delete { + file := BlockFileName(name) + if _, ok := versionOf[file]; ok { + deletes.add(file) + } + if _, ok := bySpelling[SpellingKey(name)]; ok { + for _, other := range groupFiles(name) { + deletes.add(other) + } + } + } + for _, f := range p.put { + deletes.remove(f.File) + } + for _, g := range a.ifMatch { + for _, file := range append([]string{BlockFileName(g.name)}, groupFiles(g.name)...) { + if IsBlockFileName(file) { + if v, ok := versionOf[file]; ok { + v := v + p.expected[file] = &v + } else { + p.expected[file] = nil + } + } + } + } + p.delete = deletes.items + return p +} + +// checkIfMatch compares each guard against the entry under any spelling. +func checkIfMatch(bySpelling map[string]spelled, a *normalizedApply) *Object { + mismatches := NewObject() + failed := false + for _, g := range a.ifMatch { + var actual any + if e, ok := bySpelling[SpellingKey(g.name)]; ok { + actual = e.entry.Version + } + if actual != g.version { + // `mismatches.__proto__ = …` sets the prototype in JS: no own entry. + if g.name != "__proto__" { + mismatches.Set(g.name, NewObject("expected", g.version, "actual", actual)) + } + failed = true + } + } + if failed { + return mismatches + } + return nil +} + +// resultFromVersions reports every name written or deleted, plus null for +// each other spelling the commit deleted. +func resultFromVersions(a *normalizedApply, revision string, fileVersions map[string]string, deleted []string) *Object { + versions := NewObject() + touched := map[string]bool{} + for _, op := range a.set { + if v, ok := fileVersions[BlockFileName(op.name)]; ok { + versions.Set(op.name, v) + } else { + versions.Set(op.name, nil) + } + touched[SpellingKey(op.name)] = true + } + for _, name := range a.delete { + versions.Set(name, nil) + touched[SpellingKey(name)] = true + } + for _, file := range deleted { + name := BlockNameFromFile(file) + key, _ := FullyDecodeFileName(file) + if !versions.Has(name) && touched[key] { + versions.Set(name, nil) + } + } + return NewObject("revision", revision, "versions", versions) +} + +func (h *Handler) loadSchemaForApply() (*string, *Object, error) { + stored, err := h.store.readSchema() + if err != nil { + return nil, nil, err + } + if stored == nil { + return nil, nil, nil + } + meta, err := h.parsedSchemaFor(stored) + if err != nil { + return nil, nil, err + } + v := stored.Version + return &v, meta, nil +} + +func (h *Handler) backoff(attempt int) { + ms := (float64(h.retryMinMs) + rand.Float64()*float64(max(0, h.retryMaxMs-h.retryMinMs))) * float64(attempt) + if ms > 0 { + time.Sleep(time.Duration(ms * float64(time.Millisecond))) + } +} + +func (h *Handler) blocksApply(params any) (any, error) { + a, err := normalize(params) + if err != nil { + return nil, err + } + for attempt := 1; attempt <= h.maxCommitAttempts; attempt++ { + schemaVersion, meta, err := h.loadSchemaForApply() + if err != nil { + return nil, err + } + violations, bodies := validateStatic(a, meta) + snap, content, err := h.loadCurrentContent(false, nil) + if err != nil { + return nil, err + } + bySpelling := entriesBySpelling(content) + violations = append(violations, validateAgainst(content, bySpelling, a)...) + if len(violations) > 0 { + return nil, errInvalidBlock(violations) + } + if mismatches := checkIfMatch(bySpelling, a); mismatches != nil { + return nil, errConflict(mismatches) + } + p := plan(content, bySpelling, a, bodies) + if len(p.put) == 0 && len(p.delete) == 0 { + return resultFromVersions(a, snap.Revision, nil, nil), nil + } + result, err := h.store.commit(commitAttempt{ + Put: p.put, + Delete: p.delete, + Expected: p.expected, + CheckSchema: len(a.set) > 0, + ExpectedSchemaVersion: schemaVersion, + }) + if err != nil { + return nil, err + } + if !result.Stale { + if h.onCommit != nil { + files := make([]string, 0, len(p.put)+len(p.delete)) + for _, f := range p.put { + files = append(files, f.File) + } + h.onCommit(append(files, p.delete...)) + } + return resultFromVersions(a, result.Revision, result.Versions, p.delete), nil + } + if attempt < h.maxCommitAttempts { + h.backoff(attempt) + } + } + return nil, errUnavailable("storage kept changing; gave up after "+strconv.Itoa(h.maxCommitAttempts)+" commit attempts", 250) +} diff --git a/packages/sandbox/daemon-go/internal/content/params.go b/packages/sandbox/daemon-go/internal/content/params.go index 74fc7a0622..b7739d5792 100644 --- a/packages/sandbox/daemon-go/internal/content/params.go +++ b/packages/sandbox/daemon-go/internal/content/params.go @@ -1,11 +1,12 @@ package content -// Parameter validation for the read methods, ported from params.ts (a zod 4 +// Parameter validation for the four 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 ( + "strconv" "strings" ) @@ -59,9 +60,10 @@ func (is *issues) opaque(path string, v any, nullable bool) { } var paramShapes = map[string][]string{ - "describe": {}, - "schema.get": {"ifNoneMatch"}, - "blocks.list": {"ifNoneMatch"}, + "describe": {}, + "schema.get": {"ifNoneMatch"}, + "blocks.list": {"ifNoneMatch"}, + "blocks.apply": {"set", "delete", "ifMatch"}, } // validateParams checks params for method (present is false when the request @@ -84,6 +86,34 @@ func validateParams(method string, params any, present bool) *ProtocolError { switch key { case "ifNoneMatch": is.opaque(key, v, false) + case "set": + if _, ok := asObject(v); !ok { + is.add(key, "Invalid input: expected record, received "+jsTypeName(v)) + } + case "delete": + list, ok := v.([]any) + if !ok { + is.add(key, "Invalid input: expected array, received "+jsTypeName(v)) + break + } + for i, item := range list { + if _, ok := item.(string); !ok { + is.add(key+"."+strconv.Itoa(i), "Invalid input: expected string, received "+jsTypeName(item)) + } + } + case "ifMatch": + record, ok := asObject(v) + if !ok { + is.add(key, "Invalid input: expected record, received "+jsTypeName(v)) + break + } + for _, k := range record.Keys() { + if k == "__proto__" { + continue + } + item, _ := record.Get(k) + is.opaque(key+"."+k, item, true) + } } } var unknown []string diff --git a/packages/sandbox/daemon-go/internal/content/params_test.go b/packages/sandbox/daemon-go/internal/content/params_test.go index 37cd6e78ea..570802c4d6 100644 --- a/packages/sandbox/daemon-go/internal/content/params_test.go +++ b/packages/sandbox/daemon-go/internal/content/params_test.go @@ -34,4 +34,14 @@ func TestValidateParams(t *testing.T) { 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"}`) + fail("blocks.apply", `{"set":[]}`, "set: Invalid input: expected record, received array") + fail("blocks.apply", `{"delete":[1,"a",null]}`, "delete.0: Invalid input: expected string, received number; delete.2: Invalid input: expected string, received null") + fail("blocks.apply", `{"ifMatch":{"a.b":5,"10":true,"2":null,"c":""}}`, + "ifMatch.10: Invalid input: expected string, received boolean; ifMatch.a.b: Invalid input: expected string, received number; ifMatch.c: Too small: expected string to have >=1 characters") + fail("blocks.apply", `{"set":1,"delete":2,"ifMatch":3,"q":1}`, + `set: Invalid input: expected record, received number; delete: Invalid input: expected array, received number; ifMatch: Invalid input: expected record, received number; unknown parameter "q"`) + for _, removed := range []string{"ref", "refs", "requestKey", "ifSchemaMatch", "ifUnmodifiedSince"} { + fail("blocks.apply", `{"set":{},"`+removed+`":"x"}`, `unknown parameter "`+removed+`"`) + } + ok("blocks.apply", `{"set":{"a":[]},"delete":["a"],"ifMatch":{"a":null,"__proto__":5}}`) } diff --git a/packages/sandbox/daemon-go/internal/content/rpc.go b/packages/sandbox/daemon-go/internal/content/rpc.go index 50418984bf..cc0009dd0f 100644 --- a/packages/sandbox/daemon-go/internal/content/rpc.go +++ b/packages/sandbox/daemon-go/internal/content/rpc.go @@ -10,7 +10,7 @@ import ( const maxBatchCalls = 10 -var methodNames = map[string]bool{"describe": true, "schema.get": true, "blocks.list": true} +var methodNames = map[string]bool{"describe": true, "schema.get": true, "blocks.list": true, "blocks.apply": true} var envelopeMembers = map[string]bool{"jsonrpc": true, "id": true, "method": true, "params": true} @@ -93,8 +93,10 @@ func (h *Handler) run(e *envelope) (any, error) { return h.describe() case "schema.get": return h.schemaGet(e.params) - default: + case "blocks.list": return h.blocksList(e.params) + default: + return h.blocksApply(e.params) } }