From 127047edd665e5b73d1d318bd22422b3c5b24db8 Mon Sep 17 00:00:00 2001 From: jochen Date: Thu, 8 Oct 2026 18:23:09 +0200 Subject: [PATCH] Serve a module from a runtime on its own account instead of switching users, keep bus words from bundles, and never give up a channel's work (hq ADR 0259 revision) --- merge-check.sh | 6 + node-tools/internal/bus/seat.go | 94 +++++++++-- node-tools/internal/launch/launch.go | 22 ++- node-tools/internal/launch/runas.go | 68 -------- node-tools/internal/launch/runas_test.go | 48 ------ .../internal/runtime/own_account_test.go | 56 +++++++ node-tools/internal/runtime/runtime.go | 80 ++++++---- .../internal/runtime/seat_traffic_test.go | 149 ++++++++++++++++++ 8 files changed, 364 insertions(+), 159 deletions(-) delete mode 100644 node-tools/internal/launch/runas.go delete mode 100644 node-tools/internal/launch/runas_test.go create mode 100644 node-tools/internal/runtime/own_account_test.go create mode 100644 node-tools/internal/runtime/seat_traffic_test.go diff --git a/merge-check.sh b/merge-check.sh index d04b8a5..fa55750 100644 --- a/merge-check.sh +++ b/merge-check.sh @@ -33,6 +33,9 @@ if ! command -v node >/dev/null 2>&1 || [ ! -d node_modules/@novox/mesh-sdk ]; t # What of the console launches no bundle runs anyway: how an address resolves (novox/hq issue 287), # and how a search ranks (novox/hq ADR 0245). bundleless='^(TestASeatAndAModuleOfOneNameAreEachReached|TestSearch.*|TestASearchByCommandSaysWhatItReplaces)$' + # And what of the runtime launches no bundle: how a bundle's seat traffic reaches the bus, what its + # environment holds, and what a runtime on a module's own account serves (novox/hq ADR 0259). + runtimeless='^(TestARuntimeOnAModulesOwnAccountServesThatModuleAlone|TestNoBusWordReachesABundle|TestAToolCallNeverReachesASeatsEventAcceptOrProof|TestABundlesSeatTrafficGoesThroughItsModulesBus|TestALostCompareAndSetCarriesItsCode)$' fi if command -v gcc >/dev/null 2>&1; then CGO_ENABLED=1 go test -race -count=1 $packages @@ -43,4 +46,7 @@ fi if [ -n "$bundleless" ]; then CGO_ENABLED=0 go test -count=1 -run "$bundleless" ./internal/console fi +if [ -n "${runtimeless:-}" ]; then + CGO_ENABLED=0 go test -count=1 -run "$runtimeless" ./internal/runtime +fi echo "NOT TESTED HERE: node-tools/src (TypeScript) needs @novox/mesh-sdk from the package registry" diff --git a/node-tools/internal/bus/seat.go b/node-tools/internal/bus/seat.go index b07bf20..dc5713d 100644 --- a/node-tools/internal/bus/seat.go +++ b/node-tools/internal/bus/seat.go @@ -64,6 +64,20 @@ type SeatWorker struct { // ProofTimeout bounds how long a proof waits for its answer. var ProofTimeout = 15 * time.Second +// WorkHeartbeat is how often work a bundle is still doing is said to be in progress, well inside the +// worker's ack wait. +var WorkHeartbeat = 15 * time.Second + +// WorkBackoff is how long a piece of work the bundle did not take waits before it is offered again: from +// five seconds, doubling, to ten minutes. +func WorkBackoff(delivered uint64) time.Duration { + wait := 5 * time.Second + for i := uint64(1); i < delivered && wait < 10*time.Minute; i++ { + wait *= 2 + } + return min(wait, 10*time.Minute) +} + // SubjectMatches is a subject against a permission pattern, in the server's wildcards. func SubjectMatches(pattern, subject string) bool { p, s := strings.Split(pattern, "."), strings.Split(subject, ".") @@ -122,18 +136,29 @@ func newID() string { // SeatPublish submits an accept or says an event on a seat, as the module, into the stream that keeps it, // awaited and de-duplicated by its id: the same id twice is one message. A proof is never published here. func (c *Conn) SeatPublish(module, subject string, body json.RawMessage, id string) (uint64, error) { + seq, _, err := c.SeatPublishSaid(module, subject, body, id) + return seq, err +} + +// SeatPublishSaid is SeatPublish, also saying whether the bus took it as a duplicate of one published before +// under the same id: the same message, kept once. +func (c *Conn) SeatPublishSaid(module, subject string, body json.RawMessage, id string) (uint64, bool, error) { if strings.Contains(subject, ".proof.") { - return 0, fmt.Errorf("%s is a proof, asked and answered, never kept: ask it as one", subject) + return 0, false, fmt.Errorf("%s is a proof, asked and answered, never kept: ask it as one", subject) } if err := c.mayPublish(module, subject); err != nil { - return 0, err + return 0, false, err } if len(body) == 0 || !json.Valid(body) { - return 0, fmt.Errorf("what is said on %s is JSON", subject) + return 0, false, fmt.Errorf("what is said on %s is JSON", subject) } if id == "" { id = newID() } + // **The publisher's name is part of the id the bus de-duplicates by** (security review 2026-10-08): ids + // are the publisher's own, and one module could otherwise suppress another's warrant by publishing first + // under the same id. + id = module + "." + id msg := nats.NewMsg(subject) msg.Header.Set("x-event-id", id) msg.Header.Set("x-source", module) @@ -143,9 +168,9 @@ func (c *Conn) SeatPublish(module, subject string, body json.RawMessage, id stri msg.Data = body ack, err := c.js.PublishMsg(msg, nats.MsgId(id)) if err != nil { - return 0, err + return 0, false, err } - return ack.Sequence, nil + return ack.Sequence, ack.Duplicate, nil } // SeatProve asks a proof and answers its reply: core request and reply, on no stream. @@ -263,12 +288,38 @@ func (c *Conn) SeatTake(module, consumer string, deliver func(Work) error) (func for k := range msg.Headers() { headers[k] = msg.Headers().Get(k) } - if err := deliver(Work{Worker: w.Consumer, Subject: msg.Subject(), Body: json.RawMessage(msg.Data()), Headers: headers}); err != nil { - _ = msg.NakWithDelay(NakDelay) + // Kept alive while the bundle works on it, however long that takes: the ack wait is for a holder + // that died, not one that is slow. + done := make(chan struct{}) + go func() { + tick := time.NewTicker(WorkHeartbeat) + defer tick.Stop() + for { + select { + case <-done: + return + case <-tick.C: + _ = msg.InProgress() + } + } + }() + err := deliver(Work{Worker: w.Consumer, Subject: msg.Subject(), Body: json.RawMessage(msg.Data()), Headers: headers}) + close(done) + if err != nil { + // Offered again later and later, never given up on: a channel that is away keeps its work + // (security and correctness reviews of 2026-10-08). + delivered := uint64(1) + if meta, merr := msg.Metadata(); merr == nil { + delivered = meta.NumDelivered + } + _ = msg.NakWithDelay(WorkBackoff(delivered)) return } _ = msg.Ack() - }) + }, jetstream.ConsumeErrHandler(func(_ jetstream.ConsumeContext, err error) { + // A worker deleted, or the bus refusing it, is said — never a holder that silently takes nothing. + c.Logf("[mesh-tools] %s's worker %s: %v", module, w.Consumer, err) + })) if err != nil { return nil, fmt.Errorf("reading %s's worker %s: %w", module, w.Consumer, err) } @@ -286,6 +337,20 @@ func (c *Conn) SeatRecord(module, bucket, key string) (*StateEntry, error) { if err := checkKey(key); err != nil { return nil, err } + if bucket == "" { + // The one record its membership lists, when it lists one: an asker reads the router's asks without + // knowing the router's name. + var named []string + for _, r := range t.Records { + if b, _, ok := strings.Cut(strings.TrimPrefix(r, "$JS.API.DIRECT.GET.KV_"), "."); ok { + named = append(named, b) + } + } + if len(named) != 1 { + return nil, fmt.Errorf("%s reads %d records under its name; name the one meant", module, len(named)) + } + bucket = named[0] + } subject := "$JS.API.DIRECT.GET.KV_" + bucket + ".$KV." + bucket + "." + module + "." + key if !listed(t.Records, subject) { return nil, fmt.Errorf("%s reads no record of %s under its name: its membership lists %s (novox/hq ADR 0259)", @@ -319,6 +384,10 @@ func (c *Conn) StateCreate(module, name, key string, value json.RawMessage) (uin // StateUpdate writes one key only when its revision is still the one given, and answers the new revision; // a key changed meanwhile is refused, so of two writers that read one value only the first writes. func (c *Conn) StateUpdate(module, name, key string, value json.RawMessage, revision uint64) (uint64, error) { + if revision == 0 { + // Revision zero would be read by the server as "the key is new": an update names the value it read. + return 0, fmt.Errorf("an update names the revision it read; zero is none — create the key instead") + } return c.stateWrite(module, name, key, value, func(ctx context.Context, kv jetstream.KeyValue) (uint64, error) { return kv.Update(ctx, key, value, revision) }) @@ -327,6 +396,13 @@ func (c *Conn) StateUpdate(module, name, key string, value json.RawMessage, revi // ErrStateChanged is a compare-and-set that lost: the key was made or changed by another writer first. var ErrStateChanged = errors.New("the key was written by another writer first") +// changed is ErrStateChanged naming the key, with the code a bundle is answered it by. +type changed struct{ key string } + +func (e changed) Error() string { return e.key + ": " + ErrStateChanged.Error() } +func (e changed) Is(target error) bool { return target == ErrStateChanged } +func (e changed) ErrorCode() int { return -32010 } + func (c *Conn) stateWrite(module, name, key string, value json.RawMessage, write func(context.Context, jetstream.KeyValue) (uint64, error)) (uint64, error) { s, err := c.writable(module, name) @@ -351,7 +427,7 @@ func (c *Conn) stateWrite(module, name, key string, value json.RawMessage, } rev, err := write(ctx, kv) if errors.Is(err, jetstream.ErrKeyExists) || isWrongSequence(err) { - return 0, fmt.Errorf("%s.%s: %w", name, key, ErrStateChanged) + return 0, changed{name + "." + key} } return rev, err } diff --git a/node-tools/internal/launch/launch.go b/node-tools/internal/launch/launch.go index 007e3ef..dc68ae5 100644 --- a/node-tools/internal/launch/launch.go +++ b/node-tools/internal/launch/launch.go @@ -181,12 +181,6 @@ func Start(module, entry string, env []string, mesh Bus, logf func(string, ...an start := func() (*child, error) { cmd := exec.Command(entry) cmd.Env = env - // **A module that names an account of its own runs as it** (novox/hq ADR 0259 §8): its secrets and - // state are that account's, and the operator's account — which every agent runs as — reaches - // neither. Never root, never the operator's. - if err := runAs(cmd, env); err != nil { - return nil, err - } stdin, err := cmd.StdinPipe() if err != nil { return nil, err @@ -310,7 +304,7 @@ func Start(module, entry string, env []string, mesh Bus, logf func(string, ...an } if err != nil { _ = c.write(map[string]any{"jsonrpc": "2.0", "id": m.ID, - "error": map[string]any{"code": -32000, "message": err.Error()}}) + "error": map[string]any{"code": errorCode(err), "message": err.Error()}}) return } _ = c.write(map[string]any{"jsonrpc": "2.0", "id": m.ID, "result": result}) @@ -599,3 +593,17 @@ func (c *child) deathReason() error { } return errors.New("the bundle exited") } + +// CodeChanged is the error code a bundle is answered a lost compare-and-set with (mesh-sdk's +// stdio.CodeChanged): told by its code, never by its words (novox/hq ADR 0259). +const CodeChanged = -32010 + +// errorCode is the code an error is answered to a bundle with: a lost compare-and-set's own, else the +// general one. +func errorCode(err error) int { + var coded interface{ ErrorCode() int } + if errors.As(err, &coded) { + return coded.ErrorCode() + } + return -32000 +} diff --git a/node-tools/internal/launch/runas.go b/node-tools/internal/launch/runas.go deleted file mode 100644 index 1f9725f..0000000 --- a/node-tools/internal/launch/runas.go +++ /dev/null @@ -1,68 +0,0 @@ -package launch - -import ( - "fmt" - "os/exec" - "os/user" - "strconv" - "strings" - "syscall" -) - -// RunAs is the word in a bundle's environment naming the account it runs as (novox/hq ADR 0259 §8), and -// OperatorAccount the word naming the operator's account on this machine, which the runtime is given. -const ( - RunAs = "MESH_RUN_AS" - OperatorAccount = "MESH_OPERATOR_ACCOUNT" -) - -// lookupUser is user.Lookup, a variable so a test can name accounts this machine does not have. -var lookupUser = user.Lookup - -func word(env []string, name string) string { - for _, kv := range env { - if k, v, ok := strings.Cut(kv, "="); ok && k == name { - return v - } - } - return "" -} - -// runAs starts cmd as the account its environment names, with that account's home, or leaves it as the -// runtime's own when none is named. It refuses root and the operator's account: a module asks for an -// account of its own to keep what it holds from the agents, and every agent runs as the operator. -func runAs(cmd *exec.Cmd, env []string) error { - name := word(env, RunAs) - if name == "" { - return nil - } - if operator := word(env, OperatorAccount); operator != "" && name == operator { - return fmt.Errorf("%s names the operator's account %s; a module runs as an account of its own, never "+ - "the one every agent runs as (novox/hq ADR 0259)", RunAs, name) - } - u, err := lookupUser(name) - if err != nil { - return fmt.Errorf("%s names the account %s, which this machine does not have: the module makes it with "+ - "a user resource (novox/hq ADR 0259): %w", RunAs, name, err) - } - uid, err := strconv.ParseUint(u.Uid, 10, 32) - if err != nil { - return fmt.Errorf("the account %s has no usable uid %q", name, u.Uid) - } - gid, err := strconv.ParseUint(u.Gid, 10, 32) - if err != nil { - return fmt.Errorf("the account %s has no usable gid %q", name, u.Gid) - } - if uid == 0 || name == "root" { - return fmt.Errorf("%s names root; a module that asks for an account of its own is not given root", RunAs) - } - cmd.SysProcAttr = &syscall.SysProcAttr{Credential: &syscall.Credential{Uid: uint32(uid), Gid: uint32(gid)}} - var kept []string - for _, kv := range cmd.Env { - if !strings.HasPrefix(kv, "HOME=") && !strings.HasPrefix(kv, "USER=") && !strings.HasPrefix(kv, "LOGNAME=") { - kept = append(kept, kv) - } - } - cmd.Env = append(kept, "HOME="+u.HomeDir, "USER="+name, "LOGNAME="+name) - return nil -} diff --git a/node-tools/internal/launch/runas_test.go b/node-tools/internal/launch/runas_test.go deleted file mode 100644 index 69629c9..0000000 --- a/node-tools/internal/launch/runas_test.go +++ /dev/null @@ -1,48 +0,0 @@ -package launch - -import ( - "errors" - "os/exec" - "os/user" - "strings" - "testing" -) - -// novox/hq ADR 0259 §8: a module naming an account of its own runs as it, never as root or the operator. -func TestABundleRunsAsTheAccountItNamesAndNeverRootOrTheOperator(t *testing.T) { - was := lookupUser - t.Cleanup(func() { lookupUser = was }) - lookupUser = func(name string) (*user.User, error) { - switch name { - case "telegram": - return &user.User{Username: "telegram", Uid: "961", Gid: "961", HomeDir: "/var/lib/telegram"}, nil - case "root": - return &user.User{Username: "root", Uid: "0", Gid: "0", HomeDir: "/root"}, nil - case "toor": - return &user.User{Username: "toor", Uid: "0", Gid: "0", HomeDir: "/root"}, nil - } - return nil, errors.New("unknown user") - } - cmd := exec.Command("/bin/true") - cmd.Env = []string{"HOME=/root", "PATH=/usr/bin"} - if err := runAs(cmd, []string{RunAs + "=telegram", OperatorAccount + "=jo"}); err != nil { - t.Fatal(err) - } - if c := cmd.SysProcAttr.Credential; c == nil || c.Uid != 961 || c.Gid != 961 { - t.Fatalf("not started as telegram: %+v", cmd.SysProcAttr) - } - if env := strings.Join(cmd.Env, " "); !strings.Contains(env, "HOME=/var/lib/telegram") || strings.Contains(env, "HOME=/root") { - t.Fatalf("the account's home is not its own: %s", env) - } - for name, want := range map[string]string{"jo": "operator's account", "root": "names root", "toor": "names root", - "nobody-here": "does not have"} { - err := runAs(exec.Command("/bin/true"), []string{RunAs + "=" + name, OperatorAccount + "=jo"}) - if err == nil || !strings.Contains(err.Error(), want) { - t.Errorf("%s: %v, want a refusal saying %q", name, err, want) - } - } - plain := exec.Command("/bin/true") - if err := runAs(plain, nil); err != nil || plain.SysProcAttr != nil { - t.Error("a bundle naming no account was changed") - } -} diff --git a/node-tools/internal/runtime/own_account_test.go b/node-tools/internal/runtime/own_account_test.go new file mode 100644 index 0000000..eaf433d --- /dev/null +++ b/node-tools/internal/runtime/own_account_test.go @@ -0,0 +1,56 @@ +package runtime + +import ( + "strings" + "testing" + + "github.com/novox/mesh-tools/node-tools/internal/bus" +) + +// novox/hq ADR 0259 §8: a runtime on a module's own account serves that module alone, and the machine's +// runtime never serves itself. +func TestARuntimeOnAModulesOwnAccountServesThatModuleAlone(t *testing.T) { + if got, err := ServedModulesFrom("telegram=/b/telegram/tools/telegram", "telegram"); err != nil || len(got) != 1 { + t.Fatalf("a module's own runtime could not serve it: %v", err) + } + if _, err := ServedModulesFrom("telegram=/b/t,dunst=/b/d", "telegram"); err == nil || !strings.Contains(err.Error(), "serves telegram alone") { + t.Errorf("a module's own runtime served another: %v", err) + } + if _, err := ServedModulesFrom("node-tools=/b/n", "node-tools"); err == nil { + t.Error("the machine's runtime served itself") + } +} + +// No bus word reaches a bundle (security review of 2026-10-08), not even one its module's words name. +func TestNoBusWordReachesABundle(t *testing.T) { + base := []string{"MESH_BROKER_FILE=/etc/mesh/broker", "MESH_BROKER_URL=nats://x", "MESH_TOOL_ENV={}", + "MESH_TOOL_MODULES=a=/b", "MESH_CONSOLE_LISTEN=127.0.0.1:1", "PATH=/usr/bin"} + env := strings.Join(BundleEnv(base, map[string]string{"MESH_BROKER_FILE": "/again", "OWN": "1"}, "telegram", "anchor"), "\n") + for _, never := range []string{"MESH_BROKER_FILE", "MESH_BROKER_URL", "MESH_TOOL_ENV", "MESH_TOOL_MODULES", "MESH_CONSOLE_LISTEN"} { + if strings.Contains(env, never+"=") { + t.Errorf("%s reached the bundle", never) + } + } + for _, want := range []string{"PATH=/usr/bin", "OWN=1", "MESH_MODULE=telegram", "MESH_NODE=anchor"} { + if !strings.Contains(env, want) { + t.Errorf("%s did not reach the bundle", want) + } + } +} + +// A tool call reaches only a tool's subject: whatever key a caller names, it is never a seat's event, accept +// or proof — so no tool call says a choice, a link or a code, or submits an ask, in anybody's name. +func TestAToolCallNeverReachesASeatsEventAcceptOrProof(t *testing.T) { + for _, key := range []string{"seat:intake.event.choice.telegram", "seat:intake.proof.code.telegram", + "seat:operator-channel.accept.ask.mesh-delivery", "intake.event", "seat:operator-channel.event.decided.x@anchor", + "telegram.anything", "x"} { + subject, err := bus.ToolSubject(key, "node-tools") + if err != nil { + continue + } + tokens := strings.Split(subject, ".") + if len(tokens) < 4 || tokens[3] != "tool" { + t.Errorf("%q reaches %s, which is not a tool's subject", key, subject) + } + } +} diff --git a/node-tools/internal/runtime/runtime.go b/node-tools/internal/runtime/runtime.go index c100d5d..edc1498 100644 --- a/node-tools/internal/runtime/runtime.go +++ b/node-tools/internal/runtime/runtime.go @@ -53,6 +53,13 @@ type ListedTool struct { Subjects []string `json:"subjects,omitempty"` } +// NodeRuntimeModule is the module that is a machine's runtime, serving every module on it. +const NodeRuntimeModule = "node-tools" + +// runtimeOnly are the words of the runtime's own environment no bundle is given. +var runtimeOnly = map[string]bool{ToolEnv: true, ToolModules: true, "MESH_BROKER_FILE": true, + "MESH_BROKER_URL": true, "MESH_CONSOLE_LISTEN": true} + // ServedModulesFrom reads MESH_TOOL_MODULES: `=` entries, comma-separated, several // per module. The one-module form — a bare path, or the runtime's own module — is the per-module // containers' (to-be 38 WP4c) and refused here: the node's runtime imports nothing. @@ -66,9 +73,20 @@ func ServedModulesFrom(spec, own string) ([]Served, error) { } module, path, ok := strings.Cut(entry, "=") module, path = strings.TrimSpace(module), strings.TrimSpace(path) - if !ok || module == "" || path == "" || module == own { - return nil, fmt.Errorf("%s: %q is not = of another module; the node's "+ - "runtime launches the bundles it is given and imports nothing (novox/hq ADR 0193)", ToolModules, entry) + if !ok || module == "" || path == "" { + return nil, fmt.Errorf("%s: %q is not =; the runtime launches the bundles it "+ + "is given and imports nothing (novox/hq ADR 0193)", ToolModules, entry) + } + // The node's runtime serves the machine's modules and never itself. **A runtime on a module's own + // account serves that module and nothing else** (novox/hq ADR 0259 §8): the router and a channel that + // proves its sender reach the bus on an account of their own, never the machine's runtime. + switch { + case own == NodeRuntimeModule && module == own: + return nil, fmt.Errorf("%s: %q is the runtime itself; it launches the bundles it is given and "+ + "imports nothing (novox/hq ADR 0193)", ToolModules, entry) + case own != NodeRuntimeModule && module != own: + return nil, fmt.Errorf("%s: %q is another module's; a runtime on %s's own account serves %s alone "+ + "(novox/hq ADR 0259)", ToolModules, entry, own, own) } if _, seen := by[module]; !seen { order = append(order, module) @@ -135,28 +153,7 @@ func Run(conn *bus.Conn, served []Served, envs map[string]map[string]string, log node := conn.Node() base := os.Environ() - envFor := func(module string) []string { - words := map[string]string{} - for _, kv := range base { - if k, v, ok := strings.Cut(kv, "="); ok && k != ToolEnv { - words[k] = v - } - } - for k, v := range envs[module] { - words[k] = v - } - words["MESH_SERVED_MODULE"] = module - words["MESH_MODULE"] = module - if node != "" { - words["MESH_NODE"] = node - } - out := make([]string, 0, len(words)) - for k, v := range words { - out = append(out, k+"="+v) - } - sort.Strings(out) - return out - } + envFor := func(module string) []string { return BundleEnv(base, envs[module], module, node) } // The module's events, for every child of it that subscribes (ADR 0198): one consumer per module, // bound the first time any of its children subscribes, each event handed to every child that did. @@ -710,11 +707,11 @@ func (b *moduleBus) Seat(verb string, params json.RawMessage, handOn func(string conn := b.all.conn switch verb { case "publish": - seq, err := conn.SeatPublish(b.module, asked.Subject, asked.Body, asked.ID) + seq, duplicate, err := conn.SeatPublishSaid(b.module, asked.Subject, asked.Body, asked.ID) if err != nil { return nil, err } - return json.Marshal(map[string]any{"sequence": seq}) + return json.Marshal(map[string]any{"sequence": seq, "duplicate": duplicate}) case "prove": return conn.SeatProve(b.module, asked.Subject, asked.Body) case "kinds": @@ -779,3 +776,32 @@ func (b *moduleBus) bind(key string, stop func()) { } b.all.seatBound[b.module+" "+key] = stop } + +// BundleEnv is what one module's bundle is started with: the runtime's environment without its own words, +// the module's composed words over it, and the module's and machine's names. **No bus word reaches a +// bundle** (security review of 2026-10-08): the credential and the bus's address are the runtime's, and a +// bundle reaches the bus only through the runtime. +func BundleEnv(base []string, given map[string]string, module, node string) []string { + words := map[string]string{} + for _, kv := range base { + if k, v, ok := strings.Cut(kv, "="); ok && !runtimeOnly[k] { + words[k] = v + } + } + for k, v := range given { + if !runtimeOnly[k] { + words[k] = v + } + } + words["MESH_SERVED_MODULE"] = module + words["MESH_MODULE"] = module + if node != "" { + words["MESH_NODE"] = node + } + out := make([]string, 0, len(words)) + for k, v := range words { + out = append(out, k+"="+v) + } + sort.Strings(out) + return out +} diff --git a/node-tools/internal/runtime/seat_traffic_test.go b/node-tools/internal/runtime/seat_traffic_test.go new file mode 100644 index 0000000..3986cc5 --- /dev/null +++ b/node-tools/internal/runtime/seat_traffic_test.go @@ -0,0 +1,149 @@ +package runtime + +import ( + "encoding/json" + "errors" + "strings" + "sync" + "testing" + "time" + + "github.com/nats-io/nats.go" + + "github.com/novox/mesh-tools/node-tools/internal/bus" + mt "github.com/novox/mesh-tools/node-tools/internal/meshtest" +) + +// A bundle's seat traffic through moduleBus.Seat (novox/hq ADR 0259 §3), against a real bus: publish under +// its own name with the id prefixed and a duplicate said; the record read without naming the bucket; a +// worker and a proof subject bound once however often a restarted child asks; and nothing beyond what the +// membership lists. +func TestABundlesSeatTrafficGoesThroughItsModulesBus(t *testing.T) { + mesh := mt.New(t) + nc, err := nats.Connect(mt.URL(t)) + if err != nil { + t.Fatal(err) + } + t.Cleanup(nc.Close) + js, _ := nc.JetStream() + if _, err := js.AddStream(&nats.StreamConfig{Name: "SEAT_OPERATOR_CHANNEL", Retention: nats.WorkQueuePolicy, + Subjects: []string{"mesh.seat.operator-channel.accept.>"}}); err != nil { + t.Fatal(err) + } + if _, err := js.AddConsumer("SEAT_OPERATOR_CHANNEL", &nats.ConsumerConfig{Durable: "SEAT_OPERATOR_CHANNEL_worker", + AckPolicy: nats.AckExplicitPolicy, FilterSubject: "mesh.seat.operator-channel.accept.>"}); err != nil { + t.Fatal(err) + } + kv, err := js.CreateKeyValue(&nats.KeyValueConfig{Bucket: "messenger_asks"}) + if err != nil { + t.Fatal(err) + } + _, _ = kv.Put("mesh-delivery.d-1", []byte(`{"state":"open"}`)) + asker := mt.MembershipOf("mesh-delivery", "anchor", false, nil) + asker.SeatTraffic = &bus.SeatTraffic{Publish: []string{"mesh.seat.operator-channel.accept.ask.mesh-delivery"}, + Records: []string{"$JS.API.DIRECT.GET.KV_messenger_asks.$KV.messenger_asks.mesh-delivery.>"}} + router := mt.MembershipOf("messenger", "anchor", false, nil) + router.SeatTraffic = &bus.SeatTraffic{Answers: []string{"mesh.seat.intake.proof.code.*"}, + Workers: []bus.SeatWorker{{Stream: "SEAT_OPERATOR_CHANNEL", Consumer: "SEAT_OPERATOR_CHANNEL_worker", + Filter: "mesh.seat.operator-channel.accept.>"}}} + mesh.Issue(t, asker) + mesh.Issue(t, router) + conn := connect(t, "node-tools", "anchor") + conn.Follow("mesh-delivery") + conn.Follow("messenger") + mt.Until(t, func() error { + if conn.Membership("mesh-delivery") == nil || conn.Membership("messenger") == nil { + return errors.New("not issued yet") + } + return nil + }) + all := &consumers{conn: conn, logf: t.Logf, of: map[string]*moduleEvents{}} + asking := all.forModule("mesh-delivery").(*moduleBus) + routing := all.forModule("messenger").(*moduleBus) + + publish := json.RawMessage(`{"subject":"mesh.seat.operator-channel.accept.ask.mesh-delivery","body":{"id":"d-1"},"id":"d-1"}`) + first, err := asking.Seat("publish", publish, nil) + if err != nil { + t.Fatal(err) + } + again, _ := asking.Seat("publish", publish, nil) + if !strings.Contains(string(first), `"duplicate":false`) || !strings.Contains(string(again), `"duplicate":true`) { + t.Errorf("the same id twice: %s, then %s", first, again) + } + if _, err := asking.Seat("publish", json.RawMessage(`{"subject":"mesh.seat.operator-channel.accept.ask.mesh-controller","body":{}}`), nil); err == nil { + t.Error("a module submitted under another's name") + } + got, err := asking.Seat("record", json.RawMessage(`{"key":"d-1"}`), nil) + if err != nil || !strings.Contains(string(got), `"state":"open"`) { + t.Errorf("the record read without its bucket: %s %v", got, err) + } + + var mu sync.Mutex + var taken []string + handOn := func(method string, params any) (json.RawMessage, error) { + mu.Lock() + defer mu.Unlock() + raw, _ := json.Marshal(params) + taken = append(taken, method+" "+string(raw)) + return json.RawMessage(`{}`), nil + } + for i := 0; i < 2; i++ { // a child that started again asks again + if _, err := routing.Seat("take", json.RawMessage(`{"worker":"SEAT_OPERATOR_CHANNEL_worker"}`), handOn); err != nil { + t.Fatal(err) + } + } + deadline := time.Now().Add(10 * time.Second) + for { + mu.Lock() + n := len(taken) + mu.Unlock() + if n >= 1 || time.Now().After(deadline) { + break + } + time.Sleep(50 * time.Millisecond) + } + time.Sleep(300 * time.Millisecond) + mu.Lock() + if len(taken) != 1 || !strings.Contains(taken[0], "mesh/work") || !strings.Contains(taken[0], `"x-event-id":"mesh-delivery.d-1"`) { + t.Errorf("taken %v", taken) + } + mu.Unlock() + if _, err := asking.Seat("take", json.RawMessage(`{"worker":"SEAT_OPERATOR_CHANNEL_worker"}`), handOn); err == nil { + t.Error("an asker took the router's work") + } + if _, err := asking.Seat("answer", json.RawMessage(`{"subject":"mesh.seat.intake.proof.code.*"}`), handOn); err == nil { + t.Error("an asker answered proofs") + } + if _, err := asking.Seat("anything", json.RawMessage(`{}`), nil); err == nil { + t.Error("an unknown verb was answered") + } +} + +// A lost compare-and-set is answered with its code, so the SDK knows it without reading words. +func TestALostCompareAndSetCarriesItsCode(t *testing.T) { + mesh := mt.New(t) + mesh.Bucket(t, "messenger_asks") + m := mt.MembershipOf("messenger", "anchor", false, nil) + m.State = []bus.StateIssued{{Name: "asks", Bucket: "messenger_asks", Writes: true}} + mesh.Issue(t, m) + conn := connect(t, "node-tools", "anchor") + conn.Follow("messenger") + mt.Until(t, func() error { + if conn.Membership("messenger") == nil { + return errors.New("not issued yet") + } + return nil + }) + b := (&consumers{conn: conn, logf: t.Logf, of: map[string]*moduleEvents{}}).forModule("messenger") + if _, err := b.State("create", json.RawMessage(`{"state":"asks","key":"a","value":{}}`)); err != nil { + t.Fatal(err) + } + _, err := b.State("create", json.RawMessage(`{"state":"asks","key":"a","value":{}}`)) + var coded interface{ ErrorCode() int } + if !errors.As(err, &coded) || coded.ErrorCode() != -32010 || !errors.Is(err, bus.ErrStateChanged) { + t.Errorf("a second create: %v", err) + } + if _, err := b.State("update", json.RawMessage(`{"state":"asks","key":"a","value":{},"revision":0}`)); err == nil { + t.Error("an update of revision zero was taken") + } +}