// 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..event. an event a module emits // mesh.mod..tool. a tool a module serves // mesh.seat..tool. 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 `.` 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; `.` another's; // `seat:.` a role's, with `@` 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:.", 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 }