Author SHA1 Message Date
mesh-admin 47cdfad754 Merge pull request 'Ask the bus again for a subscription it refused (hq issue 222)' (#47) from fix/issue-222-a-refused-subscription-is-asked-again into main 2026-10-03 23:20:49 +00:00
jochen 85810ccc8c Ask the bus again for a subscription it refused (novox/hq issue 222)
A push sends the machine that needs a grant and the bus's machine in the same breath, and the
runtime can subscribe the moment before the bus reloads its user list. Refused once, the
subscription stayed dead until some later membership re-served it, and a newly assigned module ran
unreachable. A refused subject the runtime answers is now asked for again for about five minutes.
2026-10-04 01:20:44 +02:00
mesh-admin 2ba7451229 Merge pull request 'A refused tool subscription is said, never fatal (hq issue 218)' (#46) from fix/a-refused-seat-subscription-is-not-fatal into main 2026-10-03 22:13:22 +00:00
jochen 6604d44372 A refused tool subscription is said, never fatal (novox/hq issue 218)
After the controller stopped granting a mesh seat to claimants that do not hold it, an image built
before the runtime followed its membership still subscribed the seat's subject, and the refusal
ended the process: ace's postgres runtime crash-looped. A subject the grants leave out now costs
that subject only, as an announcement's already did (issue 217).
2026-10-04 00:13:15 +02:00
mesh-admin 7b21440962 Merge pull request 'Serve a seat's verbs from the membership once one is issued (hq issue 218)' (#45) from fix/issue-218-the-membership-decides-the-seats into main 2026-10-03 21:49:27 +00:00
jochen 0cea8d286e Serve a seat's verbs from the membership once one is issued (novox/hq issue 218)
The runtime added every seat its start-up credential claims even after the mesh issued a
membership without it. A seat held once for the mesh is claimed on every machine running the
module, so ace's postgres announced the store seat the bus then refused it on.
2026-10-03 23:48:15 +02:00
mesh-admin 094a7d0d7c Merge pull request 'The toolchain carries the SDK the mesh last published, and esbuild (hq issue 212, ADR 0193)' (#44) from feat/the-toolchain-follows-the-sdk-and-bundles into main 2026-10-03 21:31:59 +00:00
jochen b52669577d The toolchain carries the SDK the mesh last published, and esbuild (hq issue 212, ADR 0193)
mesh-tools stands on mesh-sdk's published package: the build receives its exact version and
installs it after the package.json install, so a release is a new argument and the cached layer
cannot keep an older SDK; the planner orders the toolchain after the SDK, and every bundle after the
toolchain (issue 211). esbuild, a development dependency, is what the builder bundles each
TypeScript entrypoint and launcher into one file with.
2026-10-03 23:26:09 +02:00
mesh-admin 7bd76275f9 Merge pull request 'A refused announcement is said, never fatal (hq issue 217)' (#43) from fix/a-refused-announcement-is-not-fatal into main 2026-10-03 21:24:50 +00:00
jochen 6ba0f4dc1e A refused announcement is said, never fatal (hq issue 217)
The raw subscription that answers discovery ran its loop unguarded, so a refusal escaped as an
unhandled rejection and ended the process: every per-module container crash-looped on 2026-10-03
over a subscription that only serves the mesh seeing the runtime. It is now caught, logged, and the
runtime serves on. Tested against a bus whose permissions refuse the subject; fails without it.
2026-10-03 23:24:38 +02:00
mesh-admin bc05658772 Merge pull request 'A mesh seat's holder answers for the machine it runs on (hq ADR 0197)' (#42) from fix/a-mesh-seats-holder-answers-for-its-machine into main 2026-10-03 21:23:06 +00:00
jochen 4af59636c3 A mesh seat's holder answers for the machine it runs on (hq ADR 0197)
The controller announces the mesh-controller seat without a machine — the seat is the mesh's — and
the console then reported it as not answering on the machine it is assigned to. An announcement that
names no machine now answers for wherever its module is assigned.
2026-10-03 23:22:56 +02:00
mesh-admin 6e425a000f Merge pull request 'The node's runtime is its modules' bus: it binds each module's consumer and hands its events to the bundle (hq ADR 0198)' (#41) from feat/0198-the-runtime-is-the-bus into main 2026-10-03 21:18:33 +00:00
jochen 14b6588839 Merge remote-tracking branch 'origin/main' into feat/0198-the-runtime-is-the-bus 2026-10-03 23:16:43 +02:00
mesh-admin 490aedfd36 Merge pull request 'Announce only on the subjects the grants allow (hq ADR 0197, issue 217)' (#40) from fix/announce-only-what-the-grants-allow into main 2026-10-03 21:11:24 +00:00
jochen f11ac6441c Announce only on the subjects the grants allow (hq ADR 0197)
Both runtimes subscribed $SRV.<verb>.> as a wildcard; the grants allow the bare question and the
service's own name and instance. The bus refused the wildcard, and the TypeScript runtime treats a
refused subscription as fatal, so every per-module container crash-looped after the image rolled.
They now subscribe exactly $SRV.<verb>, $SRV.<verb>.<name> and $SRV.<verb>.<name>.<id>.
2026-10-03 22:31:02 +02:00
jochen ffe229308c The node's runtime is its modules' bus: it binds each module's consumer and hands its events to the bundle (hq ADR 0198)
A launched bundle's mesh/subscribe binds the module's own durable consumer — EVENTS, <node>_<module>,
by name as the module's own runtime bound it, so nothing is lost or replayed in the move — and every
event goes to each child of the module that subscribed as mesh/event, acknowledged only when all
answered, negatively acknowledged after a short delay when one failed or died, terminated when it is
not an event. mesh/ask calls a tool as the module. Every launched bundle is started again when it
exits, with backoff, since long-running code waits for no call. Requires SDK 0.1.6.
2026-10-03 22:29:20 +02:00
mesh-admin df4f492a72 Merge pull request 'The mesh's tools are found by address, from what announces itself on the bus (hq ADR 0195, 0197)' (#39) from feat/0197-tools-announce-themselves into main 2026-10-03 20:19:53 +00:00
jochen 66e8be0e31 The TypeScript runtime announces what it serves too, in the same services format (hq ADR 0197)
The per-module containers still run this runtime; their tools and the seats they hold (the store's,
the catalogue's) must be found by the console the same way as the node runtime's. It answers
$SRV.PING, $SRV.INFO and $SRV.STATS with one service per process, one endpoint per tool per subject
and per seat verb served, the metadata as the Go runtime writes it.
2026-10-03 22:15:48 +02:00
jochen 7722668220 Every runtime announces what it serves in the NATS services protocol; the console discovers by asking the bus (hq ADR 0197)
The Go runtime answers $SRV.PING, $SRV.INFO and $SRV.STATS (and per name and id) in the
io.nats.micro.v1 format with what it serves at the moment it is asked: one service per runtime
process, since the bus admits one reply per request from each responder, and one endpoint per tool
per subject, its metadata saying module, seat, scope, machine, description, schema and whether the
module is interchangeable. Serving is unchanged.

The console gathers one $SRV.INFO request's answers instead of asking the catalogue's roster and
each module's tools, and reads the controller's records as JSON for what should have answered: an
assignment with tools that did not announce is named, a module without tools never is. The text
parsers of node list and module list are gone. Packages share the test bus: go test -p 1.
2026-10-03 22:14:16 +02:00
jochen e8989f3cf5 Merge remote-tracking branch 'origin/main' into feat/0197-tools-announce-themselves 2026-10-03 22:08:06 +02:00
mesh-admin 915f372a85 Merge pull request 'node-tools in Go: the node's runtime, launch-only, wire-compatible with the TypeScript (hq ADR 0193)' (#38) from feat/0193-node-tools-in-go into main 2026-10-03 20:02:15 +00:00
jochen 486dad99a5 Merge remote-tracking branch 'origin/main' into feat/0193-node-tools-in-go 2026-10-03 22:02:03 +02:00
mesh-admin 65f3b68076 Merge pull request 'node-tools launches every bundle it serves; a child's emit is published as its module (hq ADR 0193)' (#36) from feat/0193-the-runtime-launches-every-bundle into main 2026-10-03 20:01:58 +00:00
mesh-admin 42e0987c64 Merge pull request 'Require SDK 0.1.5: a launched bundle names its module and emits through the runtime (hq ADR 0193)' (#37) from fix/the-toolchain-carries-sdk-0.1.5 into main 2026-10-03 19:23:59 +00:00
jochen ae6bdc9b06 Require SDK 0.1.5: a launched bundle names its module and emits through the runtime (hq ADR 0193)
The toolchain image installs from this package.json with a range; unchanged, Docker reused the
cached install and the image kept SDK 0.1.3 after 0.1.5 was published. Bundles copied that copy, so
a launched bundle registering its seat first served the seat's verbs as its own tools. Requiring
0.1.5 says what the runtime and its bundles need, and invalidates the cached layer.
2026-10-03 21:23:50 +02:00
jochen b182943c24 node-tools launches every bundle it serves; a child's emit is published as its module (hq ADR 0193)
Every served entrypoint is started as a process speaking MCP over stdio, told its module and node;
one that is not executable is refused by name. The import path, the SDK resolve hook (issue 209)
and the per-registration hand-off go. The one-module form the per-module containers use is still
imported until they move (to-be 38 WP4c). A child's mesh/publish is published as its module and
answered once accepted; a child that dies says why in its own last words. Fixtures are served
through launchers exactly as the builder writes them.
2026-10-03 21:12:15 +02:00
30 changed files with 1979 additions and 452 deletions
+11
View File
@@ -28,6 +28,14 @@ COPY node-tools/package.json ./
# credential travels no further than here. Development dependencies included: the compiler is one. # credential travels no further than here. Development dependencies included: the compiler is one.
COPY .npmrc ./.npmrc COPY .npmrc ./.npmrc
RUN npm install --no-audit --no-fund RUN npm install --no-audit --no-fund
# **The SDK this image carries is the one the mesh last published** (novox/hq issue 212). The range in
# package.json is resolved once and the layer above is cached, so a release reached no toolchain until
# that file changed. The exact version arrives as a build argument from the SDK module's published
# package (module.json `build.on`), so a new release is a new argument, this layer runs again — and
# the planner orders this module after the SDK, so a release rebuilds the toolchain and, after it,
# every bundle compiled in it (issue 211).
ARG MESH_SDK=@novox/mesh-sdk@latest
RUN npm install --no-audit --no-fund "${MESH_SDK}"
# ---- compiling: the runtime's own code built, WITHOUT the credential ------------------------- # ---- compiling: the runtime's own code built, WITHOUT the credential -------------------------
FROM ${NODE_BASE} AS compiling FROM ${NODE_BASE} AS compiling
@@ -45,6 +53,9 @@ FROM compiling AS lean
RUN npm prune --omit=dev RUN npm prune --omit=dev
# ---- toolchain: what a TypeScript bundle is compiled in --------------------------------------- # ---- toolchain: what a TypeScript bundle is compiled in ---------------------------------------
# The compiler, and esbuild (a development dependency) at /app/node_modules/esbuild: the builder
# bundles every entrypoint and launcher into one file with it (novox/hq ADR 0193), so a bundle
# carries what it imports and not this image's node_modules.
# Beside the compiler, at /app/runtime, what every TypeScript bundle runs with: the production # Beside the compiler, at /app/runtime, what every TypeScript bundle runs with: the production
# dependencies the SDK and the runtime need, and a package.json saying the compiled files are ES # dependencies the SDK and the runtime need, and a package.json saying the compiled files are ES
# modules. The builder copies this directory whole into a compiled bundle (novox/hq ADR 0188 §5), # modules. The builder copies this directory whole into a compiled bundle (novox/hq ADR 0188 §5),
+5
View File
@@ -21,6 +21,11 @@
{ {
"arg": "NODE_BASE", "arg": "NODE_BASE",
"image": "node@sha256:48e4b67d85f87bd551df43704e24d252f56cc5f8e9718841aace50f19948f0f9" "image": "node@sha256:48e4b67d85f87bd551df43704e24d252f56cc5f8e9718841aace50f19948f0f9"
},
{
"arg": "MESH_SDK",
"module": "mesh-sdk",
"artifact": "lib"
} }
] ]
}, },
+211
View File
@@ -0,0 +1,211 @@
// Package announce answers the NATS services protocol's discovery for what a runtime serves, and
// gathers the answers (novox/hq ADR 0197).
//
// A runtime does not re-serve its tools through a services library: serving is unchanged. It answers
// `$SRV.PING`, `$SRV.INFO` and `$SRV.STATS` — and the same followed by its service name, and by its
// name and id — in the format NATS's own tools read, with what it is serving at the moment it is
// asked. One service per runtime process: the bus admits one reply per request from each responder,
// so a runtime serving many modules and seats answers once, one endpoint per tool per subject, and
// says in each endpoint's metadata which module, seat, scope and machine it is.
package announce
import (
"encoding/json"
"strings"
"time"
"github.com/nats-io/nats.go/micro"
"github.com/novox/mesh-tools/node-tools/internal/bus"
)
// Version is the announced service version (semver, as the protocol requires).
const Version = "0.1.0"
// Window is how long the console gathers discovery answers: every instance answers one request, and
// how many will is what is being found out.
var Window = 750 * time.Millisecond
// Kinds of endpoint.
const (
KindTool = "tool" // a module's own tool
KindSeat = "seat" // a seat's verb, served by the module holding it
)
// Endpoint is one tool served on one subject, as it is announced.
type Endpoint struct {
Kind string
Module string // the module whose code answers
Tool string // the tool's or the verb's name
Seat string // for a seat's verb
Scope string // "mesh" or "node", for a seat's verb
Node string
Description string
Schema json.RawMessage
Interchangeable bool
Subject string
Queue string
}
// Name is the endpoint's name as the protocol allows it — letters, digits, `-` and `_` — the
// prefix and the tool joined by `__`; the metadata, not the name, is what identifies it.
func (e Endpoint) Name() string {
prefix := e.Module
if e.Kind == KindSeat {
prefix = e.Seat
}
return clean(prefix) + "__" + clean(e.Tool)
}
func clean(s string) string {
var b strings.Builder
for _, r := range s {
if r == '-' || r == '_' || (r >= 'a' && r <= 'z') || (r >= 'A' && r <= 'Z') || (r >= '0' && r <= '9') {
b.WriteRune(r)
} else {
b.WriteRune('_')
}
}
return b.String()
}
func (e Endpoint) info() micro.EndpointInfo {
md := map[string]string{
"kind": e.Kind, "module": e.Module, "tool": e.Tool, "node": e.Node,
"description": e.Description, "interchangeable": boolWord(e.Interchangeable),
}
schema := strings.TrimSpace(string(e.Schema))
if schema == "" || schema == "null" {
schema = "{}"
}
md["schema"] = schema
if e.Kind == KindSeat {
md["seat"] = e.Seat
md["scope"] = e.Scope
}
return micro.EndpointInfo{Name: e.Name(), Subject: e.Subject, QueueGroup: e.Queue, Metadata: md}
}
func boolWord(b bool) string {
if b {
return "true"
}
return "false"
}
// Service is who answers: the runtime's name, its instance, and a word about it.
type Service struct {
Name string
ID string
Description string
Metadata map[string]string
}
func (s Service) identity() micro.ServiceIdentity {
md := s.Metadata
if md == nil {
md = map[string]string{}
}
return micro.ServiceIdentity{Name: s.Name, ID: s.ID, Version: Version, Metadata: md}
}
// Info is the info_response for these endpoints.
func Info(s Service, endpoints []Endpoint) micro.Info {
out := micro.Info{ServiceIdentity: s.identity(), Type: micro.InfoResponseType,
Description: s.Description, Endpoints: []micro.EndpointInfo{}}
for _, e := range endpoints {
out.Endpoints = append(out.Endpoints, e.info())
}
return out
}
// Serve answers discovery for one service until stopped, asking `current` for its endpoints each
// time — so what is announced is what is served now, re-served memberships included.
func Serve(conn *bus.Conn, s Service, current func() []Endpoint) (func(), error) {
started := time.Now().UTC()
answer := func(subject string, _ []byte) []byte {
parts := strings.Split(subject, ".")
if len(parts) < 2 || parts[0] != "$SRV" {
return nil
}
if len(parts) >= 3 && parts[2] != s.Name {
return nil // another service's
}
if len(parts) >= 4 && parts[3] != s.ID {
return nil // another instance's
}
var v any
switch parts[1] {
case "PING":
v = micro.Ping{ServiceIdentity: s.identity(), Type: micro.PingResponseType}
case "INFO":
v = Info(s, current())
case "STATS":
st := micro.Stats{ServiceIdentity: s.identity(), Type: micro.StatsResponseType, Started: started,
Endpoints: []*micro.EndpointStats{}}
for _, e := range current() {
st.Endpoints = append(st.Endpoints, &micro.EndpointStats{Name: e.Name(), Subject: e.Subject, QueueGroup: e.Queue})
}
v = st
default:
return nil
}
body, err := json.Marshal(v)
if err != nil {
return nil
}
return body
}
var stops []func()
for _, verb := range []string{"PING", "INFO", "STATS"} {
// Exactly the questions asked of every service and of this one by name and instance — what the
// grants allow (novox/hq ADR 0197). A wildcard is refused by the bus.
for _, subject := range []string{"$SRV." + verb, "$SRV." + verb + "." + s.Name, "$SRV." + verb + "." + s.Name + "." + s.ID} {
stop, err := conn.Raw(subject, answer)
if err != nil {
for _, st := range stops {
st()
}
return func() {}, err
}
stops = append(stops, stop)
}
}
return func() {
for _, st := range stops {
st()
}
}, nil
}
// Gather asks every service on the bus what it serves and answers what came back within the window.
// An answer that is not an info_response is skipped.
func Gather(conn *bus.Conn) ([]micro.Info, error) {
raw, err := conn.Gather("$SRV.INFO", nil, Window)
if err != nil {
return nil, err
}
var out []micro.Info
for _, b := range raw {
var i micro.Info
if json.Unmarshal(b, &i) != nil || i.Type != micro.InfoResponseType {
continue
}
out = append(out, i)
}
return out, nil
}
// Endpoints reads an info_response's endpoints back into what they announce.
func Endpoints(i micro.Info) []Endpoint {
var out []Endpoint
for _, e := range i.Endpoints {
md := e.Metadata
out = append(out, Endpoint{
Kind: md["kind"], Module: md["module"], Tool: md["tool"], Seat: md["seat"], Scope: md["scope"],
Node: md["node"], Description: md["description"], Schema: json.RawMessage(md["schema"]),
Interchangeable: md["interchangeable"] == "true", Subject: e.Subject, Queue: e.QueueGroup,
})
}
return out
}
+280 -3
View File
@@ -8,6 +8,7 @@
package bus package bus
import ( import (
"context"
"crypto/sha256" "crypto/sha256"
"crypto/tls" "crypto/tls"
"crypto/x509" "crypto/x509"
@@ -22,6 +23,7 @@ import (
"time" "time"
"github.com/nats-io/nats.go" "github.com/nats-io/nats.go"
"github.com/nats-io/nats.go/jetstream"
"github.com/novox/mesh-tools/node-tools/internal/wire" "github.com/novox/mesh-tools/node-tools/internal/wire"
) )
@@ -155,6 +157,97 @@ type Conn struct {
onNew []func(Membership) onNew []func(Membership)
subs []*nats.Subscription subs []*nats.Subscription
Logf func(format string, args ...any) Logf func(format string, args ...any)
// answering is every subject this runtime answers on, so a subscription the bus refused can be
// asked for again (novox/hq issue 222).
answering map[string]*answered
// RetryAfter is how long after a refusal a subscription is asked for again, by attempt; past the
// last, it is given up and said. A variable so a test need not wait minutes.
RetryAfter []time.Duration
}
// answered is one subject the runtime answers on: what to subscribe again with, and how often it was
// refused.
type answered struct {
subject, queue string
cb nats.MsgHandler
sub *nats.Subscription
refusals int
stopped bool
}
// retryAfter is the default wait before each new attempt: a push sends the bus's user list in the
// same breath as the machine that needs it, and the bus reloads it a moment later — five minutes
// covers a slow push, and a grant that never comes is said and left.
var retryAfter = []time.Duration{2 * time.Second, 5 * time.Second, 10 * time.Second, 20 * time.Second,
30 * time.Second, 60 * time.Second, 60 * time.Second, 120 * time.Second}
var (
refusedSubject = regexp.MustCompile(`Subscription to "(\S+)"`)
refusedQueue = regexp.MustCompile(`using queue "(\S+)"`)
)
func answeringKey(subject, queue string) string { return subject + "\x00" + queue }
// refused is the bus saying no to a subscription. **Asked again, not given up** (novox/hq issue 222):
// what an account may answer is the bus's user list, written on the bus's machine, and a push that
// assigns a module somewhere sends that machine and the bus's in the same breath — the runtime can
// subscribe in the moment before the bus has reloaded. Refused once, the subscription stayed dead
// until some later membership happened to re-serve it, and the module ran unreachable meanwhile.
func (c *Conn) refused(err error) {
if !errors.Is(err, nats.ErrPermissionViolation) {
return
}
m := refusedSubject.FindStringSubmatch(err.Error())
if len(m) < 2 {
return
}
queue := ""
if q := refusedQueue.FindStringSubmatch(err.Error()); len(q) >= 2 {
queue = q[1]
}
c.mu.Lock()
a := c.answering[answeringKey(m[1], queue)]
if a == nil || a.stopped {
c.mu.Unlock()
return
}
waits := c.RetryAfter
if waits == nil {
waits = retryAfter
}
if a.refusals >= len(waits) {
c.mu.Unlock()
c.Logf("[mesh-tools] the bus still refuses %s after %d attempts; not served here until the mesh issues it again", a.subject, a.refusals)
return
}
wait := waits[a.refusals]
a.refusals++
attempt := a.refusals
c.mu.Unlock()
c.Logf("[mesh-tools] the bus refused %s; asking again in %s (attempt %d)", a.subject, wait, attempt)
time.AfterFunc(wait, func() {
c.mu.Lock()
defer c.mu.Unlock()
if a.stopped {
return
}
if a.sub != nil {
_ = a.sub.Unsubscribe()
}
var sub *nats.Subscription
var err error
if a.queue != "" {
sub, err = c.nc.QueueSubscribe(a.subject, a.queue, a.cb)
} else {
sub, err = c.nc.Subscribe(a.subject, a.cb)
}
if err != nil {
c.Logf("[mesh-tools] asking again for %s failed: %v", a.subject, err)
return
}
a.sub = sub
c.subs = append(c.subs, sub)
})
} }
// Connect dials the bus as the credential's module. A module's subjects come from its credential, // Connect dials the bus as the credential's module. A module's subjects come from its credential,
@@ -168,8 +261,14 @@ func Connect(cred Credential) (*Conn, error) {
if node == "" { if node == "" {
node = "?" node = "?"
} }
var c *Conn
opts := []nats.Option{ opts := []nats.Option{
nats.Name(node + "." + cred.Module), nats.Name(node + "." + cred.Module),
nats.ErrorHandler(func(_ *nats.Conn, _ *nats.Subscription, err error) {
if c != nil {
c.refused(err)
}
}),
// Reconnect forever: the bus restarting is an upgrade, not a reason to exit. // Reconnect forever: the bus restarting is an upgrade, not a reason to exit.
nats.MaxReconnects(-1), nats.MaxReconnects(-1),
} }
@@ -190,8 +289,8 @@ func Connect(cred Credential) (*Conn, error) {
nc.Close() nc.Close()
return nil, err return nil, err
} }
c := &Conn{nc: nc, js: js, self: cred.Module, node: cred.Node, cred: cred, c = &Conn{nc: nc, js: js, self: cred.Module, node: cred.Node, cred: cred,
issued: map[string]*Membership{}, Logf: log.Printf} issued: map[string]*Membership{}, Logf: log.Printf, answering: map[string]*answered{}}
c.Follow(cred.Module) c.Follow(cred.Module)
return c, nil return c, nil
} }
@@ -360,7 +459,20 @@ func (c *Conn) answerOn(subject, queue string, h Handler) (func(), error) {
return func() {}, err return func() {}, err
} }
c.track(sub) c.track(sub)
return func() { _ = sub.Unsubscribe() }, nil a := &answered{subject: subject, queue: queue, cb: cb, sub: sub}
c.mu.Lock()
c.answering[answeringKey(subject, queue)] = a
c.mu.Unlock()
return func() {
c.mu.Lock()
a.stopped = true
if c.answering[answeringKey(subject, queue)] == a {
delete(c.answering, answeringKey(subject, queue))
}
current := a.sub
c.mu.Unlock()
_ = current.Unsubscribe()
}, nil
} }
// nullable keeps a nil result as JSON null rather than dropping the key: the TypeScript reply always // nullable keeps a nil result as JSON null rather than dropping the key: the TypeScript reply always
@@ -512,6 +624,59 @@ func (c *Conn) PublishAs(module string, env Envelope) error {
return err return err
} }
// ServedOn is where a served module's tool is answered right now: the membership's subjects when
// issued, the derived shape otherwise — what Handle subscribes, for what announces it (ADR 0197).
func (c *Conn) ServedOn(module, tool string) []Served { return c.servedOn(module, tool) }
// Raw answers one subject with a function of the request, not a tool's reply envelope: the NATS
// services protocol's discovery subjects answer in their own format (novox/hq ADR 0197). A nil
// answer is no reply — the request was for another service.
func (c *Conn) Raw(subject string, answer func(subject string, data []byte) []byte) (func(), error) {
sub, err := c.nc.Subscribe(subject, func(msg *nats.Msg) {
if body := answer(msg.Subject, msg.Data); body != nil {
_ = msg.Respond(body)
}
})
if err != nil {
return func() {}, err
}
c.track(sub)
return func() { _ = sub.Unsubscribe() }, nil
}
// Gather publishes one request and collects every answer that arrives within the window: a
// discovery request every service instance answers (ADR 0197). It never stops early — how many will
// answer is what it is finding out.
func (c *Conn) Gather(subject string, body []byte, window time.Duration) ([][]byte, error) {
inbox := c.nc.NewRespInbox()
sub, err := c.nc.SubscribeSync(inbox)
if err != nil {
return nil, err
}
defer func() { _ = sub.Unsubscribe() }()
if err := c.nc.PublishRequest(subject, inbox, body); err != nil {
return nil, err
}
var out [][]byte
deadline := time.Now().Add(window)
for {
left := time.Until(deadline)
if left <= 0 {
return out, nil
}
msg, err := sub.NextMsg(left)
if err != nil {
if errors.Is(err, nats.ErrTimeout) {
return out, nil
}
return out, err
}
if len(msg.Data) > 0 {
out = append(out, msg.Data)
}
}
}
// Flush waits until the bus has every subscription made so far, so what is served is answerable // Flush waits until the bus has every subscription made so far, so what is served is answerable
// when this returns. // when this returns.
func (c *Conn) Flush() { _ = c.nc.Flush() } func (c *Conn) Flush() { _ = c.nc.Flush() }
@@ -566,3 +731,115 @@ func SeatToolSubject(seat, verb, scope, node string) string {
} }
return base return base
} }
// AskAs calls a tool on a module's behalf (novox/hq ADR 0198): a bare key is that module's own tool,
// `<module>.<tool>` another's, `seat:<seat>.<verb>[@<node>]` a role's — resolved through what the
// mesh issued that module to reach, as its own runtime resolved it, else the derived shape.
func (c *Conn) AskAs(module, key string, body any) (Answered, error) {
name, wanted, _ := strings.Cut(key, "@")
subject := ""
if m := c.Membership(module); m != nil {
if reach := m.Reaches[name]; len(reach) > 0 {
subject = reach[0]
if wanted != "" {
subject = ""
for _, s := range reach {
if strings.HasSuffix(s, "."+wanted) {
subject = s
}
}
}
}
}
if subject == "" {
s, err := ToolSubject(key, module)
if err != nil {
return Answered{}, err
}
subject = s
}
return c.Ask(key, body, subject)
}
// EventsStream is where every module's events land, and ConsumerOf the durable consumer the
// controller makes for a module on a machine: the same names the module's own runtime bound
// (node-tools/src/broker-nats.ts), so moving the module into the node's runtime neither loses an
// event nor sees one twice.
const EventsStream = "EVENTS"
// ConsumerOf is a module's durable consumer on a machine: `<node>_<module>`.
func ConsumerOf(node, module string) string {
if node == "" {
node = "?"
}
return node + "_" + module
}
// NakDelay is how long an event a handler failed waits before it is offered again: a transient cause
// gets another attempt, a permanent one exhausts the consumer's max-deliver rather than spinning.
var NakDelay = 5 * time.Second
// ConsumeAs reads a module's durable consumer and hands each event to deliver (novox/hq ADR 0198):
// acknowledged when deliver returns nil, negatively acknowledged after NakDelay when it returns an
// error, terminated when it is not an event at all. The consumer is the controller's to create; this
// binds to it and never makes one. Runs until stopped.
func (c *Conn) ConsumeAs(module string, deliver func(Envelope) error) (func(), error) {
durable := ConsumerOf(c.node, module)
js, err := jetstream.New(c.nc)
if err != nil {
return nil, err
}
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
defer cancel()
// Bound by name, never created and never matched against a subject: the consumer and its
// filters are the controller's (design 29 §3), exactly as the module's own runtime bound it.
consumer, err := js.Consumer(ctx, EventsStream, durable)
if err != nil {
return nil, fmt.Errorf("binding %s's consumer %s: %w", module, durable, err)
}
reading, err := consumer.Consume(func(msg jetstream.Msg) {
env, ok := envelopeOf(msg.Subject(), msg.Headers(), msg.Data())
if !ok {
// Unparseable: redelivering bytes no version of a handler can read is a loop.
_ = msg.Term()
return
}
if err := deliver(env); err != nil {
_ = msg.NakWithDelay(NakDelay)
return
}
_ = msg.Ack()
})
if err != nil {
return nil, fmt.Errorf("reading %s's consumer %s: %w", module, durable, err)
}
var once sync.Once
return func() { once.Do(reading.Stop) }, nil
}
// envelopeOf rebuilds the envelope a module sees from a delivered event: its key is the emitter and
// the event, recovered from the subject; its metadata rides as headers (novox/hq ADR 0042).
func envelopeOf(subject string, header nats.Header, data []byte) (Envelope, bool) {
if !json.Valid(data) {
return Envelope{}, false
}
headers := map[string]string{}
for k := range header {
headers[k] = header.Get(k)
}
return Envelope{Key: keyOf(subject), Node: headers["x-node"], Body: json.RawMessage(data), Headers: headers}, true
}
// keyOf is the key a module sees for an event subject: `mesh.mod.<emitter>.event.<event>` is
// `<emitter>.<event>` — the vocabulary its manifest names what it consumes in.
func keyOf(subject string) string {
before, event, found := strings.Cut(subject, ".event.")
if !found {
return subject
}
parts := strings.Split(before, ".")
if emitter := parts[len(parts)-1]; emitter != "" {
return emitter + "." + event
}
return event
}
+86
View File
@@ -0,0 +1,86 @@
package bus
import (
"encoding/json"
"fmt"
"os"
"strings"
"sync"
"testing"
"time"
"github.com/nats-io/nats.go"
)
// novox/hq issue 222: a subscription the bus refused is asked for again. A push sends the machine
// that needs a grant and the bus's machine together; the runtime may subscribe the moment before the
// bus reloads, and a refusal must not leave the module unreachable until some later membership.
func TestARefusedSubscriptionIsAskedForAgain(t *testing.T) {
url := os.Getenv("MESH_TEST_NATS")
if url == "" {
t.Skip("MESH_TEST_NATS unset")
}
c, err := Connect(Credential{URL: url, Module: "alpha", Node: "anchor"})
if err != nil {
t.Fatal(err)
}
defer c.Close()
var mu sync.Mutex
var said []string
c.Logf = func(f string, a ...any) { mu.Lock(); said = append(said, fmt.Sprintf(f, a...)); mu.Unlock() }
c.RetryAfter = []time.Duration{50 * time.Millisecond, 50 * time.Millisecond}
subject := "mesh.mod.alpha.tool.ping.anchor"
stop, err := c.answerOn(subject, "", func(json.RawMessage) (any, error) { return "pong", nil })
if err != nil {
t.Fatal(err)
}
first := c.answering[answeringKey(subject, "")].sub
// What the bus says when the grant is not there yet.
c.refused(fmt.Errorf("%w: Permissions Violation for Subscription to %q", nats.ErrPermissionViolation, subject))
deadline := time.Now().Add(2 * time.Second)
for {
c.mu.Lock()
again := c.answering[answeringKey(subject, "")].sub
c.mu.Unlock()
if again != first {
break
}
if time.Now().After(deadline) {
t.Fatal("the refused subscription was not asked for again")
}
time.Sleep(10 * time.Millisecond)
}
asker, err := nats.Connect(url)
if err != nil {
t.Fatal(err)
}
defer asker.Close()
reply, err := asker.Request(subject, []byte("{}"), 2*time.Second)
if err != nil {
t.Fatalf("the subject asked for again does not answer: %v", err)
}
if !strings.Contains(string(reply.Data), "pong") {
t.Fatalf("the subject asked for again answered %s", reply.Data)
}
// Past its attempts it is given up, and said.
for i := 0; i < 3; i++ {
c.refused(fmt.Errorf("%w: Permissions Violation for Subscription to %q", nats.ErrPermissionViolation, subject))
time.Sleep(80 * time.Millisecond)
}
mu.Lock()
gaveUp := strings.Contains(strings.Join(said, "\n"), "still refuses "+subject)
mu.Unlock()
if !gaveUp {
t.Errorf("never gave up on a subject that stays refused:\n%s", strings.Join(said, "\n"))
}
// A subject the runtime stopped answering is not asked for again.
stop()
c.refused(fmt.Errorf("%w: Permissions Violation for Subscription to %q", nats.ErrPermissionViolation, subject))
if _, still := c.answering[answeringKey(subject, "")]; still {
t.Error("a stopped subject is still tracked")
}
}
+139 -144
View File
@@ -19,12 +19,12 @@ package console
import ( import (
"encoding/json" "encoding/json"
"fmt" "fmt"
"regexp"
"sort" "sort"
"strings" "strings"
"sync" "sync"
"time" "time"
"github.com/novox/mesh-tools/node-tools/internal/announce"
"github.com/novox/mesh-tools/node-tools/internal/bus" "github.com/novox/mesh-tools/node-tools/internal/bus"
) )
@@ -159,188 +159,183 @@ func controllerOutput(conn *bus.Conn, verb string) (string, error) {
// jsonIn is the JSON document a command printed, after any lines it said first: a seat verb runs // jsonIn is the JSON document a command printed, after any lines it said first: a seat verb runs
// the controller's command, and a command may warn before it answers. // the controller's command, and a command may warn before it answers.
func jsonIn(output string) string { func jsonIn(output string) string {
if strings.HasPrefix(strings.TrimSpace(output), "{") { t := strings.TrimSpace(output)
return strings.TrimSpace(output) if strings.HasPrefix(t, "{") || strings.HasPrefix(t, "[") {
return t
} }
if i := strings.Index(output, "\n{"); i >= 0 { for _, open := range []string{"\n{", "\n["} {
return strings.TrimSpace(output[i+1:]) if i := strings.Index(output, open); i >= 0 {
return strings.TrimSpace(output[i+1:])
}
} }
return "" return ""
} }
// A machine as `node list` prints it: its name, when it was last heard from, its mode — converged // recordedModule is a module as the controller's records hold it (`module list --json`): where the
// or adopted, the only two the command prints — and its id. Lines the command says around them // mesh assigned it, and whether it declares tools — what should announce itself, and where.
// (a warning about the bus's users, "no node records yet") are not machines and are skipped. type recordedModule struct {
var nodeLine = regexp.MustCompile(`^([a-z0-9][a-z0-9-]*)\s+.*\s(converged|adopted)\s+\S+\s*$`) Module string `json:"module"`
On []string `json:"on"`
// machinesIn reads the machines from `node list`'s output. Tools bool `json:"tools"`
func machinesIn(output string) []string {
var out []string
for _, line := range strings.Split(output, "\n") {
if m := nodeLine.FindStringSubmatch(line); m != nil {
out = append(out, m[1])
}
}
sort.Strings(out)
return out
} }
// A module as `module list` prints it: name, version, how it was built, and where it runs. // recordedMachine is a machine as the controller's records hold it (`node list --json`).
var moduleLine = regexp.MustCompile(`^([a-z0-9][a-z0-9-]*)\s+\S+\s+.*?\s+on (.+)$`) type recordedMachine struct {
Name string `json:"name"`
// assignmentsIn reads, from `module list`'s output, the machines each module runs on.
func assignmentsIn(output string) map[string][]string {
out := map[string][]string{}
for _, line := range strings.Split(output, "\n") {
if line == "" || line[0] == ' ' || line[0] == '\t' {
continue
}
m := moduleLine.FindStringSubmatch(strings.TrimRight(line, " "))
if m == nil {
continue
}
on := strings.TrimSpace(m[2])
if on == "nothing" {
out[m[1]] = []string{}
continue
}
var nodes []string
for _, n := range strings.Split(on, ",") {
if n = strings.TrimSpace(n); n != "" {
nodes = append(nodes, n)
}
}
sort.Strings(nodes)
out[m[1]] = nodes
}
return out
} }
// interchangeable is whether the mesh issued the module a plain subject for this tool — one any of // indexOn asks the mesh what it holds (novox/hq ADR 0197): what answers, from every runtime's own
// its instances answers (ADR 0160): the module's own subject with no machine after it. // announcement on the bus — one `$SRV.INFO` request — and what should, from the controller's records
func interchangeable(module string, t Tool) bool { // read as JSON. Nothing is inferred from a roster and nothing is parsed from print.
plain := "mesh.mod." + module + ".tool." + t.Name
for _, s := range t.Subjects {
if s == plain {
return true
}
}
return false
}
// indexOn asks the mesh what it holds: the flat listing (catalogue, modules, the seats' tools),
// and from the controller the seats' holders, the machines and the assignments — at once.
func indexOn(conn *bus.Conn) (*index, error) { func indexOn(conn *bus.Conn) (*index, error) {
var wg sync.WaitGroup var wg sync.WaitGroup
var seatsOut, nodesOut, modulesOut string var nodesOut, modulesOut string
wg.Add(3) wg.Add(2)
go func() { defer wg.Done(); seatsOut, _ = controllerOutput(conn, "seats") }()
go func() { defer wg.Done(); nodesOut, _ = controllerOutput(conn, "nodes") }() go func() { defer wg.Done(); nodesOut, _ = controllerOutput(conn, "nodes") }()
go func() { defer wg.Done(); modulesOut, _ = controllerOutput(conn, "modules") }() go func() { defer wg.Done(); modulesOut, _ = controllerOutput(conn, "modules") }()
l, err := toolsOn(conn) infos, err := announce.Gather(conn)
wg.Wait() wg.Wait()
if err != nil { if err != nil {
return nil, err return nil, err
} }
x := &index{Modules: map[string]*moduleInfo{}, NotAnswering: l.NotAnswering, Listing: l} l := &Listing{Tools: []Tool{}, NotAnswering: []string{}}
x := &index{Modules: map[string]*moduleInfo{}, Listing: l}
// The seats: their verbs from the mesh's records, their holders from the controller. announced := map[string]map[string]bool{} // module → node → announced something
var held struct { seats := map[string]*seatInfo{}
Seats []struct { toolAt := map[string]int{} // <module>.<tool> or <seat>.<verb> → index in l.Tools
Seat string `json:"seat"` machines := map[string]bool{}
Scope string `json:"scope"` for _, info := range infos {
Holders []holder `json:"holders"` for _, e := range announce.Endpoints(info) {
} `json:"seats"` if e.Node != "" {
} machines[e.Node] = true
_ = json.Unmarshal([]byte(jsonIn(seatsOut)), &held)
holders := map[string][]holder{}
scopes := map[string]string{}
for _, s := range held.Seats {
holders[s.Seat] = s.Holders
scopes[s.Seat] = s.Scope
}
bySeat := map[string]*seatInfo{}
for _, t := range l.Tools {
if !t.Seat {
continue
}
s := bySeat[t.Module]
if s == nil {
scope := t.Scope
if scope == "" {
scope = scopes[t.Module]
} }
if scope != "node" { if announced[e.Module] == nil {
scope = "mesh" announced[e.Module] = map[string]bool{}
} }
s = &seatInfo{Seat: t.Module, Scope: scope, Holders: holders[t.Module]} announced[e.Module][e.Node] = true
bySeat[t.Module] = s switch e.Kind {
} case announce.KindSeat:
s.Verbs = append(s.Verbs, t) st := seats[e.Seat]
} if st == nil {
for _, s := range bySeat { st = &seatInfo{Seat: e.Seat, Scope: e.Scope}
x.Seats = append(x.Seats, *s) seats[e.Seat] = st
} }
sort.Slice(x.Seats, func(i, j int) bool { return x.Seats[i].Seat < x.Seats[j].Seat }) if st.Scope != "node" && e.Scope == "node" {
st.Scope = "node"
// The machines: what the controller knows, and any that hold a seat. }
seen := map[string]bool{} h := holder{Module: e.Module, Node: e.Node}
for _, n := range machinesIn(nodesOut) { if !containsHolder(st.Holders, h) {
seen[n] = true st.Holders = append(st.Holders, h)
} }
for _, s := range x.Seats { key := e.Seat + "." + e.Tool
for _, h := range s.Holders { if _, have := toolAt[key]; !have {
if h.Node != "" { toolAt[key] = len(l.Tools)
seen[h.Node] = true t := Tool{Module: e.Seat, Name: e.Tool, Description: e.Description, Input: e.Schema, Seat: true, Scope: st.Scope}
l.Tools = append(l.Tools, t)
st.Verbs = append(st.Verbs, t)
}
case announce.KindTool:
m := x.Modules[e.Module]
if m == nil {
m = &moduleInfo{Module: e.Module}
x.Modules[e.Module] = m
}
if e.Node != "" && !contains(m.On, e.Node) {
m.On = append(m.On, e.Node)
}
m.Interchangeable = m.Interchangeable || e.Interchangeable
key := e.Module + "." + e.Tool
i, have := toolAt[key]
if !have {
i = len(l.Tools)
toolAt[key] = i
l.Tools = append(l.Tools, Tool{Module: e.Module, Name: e.Tool, Description: e.Description, Input: e.Schema})
}
if !contains(l.Tools[i].Subjects, e.Subject) {
l.Tools[i].Subjects = append(l.Tools[i].Subjects, e.Subject)
}
} }
} }
} }
// A tool's subjects as a call looks them up: the plain one any instance answers first, then each
// The modules that answer tools, where they run, and whether any instance will do. // machine's.
on := assignmentsIn(modulesOut) for i := range l.Tools {
for _, t := range l.Tools { t := &l.Tools[i]
if t.Seat { if t.Seat {
continue continue
} }
m := x.Modules[t.Module] plain := "mesh.mod." + t.Module + ".tool." + t.Name
if m == nil { sort.SliceStable(t.Subjects, func(a, b int) bool {
m = &moduleInfo{Module: t.Module, On: on[t.Module]} if (t.Subjects[a] == plain) != (t.Subjects[b] == plain) {
x.Modules[t.Module] = m return t.Subjects[a] == plain
} }
m.Tools = append(m.Tools, t) return t.Subjects[a] < t.Subjects[b]
if interchangeable(t.Module, t) { })
m.Interchangeable = true }
for _, t := range l.Tools {
if !t.Seat {
x.Modules[t.Module].Tools = append(x.Modules[t.Module].Tools, t)
} }
} }
for _, m := range x.Modules { for _, m := range x.Modules {
// A module the controller could not place is placed where its own answer says it runs. sort.Strings(m.On)
if len(m.On) == 0 { }
nodes := map[string]bool{} for _, st := range seats {
for _, t := range m.Tools { sort.Slice(st.Holders, func(a, b int) bool { return st.Holders[a].Node < st.Holders[b].Node })
for _, s := range t.Subjects { x.Seats = append(x.Seats, *st)
base := "mesh.mod." + m.Module + ".tool." + t.Name + "." }
if strings.HasPrefix(s, base) { sort.Slice(x.Seats, func(i, j int) bool { return x.Seats[i].Seat < x.Seats[j].Seat })
nodes[strings.TrimPrefix(s, base)] = true sort.SliceStable(l.Tools, func(i, j int) bool {
} return l.Tools[i].Module+"."+l.Tools[i].Name < l.Tools[j].Module+"."+l.Tools[j].Name
})
// What should have answered: every assignment of a module that declares tools. Silence is named;
// a module with no tools is never a name here.
var recorded []recordedModule
if json.Unmarshal([]byte(jsonIn(modulesOut)), &recorded) == nil {
for _, m := range recorded {
if !m.Tools {
continue
}
for _, n := range m.On {
// A holder of a seat held once for the mesh announces no machine — the seat is the mesh's,
// not a machine's — so what it announced without one answers for wherever it is assigned.
if !announced[m.Module][n] && !announced[m.Module][""] {
l.NotAnswering = append(l.NotAnswering, m.Module+" on "+n)
} }
} }
for n := range nodes {
m.On = append(m.On, n)
}
sort.Strings(m.On)
} }
for _, n := range m.On { } else {
seen[n] = true l.NotAnswering = append(l.NotAnswering, "mesh-controller (its records of the modules did not answer, so what is missing cannot be said)")
}
sort.Strings(l.NotAnswering)
x.NotAnswering = l.NotAnswering
var known []recordedMachine
if json.Unmarshal([]byte(jsonIn(nodesOut)), &known) == nil {
for _, n := range known {
if n.Name != "" {
machines[n.Name] = true
}
} }
} }
for n := range seen { for n := range machines {
x.Machines = append(x.Machines, n) x.Machines = append(x.Machines, n)
} }
sort.Strings(x.Machines) sort.Strings(x.Machines)
return x, nil return x, nil
} }
func containsHolder(hs []holder, h holder) bool {
for _, x := range hs {
if x == h {
return true
}
}
return false
}
func (s *Surface) index() (*index, error) { func (s *Surface) index() (*index, error) {
s.mu.Lock() s.mu.Lock()
if s.idx != nil && time.Since(s.idxAt) <= IndexKept { if s.idx != nil && time.Since(s.idxAt) <= IndexKept {
@@ -520,7 +515,7 @@ func failure(text string) map[string]any {
func (s *Surface) discover(name string, args map[string]any) map[string]any { func (s *Surface) discover(name string, args map[string]any) map[string]any {
x, err := s.index() x, err := s.index()
if err != nil { if err != nil {
return failure(whyItFailed(catalogueModules, err)) return failure("the mesh's discovery failed: " + err.Error())
} }
str := func(k string) string { v, _ := args[k].(string); return strings.TrimSpace(v) } str := func(k string) string { v, _ := args[k].(string); return strings.TrimSpace(v) }
+51 -41
View File
@@ -4,9 +4,11 @@ import (
"encoding/json" "encoding/json"
"fmt" "fmt"
"strings" "strings"
"sync/atomic"
"testing" "testing"
"time" "time"
"github.com/novox/mesh-tools/node-tools/internal/announce"
mt "github.com/novox/mesh-tools/node-tools/internal/meshtest" mt "github.com/novox/mesh-tools/node-tools/internal/meshtest"
"github.com/novox/mesh-tools/node-tools/internal/runtime" "github.com/novox/mesh-tools/node-tools/internal/runtime"
) )
@@ -56,14 +58,20 @@ func TestTheMeshsToolsAreFoundByAddress(t *testing.T) {
} }
defer stop() defer stop()
catalogue := connect(t, "mesh-catalog", "") // Nothing may be asked for a roster or a module's `tools` any more (ADR 0197): counted.
stopCat, _ := catalogue.Handle("catalog_modules", func(json.RawMessage) (any, error) { var asked atomic.Int32
return map[string]any{"modules": []map[string]string{{"module": "alpha"}, {"module": "beta"}, {"module": "gamma"}}}, nil watcher := connect(t, "watcher", "")
}) for _, subject := range []string{"mesh.mod.*.tool.tools", "mesh.mod.*.tool.tools.*", "mesh.mod.mesh-catalog.>"} {
defer stopCat() stop, err := watcher.Raw(subject, func(string, []byte) []byte { asked.Add(1); return nil })
if err != nil {
t.Fatal(err)
}
t.Cleanup(stop)
}
// The controller, answering as its seat verbs do: the command's printed output. // The controller: its records as JSON, as its seat verbs answer them, and its own seat's verbs
controller := connect(t, "mesh-controller", "") // announced on the bus like every runtime's.
controller := connect(t, "mesh-controller", "bench")
out := func(s string) map[string]any { return map[string]any{"output": s, "ok": true} } out := func(s string) map[string]any { return map[string]any{"output": s, "ok": true} }
serve := func(verb string, answer func() any) { serve := func(verb string, answer func() any) {
stop, err := controller.HandleSubject("mesh.seat.mesh-controller.tool."+verb, func(json.RawMessage) (any, error) { stop, err := controller.HandleSubject("mesh.seat.mesh-controller.tool."+verb, func(json.RawMessage) (any, error) {
@@ -74,34 +82,29 @@ func TestTheMeshsToolsAreFoundByAddress(t *testing.T) {
} }
t.Cleanup(stop) t.Cleanup(stop)
} }
serve("tools", func() any {
return map[string]any{"seats": []map[string]any{
{"seat": "mesh-controller", "scope": "mesh", "tools": []map[string]any{
{"name": "nodes", "description": "Every machine the mesh knows.", "input": map[string]any{}}}},
{"seat": "node-shelf", "scope": "node", "tools": []map[string]any{
{"name": "list", "description": "what is on the shelf", "input": map[string]any{}},
{"name": "clear", "description": "take it all off", "input": map[string]any{}}}},
}}
})
serve("nodes", func() any { serve("nodes", func() any {
return out("bench 3m ago converged 1f2e\ndesk here converged 9a8b\n") return out("the bus's user list leaves out 2 user(s)\n" +
`[{"name":"bench","heard":"3m ago","mode":"converged","id":"1f2e"},{"name":"desk","heard":"here","mode":"converged","id":"9a8b"}]`)
}) })
serve("seats", func() any { var gammaOn atomic.Value
return out("a warning the command printed first\n" + `{ gammaOn.Store("[]")
"seats": [
{"seat": "mesh-controller", "scope": "mesh", "decision": "x", "holders": [{"module": "mesh-controller", "node": "bench"}]},
{"seat": "node-shelf", "scope": "node", "decision": "y", "holders": [{"module": "beta", "node": "desk"}]}
]
}`)
})
gammaOn := "nothing"
serve("modules", func() any { serve("modules", func() any {
return out(fmt.Sprintf("alpha 1 built 1a2b3c4d on desk\n needs container-runtime\n"+ return out(fmt.Sprintf(`[{"module":"alpha","on":["desk"],"tools":true},{"module":"beta","on":["desk"],"tools":true},`+
"beta 1 built 1a2b3c4d on desk\n"+ `{"module":"gamma","on":%s,"tools":true},{"module":"delta","on":["desk"],"tools":false},`+
"gamma 1 built 1a2b3c4d on %s\n", gammaOn)) `{"module":"epsilon","on":["bench"],"tools":true},{"module":"mesh-controller","on":["bench"],"tools":true}]`, gammaOn.Load().(string)))
}) })
catalogue.Flush() stopAnn, err := announce.Serve(controller, announce.Service{Name: "mesh-controller", ID: "bench"}, func() []announce.Endpoint {
return []announce.Endpoint{{Kind: announce.KindSeat, Module: "mesh-controller", Seat: "mesh-controller", Scope: "mesh",
// As the live controller announces: a mesh seat's holder names no machine.
Tool: "nodes", Node: "", Description: "Every machine the mesh knows.", Schema: json.RawMessage(`{}`),
Subject: "mesh.seat.mesh-controller.tool.nodes"}}
})
if err != nil {
t.Fatal(err)
}
t.Cleanup(stopAnn)
controller.Flush() controller.Flush()
watcher.Flush()
up, err := Serve(NewSurface(nodeTools, "desk.node-tools"), "127.0.0.1:0") up, err := Serve(NewSurface(nodeTools, "desk.node-tools"), "127.0.0.1:0")
if err != nil { if err != nil {
@@ -132,6 +135,14 @@ func TestTheMeshsToolsAreFoundByAddress(t *testing.T) {
!strings.Contains(overview, `"bench"`) || !strings.Contains(overview, `"desk"`) { !strings.Contains(overview, `"bench"`) || !strings.Contains(overview, `"desk"`) {
t.Errorf("overview: %s", overview) t.Errorf("overview: %s", overview)
} }
// Silence is named only where tools should have answered: epsilon declares tools on bench and
// nothing there announced it; delta declares none and is never a name.
if strings.Contains(overview, "mesh-controller on bench") {
t.Errorf("a mesh seat's holder that announced no machine is called silent on its machine:\n%s", overview)
}
if !strings.Contains(overview, "epsilon on bench") || strings.Contains(overview, "delta") {
t.Errorf("not answering: %s", overview)
}
machine, isErr := call(t, endpoint, "mesh_machine", map[string]any{"node": "desk"}) machine, isErr := call(t, endpoint, "mesh_machine", map[string]any{"node": "desk"})
if isErr || !strings.Contains(machine, "desk/node-shelf.list") || !strings.Contains(machine, "desk/beta.three") || if isErr || !strings.Contains(machine, "desk/node-shelf.list") || !strings.Contains(machine, "desk/beta.three") ||
!strings.Contains(machine, "desk/alpha.one") { !strings.Contains(machine, "desk/alpha.one") {
@@ -181,7 +192,7 @@ func TestTheMeshsToolsAreFoundByAddress(t *testing.T) {
t.Fatalf("gamma was found before it served: %s", got) t.Fatalf("gamma was found before it served: %s", got)
} }
mesh.Issue(t, mt.MembershipOf("gamma", "desk", false, nil)) mesh.Issue(t, mt.MembershipOf("gamma", "desk", false, nil))
gammaOn = "desk" gammaOn.Store(`["desk"]`)
late := connect(t, "node-tools", "desk") late := connect(t, "node-tools", "desk")
stopLate, err := runtime.Run(late, []runtime.Served{{Module: "gamma", Entrypoints: []string{mt.Fixture("env-gamma.serve.mjs")}}}, stopLate, err := runtime.Run(late, []runtime.Served{{Module: "gamma", Entrypoints: []string{mt.Fixture("env-gamma.serve.mjs")}}},
nil, (&mt.Logs{}).Logf) nil, (&mt.Logs{}).Logf)
@@ -201,22 +212,21 @@ func TestTheMeshsToolsAreFoundByAddress(t *testing.T) {
t.Errorf("a module that arrived later was not found: %s", found) t.Errorf("a module that arrived later was not found: %s", found)
} }
if n := asked.Load(); n != 0 {
t.Errorf("discovery asked a roster or a module's tools %d time(s); it asks the bus", n)
}
// The old names still answer, unannounced. // The old names still answer, unannounced.
if got, isErr := call(t, endpoint, "alpha.one", nil); isErr || !strings.Contains(got, `"alpha": 1`) { if got, isErr := call(t, endpoint, "alpha.one", nil); isErr || !strings.Contains(got, `"alpha": 1`) {
t.Errorf("an old name: %s", got) t.Errorf("an old name: %s", got)
} }
} }
func TestTheControllersPrintedListsAreRead(t *testing.T) { func TestTheControllersRecordsAreReadAsJSON(t *testing.T) {
machines := machinesIn("the bus's user list leaves out 2 user(s)\nace 2m ago converged 0c1d\n" + if got := jsonIn("a warning printed first\n[{\"name\":\"ace\"}]"); got != `[{"name":"ace"}]` {
"g14 here converged 77aa\nnovox 5s ago adopted 3e4f\n") t.Errorf("an array after a warning: %q", got)
if got := strings.Join(machines, ","); got != "ace,g14,novox" {
t.Errorf("machines: %s", got)
} }
on := assignmentsIn("baserow 1 built 7c800705 on ace\n requires postgres-database\n" + if got := jsonIn(`{"seats":[]}`); got != `{"seats":[]}` {
"confluence 1 built 7c800705 on nothing\n" + t.Errorf("an object: %q", got)
"mesh-wireguard 1 with the control plane on ace, g14, novox, shanks\n")
if strings.Join(on["baserow"], ",") != "ace" || len(on["confluence"]) != 0 || strings.Join(on["mesh-wireguard"], ",") != "ace,g14,novox,shanks" {
t.Errorf("assignments: %v", on)
} }
} }
+4 -105
View File
@@ -5,19 +5,11 @@ import (
"errors" "errors"
"fmt" "fmt"
"regexp" "regexp"
"sort"
"strings" "strings"
"sync"
"github.com/novox/mesh-tools/node-tools/internal/bus" "github.com/novox/mesh-tools/node-tools/internal/bus"
) )
const (
catalogueModules = "mesh-catalog.catalog_modules"
seatTools = "seat:mesh-controller.tools"
toolsVerb = "tools"
)
// Tool is a tool as its module — or, for a role's tool, the mesh's records — describes it. // Tool is a tool as its module — or, for a role's tool, the mesh's records — describes it.
type Tool struct { type Tool struct {
Module string Module string
@@ -67,107 +59,14 @@ func toolKey(name string, seats Seats) string {
return name return name
} }
// toolsOn asks the mesh what tools it has: the catalogue which modules it holds, each module what it // toolsOn is the flat catalogue (MESH_CONSOLE_FLAT=1): what announced itself on the bus, every
// serves, the controller's seat every role's tools — at once, so a restarting control plane hides // tool and seat verb, and the assignments with tools that did not (novox/hq ADR 0197).
// nothing else.
func toolsOn(conn *bus.Conn) (*Listing, error) { func toolsOn(conn *bus.Conn) (*Listing, error) {
type rolesAnswer struct { x, err := indexOn(conn)
Seats []struct {
Seat string `json:"seat"`
Scope string `json:"scope"`
Tools []struct {
Name string `json:"name"`
Description string `json:"description"`
Input json.RawMessage `json:"input"`
} `json:"tools"`
} `json:"seats"`
}
var wg sync.WaitGroup
var roles *rolesAnswer
wg.Add(1)
go func() {
defer wg.Done()
if got, err := conn.Ask(seatTools, map[string]any{}, ""); err == nil {
var r rolesAnswer
if json.Unmarshal(got.Result, &r) == nil {
roles = &r
}
}
}()
answered, err := conn.Ask(catalogueModules, map[string]any{}, "")
if err != nil { if err != nil {
wg.Wait()
return nil, err return nil, err
} }
var held struct { return x.Listing, nil
Modules []struct {
Module string `json:"module"`
} `json:"modules"`
}
_ = json.Unmarshal(answered.Result, &held)
names := make([]string, 0, len(held.Modules))
for _, m := range held.Modules {
if m.Module != "" {
names = append(names, m.Module)
}
}
type outcome struct {
ok bool
answer struct {
Tools *[]struct {
Name string `json:"name"`
Description string `json:"description"`
Input json.RawMessage `json:"input"`
Subjects []string `json:"subjects"`
} `json:"tools"`
Failed *string `json:"failed"`
}
}
outcomes := make([]outcome, len(names))
for i, module := range names {
wg.Add(1)
go func(i int, module string) {
defer wg.Done()
got, err := conn.Ask(module+"."+toolsVerb, map[string]any{}, "")
if err != nil {
return
}
if json.Unmarshal(got.Result, &outcomes[i].answer) == nil {
outcomes[i].ok = true
}
}(i, module)
}
wg.Wait()
l := &Listing{Tools: []Tool{}, NotAnswering: []string{}}
if roles != nil {
for _, s := range roles.Seats {
for _, t := range s.Tools {
l.Tools = append(l.Tools, Tool{Module: s.Seat, Name: t.Name, Description: t.Description,
Input: t.Input, Seat: true, Scope: s.Scope})
}
}
} else {
l.NotAnswering = append(l.NotAnswering, "mesh-controller (seat)")
}
for i, module := range names {
o := outcomes[i]
switch {
case o.ok && o.answer.Failed != nil:
l.NotAnswering = append(l.NotAnswering, fmt.Sprintf("%s (its tools bundle failed to load: %s)", module, *o.answer.Failed))
case o.ok && o.answer.Tools != nil:
for _, t := range *o.answer.Tools {
l.Tools = append(l.Tools, Tool{Module: module, Name: t.Name, Description: t.Description,
Input: t.Input, Subjects: t.Subjects})
}
default:
l.NotAnswering = append(l.NotAnswering, module)
}
}
sort.SliceStable(l.Tools, func(i, j int) bool {
return l.Tools[i].Module+"."+l.Tools[i].Name < l.Tools[j].Module+"."+l.Tools[j].Name
})
sort.Strings(l.NotAnswering)
return l, nil
} }
// callTool calls `<module>.<tool>[@<node>]`, on the subject the listing names for it when it names one. // callTool calls `<module>.<tool>[@<node>]`, on the subject the listing names for it when it names one.
+5 -12
View File
@@ -55,20 +55,13 @@ func TestTheConsoleListsAndCallsOverHTTP(t *testing.T) {
} }
defer stop() defer stop()
catalogue := connect(t, "mesh-catalog", "") // The controller's records (ADR 0197): ghost declares tools on desk and nothing announces it.
stopCat, _ := catalogue.Handle("catalog_modules", func(json.RawMessage) (any, error) {
return map[string]any{"modules": []map[string]string{{"module": "alpha"}, {"module": "beta"}, {"module": "ghost"}}}, nil
})
defer stopCat()
controller := connect(t, "mesh-controller", "") controller := connect(t, "mesh-controller", "")
stopSeat, _ := controller.HandleSubject("mesh.seat.mesh-controller.tool.tools", func(json.RawMessage) (any, error) { stopSeat, _ := controller.HandleSubject("mesh.seat.mesh-controller.tool.modules", func(json.RawMessage) (any, error) {
return map[string]any{"seats": []map[string]any{{"seat": "node-shelf", "scope": "node", "tools": []map[string]any{ return map[string]any{"ok": true, "output": `[{"module":"alpha","on":["desk"],"tools":true},` +
{"name": "list", "description": "what is on the shelf", "input": map[string]any{}}, `{"module":"beta","on":["desk"],"tools":true},{"module":"ghost","on":["desk"],"tools":true}]`}, nil
{"name": "clear", "description": "take it all off", "input": map[string]any{}},
}}}}, nil
}) })
defer stopSeat() defer stopSeat()
catalogue.Flush()
controller.Flush() controller.Flush()
flat := NewSurface(nodeTools, "desk.node-tools") flat := NewSurface(nodeTools, "desk.node-tools")
@@ -92,7 +85,7 @@ func TestTheConsoleListsAndCallsOverHTTP(t *testing.T) {
if got := strings.Join(names, ","); got != "alpha.one,alpha.two,beta.five,beta.four,beta.three,node-shelf.clear,node-shelf.list" { if got := strings.Join(names, ","); got != "alpha.one,alpha.two,beta.five,beta.four,beta.three,node-shelf.clear,node-shelf.list" {
t.Errorf("listed %s", got) t.Errorf("listed %s", got)
} }
if got := listed["_meta"].(map[string]any)["notAnswering"]; len(got.([]any)) != 1 || got.([]any)[0] != "ghost" { if got := listed["_meta"].(map[string]any)["notAnswering"]; len(got.([]any)) != 1 || got.([]any)[0] != "ghost on desk" {
t.Errorf("not answering: %v", got) t.Errorf("not answering: %v", got)
} }
for _, x := range listed["tools"].([]any) { for _, x := range listed["tools"].([]any) {
+1 -1
View File
@@ -122,7 +122,7 @@ func (s *Surface) Handle(r Request) *Reply {
} }
l, err := s.listing() l, err := s.listing()
if err != nil { if err != nil {
return refuse(r.ID, -32603, whyItFailed(catalogueModules, err)) return refuse(r.ID, -32603, "the mesh's discovery failed: "+err.Error())
} }
tools := make([]map[string]any, 0, len(l.Tools)) tools := make([]map[string]any, 0, len(l.Tools))
for _, t := range l.Tools { for _, t := range l.Tools {
+121 -11
View File
@@ -46,8 +46,26 @@ type Registration struct {
Tools []Tool Tools []Tool
} }
// Publisher publishes an event a bundle asked the runtime to emit, as the bundle's module. // Bus is what a launched bundle reaches the mesh through (novox/hq ADR 0193, ADR 0198): the runtime
type Publisher func(params json.RawMessage) error // 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.
type Bus interface {
Publish(params json.RawMessage) error
Ask(params json.RawMessage) (json.RawMessage, error)
Subscribe(deliver func(envelope json.RawMessage) error) error
}
// EventTimeout bounds how long a bundle has to handle one event before it is offered again.
var EventTimeout = 2 * time.Minute
// Restart backoff: a bundle that exits is started again at once, then after growing pauses while it
// keeps exiting, back to at once once it has run a while.
var (
RestartFirst = 500 * time.Millisecond
RestartMost = 30 * time.Second
RestartSettle = time.Minute
)
// Launched is a running bundle: what it registered, and how to stop it. // Launched is a running bundle: what it registered, and how to stop it.
type Launched struct { type Launched struct {
@@ -103,10 +121,26 @@ var stackLine = regexp.MustCompile(`^\s+at\s`)
// Start launches one bundle and learns its tools. It fails when the child cannot be started or does // Start launches one bundle and learns its tools. It fails when the child cannot be started or does
// not complete the handshake. A child that exits later is started again on its next call. // not complete the handshake. A child that exits later is started again on its next call.
func Start(module, entry string, env []string, publish Publisher, logf func(string, ...any)) (*Launched, error) { func Start(module, entry string, env []string, mesh Bus, logf func(string, ...any)) (*Launched, error) {
var mu sync.Mutex var mu sync.Mutex
var current *child var current *child
// The child that last subscribed is the one events are handed to: a restarted child subscribes
// again as its code is imported, and from then on the module's events are its.
var subscriber *child
stopped := false stopped := false
deliver := func(envelope json.RawMessage) error {
mu.Lock()
c := subscriber
mu.Unlock()
if c == nil {
return errors.New(module + "'s bundle is not running to take its events")
}
_, err := c.ask(module, "mesh/event", map[string]any{"envelope": envelope}, EventTimeout)
return err
}
var restart func(after time.Duration)
var startedAt time.Time
pause := RestartFirst
start := func() (*child, error) { start := func() (*child, error) {
cmd := exec.Command(entry) cmd := exec.Command(entry)
@@ -163,18 +197,32 @@ func Start(module, entry string, env []string, publish Publisher, logf func(stri
logf("[mesh-tools] %s's bundle said something that is not a reply: %s", module, line) logf("[mesh-tools] %s's bundle said something that is not a reply: %s", module, line)
continue continue
} }
// The bundle asks the runtime to emit (ADR 0193): published as this module, answered // The bundle asks the runtime (ADR 0193, ADR 0198): to emit, to call a tool, to hand it the
// once the bus has accepted it. Nothing else a bundle may ask. // module's events. Each on the module's behalf, answered when done; nothing else.
if m.Method != "" { if m.Method != "" {
go func(m message) { go func(m message) {
if m.Method != "mesh/publish" { var result any = map[string]any{}
var err error
switch m.Method {
case "mesh/publish":
err = mesh.Publish(m.Params)
case "mesh/ask":
var asked json.RawMessage
if asked, err = mesh.Ask(m.Params); err == nil {
result = asked
}
case "mesh/subscribe":
mu.Lock()
subscriber = c
mu.Unlock()
err = mesh.Subscribe(deliver)
default:
if len(m.ID) > 0 { if len(m.ID) > 0 {
_ = c.write(map[string]any{"jsonrpc": "2.0", "id": m.ID, _ = c.write(map[string]any{"jsonrpc": "2.0", "id": m.ID,
"error": map[string]any{"code": -32601, "message": "the runtime answers no " + m.Method + " from a bundle"}}) "error": map[string]any{"code": -32601, "message": "the runtime answers no " + m.Method + " from a bundle"}})
} }
return return
} }
err := publish(m.Params)
if len(m.ID) == 0 { if len(m.ID) == 0 {
return return
} }
@@ -183,7 +231,7 @@ func Start(module, entry string, env []string, publish Publisher, logf func(stri
"error": map[string]any{"code": -32000, "message": err.Error()}}) "error": map[string]any{"code": -32000, "message": err.Error()}})
return return
} }
_ = c.write(map[string]any{"jsonrpc": "2.0", "id": m.ID, "result": map[string]any{}}) _ = c.write(map[string]any{"jsonrpc": "2.0", "id": m.ID, "result": result})
}(m) }(m)
continue continue
} }
@@ -220,13 +268,32 @@ func Start(module, entry string, env []string, publish Publisher, logf func(stri
c.mu.Unlock() c.mu.Unlock()
close(c.dead) close(c.dead)
mu.Lock() mu.Lock()
// Only a child that was serving is brought back; one that died in its own handshake was
// never accepted, and whoever started it was told why.
wasServing := current == c
if current == c { if current == c {
current = nil current = nil
} }
if subscriber == c {
subscriber = nil
}
wasStopped := stopped wasStopped := stopped
// Started again at once if it ran a while; after a growing pause while it keeps exiting.
wait := RestartFirst
if time.Since(startedAt) < RestartSettle {
wait = pause
if pause *= 2; pause > RestartMost {
pause = RestartMost
}
} else {
pause = RestartFirst
}
mu.Unlock() mu.Unlock()
if !wasStopped { if !wasStopped && wasServing {
logf("[mesh-tools] %s; started again on its next call", why) // Every bundle stays up (ADR 0198): code that runs long — a handler, a provisioner —
// is not waiting for a call to bring it back, and a tool bundle back early costs nothing.
logf("[mesh-tools] %s; started again in %s", why, wait.Round(100*time.Millisecond))
restart(wait)
} }
}() }()
if _, err := c.ask(module, "initialize", map[string]any{"protocolVersion": Protocol, "capabilities": map[string]any{}, if _, err := c.ask(module, "initialize", map[string]any{"protocolVersion": Protocol, "capabilities": map[string]any{},
@@ -248,19 +315,62 @@ func Start(module, entry string, env []string, publish Publisher, logf func(stri
return nil, err return nil, err
} }
mu.Lock() mu.Lock()
current = fresh if current == nil {
current = fresh
startedAt = time.Now()
} else {
_ = fresh.cmd.Process.Signal(syscall.SIGTERM)
fresh = current
}
c = fresh c = fresh
mu.Unlock() mu.Unlock()
} }
return c.ask(module, method, params, timeout) return c.ask(module, method, params, timeout)
} }
restart = func(after time.Duration) {
go func() {
time.Sleep(after)
mu.Lock()
if stopped || current != nil {
mu.Unlock()
return
}
mu.Unlock()
fresh, err := start()
mu.Lock()
defer mu.Unlock()
if err != nil {
if !stopped {
wait := pause
if pause *= 2; pause > RestartMost {
pause = RestartMost
}
logf("[mesh-tools] %s's bundle did not start again: %v; trying in %s", module, err, wait)
go restart(wait)
}
return
}
if stopped {
_ = fresh.cmd.Process.Signal(syscall.SIGTERM)
return
}
if current == nil {
current = fresh
startedAt = time.Now()
} else {
_ = fresh.cmd.Process.Signal(syscall.SIGTERM)
}
}()
}
first, err := start() first, err := start()
if err != nil { if err != nil {
return nil, err return nil, err
} }
mu.Lock() mu.Lock()
current = first current = first
startedAt = time.Now()
mu.Unlock() mu.Unlock()
raw, err := asking("tools/list", map[string]any{}, HandshakeTimeout) raw, err := asking("tools/list", map[string]any{}, HandshakeTimeout)
+41
View File
@@ -1,5 +1,9 @@
// Package meshtest raises what the controller would, for tests against a real bus: the ASSIGNMENTS // Package meshtest raises what the controller would, for tests against a real bus: the ASSIGNMENTS
// and EVENTS streams, memberships issued by hand, and a fixture's path. // and EVENTS streams, memberships issued by hand, and a fixture's path.
//
// **The packages share one bus, so run them one at a time: `go test -p 1 ./...`.** Each test raises
// the streams afresh, and the console discovers every runtime that announces itself on the bus
// (novox/hq ADR 0197) — a runtime from another package's test is, correctly, found.
package meshtest package meshtest
import ( import (
@@ -141,3 +145,40 @@ func (l *Logs) Has(fragments ...string) bool {
// All is every line. // All is every line.
func (l *Logs) All() string { return strings.Join(l.lines, "\n") } func (l *Logs) All() string { return strings.Join(l.lines, "\n") }
// Consumer makes a module's durable consumer on a machine as the controller does: pull, on the
// EVENTS stream, filtered to what the module consumes, with a short ack wait so a test sees a
// redelivery in seconds rather than the mesh's minutes.
func (m *Mesh) Consumer(t *testing.T, node, module string, filters []string, ackWait time.Duration) {
t.Helper()
cfg := &nats.ConsumerConfig{Durable: node + "_" + module, AckPolicy: nats.AckExplicitPolicy,
AckWait: ackWait, MaxDeliver: 10, DeliverPolicy: nats.DeliverNewPolicy}
if len(filters) == 1 {
cfg.FilterSubject = filters[0]
} else {
cfg.FilterSubjects = filters
}
if _, err := m.js.AddConsumer("EVENTS", cfg); err != nil {
t.Fatal(err)
}
}
// Emit publishes an event as a module would, into the EVENTS stream.
func (m *Mesh) Emit(t *testing.T, subject string, body any) {
t.Helper()
data, _ := json.Marshal(body)
if _, err := m.js.Publish(subject, data); err != nil {
t.Fatal(err)
}
}
// Pending is what a module's consumer still holds: delivered and not acknowledged, and not yet
// delivered.
func (m *Mesh) Pending(t *testing.T, node, module string) (ackPending, notDelivered uint64) {
t.Helper()
info, err := m.js.ConsumerInfo("EVENTS", node+"_"+module)
if err != nil {
t.Fatal(err)
}
return uint64(info.NumAckPending), info.NumPending
}
@@ -0,0 +1,113 @@
package runtime
import (
"encoding/json"
"sort"
"strings"
"testing"
"time"
"github.com/nats-io/nats.go"
"github.com/nats-io/nats.go/micro"
"github.com/novox/mesh-tools/node-tools/internal/bus"
mt "github.com/novox/mesh-tools/node-tools/internal/meshtest"
)
// novox/hq ADR 0197: what a runtime serves, it announces — the NATS services protocol's discovery,
// read here with NATS's own types, each endpoint a subject actually served — and the announcement
// follows a re-issued membership.
func TestTheRuntimeAnnouncesWhatItServesInTheServicesProtocol(t *testing.T) {
mesh := mt.New(t)
mesh.Issue(t, mt.MembershipOf("alpha", "anchor", false, nil))
mesh.Issue(t, mt.MembershipOf("beta", "anchor", false, map[string][]string{"node-shelf": {"list", "clear"}}))
nodeTools := connect(t, "node-tools", "anchor")
stop, err := Run(nodeTools, []Served{
{"alpha", []string{mt.Fixture("many-alpha.serve.mjs")}},
{"beta", []string{mt.Fixture("many-beta.serve.mjs")}},
}, nil, (&mt.Logs{}).Logf)
if err != nil {
t.Fatal(err)
}
defer stop()
nc, err := nats.Connect(mt.URL(t))
if err != nil {
t.Fatal(err)
}
defer nc.Close()
info := func() micro.Info {
t.Helper()
msg, err := nc.Request("$SRV.INFO.node-tools.anchor", nil, 2*time.Second)
if err != nil {
t.Fatal(err)
}
var i micro.Info
if err := json.Unmarshal(msg.Data, &i); err != nil {
t.Fatal(err)
}
return i
}
subjects := func(i micro.Info) string {
var out []string
for _, e := range i.Endpoints {
out = append(out, e.Subject+"|"+e.QueueGroup+"|"+e.Metadata["kind"]+"|"+e.Metadata["module"]+"|"+e.Metadata["seat"]+"|"+e.Metadata["scope"]+"|"+e.Metadata["interchangeable"])
}
sort.Strings(out)
return strings.Join(out, "\n")
}
got := info()
if got.Type != micro.InfoResponseType || got.Name != "node-tools" || got.ID != "anchor" || got.Version == "" {
t.Errorf("identity: %+v", got.ServiceIdentity)
}
want := strings.Join([]string{
"mesh.mod.alpha.tool.one.anchor||tool|alpha|||false",
"mesh.mod.alpha.tool.two.anchor||tool|alpha|||false",
"mesh.mod.beta.tool.five.anchor||tool|beta|||false",
"mesh.mod.beta.tool.four.anchor||tool|beta|||false",
"mesh.mod.beta.tool.three.anchor||tool|beta|||false",
"mesh.seat.node-shelf.tool.clear.anchor||seat|beta|node-shelf|node|false",
"mesh.seat.node-shelf.tool.list.anchor||seat|beta|node-shelf|node|false",
}, "\n")
if s := subjects(got); s != want {
t.Errorf("announced:\n%s\nwant:\n%s", s, want)
}
for _, e := range got.Endpoints {
if e.Metadata["node"] != "anchor" || e.Metadata["description"] == "" || !json.Valid([]byte(e.Metadata["schema"])) {
t.Errorf("endpoint metadata: %+v", e)
}
}
// Every subject announced is answered.
asker := connect(t, "console", "workstation")
for _, e := range got.Endpoints {
if _, err := asker.Ask("", map[string]any{}, e.Subject); err != nil {
t.Errorf("announced %s and does not answer it: %v", e.Subject, err)
}
}
// PING answers with the same identity, and a request for another service is not answered.
if msg, err := nc.Request("$SRV.PING", nil, 2*time.Second); err != nil || !strings.Contains(string(msg.Data), micro.PingResponseType) {
t.Errorf("ping: %v %v", msg, err)
}
if _, err := nc.Request("$SRV.INFO.somebody-else", nil, 300*time.Millisecond); err == nil {
t.Error("answered a request for another service")
}
// alpha is re-issued a plain subject: the announcement says so, and says it is interchangeable.
mesh.Issue(t, mt.MembershipOf("alpha", "anchor", true, nil))
var after string
for i := 0; i < 40; i++ {
after = subjects(info())
if strings.Contains(after, "mesh.mod.alpha.tool.one|serve.alpha|tool|alpha|||true") {
break
}
time.Sleep(100 * time.Millisecond)
}
if !strings.Contains(after, "mesh.mod.alpha.tool.one|serve.alpha|tool|alpha|||true") ||
!strings.Contains(after, "mesh.mod.alpha.tool.one.anchor||tool|alpha|||true") {
t.Errorf("after a re-issued membership:\n%s", after)
}
_ = bus.Served{}
}
+157
View File
@@ -0,0 +1,157 @@
package runtime
import (
"fmt"
"os"
"path/filepath"
"strings"
"sync"
"testing"
"time"
"github.com/novox/mesh-tools/node-tools/internal/bus"
"github.com/novox/mesh-tools/node-tools/internal/launch"
mt "github.com/novox/mesh-tools/node-tools/internal/meshtest"
)
// A logger safe to write from the runtime's goroutines.
type lines struct {
mu sync.Mutex
l []string
}
func (s *lines) logf(format string, args ...any) {
s.mu.Lock()
defer s.mu.Unlock()
s.l = append(s.l, strings.TrimSpace(sprintf(format, args...)))
}
func (s *lines) all() string {
s.mu.Lock()
defer s.mu.Unlock()
return strings.Join(s.l, "\n")
}
func sprintf(format string, args ...any) string { return fmtSprintf(format, args...) }
func countLines(t *testing.T, path, want string) int {
t.Helper()
b, _ := os.ReadFile(path)
n := 0
for _, l := range strings.Split(string(b), "\n") {
if l == want {
n++
}
}
return n
}
// novox/hq ADR 0198: a module's long-running code is launched by the node's runtime, and the runtime
// is its bus — it binds the module's own consumer, hands each event to the bundle, and acknowledges it
// only when the bundle has taken it; one the bundle failed or died on is offered again.
func TestTheRuntimeHandsAModulesEventsToItsBundleAndAcknowledgesThemOnlyWhenTaken(t *testing.T) {
mesh := mt.New(t)
restore := bus.NakDelay
bus.NakDelay = 300 * time.Millisecond
t.Cleanup(func() { bus.NakDelay = restore })
firstPause := launch.RestartFirst
launch.RestartFirst = 100 * time.Millisecond
t.Cleanup(func() { launch.RestartFirst = firstPause })
mesh.Issue(t, mt.MembershipOf("watcher", "anchor", false, nil))
mesh.Issue(t, mt.MembershipOf("beta", "anchor", true, nil))
// The controller's consumer for the module, as its own runtime bound it: anchor_watcher on EVENTS.
mesh.Consumer(t, "anchor", "watcher", []string{"mesh.mod.alpha.event.>"}, 2*time.Second)
nodeTools := connect(t, "node-tools", "anchor")
asker := connect(t, "console", "workstation")
dir := t.TempDir()
log := filepath.Join(dir, "watch.log")
said := &lines{}
stop, err := Run(nodeTools, []Served{
{Module: "watcher", Entrypoints: []string{mt.Fixture("watcher.serve.mjs")}},
{Module: "beta", Entrypoints: []string{mt.Fixture("many-beta.serve.mjs")}},
}, map[string]map[string]string{"watcher": {"WATCH_LOG": log}}, said.logf)
if err != nil {
t.Fatal(err)
}
t.Cleanup(stop)
if !strings.Contains(said.all(), "watcher's events arrive on its consumer anchor_watcher") {
t.Fatalf("the module's consumer was not bound:\n%s", said.all())
}
// Handled once, acknowledged once.
mesh.Emit(t, "mesh.mod.alpha.event.happened", map[string]any{"n": 1})
mt.Until(t, func() error {
if countLines(t, log, "handled 1") != 1 {
return errorf("event 1 not handled yet")
}
return nil
})
// The handler fails the first time: not acknowledged, offered again, then handled.
mesh.Emit(t, "mesh.mod.alpha.event.happened", map[string]any{"n": 2, "fail": true})
// The bundle dies on the first offer: not acknowledged, the bundle is started again and the
// event is offered again once its ack wait has passed.
mesh.Emit(t, "mesh.mod.alpha.event.happened", map[string]any{"n": 3, "die": true})
deadline := time.Now().Add(20 * time.Second)
for time.Now().Before(deadline) {
if countLines(t, log, "handled 2") == 1 && countLines(t, log, "handled 3") == 1 {
break
}
time.Sleep(200 * time.Millisecond)
}
if countLines(t, log, "handled 2") != 1 || countLines(t, log, "handled 3") != 1 {
b, _ := os.ReadFile(log)
t.Fatalf("a failed or interrupted event was not offered again and handled once:\n%s\nruntime:\n%s", b, said.all())
}
if countLines(t, log, "started") < 2 {
t.Errorf("the bundle that died was not started again: %s", said.all())
}
mt.Until(t, func() error {
ack, notYet := mesh.Pending(t, "anchor", "watcher")
if ack != 0 || notYet != 0 {
return errorf("still pending: %d unacknowledged, %d undelivered", ack, notYet)
}
return nil
})
if countLines(t, log, "handled 1") != 1 {
t.Errorf("event 1 was handled more than once")
}
// A tool asks another module's tool through the runtime, as the module.
got, err := call(t, asker, "watcher.relay@anchor", map[string]any{})
if err != nil {
t.Fatal(err)
}
same(t, got, `{"beta":3}`)
}
// A long-running bundle that exits is started again, without waiting for a call (ADR 0198).
func TestALongRunningBundleThatExitsIsStartedAgain(t *testing.T) {
mesh := mt.New(t)
firstPause := launch.RestartFirst
launch.RestartFirst = 100 * time.Millisecond
t.Cleanup(func() { launch.RestartFirst = firstPause })
mesh.Issue(t, mt.MembershipOf("flaky", "anchor", false, nil))
nodeTools := connect(t, "node-tools", "anchor")
log := filepath.Join(t.TempDir(), "flaky.log")
said := &lines{}
stop, err := Run(nodeTools, []Served{{Module: "flaky", Entrypoints: []string{mt.Fixture("flaky.serve.mjs")}}},
map[string]map[string]string{"flaky": {"FLAKY_LOG": log}}, said.logf)
if err != nil {
t.Fatal(err)
}
t.Cleanup(stop)
mt.Until(t, func() error {
if countLines(t, log, "started") < 3 {
return errorf("started %d time(s)", countLines(t, log, "started"))
}
return nil
})
if !strings.Contains(said.all(), "flaky's bundle exited (3)") || !strings.Contains(said.all(), "started again in") {
t.Errorf("the restart was not said:\n%s", said.all())
}
}
func fmtSprintf(format string, args ...any) string { return fmt.Sprintf(format, args...) }
func errorf(format string, args ...any) error { return fmt.Errorf(format, args...) }
+209 -7
View File
@@ -14,6 +14,7 @@ import (
"strings" "strings"
"sync" "sync"
"github.com/novox/mesh-tools/node-tools/internal/announce"
"github.com/novox/mesh-tools/node-tools/internal/bus" "github.com/novox/mesh-tools/node-tools/internal/bus"
"github.com/novox/mesh-tools/node-tools/internal/launch" "github.com/novox/mesh-tools/node-tools/internal/launch"
) )
@@ -156,6 +157,10 @@ func Run(conn *bus.Conn, served []Served, envs map[string]map[string]string, log
return out return out
} }
// 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.
events := &consumers{conn: conn, logf: logf, of: map[string]*moduleEvents{}}
failed := map[string]string{} failed := map[string]string{}
var registrations []registration var registrations []registration
var stops []func() var stops []func()
@@ -171,13 +176,7 @@ func Run(conn *bus.Conn, served []Served, envs map[string]map[string]string, log
fail(path + " is not executable; a bundle the runtime serves is started, never imported, and its build makes it executable (novox/hq ADR 0193)") fail(path + " is not executable; a bundle the runtime serves is started, never imported, and its build makes it executable (novox/hq ADR 0193)")
continue continue
} }
child, err := launch.Start(module, path, envFor(module), func(params json.RawMessage) error { child, err := launch.Start(module, path, envFor(module), events.forModule(module), logf)
var env bus.Envelope
if err := json.Unmarshal(params, &env); err != nil {
return fmt.Errorf("not an event envelope: %w", err)
}
return conn.PublishAs(module, env)
}, logf)
if err != nil { if err != nil {
fail(err.Error()) fail(err.Error())
continue continue
@@ -208,6 +207,7 @@ func Run(conn *bus.Conn, served []Served, envs map[string]map[string]string, log
} }
} }
stopAll := func() { stopAll := func() {
events.stopAll()
for i := len(stops) - 1; i >= 0; i-- { for i := len(stops) - 1; i >= 0; i-- {
stops[i]() stops[i]()
} }
@@ -276,10 +276,83 @@ func Run(conn *bus.Conn, served []Served, envs map[string]map[string]string, log
} }
logf("%s", line) logf("%s", line)
stops = append(stops, serveSeats(conn, modules, registrations, logf)) stops = append(stops, serveSeats(conn, modules, registrations, logf))
// **What it serves, it announces** (novox/hq ADR 0197): the NATS services protocol's discovery,
// answered with what is served at the moment it is asked — re-served memberships included.
announced, err := announce.Serve(conn, announce.Service{
Name: conn.Module(), ID: instanceOf(conn),
Description: "the mesh's tool runtime on " + node + ": every assigned module's tools and the seats they hold",
Metadata: map[string]string{"node": node},
}, func() []announce.Endpoint { return endpointsOf(conn, own, registrations, modules) })
if err != nil {
logf("[mesh-tools] cannot announce what it serves: %v", err)
} else {
stops = append(stops, announced)
}
conn.Flush() conn.Flush()
return stopAll, nil return stopAll, nil
} }
// instanceOf is this runtime's instance on the bus: its machine, which is what tells two instances of
// one service apart; the connection's module where it has no machine.
func instanceOf(conn *bus.Conn) string {
if n := conn.Node(); n != "" {
return n
}
return conn.Module()
}
// endpointsOf is everything this runtime serves now: each served module's tools on every subject the
// mesh issued for them, and each held seat's verbs on the seat's subject (ADR 0197).
func endpointsOf(conn *bus.Conn, own []registration, registrations []registration, modules []string) []announce.Endpoint {
node := conn.Node()
var out []announce.Endpoint
for _, r := range own {
for _, t := range r.tools {
served := conn.ServedOn(r.module, t.Name)
interchangeable := false
for _, s := range served {
interchangeable = interchangeable || s.Subject == "mesh.mod."+r.module+".tool."+t.Name
}
for _, s := range served {
out = append(out, announce.Endpoint{Kind: announce.KindTool, Module: r.module, Tool: t.Name,
Node: node, Description: t.Description, Schema: t.Input, Interchangeable: interchangeable,
Subject: s.Subject, Queue: s.Queue})
}
}
}
impl := map[string]map[string]launch.Tool{}
for _, r := range registrations {
if impl[r.module] == nil {
impl[r.module] = map[string]launch.Tool{}
}
for _, t := range r.tools {
impl[r.module][t.Name] = t
}
}
have := map[string]bool{}
for _, module := range modules {
m := conn.Membership(module)
if m == nil {
continue
}
for _, v := range m.Seats {
t, ok := impl[v.Seat][v.Verb]
if !ok || have[v.Subject] {
continue
}
have[v.Subject] = true
scope := "mesh"
if node != "" && strings.HasSuffix(v.Subject, "."+node) {
scope = "node"
}
out = append(out, announce.Endpoint{Kind: announce.KindSeat, Module: module, Tool: v.Verb, Seat: v.Seat,
Scope: scope, Node: node, Description: t.Description, Schema: t.Input, Subject: v.Subject})
}
}
return out
}
func orNone(names []string) string { func orNone(names []string) string {
if len(names) == 0 { if len(names) == 0 {
return "(none)" return "(none)"
@@ -376,3 +449,132 @@ func serveSeats(conn *bus.Conn, modules []string, registrations []registration,
stops = nil stops = nil
} }
} }
// consumers holds, per module the runtime serves, the one durable consumer its events arrive on and
// the children its events are handed to (novox/hq ADR 0198).
type consumers struct {
conn *bus.Conn
logf func(string, ...any)
mu sync.Mutex
of map[string]*moduleEvents
}
type moduleEvents struct {
stop func()
delivers []*func(json.RawMessage) error
}
// stopAll unbinds every module's consumer.
func (c *consumers) stopAll() {
c.mu.Lock()
defer c.mu.Unlock()
for _, m := range c.of {
if m.stop != nil {
m.stop()
}
}
c.of = map[string]*moduleEvents{}
}
// forModule is the bus one launched bundle of a module reaches the mesh through.
func (c *consumers) forModule(module string) launch.Bus {
return &moduleBus{all: c, module: module}
}
type moduleBus struct {
all *consumers
module string
mu sync.Mutex
deliver *func(json.RawMessage) error
}
// Publish emits an event as the module (ADR 0193).
func (b *moduleBus) Publish(params json.RawMessage) error {
var env bus.Envelope
if err := json.Unmarshal(params, &env); err != nil {
return fmt.Errorf("not an event envelope: %w", err)
}
return b.all.conn.PublishAs(b.module, env)
}
// Ask calls a tool as the module: `{key, body}`, answered with the tool's result (ADR 0198).
func (b *moduleBus) Ask(params json.RawMessage) (json.RawMessage, error) {
var asked struct {
Key string `json:"key"`
Body json.RawMessage `json:"body"`
}
if err := json.Unmarshal(params, &asked); err != nil || asked.Key == "" {
return nil, fmt.Errorf("mesh/ask names no tool: {key, body}")
}
body := any(asked.Body)
if len(asked.Body) == 0 {
body = map[string]any{}
}
answered, err := b.all.conn.AskAs(b.module, asked.Key, body)
if err != nil {
return nil, err
}
if len(answered.Result) == 0 {
return json.RawMessage("null"), nil
}
return answered.Result, nil
}
// Subscribe hands this bundle the module's events. The consumer is bound once per module; each of
// the module's children that subscribed is handed every event, and the event is acknowledged only
// when all of them took it — one consumer split between two readers would give each half.
func (b *moduleBus) Subscribe(deliver func(json.RawMessage) error) error {
b.mu.Lock()
if b.deliver == nil {
d := deliver
b.deliver = &d
} else {
*b.deliver = deliver
}
mine := b.deliver
b.mu.Unlock()
c := b.all
c.mu.Lock()
defer c.mu.Unlock()
m := c.of[b.module]
if m == nil {
m = &moduleEvents{}
c.of[b.module] = m
}
listed := false
for _, d := range m.delivers {
if d == mine {
listed = true
}
}
if !listed {
m.delivers = append(m.delivers, mine)
}
if m.stop != nil {
return nil
}
module := b.module
stop, err := c.conn.ConsumeAs(module, func(env bus.Envelope) error {
raw, err := json.Marshal(env)
if err != nil {
return err
}
c.mu.Lock()
targets := append([]*func(json.RawMessage) error(nil), c.of[module].delivers...)
c.mu.Unlock()
for _, d := range targets {
if err := (*d)(raw); err != nil {
c.logf("[mesh-tools] %s did not take %s: %v; offered again", module, env.Key, err)
return err
}
}
return nil
})
if err != nil {
return err
}
m.stop = stop
c.logf("[mesh-tools] %s's events arrive on its consumer %s", module, bus.ConsumerOf(c.conn.Node(), module))
return nil
}
+3 -2
View File
@@ -13,11 +13,12 @@
"test": "node --test --test-concurrency=1 --experimental-strip-types 'test/*.test.ts'" "test": "node --test --test-concurrency=1 --experimental-strip-types 'test/*.test.ts'"
}, },
"dependencies": { "dependencies": {
"@novox/mesh-sdk": "^0.1.0", "@novox/mesh-sdk": "^0.1.6",
"nats": "^2.29.0" "nats": "^2.29.0"
}, },
"devDependencies": { "devDependencies": {
"@types/node": "^22.20.1", "@types/node": "^22.20.1",
"typescript": "^5.9.3" "typescript": "^5.9.3",
"esbuild": "^0.25.0"
} }
} }
+139
View File
@@ -0,0 +1,139 @@
// What a runtime serves, it announces (novox/hq ADR 0197): the NATS services protocol's discovery —
// `$SRV.PING`, `$SRV.INFO`, `$SRV.STATS`, and the same followed by the service's name and its id —
// answered in the io.nats.micro.v1 format with what is served at the moment of the request. Serving is
// unchanged; this only says what is served. One service per runtime process, because the bus admits
// one reply per request from each responder: one endpoint per tool per subject, its metadata saying
// which module, seat, scope and machine it is. The same shape the Go runtime answers.
import { StringCodec } from "nats";
import { asSchema } from "@novox/mesh-sdk/stdio";
import type { ToolDefinition } from "@novox/mesh-sdk/tools";
import type { RuntimeBroker } from "./broker-nats.js";
const sc = StringCodec();
export const VERSION = "0.1.0";
export const INFO_RESPONSE = "io.nats.micro.v1.info_response";
export const PING_RESPONSE = "io.nats.micro.v1.ping_response";
export const STATS_RESPONSE = "io.nats.micro.v1.stats_response";
/** One tool served on one subject, as announced. */
export interface Endpoint {
kind: "tool" | "seat";
module: string;
tool: string;
seat?: string;
scope?: "mesh" | "node";
node: string;
description: string;
schema: unknown;
interchangeable: boolean;
subject: string;
queue?: string;
}
/** A seat's verb as the runtime serves it, with the definition that answers it. */
export interface ServedSeatVerb {
seat: string;
verb: string;
subject: string;
holder: string;
tool: ToolDefinition;
}
export interface Service {
name: string;
id: string;
description: string;
metadata: Record<string, string>;
}
/** The endpoint's name as the protocol allows it; the metadata, not the name, identifies it. */
function nameOf(e: Endpoint): string {
const clean = (s: string) => s.replace(/[^A-Za-z0-9_-]/g, "_");
return `${clean(e.kind === "seat" ? e.seat ?? "" : e.module)}__${clean(e.tool)}`;
}
/** The info_response for these endpoints. */
export function info(s: Service, endpoints: Endpoint[]): Record<string, unknown> {
return {
name: s.name, id: s.id, version: VERSION, metadata: s.metadata, type: INFO_RESPONSE, description: s.description,
endpoints: endpoints.map((e) => {
const metadata: Record<string, string> = {
kind: e.kind, module: e.module, tool: e.tool, node: e.node, description: e.description,
schema: JSON.stringify(e.schema ?? {}), interchangeable: e.interchangeable ? "true" : "false",
};
if (e.kind === "seat") {
metadata.seat = e.seat ?? "";
metadata.scope = e.scope ?? "mesh";
}
return { name: nameOf(e), subject: e.subject, queue_group: e.queue ?? "", metadata };
}),
};
}
/** Everything served now: each module's tools on every subject issued for them, and each held
* seat's verbs on the seat's subject. */
export function endpointsOf(
broker: RuntimeBroker,
own: { module: string; tools: ToolDefinition[] }[],
seats: ServedSeatVerb[],
): Endpoint[] {
const node = broker.node ?? "";
const out: Endpoint[] = [];
for (const { module, tools } of own) {
for (const t of tools) {
const served = broker.servedOn ? broker.servedOn(module, t.name) : [];
const interchangeable = served.some((s) => s.subject === `mesh.mod.${module}.tool.${t.name}`);
for (const s of served) {
out.push({ kind: "tool", module, tool: t.name, node, description: t.description, schema: asSchema(t.input),
interchangeable, subject: s.subject, queue: s.queue });
}
}
}
for (const v of seats) {
out.push({ kind: "seat", module: v.holder, tool: v.verb, seat: v.seat,
scope: node && v.subject.endsWith(`.${node}`) ? "node" : "mesh", node, description: v.tool.description,
schema: asSchema(v.tool.input), interchangeable: false, subject: v.subject });
}
return out;
}
/** Answer discovery for one service until stopped. */
export function announce(broker: RuntimeBroker, s: Service, current: () => Endpoint[]): () => void {
const started = new Date().toISOString();
const identity = { name: s.name, id: s.id, version: VERSION, metadata: s.metadata };
const answer = (subject: string): Uint8Array | undefined => {
const parts = subject.split(".");
if (parts[0] !== "$SRV" || parts.length < 2) return undefined;
if (parts.length >= 3 && parts[2] !== s.name) return undefined; // another service's
if (parts.length >= 4 && parts[3] !== s.id) return undefined; // another instance's
let v: unknown;
switch (parts[1]) {
case "PING":
v = { ...identity, type: PING_RESPONSE };
break;
case "INFO":
v = info(s, current());
break;
case "STATS":
v = { ...identity, type: STATS_RESPONSE, started, endpoints: current().map((e) => ({
name: nameOf(e), subject: e.subject, queue_group: e.queue ?? "", num_requests: 0, num_errors: 0,
last_error: "", processing_time: 0, average_processing_time: 0 })) };
break;
default:
return undefined;
}
return sc.encode(JSON.stringify(v));
};
const stops: (() => void)[] = [];
// Exactly the questions asked of every service and of this one by name and instance — what the
// grants allow (novox/hq ADR 0197). A wildcard is refused by the bus, and a refused subscription
// ends a runtime: on 2026-10-03 it crash-looped every container that announced itself.
for (const verb of ["PING", "INFO", "STATS"]) {
for (const subject of [`$SRV.${verb}`, `$SRV.${verb}.${s.name}`, `$SRV.${verb}.${s.name}.${s.id}`]) {
stops.push(broker.raw!(subject, (subj) => answer(subj)));
}
}
return () => stops.forEach((stop) => stop());
}
+49 -10
View File
@@ -97,6 +97,13 @@ export interface RuntimeBroker extends Broker {
serving(): string[]; serving(): string[];
/** The module this connection is: what its credential named, and what a bare key serves as. */ /** The module this connection is: what its credential named, and what a bare key serves as. */
readonly module: string; readonly module: string;
/** Where a served module's tool is answered right now (ADR 0197: what it serves, it announces). */
servedOn?(module: string, tool: string): { subject: string; queue?: string }[];
/** Answer a subject in a format of its own, not a tool's reply envelope — the NATS services
* protocol's discovery (novox/hq ADR 0197). An undefined answer is no reply. */
raw?(subject: string, answer: (subject: string, data: Uint8Array) => Uint8Array | undefined): () => void;
/** The machine this connection serves on, when its credential names one. */
readonly node?: string;
} }
/** /**
@@ -270,17 +277,27 @@ export async function connectNats(
const sub = conn.subscribe(subject, queue ? { queue } : {}); const sub = conn.subscribe(subject, queue ? { queue } : {});
subs.push(sub); subs.push(sub);
void (async () => { void (async () => {
for await (const msg of sub) { // **A refused subscription is said, never fatal** (novox/hq 04-ISSUES/218, as 217 for the
let reply: { result?: Res; error?: string; node?: string }; // announcements). The grants are the mesh's word on what this account may answer; a subject
try { // they leave out — a seat claimed here and held elsewhere — costs that subject, never the
reply = { result: await handler(JSON.parse(sc.decode(msg.data)) as Req) }; // module's other tools, its handlers and its provisioning. Unhandled, the refusal ended the
} catch (err) { // process and a module's runtime crash-looped on 2026-10-04.
// The caller is told, rather than left to time out: a handler that threw is a try {
// different failure from a tool nobody serves, and only one of them is worth retrying. for await (const msg of sub) {
reply = { error: err instanceof Error ? err.message : String(err) }; let reply: { result?: Res; error?: string; node?: string };
try {
reply = { result: await handler(JSON.parse(sc.decode(msg.data)) as Req) };
} catch (err) {
// The caller is told, rather than left to time out: a handler that threw is a
// different failure from a tool nobody serves, and only one of them is worth retrying.
reply = { error: err instanceof Error ? err.message : String(err) };
}
if (node) reply.node = node;
msg.respond(sc.encode(JSON.stringify(reply)));
} }
if (node) reply.node = node; } catch (err) {
msg.respond(sc.encode(JSON.stringify(reply))); console.log(`[mesh-tools] the bus refused ${subject}: ${err instanceof Error ? err.message : String(err)}; ` +
"not served here, and the rest serves on");
} }
})(); })();
return () => sub.unsubscribe(); return () => sub.unsubscribe();
@@ -348,6 +365,28 @@ export async function connectNats(
}, },
module: self, module: self,
node,
servedOn,
raw(subject: string, answer: (subject: string, data: Uint8Array) => Uint8Array | undefined): () => void {
const sub = conn.subscribe(subject);
subs.push(sub);
void (async () => {
// **A refusal here is said, never fatal** (novox/hq 04-ISSUES/217). This serves discovery —
// what the runtime says about itself — not the work; a bus that refuses it costs the mesh
// seeing this runtime, not the runtime's tools and handlers. Unhandled, the refusal ended the
// process and every per-module container crash-looped on 2026-10-03.
try {
for await (const msg of sub) {
const body = answer(msg.subject, msg.data);
if (body) msg.respond(body);
}
} catch (err) {
console.log(`[mesh-tools] the bus refused ${subject}: ${err instanceof Error ? err.message : String(err)}; ` +
"discovery will not see this runtime there, and it serves on");
}
})();
return () => sub.unsubscribe();
},
follow, follow,
serving: () => [...issued.keys()], serving: () => [...issued.keys()],
membership: (module?: string) => issued.get(module ?? self), membership: (module?: string) => issued.get(module ?? self),
+27 -4
View File
@@ -13,6 +13,8 @@
import { spawn, type ChildProcess } from "node:child_process"; import { spawn, type ChildProcess } from "node:child_process";
import { accessSync, constants } from "node:fs"; import { accessSync, constants } from "node:fs";
import type { ToolDefinition } from "@novox/mesh-sdk/tools"; import type { ToolDefinition } from "@novox/mesh-sdk/tools";
import { broker, type Envelope } from "@novox/mesh-sdk/messaging";
import { atWork } from "./broker-nats.js";
/** The protocol version this speaks; a bundle says the same. */ /** The protocol version this speaks; a bundle says the same. */
export const PROTOCOL = "2025-03-26"; export const PROTOCOL = "2025-03-26";
@@ -71,29 +73,50 @@ export async function launch(module: string, entry: string, env: NodeJS.ProcessE
const line = buffered.slice(0, at).trim(); const line = buffered.slice(0, at).trim();
buffered = buffered.slice(at + 1); buffered = buffered.slice(at + 1);
if (!line) continue; if (!line) continue;
let reply: { id?: number; result?: unknown; error?: { message?: string } }; let reply: { id?: number | string; method?: string; params?: unknown; result?: unknown; error?: { message?: string } };
try { try {
reply = JSON.parse(line); reply = JSON.parse(line);
} catch { } catch {
console.log(`[mesh-tools] ${module}'s bundle said something that is not a reply: ${line.slice(0, 120)}`); console.log(`[mesh-tools] ${module}'s bundle said something that is not a reply: ${line.slice(0, 120)}`);
continue; continue;
} }
// **The bundle asks the runtime to emit** (novox/hq ADR 0193): published on the bus as this
// module, and answered once the bus has accepted it, so the tool's emit means what it means
// in-process. Nothing else a bundle may ask.
if (typeof reply.method === "string") {
const id = reply.id;
const answer = (m: Record<string, unknown>) => proc.stdin!.write(JSON.stringify({ jsonrpc: "2.0", id, ...m }) + "\n");
if (reply.method !== "mesh/publish") {
if (id !== undefined) answer({ error: { code: -32601, message: `the runtime answers no ${reply.method} from a bundle` } });
continue;
}
atWork.run({ module }, () => broker().publish(reply.params as Envelope<unknown>))
.then(() => { if (id !== undefined) answer({ result: {} }); })
.catch((err: unknown) => { if (id !== undefined) answer({ error: { code: -32000, message: err instanceof Error ? err.message : String(err) } }); });
continue;
}
const waiting = typeof reply.id === "number" ? pending.get(reply.id) : undefined; const waiting = typeof reply.id === "number" ? pending.get(reply.id) : undefined;
if (!waiting) continue; if (!waiting) continue;
pending.delete(reply.id!); pending.delete(reply.id as number);
clearTimeout(waiting.timer); clearTimeout(waiting.timer);
if (reply.error) waiting.reject(new Error(reply.error.message ?? "the bundle refused the request")); if (reply.error) waiting.reject(new Error(reply.error.message ?? "the bundle refused the request"));
else waiting.resolve(reply.result); else waiting.resolve(reply.result);
} }
}); });
// stderr is the bundle's log; kept under the module's name so a fault reads where it belongs. // stderr is the bundle's log; kept under the module's name so a fault reads where it belongs.
// The last thing it said is kept, so a bundle that dies says why in its own words, not by code.
let lastSaid = "";
proc.stderr!.on("data", (chunk: Buffer) => { proc.stderr!.on("data", (chunk: Buffer) => {
for (const line of chunk.toString("utf8").split("\n")) if (line.trim()) console.log(`[${module}] ${line}`); for (const line of chunk.toString("utf8").split("\n")) {
if (!line.trim()) continue;
console.log(`[${module}] ${line}`);
if (/\S/.test(line) && !/^\s+at\s/.test(line) && !/^Node\.js v/.test(line)) lastSaid = line.trim();
}
}); });
const exited = new Promise<never>((_, reject) => { const exited = new Promise<never>((_, reject) => {
proc.once("error", (err) => reject(err)); proc.once("error", (err) => reject(err));
proc.once("exit", (code, signal) => { proc.once("exit", (code, signal) => {
const why = `${module}'s bundle exited (${signal ?? code})`; const why = `${module}'s bundle exited (${signal ?? code})` + (lastSaid ? `: ${lastSaid}` : "");
for (const [id, p] of pending) { for (const [id, p] of pending) {
pending.delete(id); pending.delete(id);
clearTimeout(p.timer); clearTimeout(p.timer);
+59 -102
View File
@@ -9,16 +9,14 @@
// and the per-module shape is the list with one entry. A bundle that fails to import is named — // and the per-module shape is the list with one entry. A bundle that fails to import is named —
// in the log and in what `tools` answers for it — and the others serve. // in the log and in what `tools` answers for it — and the others serve.
import { fileURLToPath, pathToFileURL } from "node:url"; import { pathToFileURL } from "node:url";
import { dirname, join, resolve } from "node:path"; import { resolve } from "node:path";
import { existsSync, readFileSync } from "node:fs";
import { registerHooks } from "node:module";
import { useBroker } from "@novox/mesh-sdk/messaging"; import { useBroker } from "@novox/mesh-sdk/messaging";
import * as sdkTools from "@novox/mesh-sdk/tools";
import { collectTools, toolKey, type ToolDefinition } from "@novox/mesh-sdk/tools"; import { collectTools, toolKey, type ToolDefinition } from "@novox/mesh-sdk/tools";
import type { Broker } from "@novox/mesh-sdk/messaging"; import type { Broker } from "@novox/mesh-sdk/messaging";
import { atWork, seatToolSubject, type Credential, type RuntimeBroker } from "./broker-nats.js"; import { atWork, seatToolSubject, type Credential, type RuntimeBroker } from "./broker-nats.js";
import { launch, launches } from "./launch.js"; import { launch, launches } from "./launch.js";
import { announce, endpointsOf, type ServedSeatVerb } from "./announce.js";
/** /**
* The one verb every module's runtime answers for it (novox/hq ADR 0152, design 34 §3): the * The one verb every module's runtime answers for it (novox/hq ADR 0152, design 34 §3): the
@@ -115,6 +113,7 @@ export async function runTools(opts: RuntimeOptions): Promise<() => void> {
for (const s of opts.serves ?? []) { for (const s of opts.serves ?? []) {
served.set(s.module, [...(served.get(s.module) ?? []), ...s.entrypoints]); served.set(s.module, [...(served.get(s.module) ?? []), ...s.entrypoints]);
} }
const ownEntrypoints = opts.moduleEntrypoints ?? [];
if (opts.moduleEntrypoints?.length) { if (opts.moduleEntrypoints?.length) {
if (!self) { if (!self) {
throw new Error( throw new Error(
@@ -143,34 +142,38 @@ export async function runTools(opts: RuntimeOptions): Promise<() => void> {
const envs = opts.envs ?? new Map<string, Record<string, string>>(); const envs = opts.envs ?? new Map<string, Record<string, string>>();
const envFor = (module: string): NodeJS.ProcessEnv => ({ ...process.env, ...(envs.get(module) ?? {}) }); const envFor = (module: string): NodeJS.ProcessEnv => ({ ...process.env, ...(envs.get(module) ?? {}) });
// Import each bundle, guarded (ADR 0175: one faulty bundle must not take the node's tools down). // Every bundle this runtime serves is launched as a process and spoken to over MCP on stdio (ADR
// Importing the entrypoint runs its registerModuleTools(...) — that is the whole handshake — and // 0188, ADR 0193): given the runtime's words and its module's own, and told the module it serves
// the registrations it adds are the ones that appear after it, which is how each is attributed // it as. The runtime knows no language; an entrypoint that is not executable was not built to be
// to the module whose bundle made it. // served, and is refused by name. A bundle that fails to start is named, and the others serve
// A bundle that is not plain JavaScript — or is marked executable — is launched as a process // (ADR 0175: one faulty bundle must not take the node's tools down).
// and spoken to over MCP on stdio instead (ADR 0188); what it lists is registered the same way. //
// Every bundle's import of the SDK resolves to this runtime's copy (04-ISSUES/209): one registry // The one-module form — the credential's own module's entrypoints, which the per-module containers
// of tools, one broker. Installed before the first bundle is imported. // still use for their event handlers and provisioners until they move (to-be 38 WP4c) — is imported
oneSdk(); // into this process as before: one module, one SDK, its container's own environment.
// The machine this runtime serves, which an event a launched tool emits is stamped with.
const node = opts.credential?.node ?? (typeof (runtime as unknown as { node?: unknown }).node === "string" ? (runtime as unknown as { node: string }).node : undefined);
const failed = new Map<string, string>(); const failed = new Map<string, string>();
const owner: string[] = []; // registration index → the module whose bundle registered it
const launched: { module: string; owner: string; tools: ToolDefinition[] }[] = []; const launched: { module: string; owner: string; tools: ToolDefinition[] }[] = [];
const children: Array<() => void> = []; const children: Array<() => void> = [];
for (const [module, entrypoints] of served) { for (const [module, entrypoints] of served) {
const imported = module === self && ownEntrypoints.length > 0;
for (const entry of entrypoints) { for (const entry of entrypoints) {
const path = resolve(entry); const path = resolve(entry);
try { try {
if (launches(path)) { if (imported) {
// Told the module it serves it as, so a seat's verbs are the seat's (ADR 0193). await import(pathToFileURL(path).href);
const child = await launch(module, path, { ...envFor(module), MESH_SERVED_MODULE: module });
children.push(child.stop);
for (const r of child.registrations) launched.push({ ...r, owner: module });
continue; continue;
} }
const before = collectTools().length; if (!launches(path)) {
await import(pathToFileURL(path).href); throw new Error(`${path} is not executable; a bundle the runtime serves is started, never imported, and its build makes it executable (novox/hq ADR 0193)`);
const after = collectTools().length; }
for (let i = before; i < after; i++) owner[i] = module; const child = await launch(module, path, {
...envFor(module), MESH_SERVED_MODULE: module, MESH_MODULE: module,
...(node ? { MESH_NODE: node } : {}),
});
children.push(child.stop);
for (const r of child.registrations) launched.push({ ...r, owner: module });
} catch (err) { } catch (err) {
const why = err instanceof Error ? err.message : String(err); const why = err instanceof Error ? err.message : String(err);
failed.set(module, why); failed.set(module, why);
@@ -187,16 +190,8 @@ export async function runTools(opts: RuntimeOptions): Promise<() => void> {
// out rather than fatal — on 2026-10-01 the credential of a module that had just learned to // out rather than fatal — on 2026-10-01 the credential of a module that had just learned to
// implement a seat did not yet name the claim, and the whole runtime restarted for it. // implement a seat did not yet name the claim, and the whole runtime restarted for it.
const claimed = seatsClaimed(served.keys(), self, opts.credential, runtime); const claimed = seatsClaimed(served.keys(), self, opts.credential, runtime);
// Each registration's contributor is given its own module's environment and no other's (ADR 0192):
// the runtime's own words, and over them what the mesh composed for the module whose bundle made
// the registration. An SDK too old to ask per registration cannot do that; said, not hidden.
const each = (sdkTools as { collectToolsEach?: (f: (module: string, i: number) => NodeJS.ProcessEnv) => { module: string; tools: ToolDefinition[] }[] }).collectToolsEach;
if (!each && envs.size > 0) {
console.log(`[mesh-tools] this runtime's SDK cannot give each bundle its own environment; ${[...envs.keys()].join(", ")} serve with the runtime's words only (novox/hq ADR 0192)`);
}
const collected = each ? each((module, i) => envFor(owner[i] ?? self ?? module)) : collectTools();
const registrations = [ const registrations = [
...collected.map((r, i) => ({ ...r, owner: owner[i] ?? self ?? r.module })), ...collectTools().map((r) => ({ ...r, owner: self ?? r.module })),
...launched, ...launched,
]; ];
const ownRegistrations = registrations.filter(({ module, owner: by }) => { const ownRegistrations = registrations.filter(({ module, owner: by }) => {
@@ -265,7 +260,16 @@ export async function runTools(opts: RuntimeOptions): Promise<() => void> {
console.log(`[mesh-tools] serving ${names.length} tool(s) for ${served.size} module(s): ${names.join(", ") || "(none)"}` + console.log(`[mesh-tools] serving ${names.length} tool(s) for ${served.size} module(s): ${names.join(", ") || "(none)"}` +
(failed.size ? `; not serving ${[...failed.keys()].join(", ")}, whose bundle(s) failed to load` : "")); (failed.size ? `; not serving ${[...failed.keys()].join(", ")}, whose bundle(s) failed to load` : ""));
stops.push(await serveClaimedSeats(runtime, [...served.keys()], self, opts.credential, registrations)); const seats = await serveClaimedSeats(runtime, [...served.keys()], self, opts.credential, registrations);
stops.push(seats.stop);
// **What it serves, it announces** (novox/hq ADR 0197), asked at the moment of the request.
if (typeof runtime.raw === "function") {
stops.push(announce(runtime, {
name: self ?? "runtime", id: runtime.node ?? self ?? "runtime",
description: `the tool runtime of ${self ?? "a module"}${runtime.node ? ` on ${runtime.node}` : ""}`,
metadata: runtime.node ? { node: runtime.node } : {},
}, () => endpointsOf(runtime, ownRegistrations, seats.serving())));
}
return () => stop(); return () => stop();
} }
@@ -309,11 +313,17 @@ async function serveClaimedSeats(
self: string | undefined, self: string | undefined,
credential: Credential | undefined, credential: Credential | undefined,
registrations: { module: string; owner: string; tools: ToolDefinition[] }[], registrations: { module: string; owner: string; tools: ToolDefinition[] }[],
): Promise<() => void> { ): Promise<{ stop: () => void; serving: () => ServedSeatVerb[] }> {
if (typeof broker.handleSubject !== "function") return () => {}; if (typeof broker.handleSubject !== "function") return { stop: () => {}, serving: () => [] };
// A seat's verbs are the role's, not the software's (ADR 0159): implemented under the seat's // A seat's verbs are the role's, not the software's (ADR 0159): implemented under the seat's
// name — `registerModuleTools("mesh-store", …)` — and never confused with the module's own tools. // name — `registerModuleTools("mesh-store", …)` — and never confused with the module's own tools.
const implementations = new Map<string, Map<string, (args: Record<string, unknown>) => Promise<unknown>>>(); const implementations = new Map<string, Map<string, (args: Record<string, unknown>) => Promise<unknown>>>();
const definitions = new Map<string, Map<string, ToolDefinition>>();
for (const { module, tools } of registrations) {
const defs = definitions.get(module) ?? new Map<string, ToolDefinition>();
for (const t of tools) defs.set(t.name, t);
definitions.set(module, defs);
}
for (const { module, owner, tools } of registrations) { for (const { module, owner, tools } of registrations) {
const verbs = implementations.get(module) ?? new Map<string, (args: Record<string, unknown>) => Promise<unknown>>(); const verbs = implementations.get(module) ?? new Map<string, (args: Record<string, unknown>) => Promise<unknown>>();
for (const t of tools) verbs.set(t.name, (args) => atWork.run({ module: owner }, () => t.run(args))); for (const t of tools) verbs.set(t.name, (args) => atWork.run({ module: owner }, () => t.run(args)));
@@ -332,13 +342,16 @@ async function serveClaimedSeats(
for (const module of served) { for (const module of served) {
const m = typeof broker.membership === "function" ? broker.membership(module) : undefined; const m = typeof broker.membership === "function" ? broker.membership(module) : undefined;
for (const s of m?.seats ?? []) add({ seat: s.seat, verb: s.verb, subject: s.subject, holder: module }); for (const s of m?.seats ?? []) add({ seat: s.seat, verb: s.verb, subject: s.subject, holder: module });
// The credential's claims, for the module's own runtime: where the mesh issued the verb when // The credential's claims, for the module's own runtime, ONLY until the mesh issues a
// it has; the derived shape until then. // membership (novox/hq issue 218). The credential names what the module claims; the membership
if (module !== self) continue; // names what it holds on this machine. A seat held once for the mesh is claimed by every
// machine running the module and held by one, so once a membership exists it decides: a claim
// it leaves out is not held here, and serving it anyway announced the seat from a machine the
// bus then refused it on.
if (module !== self || m) continue;
for (const claim of credential?.claims ?? []) { for (const claim of credential?.claims ?? []) {
for (const verb of claim.serves ?? []) { for (const verb of claim.serves ?? []) {
const subject = m?.seats?.find((s) => s.seat === claim.seat && s.verb === verb)?.subject const subject = seatToolSubject(claim.seat, verb, claim.scope, credential?.node);
?? seatToolSubject(claim.seat, verb, claim.scope, credential?.node);
add({ seat: claim.seat, verb, subject, holder: module }); add({ seat: claim.seat, verb, subject, holder: module });
} }
} }
@@ -347,81 +360,25 @@ async function serveClaimedSeats(
}; };
let stops: (() => void)[] = []; let stops: (() => void)[] = [];
let servingNow: ServedSeatVerb[] = [];
const serve = async (): Promise<void> => { const serve = async (): Promise<void> => {
stops.forEach((s) => s()); stops.forEach((s) => s());
stops = []; stops = [];
servingNow = [];
for (const v of wanted()) { for (const v of wanted()) {
const run = implementations.get(v.seat)?.get(v.verb); const run = implementations.get(v.seat)?.get(v.verb);
const tool = definitions.get(v.seat)?.get(v.verb);
if (!run) { if (!run) {
console.log(`[mesh-tools] ${v.holder} claims ${v.seat} and implements no ${v.verb}, which that seat promises; not served`); console.log(`[mesh-tools] ${v.holder} claims ${v.seat} and implements no ${v.verb}, which that seat promises; not served`);
continue; continue;
} }
stops.push(await broker.handleSubject(v.subject, run)); stops.push(await broker.handleSubject(v.subject, run));
if (tool) servingNow.push({ ...v, tool });
console.log(`[mesh-tools] serving ${v.seat}'s ${v.verb} on ${v.subject}, admitted where ${v.holder} holds the seat`); console.log(`[mesh-tools] serving ${v.seat}'s ${v.verb} on ${v.subject}, admitted where ${v.holder} holds the seat`);
} }
}; };
await serve(); await serve();
// A membership issued to any served module may add, move or withdraw a seat's verbs. // A membership issued to any served module may add, move or withdraw a seat's verbs.
if (typeof broker.onMembership === "function") broker.onMembership(() => void serve()); if (typeof broker.onMembership === "function") broker.onMembership(() => void serve());
return () => stops.forEach((s) => s()); return { stop: () => stops.forEach((s) => s()), serving: () => servingNow };
}
let sdkHooked = false;
const SDK = "@novox/mesh-sdk";
/**
* The one SDK in a node's runtime (novox/hq ADR 0175, 04-ISSUES/209).
*
* A bundle carries its own dependencies — the toolchain copies them in so a bundle starts anywhere
* (to-be 38 WP3) — and among them is a copy of the SDK. Imported in this process, that copy would be
* a second SDK: its own registry of tools, its own broker handle. A bundle calling registerModuleTools
* through it registers into a list this runtime never reads, and its tools are silently not served.
* So every import of the SDK, from whichever bundle, is resolved as if this runtime had written it:
* one registry, one broker — the runtime's. Everything else a bundle carries resolves from the
* bundle's own tree, as before. A launched bundle (ADR 0188) is another process and is untouched.
*
* Installed once, in-thread, before the first bundle is imported; the hook sees every import after,
* `require` included. A bundle whose own copy is another version than the runtime's is said once,
* so a tool failing against the runtime's SDK points at the bundle rather than at the runtime.
*/
function oneSdk(): void {
if (sdkHooked) return;
registerHooks({
resolve(specifier, context, next) {
if (specifier === SDK || specifier.startsWith(SDK + "/")) {
if (context.parentURL) sayOtherSdk(context.parentURL);
return next(specifier, { ...context, parentURL: import.meta.url });
}
return next(specifier, context);
},
});
sdkHooked = true;
}
const sdkSaid = new Set<string>();
/** The version of the SDK copy nearest a file, by its package.json, or nothing when the file has none above it. */
function sdkVersionNear(fileURL: string): { dir: string; version: string } | undefined {
let dir = dirname(fileURLToPath(fileURL));
for (;;) {
const pkg = join(dir, "node_modules", SDK, "package.json");
if (existsSync(pkg)) {
try {
return { dir, version: String((JSON.parse(readFileSync(pkg, "utf8")) as { version?: string }).version ?? "?") };
} catch {
return { dir, version: "?" };
}
}
const up = dirname(dir);
if (up === dir) return undefined;
dir = up;
}
}
function sayOtherSdk(parentURL: string): void {
if (!parentURL.startsWith("file:")) return;
const own = sdkVersionNear(import.meta.url);
const theirs = sdkVersionNear(parentURL);
if (!theirs || !own || theirs.dir === own.dir || sdkSaid.has(theirs.dir)) return;
sdkSaid.add(theirs.dir);
if (theirs.version !== own.version) {
console.log(`[mesh-tools] ${theirs.dir} carries ${SDK} ${theirs.version}; this runtime's is ${own.version}, and the bundle speaks to the runtime's`);
}
} }
+69
View File
@@ -0,0 +1,69 @@
/**
* What a runtime serves, it announces (novox/hq ADR 0197): the TypeScript runtime the per-module
* containers still run answers the NATS services protocol's discovery in the same shape as the Go
* tool runtime — its module's tools on every subject issued, and the seat verbs it serves.
*
* MESH_TEST_NATS=nats://127.0.0.1:14232 node --test --experimental-strip-types test/announce.test.ts
*/
import assert from "node:assert/strict";
import { test } from "node:test";
import { fileURLToPath } from "node:url";
import { connect, StringCodec } from "nats";
import { resetTools } from "@novox/mesh-sdk/tools";
import { connectNats, membershipSubject } from "../dist/broker-nats.js";
import { runTools } from "../dist/runtime.js";
const url = process.env.MESH_TEST_NATS;
const fixture = (name: string) => fileURLToPath(new URL(`./fixtures/${name}`, import.meta.url));
const sc = StringCodec();
test("the runtime answers $SRV.INFO with what it serves, in the services protocol's format", async (t) => {
if (!url) return t.skip("MESH_TEST_NATS unset");
resetTools();
const nc = await connect({ servers: url });
const jsm = await nc.jetstreamManager();
try {
await jsm.streams.delete("ASSIGNMENTS");
} catch {
// none yet
}
await jsm.streams.add({ name: "ASSIGNMENTS", subjects: ["mesh.assignment.>"], max_msgs_per_subject: 1, allow_direct: true } as never);
await nc.jetstream().publish(membershipSubject("anchor", "shop"), sc.encode(JSON.stringify({
node: "anchor", module: "shop",
serves: [{ subject: "mesh.mod.shop.tool.{tool}.anchor" }, { subject: "mesh.mod.shop.tool.{tool}", queue: "serve.shop" }],
emits: "mesh.mod.shop.event.{event}", tools: "mesh.mod.shop.tool.tools",
})));
const shop = await connectNats({ url, node: "anchor", module: "shop" });
let stop = () => {};
try {
stop = await runTools({ broker: shop, moduleEntrypoints: [fixture("shop-tools.mjs")] });
const msg = await nc.request("$SRV.INFO.shop.anchor", sc.encode(""), { timeout: 2000 });
const info = JSON.parse(sc.decode(msg.data)) as {
type: string; name: string; id: string; version: string;
endpoints: { name: string; subject: string; queue_group: string; metadata: Record<string, string> }[];
};
assert.equal(info.type, "io.nats.micro.v1.info_response");
assert.equal(info.name, "shop");
assert.equal(info.id, "anchor");
assert.ok(info.version);
const price = info.endpoints.filter((e) => e.metadata.tool === "price").map((e) => `${e.subject}|${e.queue_group}`).sort();
assert.deepEqual(price, ["mesh.mod.shop.tool.price.anchor|", "mesh.mod.shop.tool.price|serve.shop"]);
const one = info.endpoints.find((e) => e.metadata.tool === "price")!;
assert.equal(one.metadata.kind, "tool");
assert.equal(one.metadata.module, "shop");
assert.equal(one.metadata.node, "anchor");
assert.equal(one.metadata.interchangeable, "true");
assert.ok(JSON.parse(one.metadata.schema).type === "object");
// Ping answers with the same identity; another service's request is not answered.
const ping = JSON.parse(sc.decode((await nc.request("$SRV.PING", sc.encode(""), { timeout: 2000 })).data));
assert.equal(ping.type, "io.nats.micro.v1.ping_response");
await assert.rejects(nc.request("$SRV.INFO.somebody-else", sc.encode(""), { timeout: 300 }));
} finally {
stop();
await shop.close();
await nc.close();
resetTools();
}
});
+7
View File
@@ -0,0 +1,7 @@
#!/usr/bin/env node
// A long-running bundle that keeps exiting (novox/hq ADR 0198): the runtime starts it again each time.
import { appendFileSync } from "node:fs";
import { serveStdio } from "@novox/mesh-sdk/stdio";
appendFileSync(process.env.FLAKY_LOG, "started\n");
setTimeout(() => process.exit(3), 300);
await serveStdio("flaky", []);
+10
View File
@@ -0,0 +1,10 @@
# A bus that refuses one subject: what the mesh's grants do to a subject they do not name (issue 217).
authorization {
users = [
{ user: "runtime", password: "runtime", permissions: {
publish: { allow: [">"] }
subscribe: { allow: [">"], deny: ["$SRV.PING.>", "mesh.seat.held-elsewhere.>"] }
} }
]
}
jetstream: enabled
+13
View File
@@ -0,0 +1,13 @@
// A copy for the issue 218 test: a fixture registers once per process, at import.
// A module that is both software and a role: postgres's own tools under its name, and its
// implementation of the mesh-store seat's verbs under the seat's (novox/hq ADR 0159, 0160).
import { registerModuleTools } from "@novox/mesh-sdk/tools";
registerModuleTools("postgres", () => [
{ name: "postgres_create_database", description: "make one", input: {}, run: async () => ({ made: true }) },
{ name: "databases", description: "postgres's own listing", input: {}, run: async () => ({ software: "postgres" }) },
]);
registerModuleTools("mesh-store", () => [
{ name: "databases", description: "what the store holds", input: {}, run: async () => ({ seat: "mesh-store" }) },
]);
+24
View File
@@ -0,0 +1,24 @@
// A module whose code runs long (novox/hq ADR 0198): it subscribes as it is imported, handles each
// event once it can, fails the first time it is asked to, dies the first time it is told to, and one
// tool asks another module's tool through the runtime. What it did is written to WATCH_LOG.
import { appendFileSync, existsSync, writeFileSync } from "node:fs";
import { on } from "@novox/mesh-sdk/events";
import { broker } from "@novox/mesh-sdk/messaging";
import { registerModuleTools } from "@novox/mesh-sdk/tools";
const log = process.env.WATCH_LOG;
const once = (mark) => {
const file = `${log}.${mark}`;
if (existsSync(file)) return false;
writeFileSync(file, "1");
return true;
};
appendFileSync(log, "started\n");
await on("alpha.happened", async (e) => {
if (e.body.fail && once(`fail-${e.body.n}`)) throw new Error(`not yet ${e.body.n}`);
if (e.body.die && once(`die-${e.body.n}`)) process.exit(7);
appendFileSync(log, `handled ${e.body.n}\n`);
});
registerModuleTools("watcher", () => [
{ name: "relay", description: "asks beta", input: {}, run: async () => broker().request("beta.three", {}) },
]);
+5
View File
@@ -0,0 +1,5 @@
#!/usr/bin/env node
// The launcher the builder writes beside an entrypoint (novox/hq ADR 0193), for watcher.mjs.
import { serveRegisteredOverStdio } from "@novox/mesh-sdk/stdio";
await import("./watcher.mjs");
await serveRegisteredOverStdio();
+30
View File
@@ -140,6 +140,36 @@ test("a seat's verbs are implemented under the seat's name, served where issued,
} }
}); });
// novox/hq issue 218: the store seat is claimed by every machine running postgres and held by one.
// Where the mesh issued a membership without the seat, the claim in the credential serves nothing:
// the module's own tools answer, the seat's verbs do not, and the runtime does not announce them.
test("a claimant the membership does not make the holder serves none of the seat's verbs", async (t) => {
if (!url) return t.skip("MESH_TEST_NATS unset");
resetTools();
const stream = await anAssignmentsStream();
await stream.issue({
node: "elsewhere",
module: "postgres",
serves: [{ subject: "mesh.mod.postgres.tool.{tool}.elsewhere" }],
emits: "mesh.mod.postgres.event.{event}",
tools: "mesh.mod.postgres.tool.tools",
});
const credential = { url, node: "elsewhere", module: "postgres", claims: [{ seat: "mesh-store", scope: "mesh", serves: ["databases", "query"] }] };
const pg = await connectNats(credential);
const asker = await connectNats({ url, module: "console", node: "workstation" });
let stop = () => {};
try {
stop = await runTools({ broker: pg, credential, moduleEntrypoints: [fixture("store-claimant.mjs")] });
assert.deepEqual((await callTool(asker, "postgres.databases@elsewhere", {})).result, { software: "postgres" });
await assert.rejects(callTool(asker, "seat:mesh-store.databases", {}), /no responders|503/i, "the store is not held here");
} finally {
stop();
await asker.close();
await pg.close();
await stream.close();
}
});
test("a module named like its seat registers once, and answers as the module and as the seat", async (t) => { test("a module named like its seat registers once, and answers as the module and as the seat", async (t) => {
if (!url) return t.skip("MESH_TEST_NATS unset"); if (!url) return t.skip("MESH_TEST_NATS unset");
resetTools(); resetTools();
+41 -9
View File
@@ -122,9 +122,9 @@ test("the node's runtime serves five modules' bundles on one credential — two
broker: nodeTools, broker: nodeTools,
credential, credential,
serves: [ serves: [
{ module: "alpha", entrypoints: [fixture("many-alpha.mjs")] }, { module: "alpha", entrypoints: [fixture("many-alpha.serve.mjs")] },
{ module: "beta", entrypoints: [fixture("many-beta.mjs")] }, { module: "beta", entrypoints: [fixture("many-beta.serve.mjs")] },
{ module: "gamma", entrypoints: [fixture("many-broken.mjs")] }, { module: "gamma", entrypoints: [fixture("many-broken.serve.mjs")] },
{ module: "delta", entrypoints: [fixture("many-delta.py")] }, { module: "delta", entrypoints: [fixture("many-delta.py")] },
{ module: "epsilon", entrypoints: [fixture("many-epsilon.mjs")] }, { module: "epsilon", entrypoints: [fixture("many-epsilon.mjs")] },
], ],
@@ -132,7 +132,7 @@ test("the node's runtime serves five modules' bundles on one credential — two
console.log = log; console.log = log;
assert.deepEqual(nodeTools.serving().sort(), ["alpha", "beta", "delta", "epsilon", "gamma", "node-tools"]); assert.deepEqual(nodeTools.serving().sort(), ["alpha", "beta", "delta", "epsilon", "gamma", "node-tools"]);
assert.ok(said.some((s) => /the operator's account here is somebody \(home \/home\/somebody\)/.test(s)), said.join("\n")); assert.ok(said.some((s) => /the operator's account here is somebody \(home \/home\/somebody\)/.test(s)), said.join("\n"));
assert.ok(said.some((s) => /gamma's bundle .*many-broken\.mjs failed to load: gamma's bundle cannot find its client; its tools are not served here/.test(s)), said.join("\n")); assert.ok(said.some((s) => /gamma's bundle .*many-broken\.serve\.mjs failed to load: gamma's bundle exited \(1\): Error: gamma's bundle cannot find its client; its tools are not served here/.test(s)), said.join("\n"));
assert.ok(said.some((s) => /serving 8 tool\(s\) for 5 module\(s\): alpha\.one, alpha\.two, beta\.three, beta\.four, beta\.five, delta\.greet, delta\.die, epsilon\.seven; not serving gamma/.test(s)), said.join("\n")); assert.ok(said.some((s) => /serving 8 tool\(s\) for 5 module\(s\): alpha\.one, alpha\.two, beta\.three, beta\.four, beta\.five, delta\.greet, delta\.die, epsilon\.seven; not serving gamma/.test(s)), said.join("\n"));
// Five tools answer, each where its module's membership says: alpha anywhere and here, beta here only. // Five tools answer, each where its module's membership says: alpha anywhere and here, beta here only.
@@ -163,7 +163,7 @@ test("the node's runtime serves five modules' bundles on one credential — two
// `tools` answers for each: what alpha and beta serve, and why gamma serves nothing. // `tools` answers for each: what alpha and beta serve, and why gamma serves nothing.
const gamma = await asker.request<Record<string, never>, { module: string; tools: unknown[]; failed?: string }>("gamma.tools@anchor", {}); const gamma = await asker.request<Record<string, never>, { module: string; tools: unknown[]; failed?: string }>("gamma.tools@anchor", {});
assert.deepEqual(gamma, { module: "gamma", tools: [], failed: "gamma's bundle cannot find its client" }); assert.deepEqual(gamma, { module: "gamma", tools: [], failed: "gamma's bundle exited (1): Error: gamma's bundle cannot find its client" });
const beta = await asker.request<Record<string, never>, { tools: { name: string; subjects?: string[] }[] }>("beta.tools@anchor", {}); const beta = await asker.request<Record<string, never>, { tools: { name: string; subjects?: string[] }[] }>("beta.tools@anchor", {});
assert.deepEqual(beta.tools.map((x) => x.name), ["three", "four", "five"]); assert.deepEqual(beta.tools.map((x) => x.name), ["three", "four", "five"]);
assert.deepEqual(beta.tools[0]!.subjects, ["mesh.mod.beta.tool.three.anchor"]); assert.deepEqual(beta.tools[0]!.subjects, ["mesh.mod.beta.tool.three.anchor"]);
@@ -174,7 +174,7 @@ test("the node's runtime serves five modules' bundles on one credential — two
try { try {
const have = await toolsOn(asker); const have = await toolsOn(asker);
assert.deepEqual(have.tools.map((x) => `${x.module}.${x.name}`), ["alpha.one", "alpha.two", "beta.five", "beta.four", "beta.three", "delta.die", "delta.greet", "epsilon.seven"]); assert.deepEqual(have.tools.map((x) => `${x.module}.${x.name}`), ["alpha.one", "alpha.two", "beta.five", "beta.four", "beta.three", "delta.die", "delta.greet", "epsilon.seven"]);
assert.deepEqual(have.notAnswering, ["gamma (its tools bundle failed to load: gamma's bundle cannot find its client)", "mesh-controller (seat)"]); assert.deepEqual(have.notAnswering, ["gamma (its tools bundle failed to load: gamma's bundle exited (1): Error: gamma's bundle cannot find its client)", "mesh-controller (seat)"]);
} finally { } finally {
await catalogue.close(); await catalogue.close();
} }
@@ -215,6 +215,8 @@ test("a bundle carrying its own copy of the SDK registers into the runtime's reg
writeFileSync(join(dir, "node_modules", "zeta-flavour", "package.json"), '{"name":"zeta-flavour","type":"module","main":"index.js"}\n'); writeFileSync(join(dir, "node_modules", "zeta-flavour", "package.json"), '{"name":"zeta-flavour","type":"module","main":"index.js"}\n');
writeFileSync(join(dir, "node_modules", "zeta-flavour", "index.js"), 'export const flavour = "the bundle\'s own";\n'); writeFileSync(join(dir, "node_modules", "zeta-flavour", "index.js"), 'export const flavour = "the bundle\'s own";\n');
writeFileSync(join(dir, "package.json"), '{"type":"module","private":true}\n'); writeFileSync(join(dir, "package.json"), '{"type":"module","private":true}\n');
writeFileSync(join(dir, "index.serve.mjs"),
'#!/usr/bin/env node\nimport { serveRegisteredOverStdio } from "@novox/mesh-sdk/stdio";\nawait import("./index.js");\nawait serveRegisteredOverStdio();\n', { mode: 0o755 });
writeFileSync(join(dir, "index.js"), writeFileSync(join(dir, "index.js"),
'import { registerModuleTools } from "@novox/mesh-sdk/tools";\n' + 'import { registerModuleTools } from "@novox/mesh-sdk/tools";\n' +
'import { flavour } from "zeta-flavour";\n' + 'import { flavour } from "zeta-flavour";\n' +
@@ -228,7 +230,7 @@ test("a bundle carrying its own copy of the SDK registers into the runtime's reg
const asker = await connectNats({ url, module: "console", node: "workstation" }); const asker = await connectNats({ url, module: "console", node: "workstation" });
closing.push(() => asker.close()); closing.push(() => asker.close());
console.log = (...a: unknown[]) => said.push(a.join(" ")); console.log = (...a: unknown[]) => said.push(a.join(" "));
stop = await runTools({ broker: nodeTools, credential, serves: [{ module: "zeta", entrypoints: [join(dir, "index.js")] }] }); stop = await runTools({ broker: nodeTools, credential, serves: [{ module: "zeta", entrypoints: [join(dir, "index.serve.mjs")] }] });
console.log = log; console.log = log;
assert.ok(said.some((s) => /serving 1 tool\(s\) for 1 module\(s\): zeta\.probe/.test(s)), said.join("\n")); assert.ok(said.some((s) => /serving 1 tool\(s\) for 1 module\(s\): zeta\.probe/.test(s)), said.join("\n"));
// The SDK is the runtime's (the registration arrived); the bundle's other dependency is its own. // The SDK is the runtime's (the registration arrived); the bundle's other dependency is its own.
@@ -266,8 +268,8 @@ test("each bundle is given its own environment and none of another's, imported o
stop = await runTools({ stop = await runTools({
broker: nodeTools, credential, envs, broker: nodeTools, credential, envs,
serves: [ serves: [
{ module: "gamma", entrypoints: [fixture("env-gamma.mjs")] }, { module: "gamma", entrypoints: [fixture("env-gamma.serve.mjs")] },
{ module: "delta", entrypoints: [fixture("env-delta.mjs")] }, { module: "delta", entrypoints: [fixture("env-delta.serve.mjs")] },
{ module: "zeta", entrypoints: [fixture("env-zeta.mjs")] }, { module: "zeta", entrypoints: [fixture("env-zeta.mjs")] },
], ],
}); });
@@ -316,3 +318,33 @@ test("a launched bundle is told the module it serves, so its seat's verbs stay t
resetTools(); resetTools();
} }
}); });
test("an entrypoint that is not executable is refused by name, and the others serve (ADR 0193)", async (t) => {
if (!url) return t.skip("MESH_TEST_NATS unset");
resetTools();
const mesh = await aMesh();
for (const m of ["alpha", "plain"]) await mesh.issue(membershipOf(m, "anchor"));
const credential = { url, node: "anchor", module: "node-tools" };
const nodeTools = await connectNats(credential);
const asker = await connectNats({ url, module: "console", node: "workstation" });
const said: string[] = [];
const log = console.log;
let stop = () => {};
try {
console.log = (...a: unknown[]) => said.push(a.join(" "));
stop = await runTools({ broker: nodeTools, credential, serves: [
{ module: "alpha", entrypoints: [fixture("many-alpha.serve.mjs")] },
{ module: "plain", entrypoints: [fixture("many-alpha.mjs")] },
] });
console.log = log;
assert.ok(said.some((s) => /plain's bundle .*many-alpha\.mjs failed to load: .* is not executable; a bundle the runtime serves is started, never imported/.test(s)), said.join("\n"));
assert.deepEqual((await callTool(asker, "alpha.one@anchor", {})).result, { alpha: 1 });
} finally {
console.log = log;
stop();
await asker.close();
await nodeTools.close();
await mesh.close();
resetTools();
}
});
+1 -1
View File
@@ -36,7 +36,7 @@ test("as node-tools, serve loads the bundles and is the console on loopback", as
env: { env: {
...process.env, ...process.env,
MESH_BROKER_FILE: credential, MESH_BROKER_FILE: credential,
MESH_TOOL_MODULES: `alpha=${fixture("many-alpha.mjs")}`, MESH_TOOL_MODULES: `alpha=${fixture("many-alpha.serve.mjs")}`,
MESH_CONSOLE_LISTEN: "127.0.0.1:0", MESH_CONSOLE_LISTEN: "127.0.0.1:0",
}, },
stdio: ["ignore", "pipe", "pipe"], stdio: ["ignore", "pipe", "pipe"],
+68
View File
@@ -0,0 +1,68 @@
/**
* A refused announcement is said and never fatal (novox/hq 04-ISSUES/217): against a bus whose
* permissions refuse one discovery subject, the runtime's raw subscription is refused, logged, and
* the process keeps serving — its tools still answer.
*
* docker run -d --rm --name t -p 14233:4222 -v $PWD/test/fixtures/refusing-nats.conf:/c.conf nats:2.10-alpine -c /c.conf
* MESH_TEST_REFUSING_NATS=nats://127.0.0.1:14233 node --test --experimental-strip-types test/refused.test.ts
*/
import assert from "node:assert/strict";
import { test } from "node:test";
import { connectNats } from "../dist/broker-nats.js";
import { callTool } from "../dist/client.js";
const url = process.env.MESH_TEST_REFUSING_NATS;
test("a refused discovery subscription is logged and the runtime serves on", async (t) => {
if (!url) return t.skip("MESH_TEST_REFUSING_NATS unset");
const bus = await connectNats({ url, user: "runtime", password: "runtime", module: "alpha", node: "anchor" });
const said: string[] = [];
const log = console.log;
console.log = (...a: unknown[]) => said.push(a.join(" "));
const crashed: unknown[] = [];
const onRejection = (e: unknown) => crashed.push(e);
process.on("unhandledRejection", onRejection);
try {
(bus as unknown as { raw: (s: string, f: () => Uint8Array | undefined) => () => void }).raw("$SRV.PING.>", () => undefined);
const stop = await bus.handle("alpha.ping", async () => ({ pong: true }));
for (let i = 0; i < 50 && !said.some((s) => s.includes("the bus refused $SRV.PING.>")); i++) await new Promise((r) => setTimeout(r, 50));
console.log = log;
assert.ok(said.some((s) => /the bus refused \$SRV\.PING\.>.*serves on/.test(s)), said.join("\n"));
assert.equal(crashed.length, 0, `the refusal escaped: ${String(crashed[0])}`);
stop();
} finally {
console.log = log;
process.off("unhandledRejection", onRejection);
await bus.close();
}
});
// novox/hq issue 218: a tool subject the grants leave out — a seat claimed here and held elsewhere —
// is refused, said, and the module's other tools still answer.
test("a refused tool subscription is logged and the module's other tools answer", async (t) => {
if (!url) return t.skip("MESH_TEST_REFUSING_NATS unset");
const bus = await connectNats({ url, user: "runtime", password: "runtime", module: "alpha", node: "anchor" });
const asker = await connectNats({ url, user: "runtime", password: "runtime", module: "console", node: "workstation" });
const said: string[] = [];
const log = console.log;
console.log = (...a: unknown[]) => said.push(a.join(" "));
const crashed: unknown[] = [];
const onRejection = (e: unknown) => crashed.push(e);
process.on("unhandledRejection", onRejection);
try {
const refused = await bus.handleSubject!("mesh.seat.held-elsewhere.tool.databases", async () => ({ seat: true }));
const stop = await bus.handle("alpha.ping", async () => ({ pong: true }));
for (let i = 0; i < 50 && !said.some((s) => s.includes("the bus refused mesh.seat.held-elsewhere")); i++) await new Promise((r) => setTimeout(r, 50));
console.log = log;
assert.ok(said.some((s) => /the bus refused mesh\.seat\.held-elsewhere\.tool\.databases.*serves on/.test(s)), said.join("\n"));
assert.equal(crashed.length, 0, `the refusal escaped: ${String(crashed[0])}`);
assert.deepEqual((await callTool(asker, "alpha.ping@anchor", {})).result, { pong: true });
stop();
refused();
} finally {
console.log = log;
process.off("unhandledRejection", onRejection);
await asker.close();
await bus.close();
}
});