From 6cba894f01e427eb4f75c20820ed5fce43b54ed1 Mon Sep 17 00:00:00 2001 From: jochen Date: Sun, 4 Oct 2026 02:45:07 +0200 Subject: [PATCH] The runtime serves a bundle its module's state (novox/hq ADR 0201) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit mesh/state.get, put, delete, keys and watch on the stdio channel, from the buckets the membership issues. A watch hands the current values without deletions, then every change, and is answered once the current values are delivered. Refused with the reason: state not issued, a reader's write, a value with a credential-named field — the bus alone would answer a refused write with a timeout. --- node-tools/internal/bus/bus.go | 2 + node-tools/internal/bus/state.go | 357 ++++++++++++++++++++++ node-tools/internal/bus/state_test.go | 23 ++ node-tools/internal/launch/launch.go | 38 +++ node-tools/internal/meshtest/meshtest.go | 12 + node-tools/internal/runtime/runtime.go | 66 ++++ node-tools/internal/runtime/state_test.go | 137 +++++++++ node-tools/test/fixtures/state-keeper.mjs | 58 ++++ 8 files changed, 693 insertions(+) create mode 100644 node-tools/internal/bus/state.go create mode 100644 node-tools/internal/bus/state_test.go create mode 100644 node-tools/internal/runtime/state_test.go create mode 100755 node-tools/test/fixtures/state-keeper.mjs diff --git a/node-tools/internal/bus/bus.go b/node-tools/internal/bus/bus.go index d1a7af3..5854cdd 100644 --- a/node-tools/internal/bus/bus.go +++ b/node-tools/internal/bus/bus.go @@ -60,6 +60,8 @@ type Membership struct { Emits string `json:"emits"` Reaches map[string][]string `json:"reaches,omitempty"` Tools string `json:"tools"` + // State is every bucket the module's code may reach, by the name it uses (novox/hq ADR 0201). + State []StateIssued `json:"state,omitempty"` } // Served is an address a tool is answered on; `{tool}` stands for the tool's name. diff --git a/node-tools/internal/bus/state.go b/node-tools/internal/bus/state.go new file mode 100644 index 0000000..deb39a9 --- /dev/null +++ b/node-tools/internal/bus/state.go @@ -0,0 +1,357 @@ +package bus + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "regexp" + "sort" + "strings" + "sync" + "time" + + "github.com/nats-io/nats.go/jetstream" +) + +// A module's state on the bus (novox/hq ADR 0201): key-value buckets the controller creates from what +// the module declared, and issues to each assignment in its membership by the name the module uses — +// its own state by the local name, another module's as `.`. +// +// **The runtime keeps each module to its own buckets.** One account per machine carries every module +// on it, so the bus enforces only the union; and a write the bus refuses reaches the writer as a +// timeout, not a refusal (measured, novox/hq research 024). So what a module may reach is decided here, +// from its membership, and refused with the reason before anything is sent. + +// StateIssued is one bucket an assignment may reach (ADR 0201). +type StateIssued struct { + Name string `json:"name"` + Bucket string `json:"bucket"` + Writes bool `json:"writes,omitempty"` +} + +// StateEntry is one key's current value. +type StateEntry struct { + Key string `json:"key"` + Value json.RawMessage `json:"value"` + Revision uint64 `json:"revision"` +} + +// StateChange is one change a watch delivers: a key put or deleted. +type StateChange struct { + State string `json:"state"` + Key string `json:"key"` + Op string `json:"op"` + Value json.RawMessage `json:"value,omitempty"` + Revision uint64 `json:"revision"` + // Current is true for a value that was there when the watch began, false for a change since. + Current bool `json:"current"` +} + +// StateTimeout bounds one state operation on the bus. +var StateTimeout = 10 * time.Second + +// stateKey is a key the bus can hold: letters, digits and `-/_=.`, no leading or trailing dot. A +// module's convention for naming a machine in a key (`all.`, `.`) fits it. +var stateKey = regexp.MustCompile(`^[-/_=a-zA-Z0-9]+(\.[-/_=a-zA-Z0-9]+)*$`) + +// issuedState is the bucket a module may reach by a name, and whether it may write it — or the reason +// it may not reach it at all. +func (c *Conn) issuedState(module, name string) (StateIssued, error) { + m := c.Membership(module) + if m == nil { + return StateIssued{}, fmt.Errorf("%s has no membership issued on %s yet, so no state of it is reachable "+ + "until the mesh issues one (novox/hq ADR 0201)", module, c.node) + } + var names []string + for _, s := range m.State { + if s.Name == name { + return s, nil + } + names = append(names, s.Name) + } + sort.Strings(names) + issued := "none" + if len(names) > 0 { + issued = strings.Join(names, ", ") + } + return StateIssued{}, fmt.Errorf("%s keeps and reads no state called %q: it declares what it keeps under "+ + "`state` and what it reads under `reads` as ., and was issued: %s (novox/hq ADR 0201)", + module, name, issued) +} + +func (c *Conn) bucket(ctx context.Context, s StateIssued) (jetstream.KeyValue, error) { + js, err := jetstream.New(c.nc) + if err != nil { + return nil, err + } + kv, err := js.KeyValue(ctx, s.Bucket) + if err != nil { + if errors.Is(err, jetstream.ErrBucketNotFound) { + return nil, fmt.Errorf("the state %q is issued and its bucket is not on the bus yet: the controller "+ + "creates it from the catalogue on its next raise", s.Name) + } + return nil, err + } + return kv, nil +} + +func checkKey(key string) error { + if !stateKey.MatchString(key) { + return fmt.Errorf("%q is not a key the bus can hold: letters, digits and -/_=, in dot-separated "+ + "names", key) + } + return nil +} + +// StateGet is one key's current value, or nil when it has none. +func (c *Conn) StateGet(module, name, key string) (*StateEntry, error) { + s, err := c.issuedState(module, name) + if err != nil { + return nil, err + } + if err := checkKey(key); err != nil { + return nil, err + } + ctx, cancel := context.WithTimeout(context.Background(), StateTimeout) + defer cancel() + kv, err := c.bucket(ctx, s) + if err != nil { + return nil, err + } + e, err := kv.Get(ctx, key) + if errors.Is(err, jetstream.ErrKeyNotFound) { + return nil, nil + } + if err != nil { + return nil, err + } + return &StateEntry{Key: e.Key(), Value: valueOf(e.Value()), Revision: e.Revision()}, nil +} + +// StatePut writes one key, as the module, where the module keeps the state. It answers the revision. +func (c *Conn) StatePut(module, name, key string, value json.RawMessage) (uint64, error) { + s, err := c.writable(module, name) + if err != nil { + return 0, err + } + if err := checkKey(key); err != nil { + return 0, err + } + if len(value) == 0 || !json.Valid(value) { + return 0, fmt.Errorf("a state value is JSON") + } + if field := credentialField(value); field != "" { + return 0, fmt.Errorf("%s's %s.%s carries a field %q, which names a credential: no secret is kept in state, "+ + "sealed or not — a bucket is a stream, and a machine joining a year later reads it whole. Name the "+ + "secret and fetch it on request/reply (novox/hq ADR 0201, design 32 §10)", module, name, key, field) + } + ctx, cancel := context.WithTimeout(context.Background(), StateTimeout) + defer cancel() + kv, err := c.bucket(ctx, s) + if err != nil { + return 0, err + } + return kv.Put(ctx, key, value) +} + +// StateDelete removes one key, as the module, where the module keeps the state. A key that was not +// there is not an error: what is asked for is that it is gone. +func (c *Conn) StateDelete(module, name, key string) error { + s, err := c.writable(module, name) + if err != nil { + return err + } + if err := checkKey(key); err != nil { + return err + } + ctx, cancel := context.WithTimeout(context.Background(), StateTimeout) + defer cancel() + kv, err := c.bucket(ctx, s) + if err != nil { + return err + } + return kv.Delete(ctx, key) +} + +// StateKeys is every key with a current value, sorted. +func (c *Conn) StateKeys(module, name string) ([]string, error) { + s, err := c.issuedState(module, name) + if err != nil { + return nil, err + } + ctx, cancel := context.WithTimeout(context.Background(), StateTimeout) + defer cancel() + kv, err := c.bucket(ctx, s) + if err != nil { + return nil, err + } + lister, err := kv.ListKeys(ctx) + if err != nil { + return nil, err + } + defer func() { _ = lister.Stop() }() + keys := []string{} + for k := range lister.Keys() { + keys = append(keys, k) + } + sort.Strings(keys) + return keys, nil +} + +func (c *Conn) writable(module, name string) (StateIssued, error) { + s, err := c.issuedState(module, name) + if err != nil { + return s, err + } + if !s.Writes { + owner := name + if dot := strings.LastIndex(name, "."); dot > 0 { + owner = name[:dot] + } + return s, fmt.Errorf("%s reads %s and does not keep it: only %s's own instances write it (novox/hq ADR 0201)", + module, name, owner) + } + return s, nil +} + +// StateWatch hands deliver the current value of every key matching the pattern — none that is +// deleted — and then every change, in order (ADR 0201). It returns once the current values are +// delivered; deliver is called from one goroutine, one change at a time, and an error from it is said +// and the watch goes on: state is not a queue, and the next change, or the next start, reads it again. +// The pattern is a key, with `*` for one name and `**` for the rest; empty is every key. +func (c *Conn) StateWatch(module, name, pattern string, deliver func(StateChange) error) (stop func(), err error) { + s, err := c.issuedState(module, name) + if err != nil { + return nil, err + } + filter := ">" + if pattern != "" { + parts := strings.Split(pattern, ".") + for i, p := range parts { + if p == "**" { + if i != len(parts)-1 { + return nil, fmt.Errorf("%q: `**` stands for the rest of a key, so it comes last", pattern) + } + parts[i] = ">" + continue + } + if p != "*" && !stateKey.MatchString(p) { + return nil, fmt.Errorf("%q is not a key pattern: names, `*` for one and `**` for the rest", pattern) + } + } + filter = strings.Join(parts, ".") + } + ctx, cancel := context.WithTimeout(context.Background(), StateTimeout) + kv, err := c.bucket(ctx, s) + cancel() + if err != nil { + return nil, err + } + watching, stopWatching := context.WithCancel(context.Background()) + w, err := kv.Watch(watching, filter) + if err != nil { + stopWatching() + return nil, err + } + var once sync.Once + stop = func() { + once.Do(func() { + _ = w.Stop() + stopWatching() + }) + } + current := make(chan struct{}) + go func() { + initial := true + for e := range w.Updates() { + if e == nil { + // The end of what was there when the watch began. + if initial { + initial = false + close(current) + } + continue + } + change := StateChange{State: name, Key: e.Key(), Revision: e.Revision(), Current: initial} + switch e.Operation() { + case jetstream.KeyValuePut: + change.Op = "put" + change.Value = valueOf(e.Value()) + default: + // A deletion among the current values is a key that is not there: not handed over + // (measured, research 024 — the server sends its marker among the initial values). + if initial { + continue + } + change.Op = "delete" + } + if err := deliver(change); err != nil { + c.Logf("[mesh-tools] %s did not take %s %s of %s: %v; its next change, or its next start, reads it again", + module, change.Op, change.Key, name, err) + } + } + if initial { + close(current) + } + }() + select { + case <-current: + return stop, nil + case <-time.After(StateTimeout): + stop() + return nil, fmt.Errorf("the current values of %s did not arrive in %s", name, StateTimeout) + } +} + +func valueOf(raw []byte) json.RawMessage { + if len(raw) == 0 || !json.Valid(raw) { + b, _ := json.Marshal(string(raw)) + return b + } + return json.RawMessage(raw) +} + +// credentialEndings are the ends of a field name that say its value is a credential. A guard against +// the ordinary mistake, not a determined one: a sealed value is plain text to anything inspecting it, +// so the rule is checked where it can be and said to be partial (ADR 0201). +var credentialEndings = []string{"password", "passwd", "secret", "token", "credential", "credentials", + "authorization", "apikey", "privatekey", "accesskey", "cookie"} + +// credentialField is the first field anywhere in a JSON value whose name says it is a credential, or "". +func credentialField(value json.RawMessage) string { + var v any + if json.Unmarshal(value, &v) != nil { + return "" + } + var walk func(any) string + walk = func(v any) string { + switch t := v.(type) { + case map[string]any: + keys := make([]string, 0, len(t)) + for k := range t { + keys = append(keys, k) + } + sort.Strings(keys) + for _, k := range keys { + norm := strings.NewReplacer("-", "", "_", "", " ", "").Replace(strings.ToLower(k)) + for _, end := range credentialEndings { + if strings.HasSuffix(norm, end) { + return k + } + } + if found := walk(t[k]); found != "" { + return found + } + } + case []any: + for _, x := range t { + if found := walk(x); found != "" { + return found + } + } + } + return "" + } + return walk(v) +} diff --git a/node-tools/internal/bus/state_test.go b/node-tools/internal/bus/state_test.go new file mode 100644 index 0000000..c19be0b --- /dev/null +++ b/node-tools/internal/bus/state_test.go @@ -0,0 +1,23 @@ +package bus + +import ( + "encoding/json" + "testing" +) + +// A value naming a credential anywhere in it is found (novox/hq ADR 0201) — the guard against the +// ordinary mistake — and a value that only mentions tokens as a count is not. +func TestACredentialNamedFieldIsFoundAnywhereInAValue(t *testing.T) { + for value, want := range map[string]string{ + `{"url":"http://x","headers":{"Authorization":"Bearer s"}}`: "Authorization", + `[{"name":"a","env":{"API_KEY":"s"}}]`: "API_KEY", + `{"refresh_token":"s"}`: "refresh_token", + `{"clientSecret":"s"}`: "clientSecret", + `{"maxTokens":4096,"name":"a","generation":3}`: "", + `"just a string"`: "", + } { + if got := credentialField(json.RawMessage(value)); got != want { + t.Errorf("%s: found %q, want %q", value, got, want) + } + } +} diff --git a/node-tools/internal/launch/launch.go b/node-tools/internal/launch/launch.go index f96a4fb..22a39f5 100644 --- a/node-tools/internal/launch/launch.go +++ b/node-tools/internal/launch/launch.go @@ -50,10 +50,16 @@ type Registration struct { // publishes, asks and subscribes on the module's behalf. Subscribe is asked each time the module's // code subscribes; the runtime binds the module's consumer once and calls deliver for every event, // acknowledging it on the bus when deliver returns nil. +// +// State and Watch reach the module's state (novox/hq ADR 0201): State answers `get`, `put`, `delete` +// and `keys`; Watch hands deliver the current values and then every change, and returns once the +// current values are delivered — the watch is the child's, and stops when the child does. type Bus interface { Publish(params json.RawMessage) error Ask(params json.RawMessage) (json.RawMessage, error) Subscribe(deliver func(envelope json.RawMessage) error) error + State(verb string, params json.RawMessage) (json.RawMessage, error) + Watch(params json.RawMessage, deliver func(change json.RawMessage) error) (stop func(), err error) } // EventTimeout bounds how long a bundle has to handle one event before it is offered again. @@ -104,6 +110,9 @@ type child struct { next int64 dead chan struct{} why error + // watches are the state watches this child asked for, stopped when it exits (ADR 0201): a child + // that starts again watches again, as its code runs again. + watches []func() } func (c *child) write(m any) error { @@ -216,6 +225,30 @@ func Start(module, entry string, env []string, mesh Bus, logf func(string, ...an subscriber = c mu.Unlock() err = mesh.Subscribe(deliver) + case "mesh/state.get", "mesh/state.put", "mesh/state.delete", "mesh/state.keys": + var answered json.RawMessage + if answered, err = mesh.State(strings.TrimPrefix(m.Method, "mesh/state."), m.Params); err == nil { + result = answered + } + case "mesh/state.watch": + // Each change is asked of this child, in order; the watch is answered once the + // current values have been handed over, so a bundle that awaits it has the whole + // of the state before it goes on (ADR 0201). + var stop func() + stop, err = mesh.Watch(m.Params, func(change json.RawMessage) error { + _, err := c.ask(module, "mesh/state", map[string]any{"change": change}, EventTimeout) + return err + }) + if err == nil { + c.mu.Lock() + if c.why != nil { + c.mu.Unlock() + stop() + } else { + c.watches = append(c.watches, stop) + c.mu.Unlock() + } + } default: if len(m.ID) > 0 { _ = c.write(map[string]any{"jsonrpc": "2.0", "id": m.ID, @@ -265,7 +298,12 @@ func Start(module, entry string, env []string, mesh Bus, logf func(string, ...an saidMu.Unlock() c.mu.Lock() c.why = errors.New(why) + watches := c.watches + c.watches = nil c.mu.Unlock() + for _, stop := range watches { + stop() + } close(c.dead) mu.Lock() // Only a child that was serving is brought back; one that died in its own handshake was diff --git a/node-tools/internal/meshtest/meshtest.go b/node-tools/internal/meshtest/meshtest.go index 337d3d4..f5bc1bf 100644 --- a/node-tools/internal/meshtest/meshtest.go +++ b/node-tools/internal/meshtest/meshtest.go @@ -182,3 +182,15 @@ func (m *Mesh) Pending(t *testing.T, node, module string) (ackPending, notDelive } return uint64(info.NumAckPending), info.NumPending } + +// Bucket makes a module's state afresh as the controller does from the catalogue (novox/hq ADR 0201), +// and answers it for writing what is there before a bundle starts. +func (m *Mesh) Bucket(t *testing.T, name string) nats.KeyValue { + t.Helper() + _ = m.js.DeleteKeyValue(name) + kv, err := m.js.CreateKeyValue(&nats.KeyValueConfig{Bucket: name, History: 1, MaxValueSize: 256 * 1024}) + if err != nil { + t.Fatal(err) + } + return kv +} diff --git a/node-tools/internal/runtime/runtime.go b/node-tools/internal/runtime/runtime.go index e2c3524..5293f82 100644 --- a/node-tools/internal/runtime/runtime.go +++ b/node-tools/internal/runtime/runtime.go @@ -578,3 +578,69 @@ func (b *moduleBus) Subscribe(deliver func(json.RawMessage) error) error { c.logf("[mesh-tools] %s's events arrive on its consumer %s", module, bus.ConsumerOf(c.conn.Node(), module)) return nil } + +// stateAsked is what a bundle names when it reaches its state (ADR 0201): the state by the name its +// module uses, a key, and for a put the value. +type stateAsked struct { + State string `json:"state"` + Key string `json:"key"` + Value json.RawMessage `json:"value"` +} + +// State answers a bundle's `get`, `put`, `delete` and `keys` on its module's state (ADR 0201). The +// runtime refuses, with the reason, a state the module was not issued and a write to one it only reads. +func (b *moduleBus) State(verb string, params json.RawMessage) (json.RawMessage, error) { + var asked stateAsked + if err := json.Unmarshal(params, &asked); err != nil || asked.State == "" { + return nil, fmt.Errorf("mesh/state.%s names no state: {state, key, value}", verb) + } + conn := b.all.conn + var answer any + switch verb { + case "get": + entry, err := conn.StateGet(b.module, asked.State, asked.Key) + if err != nil { + return nil, err + } + if entry == nil { + return json.RawMessage("null"), nil + } + answer = entry + case "put": + revision, err := conn.StatePut(b.module, asked.State, asked.Key, asked.Value) + if err != nil { + return nil, err + } + answer = map[string]any{"revision": revision} + case "delete": + if err := conn.StateDelete(b.module, asked.State, asked.Key); err != nil { + return nil, err + } + answer = map[string]any{} + case "keys": + keys, err := conn.StateKeys(b.module, asked.State) + if err != nil { + return nil, err + } + answer = keys + default: + return nil, fmt.Errorf("the runtime answers no mesh/state.%s", verb) + } + return json.Marshal(answer) +} + +// Watch hands a bundle its module's state as it is and as it changes (ADR 0201): `{state, key}`, the +// key a pattern with `*` and `**`, empty for every key. +func (b *moduleBus) Watch(params json.RawMessage, deliver func(json.RawMessage) error) (func(), error) { + var asked stateAsked + if err := json.Unmarshal(params, &asked); err != nil || asked.State == "" { + return nil, fmt.Errorf("mesh/state.watch names no state: {state, key}") + } + return b.all.conn.StateWatch(b.module, asked.State, asked.Key, func(change bus.StateChange) error { + raw, err := json.Marshal(change) + if err != nil { + return err + } + return deliver(raw) + }) +} diff --git a/node-tools/internal/runtime/state_test.go b/node-tools/internal/runtime/state_test.go new file mode 100644 index 0000000..f89bb71 --- /dev/null +++ b/node-tools/internal/runtime/state_test.go @@ -0,0 +1,137 @@ +package runtime + +import ( + "os" + "path/filepath" + "strings" + "testing" + + "github.com/novox/mesh-tools/node-tools/internal/bus" + mt "github.com/novox/mesh-tools/node-tools/internal/meshtest" +) + +// novox/hq ADR 0201: a module keeps its current state in buckets it declares, and its code reaches +// them through the runtime. The owner's instances write and read; a reader only reads; a watch hands +// the current values — none that is deleted — and then every change; the runtime refuses, with the +// reason, what the module was not issued, a write to state it only reads, and a value naming a +// credential. +func TestABundleKeepsAndWatchesItsStateThroughTheRuntime(t *testing.T) { + mesh := mt.New(t) + kv := mesh.Bucket(t, "keeper_servers") + // What was there before anything started: one value, and one key since deleted. + if _, err := kv.Put("all.one", []byte(`{"url":"http://one"}`)); err != nil { + t.Fatal(err) + } + if _, err := kv.Put("all.gone", []byte(`{}`)); err != nil { + t.Fatal(err) + } + if err := kv.Delete("all.gone"); err != nil { + t.Fatal(err) + } + + keeper := mt.MembershipOf("keeper", "anchor", false, nil) + keeper.State = []bus.StateIssued{{Name: "servers", Bucket: "keeper_servers", Writes: true}} + peeker := mt.MembershipOf("peeker", "anchor", false, nil) + peeker.State = []bus.StateIssued{{Name: "keeper.servers", Bucket: "keeper_servers"}} + mesh.Issue(t, keeper) + mesh.Issue(t, peeker) + + nodeTools := connect(t, "node-tools", "anchor") + asker := connect(t, "console", "workstation") + dir := t.TempDir() + keeperLog, peekerLog := filepath.Join(dir, "keeper.log"), filepath.Join(dir, "peeker.log") + said := &lines{} + stop, err := Run(nodeTools, []Served{ + {Module: "keeper", Entrypoints: []string{mt.Fixture("state-keeper.mjs")}}, + {Module: "peeker", Entrypoints: []string{mt.Fixture("state-keeper.mjs")}}, + }, map[string]map[string]string{ + "keeper": {"STATE_NAME": "servers", "STATE_LOG": keeperLog}, + "peeker": {"STATE_NAME": "keeper.servers", "STATE_LOG": peekerLog}, + }, said.logf) + if err != nil { + t.Fatal(err) + } + t.Cleanup(stop) + + // Both read the whole current state at start: the value, not the deleted key, then the end of it. + for _, log := range []string{keeperLog, peekerLog} { + waitFor(t, log, "current") + lines := read(t, log) + if lines[0] != `was put all.one {"url":"http://one"}` || lines[1] != "current" || strings.Contains(strings.Join(lines, "\n"), "all.gone") { + t.Fatalf("the current state was not handed over as it is:\n%s\nruntime:\n%s", strings.Join(lines, "\n"), said.all()) + } + } + + // The owner writes; every watch sees the change. + got, err := call(t, asker, "keeper.put@anchor", map[string]any{"state": "servers", "key": "anchor.two", "value": map[string]any{"url": "http://two"}}) + if err != nil { + t.Fatal(err) + } + if !strings.Contains(got, `"revision"`) { + t.Fatalf("a put answered %s", got) + } + waitFor(t, keeperLog, `now put anchor.two {"url":"http://two"}`) + waitFor(t, peekerLog, `now put anchor.two {"url":"http://two"}`) + + // Get and keys, by the owner and by the reader. + got, err = call(t, asker, "peeker.get@anchor", map[string]any{"state": "keeper.servers", "key": "anchor.two"}) + if err != nil { + t.Fatal(err) + } + if !strings.Contains(got, `"value":{"url":"http://two"}`) { + t.Fatalf("the reader's get answered %s", got) + } + got, err = call(t, asker, "peeker.get@anchor", map[string]any{"state": "keeper.servers", "key": "nothing.here"}) + if err != nil || got != "null" { + t.Fatalf("a key with no value answered %s, %v", got, err) + } + got, err = call(t, asker, "keeper.keys@anchor", map[string]any{"state": "servers"}) + if err != nil { + t.Fatal(err) + } + same(t, got, `["all.one","anchor.two"]`) + + // Refused, with the reason, before anything is sent. + for _, c := range []struct { + key, why string + args map[string]any + }{ + {"peeker.put@anchor", "peeker reads keeper.servers and does not keep it", + map[string]any{"state": "keeper.servers", "key": "x", "value": 1}}, + {"keeper.get@anchor", `keeper keeps and reads no state called "other"`, + map[string]any{"state": "other", "key": "x"}}, + {"keeper.put@anchor", `carries a field "accessToken", which names a credential`, + map[string]any{"state": "servers", "key": "x", "value": map[string]any{"headers": map[string]any{"accessToken": "s"}}}}, + {"keeper.put@anchor", "is not a key the bus can hold", + map[string]any{"state": "servers", "key": "a..b", "value": 1}}, + } { + if _, err := call(t, asker, c.key, c.args); err == nil || !strings.Contains(err.Error(), c.why) { + t.Errorf("%s %v: want refused for %q, got %v", c.key, c.args, c.why, err) + } + } + + // A delete reaches every watch. + if _, err := call(t, asker, "keeper.del@anchor", map[string]any{"state": "servers", "key": "all.one"}); err != nil { + t.Fatal(err) + } + waitFor(t, peekerLog, "now delete all.one") + waitFor(t, keeperLog, "now delete all.one") +} + +func read(t *testing.T, path string) []string { + t.Helper() + b, _ := os.ReadFile(path) + return strings.Split(strings.TrimSpace(string(b)), "\n") +} + +func waitFor(t *testing.T, path, line string) { + t.Helper() + mt.Until(t, func() error { + for _, l := range read(t, path) { + if l == line { + return nil + } + } + return errorf("%s has no %q yet: %s", filepath.Base(path), line, strings.Join(read(t, path), " | ")) + }) +} diff --git a/node-tools/test/fixtures/state-keeper.mjs b/node-tools/test/fixtures/state-keeper.mjs new file mode 100755 index 0000000..3260b7f --- /dev/null +++ b/node-tools/test/fixtures/state-keeper.mjs @@ -0,0 +1,58 @@ +#!/usr/bin/env node +// A bundle that keeps and reads state through the runtime (novox/hq ADR 0201), written against the +// protocol with no SDK: on start it watches STATE_NAME and writes every change it is handed to +// STATE_LOG; its tools put, get, delete and list keys of whichever state they name. +import { appendFileSync } from "node:fs"; +import { createInterface } from "node:readline"; + +const say = (m) => process.stdout.write(JSON.stringify({ jsonrpc: "2.0", ...m }) + "\n"); +let next = 1; +const waiting = new Map(); +const ask = (method, params) => new Promise((resolve, reject) => { + const id = `keeper-${next++}`; + waiting.set(id, { resolve, reject }); + say({ id, method, params }); +}); +const log = (line) => appendFileSync(process.env.STATE_LOG, line + "\n"); +const text = (v) => ({ content: [{ type: "text", text: JSON.stringify(v ?? null) }] }); +const tool = (name) => ({ name, description: name, inputSchema: { type: "object", properties: {} } }); +const TOOLS = ["put", "get", "del", "keys"].map(tool); + +createInterface({ input: process.stdin }).on("line", async (line) => { + const m = JSON.parse(line); + if (m.method === undefined && waiting.has(m.id)) { + const w = waiting.get(m.id); waiting.delete(m.id); + m.error ? w.reject(new Error(m.error.message)) : w.resolve(m.result); + return; + } + const p = m.params ?? {}; + switch (m.method) { + case "initialize": + say({ id: m.id, result: { protocolVersion: "2025-03-26", capabilities: { tools: {} }, serverInfo: { name: "keeper", version: "1" } } }); + // Watched once the runtime has the handshake: the current values, then every change. + ask("mesh/state.watch", { state: process.env.STATE_NAME }) + .then(() => log("current"), (e) => log("watch refused: " + e.message)); + return; + case "tools/list": + say({ id: m.id, result: { tools: TOOLS } }); + return; + case "mesh/state": { + const c = p.change; + log(`${c.current ? "was" : "now"} ${c.op} ${c.key}${c.value !== undefined ? " " + JSON.stringify(c.value) : ""}`); + say({ id: m.id, result: {} }); + return; + } + case "tools/call": { + const a = p.arguments ?? {}; + const verb = { put: "put", get: "get", del: "delete", keys: "keys" }[p.name]; + try { + say({ id: m.id, result: text(await ask(`mesh/state.${verb}`, { state: a.state, key: a.key, value: a.value })) }); + } catch (e) { + say({ id: m.id, result: { content: [{ type: "text", text: e.message }], isError: true } }); + } + return; + } + default: + if (m.id !== undefined) say({ id: m.id, error: { code: -32601, message: "no " + m.method } }); + } +});