The runtime serves a bundle its module's state (novox/hq ADR 0201)

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.
This commit is contained in:
jochen
2026-10-04 02:45:07 +02:00
parent 47cdfad754
commit 6cba894f01
8 changed files with 693 additions and 0 deletions
+2
View File
@@ -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.
+357
View File
@@ -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 `<module>.<name>`.
//
// **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.<server>`, `<machine>.<server>`) 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 <module>.<name>, 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)
}
+23
View File
@@ -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)
}
}
}
+38
View File
@@ -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
+12
View File
@@ -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
}
+66
View File
@@ -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)
})
}
+137
View File
@@ -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), " | "))
})
}