The runtime knows no language, so nothing ties it to Node.js. This ports its serve mode — the pinned bus connection and patient connect, following memberships, launching every served bundle over MCP on stdio with its own environment, a child's emit published as its module, each tool, the tools verb and seat verbs served where the mesh issued them, and the console on loopback — to one static binary. Same subjects, request and reply bodies, event headers and MCP answers. The TypeScript stays: it is still the runtime inside the per-module containers until WP4c. Tests run against a real bus and share the TypeScript fixtures.
569 lines
17 KiB
Go
569 lines
17 KiB
Go
// Package bus is the node runtime's connection to the mesh bus, on NATS (novox/hq design 25, design
|
|
// 29, ADR 0160, ADR 0175). It is the Go port of node-tools' broker-nats.ts, and speaks the same
|
|
// wire: the same subjects, the same JSON request and reply bodies, the same event headers.
|
|
//
|
|
// mesh.mod.<module>.event.<type> an event a module emits
|
|
// mesh.mod.<module>.tool.<tool> a tool a module serves
|
|
// mesh.seat.<seat>.tool.<verb> a role's verb, answered by whoever holds the seat
|
|
package bus
|
|
|
|
import (
|
|
"crypto/sha256"
|
|
"crypto/tls"
|
|
"crypto/x509"
|
|
"encoding/hex"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"log"
|
|
"regexp"
|
|
"strings"
|
|
"sync"
|
|
"time"
|
|
|
|
"github.com/nats-io/nats.go"
|
|
|
|
"github.com/novox/mesh-tools/node-tools/internal/wire"
|
|
)
|
|
|
|
// RequestTimeout is how long a call waits for its answer — what modules already expect.
|
|
const RequestTimeout = 30 * time.Second
|
|
|
|
const assignmentsStream = "ASSIGNMENTS"
|
|
|
|
// Credential is a broker credential as the mesh delivers it (novox/hq ADR 0120).
|
|
type Credential struct {
|
|
URL string `json:"url"`
|
|
Fingerprint string `json:"fingerprint,omitempty"`
|
|
Node string `json:"node,omitempty"`
|
|
Module string `json:"module,omitempty"`
|
|
User string `json:"user,omitempty"`
|
|
Password string `json:"password,omitempty"`
|
|
Claims []Claim `json:"claims,omitempty"`
|
|
}
|
|
|
|
// Claim is a seat a module claims, with the verbs it promises (novox/hq ADR 0159).
|
|
type Claim struct {
|
|
Seat string `json:"seat"`
|
|
Scope string `json:"scope,omitempty"`
|
|
Serves []string `json:"serves,omitempty"`
|
|
}
|
|
|
|
// Membership is what the mesh issued one assignment (novox/hq ADR 0160).
|
|
type Membership struct {
|
|
Node string `json:"node"`
|
|
Module string `json:"module"`
|
|
Serves []Served `json:"serves"`
|
|
Seats []SeatVerb `json:"seats,omitempty"`
|
|
Emits string `json:"emits"`
|
|
Reaches map[string][]string `json:"reaches,omitempty"`
|
|
Tools string `json:"tools"`
|
|
}
|
|
|
|
// Served is an address a tool is answered on; `{tool}` stands for the tool's name.
|
|
type Served struct {
|
|
Subject string `json:"subject"`
|
|
Queue string `json:"queue,omitempty"`
|
|
}
|
|
|
|
// SeatVerb is one verb of a seat a module holds, where it is answered.
|
|
type SeatVerb struct {
|
|
Seat string `json:"seat"`
|
|
Verb string `json:"verb"`
|
|
Subject string `json:"subject"`
|
|
}
|
|
|
|
// Envelope is an event as a module emits it: the body is the payload, the metadata rides as headers
|
|
// (novox/hq ADR 0042).
|
|
type Envelope struct {
|
|
Key string `json:"key"`
|
|
Node string `json:"node,omitempty"`
|
|
Body json.RawMessage `json:"body"`
|
|
Headers map[string]string `json:"headers,omitempty"`
|
|
}
|
|
|
|
// Answered is a call's result and the machine that gave it (novox/hq ADR 0159).
|
|
type Answered struct {
|
|
Result json.RawMessage
|
|
Node string
|
|
}
|
|
|
|
// Handler answers one request body.
|
|
type Handler func(body json.RawMessage) (any, error)
|
|
|
|
// ErrPin is a bus whose certificate is not the one the mesh pinned: final, never retried.
|
|
var ErrPin = errors.New("the bus's certificate does not match the pin")
|
|
|
|
// Fatal is why a connection failure is final rather than "not yet", or "" when waiting may fix it —
|
|
// the same classification the TypeScript runtime makes.
|
|
func Fatal(err error) string {
|
|
if err == nil {
|
|
return ""
|
|
}
|
|
msg := err.Error()
|
|
if errors.Is(err, ErrPin) || strings.Contains(msg, ErrPin.Error()) {
|
|
return "the bus's certificate does not match the pin"
|
|
}
|
|
if regexp.MustCompile(`(?i)invalid url|no servers available for connection: .*url`).MatchString(msg) ||
|
|
strings.Contains(msg, "nats: invalid url") {
|
|
return "the bus address is not a usable URL"
|
|
}
|
|
if regexp.MustCompile(`(?i)authorization violation|user authentication expired|permissions violation`).MatchString(msg) {
|
|
return "the bus refused this account"
|
|
}
|
|
return ""
|
|
}
|
|
|
|
func normalizeFingerprint(f string) string {
|
|
f = strings.TrimSpace(f)
|
|
if len(f) > 7 && strings.EqualFold(f[:7], "sha256:") {
|
|
f = f[7:]
|
|
}
|
|
return strings.ToLower(strings.ReplaceAll(f, ":", ""))
|
|
}
|
|
|
|
// pinned accepts exactly the certificate with this SHA-256 and no other. The pin is the only check:
|
|
// the bus's certificate names the seat, not the address a machine dials it by.
|
|
func pinned(want string) *tls.Config {
|
|
want = normalizeFingerprint(want)
|
|
return &tls.Config{
|
|
InsecureSkipVerify: true, //nolint:gosec // replaced by the pin, which is stricter
|
|
MinVersion: tls.VersionTLS12,
|
|
VerifyPeerCertificate: func(raw [][]byte, _ [][]*x509.Certificate) error {
|
|
if len(raw) == 0 {
|
|
return fmt.Errorf("%w: it presented none", ErrPin)
|
|
}
|
|
sum := sha256.Sum256(raw[0])
|
|
if got := hex.EncodeToString(sum[:]); got != want {
|
|
return fmt.Errorf("%w: it presented %s, not the pinned %s", ErrPin, got, want)
|
|
}
|
|
return nil
|
|
},
|
|
}
|
|
}
|
|
|
|
// Conn is the runtime's connection: the module it is, the memberships it follows, the subjects it
|
|
// answers.
|
|
type Conn struct {
|
|
nc *nats.Conn
|
|
js nats.JetStreamContext
|
|
self string
|
|
node string
|
|
cred Credential
|
|
mu sync.Mutex
|
|
issued map[string]*Membership // module → membership; present with nil = followed, none issued
|
|
onNew []func(Membership)
|
|
subs []*nats.Subscription
|
|
Logf func(format string, args ...any)
|
|
}
|
|
|
|
// Connect dials the bus as the credential's module. A module's subjects come from its credential,
|
|
// never from its calls (ADR 0074).
|
|
func Connect(cred Credential) (*Conn, error) {
|
|
if cred.Module == "" {
|
|
return nil, errors.New("a broker credential with no module: the runtime derives its subjects " +
|
|
"from the account the mesh issued, and cannot guess which module it is")
|
|
}
|
|
node := cred.Node
|
|
if node == "" {
|
|
node = "?"
|
|
}
|
|
opts := []nats.Option{
|
|
nats.Name(node + "." + cred.Module),
|
|
// Reconnect forever: the bus restarting is an upgrade, not a reason to exit.
|
|
nats.MaxReconnects(-1),
|
|
}
|
|
if cred.User != "" {
|
|
opts = append(opts, nats.UserInfo(cred.User, cred.Password),
|
|
// Its own inbox: every user's inbox is private to it (design 25 §4).
|
|
nats.CustomInboxPrefix("_INBOX."+cred.User))
|
|
}
|
|
if strings.TrimSpace(cred.Fingerprint) != "" {
|
|
opts = append(opts, nats.Secure(pinned(cred.Fingerprint)))
|
|
}
|
|
nc, err := nats.Connect(cred.URL, opts...)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
js, err := nc.JetStream()
|
|
if err != nil {
|
|
nc.Close()
|
|
return nil, err
|
|
}
|
|
c := &Conn{nc: nc, js: js, self: cred.Module, node: cred.Node, cred: cred,
|
|
issued: map[string]*Membership{}, Logf: log.Printf}
|
|
c.Follow(cred.Module)
|
|
return c, nil
|
|
}
|
|
|
|
// Module is what this connection is.
|
|
func (c *Conn) Module() string { return c.self }
|
|
|
|
// Node is the machine this connection's account is scoped to.
|
|
func (c *Conn) Node() string { return c.node }
|
|
|
|
// Credential is what this connection was opened with.
|
|
func (c *Conn) Credential() Credential { return c.cred }
|
|
|
|
// MembershipSubject is the one address a runtime derives for an assignment (ADR 0160).
|
|
func MembershipSubject(node, module string) string {
|
|
return "mesh.assignment." + node + "." + module
|
|
}
|
|
|
|
// Follow reads a module's membership on this machine once and follows it live, so its tools are
|
|
// served where the mesh issued them (ADR 0175).
|
|
func (c *Conn) Follow(module string) {
|
|
c.mu.Lock()
|
|
if c.node == "" {
|
|
c.mu.Unlock()
|
|
return
|
|
}
|
|
if _, has := c.issued[module]; has {
|
|
c.mu.Unlock()
|
|
return
|
|
}
|
|
c.issued[module] = nil
|
|
c.mu.Unlock()
|
|
subject := MembershipSubject(c.node, module)
|
|
// The subject-addressed direct get: the one address the mesh grants this account on the
|
|
// stream's API.
|
|
if got, err := c.nc.Request("$JS.API.DIRECT.GET."+assignmentsStream+"."+subject, nil, 5*time.Second); err == nil {
|
|
if got.Header.Get("Status") == "" && len(got.Data) > 0 {
|
|
var m Membership
|
|
if json.Unmarshal(got.Data, &m) == nil {
|
|
c.mu.Lock()
|
|
c.issued[module] = &m
|
|
c.mu.Unlock()
|
|
}
|
|
}
|
|
}
|
|
if c.Membership(module) == nil {
|
|
c.Logf("[mesh-tools] no membership issued for %s on %s yet; serving the derived shape until one arrives", module, c.node)
|
|
}
|
|
sub, err := c.nc.Subscribe(subject, func(msg *nats.Msg) {
|
|
var m Membership
|
|
if err := json.Unmarshal(msg.Data, &m); err != nil {
|
|
c.Logf("[mesh-tools] a membership arrived that is not one: %v", err)
|
|
return
|
|
}
|
|
c.mu.Lock()
|
|
c.issued[module] = &m
|
|
handlers := append([]func(Membership){}, c.onNew...)
|
|
c.mu.Unlock()
|
|
c.Logf("[mesh-tools] %s on %s was issued a new membership; re-serving on it", module, c.node)
|
|
for _, h := range handlers {
|
|
h(m)
|
|
}
|
|
_ = c.nc.Flush()
|
|
})
|
|
if err == nil {
|
|
c.track(sub)
|
|
}
|
|
}
|
|
|
|
// Membership is what the mesh issued a module here, or nil when nothing has been issued.
|
|
func (c *Conn) Membership(module string) *Membership {
|
|
c.mu.Lock()
|
|
defer c.mu.Unlock()
|
|
return c.issued[module]
|
|
}
|
|
|
|
// Following says whether this connection follows a module's membership.
|
|
func (c *Conn) Following(module string) bool {
|
|
c.mu.Lock()
|
|
defer c.mu.Unlock()
|
|
_, has := c.issued[module]
|
|
return has
|
|
}
|
|
|
|
// Serving is every module whose membership this connection follows.
|
|
func (c *Conn) Serving() []string {
|
|
c.mu.Lock()
|
|
defer c.mu.Unlock()
|
|
out := make([]string, 0, len(c.issued))
|
|
for m := range c.issued {
|
|
out = append(out, m)
|
|
}
|
|
return out
|
|
}
|
|
|
|
// OnMembership is called with every new membership any followed module is issued.
|
|
func (c *Conn) OnMembership(h func(Membership)) {
|
|
c.mu.Lock()
|
|
c.onNew = append(c.onNew, h)
|
|
c.mu.Unlock()
|
|
}
|
|
|
|
func (c *Conn) track(s *nats.Subscription) {
|
|
c.mu.Lock()
|
|
c.subs = append(c.subs, s)
|
|
c.mu.Unlock()
|
|
}
|
|
|
|
// servedOn is where a served module's tool is answered: its membership's subjects when issued, the
|
|
// derived shape otherwise (the shape the mesh issues on day one).
|
|
func (c *Conn) servedOn(module, tool string) []Served {
|
|
if m := c.Membership(module); m != nil {
|
|
out := make([]Served, 0, len(m.Serves)+1)
|
|
for _, s := range m.Serves {
|
|
out = append(out, Served{Subject: strings.ReplaceAll(s.Subject, "{tool}", tool), Queue: s.Queue})
|
|
}
|
|
if tool == "tools" && m.Tools != "" {
|
|
found := false
|
|
for _, s := range out {
|
|
found = found || s.Subject == m.Tools
|
|
}
|
|
if !found {
|
|
out = append([]Served{{Subject: m.Tools, Queue: "serve." + module}}, out...)
|
|
}
|
|
}
|
|
return out
|
|
}
|
|
base := "mesh.mod." + module + ".tool." + tool
|
|
out := []Served{{Subject: base, Queue: "serve." + module}}
|
|
if c.node != "" {
|
|
out = append(out, Served{Subject: base + "." + c.node})
|
|
}
|
|
return out
|
|
}
|
|
|
|
type reply struct {
|
|
Result any `json:"result,omitempty"`
|
|
Error string `json:"error,omitempty"`
|
|
Node string `json:"node,omitempty"`
|
|
}
|
|
|
|
// answerOn answers one subject with one handler, and says which machine answered (ADR 0159).
|
|
func (c *Conn) answerOn(subject, queue string, h Handler) (func(), error) {
|
|
cb := func(msg *nats.Msg) {
|
|
go func() {
|
|
var r reply
|
|
result, err := h(json.RawMessage(msg.Data))
|
|
if err != nil {
|
|
r.Error = err.Error()
|
|
} else {
|
|
r.Result = nullable(result)
|
|
}
|
|
r.Node = c.node
|
|
body, _ := wire.Marshal(r)
|
|
_ = msg.Respond(body)
|
|
}()
|
|
}
|
|
var sub *nats.Subscription
|
|
var err error
|
|
if queue != "" {
|
|
sub, err = c.nc.QueueSubscribe(subject, queue, cb)
|
|
} else {
|
|
sub, err = c.nc.Subscribe(subject, cb)
|
|
}
|
|
if err != nil {
|
|
return func() {}, err
|
|
}
|
|
c.track(sub)
|
|
return func() { _ = sub.Unsubscribe() }, nil
|
|
}
|
|
|
|
// nullable keeps a nil result as JSON null rather than dropping the key: the TypeScript reply always
|
|
// carries `result` when the handler did not throw.
|
|
func nullable(v any) any {
|
|
if v == nil {
|
|
return json.RawMessage("null")
|
|
}
|
|
return v
|
|
}
|
|
|
|
// HandleSubject answers a subject outright: a seat's verb where the mesh issued it.
|
|
func (c *Conn) HandleSubject(subject string, h Handler) (func(), error) {
|
|
return c.answerOn(subject, "", h)
|
|
}
|
|
|
|
// Handle serves `<module>.<tool>` where the mesh issued that module, and follows its membership:
|
|
// when a new one arrives, it serves where it now says and stops where it no longer does.
|
|
func (c *Conn) Handle(key string, h Handler) (func(), error) {
|
|
if strings.HasPrefix(key, "seat:") {
|
|
subject, err := ToolSubject(key, c.self)
|
|
if err != nil {
|
|
return func() {}, err
|
|
}
|
|
return c.answerOn(subject, "", h)
|
|
}
|
|
module, tool := c.self, key
|
|
if dot := strings.Index(key, "."); dot >= 0 {
|
|
module, tool = key[:dot], key[dot+1:]
|
|
}
|
|
if module != c.self && !c.Following(module) {
|
|
return func() {}, fmt.Errorf("%s cannot serve %s: a module serves its own tools, and a runtime "+
|
|
"those of the modules it follows", c.self, key)
|
|
}
|
|
var mu sync.Mutex
|
|
var stops []func()
|
|
serve := func() {
|
|
mu.Lock()
|
|
defer mu.Unlock()
|
|
for _, s := range stops {
|
|
s()
|
|
}
|
|
stops = nil
|
|
for _, s := range c.servedOn(module, tool) {
|
|
if stop, err := c.answerOn(s.Subject, s.Queue, h); err == nil {
|
|
stops = append(stops, stop)
|
|
} else {
|
|
c.Logf("[mesh-tools] cannot serve %s on %s: %v", key, s.Subject, err)
|
|
}
|
|
}
|
|
}
|
|
serve()
|
|
c.OnMembership(func(m Membership) {
|
|
if m.Module == module {
|
|
serve()
|
|
}
|
|
})
|
|
return func() {
|
|
mu.Lock()
|
|
defer mu.Unlock()
|
|
for _, s := range stops {
|
|
s()
|
|
}
|
|
stops = nil
|
|
}, nil
|
|
}
|
|
|
|
// reachedAt is where a call by key goes: a subject this connection's own membership says it
|
|
// reaches — the machine's when named — else the derived shape.
|
|
func (c *Conn) reachedAt(key string) (string, error) {
|
|
name, wanted, _ := strings.Cut(key, "@")
|
|
if m := c.Membership(c.self); m != nil {
|
|
if reach := m.Reaches[name]; len(reach) > 0 {
|
|
if wanted == "" {
|
|
return reach[0], nil
|
|
}
|
|
for _, s := range reach {
|
|
if strings.HasSuffix(s, "."+wanted) {
|
|
return s, nil
|
|
}
|
|
}
|
|
}
|
|
}
|
|
return ToolSubject(key, c.self)
|
|
}
|
|
|
|
// Ask calls a tool by key — or on a subject the mesh listed for it — and learns which machine
|
|
// answered. Core request/reply: a tool call is never persisted (design 25 §3).
|
|
func (c *Conn) Ask(key string, body any, on string) (Answered, error) {
|
|
subject := on
|
|
if subject == "" {
|
|
s, err := c.reachedAt(key)
|
|
if err != nil {
|
|
return Answered{}, err
|
|
}
|
|
subject = s
|
|
}
|
|
data, err := wire.Marshal(body)
|
|
if err != nil {
|
|
return Answered{}, err
|
|
}
|
|
msg, err := c.nc.Request(subject, data, RequestTimeout)
|
|
if err != nil {
|
|
if errors.Is(err, nats.ErrNoResponders) {
|
|
return Answered{}, errors.New("503 no responders")
|
|
}
|
|
if errors.Is(err, nats.ErrTimeout) {
|
|
return Answered{}, errors.New("timeout")
|
|
}
|
|
return Answered{}, err
|
|
}
|
|
var r struct {
|
|
Result json.RawMessage `json:"result"`
|
|
Error string `json:"error"`
|
|
Node string `json:"node"`
|
|
}
|
|
if err := json.Unmarshal(msg.Data, &r); err != nil {
|
|
return Answered{}, err
|
|
}
|
|
if r.Error != "" {
|
|
return Answered{}, errors.New(r.Error)
|
|
}
|
|
return Answered{Result: r.Result, Node: r.Node}, nil
|
|
}
|
|
|
|
// PublishAs emits an event as a module: published into JetStream and awaited, de-duplicated by its
|
|
// own id (ADR 0042). The body is the payload; the metadata rides as headers.
|
|
func (c *Conn) PublishAs(module string, env Envelope) error {
|
|
msg := nats.NewMsg("mesh.mod." + module + ".event." + env.Key)
|
|
for k, v := range env.Headers {
|
|
msg.Header.Set(k, v)
|
|
}
|
|
if env.Headers["content-type"] == "" {
|
|
msg.Header.Set("content-type", "application/json")
|
|
}
|
|
if env.Node != "" {
|
|
msg.Header.Set("x-node", env.Node)
|
|
}
|
|
body := env.Body
|
|
if len(body) == 0 {
|
|
body = json.RawMessage("null")
|
|
}
|
|
msg.Data = body
|
|
var opts []nats.PubOpt
|
|
if id := env.Headers["x-event-id"]; id != "" {
|
|
opts = append(opts, nats.MsgId(id))
|
|
}
|
|
_, err := c.js.PublishMsg(msg, opts...)
|
|
return err
|
|
}
|
|
|
|
// Flush waits until the bus has every subscription made so far, so what is served is answerable
|
|
// when this returns.
|
|
func (c *Conn) Flush() { _ = c.nc.Flush() }
|
|
|
|
// Close unsubscribes everything and drains, so an in-flight reply is finished rather than dropped.
|
|
func (c *Conn) Close() {
|
|
c.mu.Lock()
|
|
subs := c.subs
|
|
c.subs = nil
|
|
c.mu.Unlock()
|
|
for _, s := range subs {
|
|
_ = s.Unsubscribe()
|
|
}
|
|
_ = c.nc.Drain()
|
|
}
|
|
|
|
// ToolSubject is a tool's subject. A bare name is this module's own; `<module>.<tool>` another's;
|
|
// `seat:<seat>.<verb>` a role's, with `@<node>` for a node-scoped seat (design 33 §4).
|
|
func ToolSubject(key, self string) (string, error) {
|
|
if strings.HasPrefix(key, "seat:") {
|
|
rest := strings.TrimPrefix(key, "seat:")
|
|
dot := strings.Index(rest, ".")
|
|
if dot < 0 {
|
|
return "", fmt.Errorf("%q names a seat and no verb: seat:<seat>.<verb>", key)
|
|
}
|
|
seat := rest[:dot]
|
|
verb, node, _ := strings.Cut(rest[dot+1:], "@")
|
|
if node != "" {
|
|
return "mesh.seat." + seat + ".tool." + verb + "." + node, nil
|
|
}
|
|
return "mesh.seat." + seat + ".tool." + verb, nil
|
|
}
|
|
name, node, _ := strings.Cut(key, "@")
|
|
var base string
|
|
if dot := strings.Index(name, "."); dot < 0 {
|
|
base = "mesh.mod." + self + ".tool." + name
|
|
} else {
|
|
base = "mesh.mod." + name[:dot] + ".tool." + name[dot+1:]
|
|
}
|
|
if node != "" {
|
|
return base + "." + node, nil
|
|
}
|
|
return base, nil
|
|
}
|
|
|
|
// SeatToolSubject is a seat's verb as its holder serves it: flat for a mesh seat, carrying the
|
|
// machine for a node-scoped one.
|
|
func SeatToolSubject(seat, verb, scope, node string) string {
|
|
base := "mesh.seat." + seat + ".tool." + verb
|
|
if scope == "node" && node != "" {
|
|
return base + "." + node
|
|
}
|
|
return base
|
|
}
|