Compare commits
25
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
2ba7451229 | ||
|
|
6604d44372 | ||
|
|
7b21440962 | ||
|
|
0cea8d286e | ||
|
|
094a7d0d7c | ||
|
|
b52669577d | ||
|
|
7bd76275f9 | ||
|
|
6ba0f4dc1e | ||
|
|
bc05658772 | ||
|
|
4af59636c3 | ||
|
|
6e425a000f | ||
|
|
14b6588839 | ||
|
|
490aedfd36 | ||
|
|
f11ac6441c | ||
|
|
ffe229308c | ||
|
|
df4f492a72 | ||
|
|
66e8be0e31 | ||
|
|
7722668220 | ||
|
|
e8989f3cf5 | ||
|
|
915f372a85 | ||
|
|
486dad99a5 | ||
|
|
65f3b68076 | ||
|
|
42e0987c64 | ||
|
|
ae6bdc9b06 | ||
|
|
b182943c24 |
+11
@@ -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),
|
||||||
|
|||||||
@@ -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"
|
||||||
}
|
}
|
||||||
]
|
]
|
||||||
},
|
},
|
||||||
|
|||||||
@@ -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, µ.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
|
||||||
|
}
|
||||||
@@ -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"
|
||||||
)
|
)
|
||||||
@@ -512,6 +514,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 +621,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
|
||||||
|
}
|
||||||
|
|||||||
@@ -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) }
|
||||||
|
|
||||||
|
|||||||
@@ -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)
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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.
|
||||||
|
|||||||
@@ -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) {
|
||||||
|
|||||||
@@ -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 {
|
||||||
|
|||||||
@@ -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)
|
||||||
|
|||||||
@@ -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{}
|
||||||
|
}
|
||||||
@@ -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...) }
|
||||||
@@ -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
|
||||||
|
}
|
||||||
|
|||||||
@@ -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"
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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());
|
||||||
|
}
|
||||||
@@ -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),
|
||||||
|
|||||||
@@ -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
@@ -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`);
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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
@@ -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
@@ -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
@@ -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" }) },
|
||||||
|
]);
|
||||||
Vendored
+24
@@ -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
@@ -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();
|
||||||
@@ -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();
|
||||||
|
|||||||
@@ -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();
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|||||||
@@ -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"],
|
||||||
|
|||||||
@@ -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();
|
||||||
|
}
|
||||||
|
});
|
||||||
Reference in New Issue
Block a user