mesh/merge-gate pass: builds mesh-tools, node-tools → ace, g14, novox, shanks; no bus step; every machine composes with the change as it did without (4 of …
mesh/repo-check pass: THE CHANGE ALTERS ITS OWN CHECK (merge-check.sh): main's version judged it; the change's judges the pull requests after it merges; it…
mesh/delivery ready: it delivers once merged
mesh/delivery-group group feat/asks-answered-on-any-channel rejected: a member's own check failed
440 lines
16 KiB
Go
440 lines
16 KiB
Go
package bus
|
|
|
|
import (
|
|
"context"
|
|
"crypto/rand"
|
|
"encoding/hex"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"strings"
|
|
"sync"
|
|
"time"
|
|
|
|
"github.com/nats-io/nats.go"
|
|
"github.com/nats-io/nats.go/jetstream"
|
|
)
|
|
|
|
// A bundle's traffic on seats that name their caller or their kind (novox/hq ADR 0259 §3): submitting an
|
|
// accept and saying an event under its own name or kind, taking a held seat's work from its worker, asking a
|
|
// proof and answering one, and reading a record kept for it.
|
|
//
|
|
// **The runtime keeps each module to its own name and kind.** One account per machine carries every module
|
|
// on it, so the bus enforces only the union; what one module may reach is its membership's `seat-traffic`,
|
|
// and anything else is refused here, with the reason, before it is sent. A bundle reaches this only through
|
|
// its own standard input and output: no tool call does.
|
|
|
|
// SeatTraffic is what a membership lists of a module's seat traffic (the controller's broker.SeatTraffic).
|
|
type SeatTraffic struct {
|
|
Publish []string `json:"publish,omitempty"`
|
|
Subscribe []string `json:"subscribe,omitempty"`
|
|
Answers []string `json:"answers,omitempty"`
|
|
Workers []SeatWorker `json:"workers,omitempty"`
|
|
Records []string `json:"records,omitempty"`
|
|
// Kinds are the holders of every kinded bench the module uses or watches, with what each promises: the
|
|
// controller's record of their claims (novox/hq ADR 0259 §5).
|
|
Kinds []KindHeld `json:"kinds,omitempty"`
|
|
}
|
|
|
|
// KindHeld is one kind of a kinded bench and who holds it.
|
|
type KindHeld struct {
|
|
Seat string `json:"seat"`
|
|
Kind string `json:"kind"`
|
|
Module string `json:"module"`
|
|
Node string `json:"node"`
|
|
Capabilities []string `json:"capabilities,omitempty"`
|
|
}
|
|
|
|
// SeatKinds is the kinds the module's membership names, as the controller issued them.
|
|
func (c *Conn) SeatKinds(module string) ([]KindHeld, error) {
|
|
t, err := c.trafficOf(module)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
return append([]KindHeld{}, t.Kinds...), nil
|
|
}
|
|
|
|
// SeatWorker is one work queue a holder takes from.
|
|
type SeatWorker struct {
|
|
Stream string `json:"Stream"`
|
|
Consumer string `json:"Consumer"`
|
|
Filter string `json:"Filter"`
|
|
}
|
|
|
|
// ProofTimeout bounds how long a proof waits for its answer.
|
|
var ProofTimeout = 15 * time.Second
|
|
|
|
// WorkHeartbeat is how often work a bundle is still doing is said to be in progress, well inside the
|
|
// worker's ack wait.
|
|
var WorkHeartbeat = 15 * time.Second
|
|
|
|
// WorkBackoff is how long a piece of work the bundle did not take waits before it is offered again: from
|
|
// five seconds, doubling, to ten minutes.
|
|
func WorkBackoff(delivered uint64) time.Duration {
|
|
wait := 5 * time.Second
|
|
for i := uint64(1); i < delivered && wait < 10*time.Minute; i++ {
|
|
wait *= 2
|
|
}
|
|
return min(wait, 10*time.Minute)
|
|
}
|
|
|
|
// SubjectMatches is a subject against a permission pattern, in the server's wildcards.
|
|
func SubjectMatches(pattern, subject string) bool {
|
|
p, s := strings.Split(pattern, "."), strings.Split(subject, ".")
|
|
for i, part := range p {
|
|
if part == ">" {
|
|
return len(s) > i
|
|
}
|
|
if i >= len(s) || (part != "*" && part != s[i]) {
|
|
return false
|
|
}
|
|
}
|
|
return len(p) == len(s)
|
|
}
|
|
|
|
func (c *Conn) trafficOf(module string) (*SeatTraffic, error) {
|
|
m := c.Membership(module)
|
|
if m == nil {
|
|
return nil, fmt.Errorf("%s has no membership issued on %s yet, so it reaches no seat (novox/hq ADR 0259)",
|
|
module, c.node)
|
|
}
|
|
if m.SeatTraffic == nil {
|
|
return nil, fmt.Errorf("%s holds, uses and watches no seat that names its caller or its kind (novox/hq "+
|
|
"ADR 0259): its manifest says so under `uses`, `claims` and `consumes`", module)
|
|
}
|
|
return m.SeatTraffic, nil
|
|
}
|
|
|
|
func listed(patterns []string, subject string) bool {
|
|
for _, p := range patterns {
|
|
if SubjectMatches(p, subject) {
|
|
return true
|
|
}
|
|
}
|
|
return false
|
|
}
|
|
|
|
// mayPublish refuses a subject the module's seat traffic does not list. A wildcard is never published.
|
|
func (c *Conn) mayPublish(module, subject string) error {
|
|
t, err := c.trafficOf(module)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if strings.ContainsAny(subject, "*>") || !listed(t.Publish, subject) {
|
|
return fmt.Errorf("%s may not publish %s: a module submits and says only under its own name or kind, "+
|
|
"and its membership lists %s (novox/hq ADR 0259)", module, subject, strings.Join(t.Publish, ", "))
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func newID() string {
|
|
var b [16]byte
|
|
_, _ = rand.Read(b[:])
|
|
return hex.EncodeToString(b[:])
|
|
}
|
|
|
|
// SeatPublish submits an accept or says an event on a seat, as the module, into the stream that keeps it,
|
|
// awaited and de-duplicated by its id: the same id twice is one message. A proof is never published here.
|
|
func (c *Conn) SeatPublish(module, subject string, body json.RawMessage, id string) (uint64, error) {
|
|
seq, _, err := c.SeatPublishSaid(module, subject, body, id)
|
|
return seq, err
|
|
}
|
|
|
|
// SeatPublishSaid is SeatPublish, also saying whether the bus took it as a duplicate of one published before
|
|
// under the same id: the same message, kept once.
|
|
func (c *Conn) SeatPublishSaid(module, subject string, body json.RawMessage, id string) (uint64, bool, error) {
|
|
if strings.Contains(subject, ".proof.") {
|
|
return 0, false, fmt.Errorf("%s is a proof, asked and answered, never kept: ask it as one", subject)
|
|
}
|
|
if err := c.mayPublish(module, subject); err != nil {
|
|
return 0, false, err
|
|
}
|
|
if len(body) == 0 || !json.Valid(body) {
|
|
return 0, false, fmt.Errorf("what is said on %s is JSON", subject)
|
|
}
|
|
if id == "" {
|
|
id = newID()
|
|
}
|
|
// **The publisher's name is part of the id the bus de-duplicates by** (security review 2026-10-08): ids
|
|
// are the publisher's own, and one module could otherwise suppress another's warrant by publishing first
|
|
// under the same id.
|
|
id = module + "." + id
|
|
msg := nats.NewMsg(subject)
|
|
msg.Header.Set("x-event-id", id)
|
|
msg.Header.Set("x-source", module)
|
|
msg.Header.Set("x-node", c.node)
|
|
msg.Header.Set("x-time", time.Now().UTC().Format(time.RFC3339Nano))
|
|
msg.Header.Set("content-type", "application/json")
|
|
msg.Data = body
|
|
ack, err := c.js.PublishMsg(msg, nats.MsgId(id))
|
|
if err != nil {
|
|
return 0, false, err
|
|
}
|
|
return ack.Sequence, ack.Duplicate, nil
|
|
}
|
|
|
|
// SeatProve asks a proof and answers its reply: core request and reply, on no stream.
|
|
func (c *Conn) SeatProve(module, subject string, body json.RawMessage) (json.RawMessage, error) {
|
|
if !strings.Contains(subject, ".proof.") {
|
|
return nil, fmt.Errorf("%s is not a proof", subject)
|
|
}
|
|
if err := c.mayPublish(module, subject); err != nil {
|
|
return nil, err
|
|
}
|
|
reply, err := c.nc.Request(subject, body, ProofTimeout)
|
|
if errors.Is(err, nats.ErrNoResponders) {
|
|
return nil, fmt.Errorf("nothing answers %s: the module watching the seat is not running", subject)
|
|
}
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
if !json.Valid(reply.Data) {
|
|
return nil, fmt.Errorf("the answer to %s is not JSON", subject)
|
|
}
|
|
return json.RawMessage(reply.Data), nil
|
|
}
|
|
|
|
// SeatAnswer answers every proof on a subject the module's seat traffic lists among its answers, with what
|
|
// answer returns. A returned error is answered as `{"error": "<words>"}`, never as silence.
|
|
func (c *Conn) SeatAnswer(module, subject string, answer func(subject string, body json.RawMessage) (json.RawMessage, error)) (func(), error) {
|
|
t, err := c.trafficOf(module)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
allowed := false
|
|
for _, a := range t.Answers {
|
|
allowed = allowed || a == subject
|
|
}
|
|
if !allowed {
|
|
return nil, fmt.Errorf("%s does not answer %s: its membership lists %s (novox/hq ADR 0259)", module,
|
|
subject, strings.Join(t.Answers, ", "))
|
|
}
|
|
sub, err := c.nc.Subscribe(subject, func(msg *nats.Msg) {
|
|
if msg.Reply == "" {
|
|
return
|
|
}
|
|
body := json.RawMessage(msg.Data)
|
|
if !json.Valid(body) {
|
|
_ = msg.Respond([]byte(`{"error":"a proof is JSON"}`))
|
|
return
|
|
}
|
|
out, err := answer(msg.Subject, body)
|
|
if err != nil {
|
|
b, _ := json.Marshal(map[string]string{"error": err.Error()})
|
|
_ = msg.Respond(b)
|
|
return
|
|
}
|
|
if len(out) == 0 {
|
|
out = json.RawMessage(`{}`)
|
|
}
|
|
if err := msg.Respond(out); err != nil {
|
|
c.Logf("[mesh-tools] %s could not answer the proof on %s: %v", module, msg.Subject, err)
|
|
}
|
|
})
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
c.track(sub)
|
|
var once sync.Once
|
|
return func() { once.Do(func() { _ = sub.Unsubscribe() }) }, nil
|
|
}
|
|
|
|
// Work is one message a holder took from a seat's work queue.
|
|
type Work struct {
|
|
Worker string `json:"worker"`
|
|
Subject string `json:"subject"`
|
|
Body json.RawMessage `json:"body"`
|
|
Headers map[string]string `json:"headers,omitempty"`
|
|
}
|
|
|
|
// SeatTake reads a worker the module's seat traffic lists and hands each message to deliver: acknowledged
|
|
// when deliver returns nil, offered again after NakDelay when it returns an error, ended when it is not
|
|
// JSON. The worker is the controller's to create; this binds to it by name and never makes one.
|
|
func (c *Conn) SeatTake(module, consumer string, deliver func(Work) error) (func(), error) {
|
|
t, err := c.trafficOf(module)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
var w *SeatWorker
|
|
for i := range t.Workers {
|
|
if t.Workers[i].Consumer == consumer {
|
|
w = &t.Workers[i]
|
|
}
|
|
}
|
|
if w == nil {
|
|
var names []string
|
|
for _, x := range t.Workers {
|
|
names = append(names, x.Consumer)
|
|
}
|
|
return nil, fmt.Errorf("%s takes no work from %s: its membership lists %s (novox/hq ADR 0259)", module,
|
|
consumer, strings.Join(names, ", "))
|
|
}
|
|
js, err := jetstream.New(c.nc)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
|
|
defer cancel()
|
|
bound, err := js.Consumer(ctx, w.Stream, w.Consumer)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("binding %s's worker %s on %s: %w", module, w.Consumer, w.Stream, err)
|
|
}
|
|
reading, err := bound.Consume(func(msg jetstream.Msg) {
|
|
if !json.Valid(msg.Data()) {
|
|
_ = msg.Term()
|
|
return
|
|
}
|
|
headers := map[string]string{}
|
|
for k := range msg.Headers() {
|
|
headers[k] = msg.Headers().Get(k)
|
|
}
|
|
// Kept alive while the bundle works on it, however long that takes: the ack wait is for a holder
|
|
// that died, not one that is slow.
|
|
done := make(chan struct{})
|
|
go func() {
|
|
tick := time.NewTicker(WorkHeartbeat)
|
|
defer tick.Stop()
|
|
for {
|
|
select {
|
|
case <-done:
|
|
return
|
|
case <-tick.C:
|
|
_ = msg.InProgress()
|
|
}
|
|
}
|
|
}()
|
|
err := deliver(Work{Worker: w.Consumer, Subject: msg.Subject(), Body: json.RawMessage(msg.Data()), Headers: headers})
|
|
close(done)
|
|
if err != nil {
|
|
// Offered again later and later, never given up on: a channel that is away keeps its work
|
|
// (security and correctness reviews of 2026-10-08).
|
|
delivered := uint64(1)
|
|
if meta, merr := msg.Metadata(); merr == nil {
|
|
delivered = meta.NumDelivered
|
|
}
|
|
_ = msg.NakWithDelay(WorkBackoff(delivered))
|
|
return
|
|
}
|
|
_ = msg.Ack()
|
|
}, jetstream.ConsumeErrHandler(func(_ jetstream.ConsumeContext, err error) {
|
|
// A worker deleted, or the bus refusing it, is said — never a holder that silently takes nothing.
|
|
c.Logf("[mesh-tools] %s's worker %s: %v", module, w.Consumer, err)
|
|
}))
|
|
if err != nil {
|
|
return nil, fmt.Errorf("reading %s's worker %s: %w", module, w.Consumer, err)
|
|
}
|
|
var once sync.Once
|
|
return func() { once.Do(reading.Stop) }, nil
|
|
}
|
|
|
|
// SeatRecord reads one key of a record a seat keeps for the module, under its own name: the bucket's
|
|
// current value of `<module>.<key>`, or nil when there is none.
|
|
func (c *Conn) SeatRecord(module, bucket, key string) (*StateEntry, error) {
|
|
t, err := c.trafficOf(module)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
if err := checkKey(key); err != nil {
|
|
return nil, err
|
|
}
|
|
if bucket == "" {
|
|
// The one record its membership lists, when it lists one: an asker reads the router's asks without
|
|
// knowing the router's name.
|
|
var named []string
|
|
for _, r := range t.Records {
|
|
if b, _, ok := strings.Cut(strings.TrimPrefix(r, "$JS.API.DIRECT.GET.KV_"), "."); ok {
|
|
named = append(named, b)
|
|
}
|
|
}
|
|
if len(named) != 1 {
|
|
return nil, fmt.Errorf("%s reads %d records under its name; name the one meant", module, len(named))
|
|
}
|
|
bucket = named[0]
|
|
}
|
|
subject := "$JS.API.DIRECT.GET.KV_" + bucket + ".$KV." + bucket + "." + module + "." + key
|
|
if !listed(t.Records, subject) {
|
|
return nil, fmt.Errorf("%s reads no record of %s under its name: its membership lists %s (novox/hq ADR 0259)",
|
|
module, bucket, strings.Join(t.Records, ", "))
|
|
}
|
|
reply, err := c.nc.Request(subject, nil, StateTimeout)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
if status := reply.Header.Get("Status"); status == "404" {
|
|
return nil, nil
|
|
} else if status != "" {
|
|
return nil, fmt.Errorf("the record %s.%s could not be read: %s %s", module, key, status, reply.Header.Get("Description"))
|
|
}
|
|
if op := reply.Header.Get("KV-Operation"); op == "DEL" || op == "PURGE" {
|
|
return nil, nil
|
|
}
|
|
var revision uint64
|
|
fmt.Sscan(reply.Header.Get("Nats-Sequence"), &revision)
|
|
return &StateEntry{Key: module + "." + key, Value: valueOf(reply.Data), Revision: revision}, nil
|
|
}
|
|
|
|
// StateCreate writes one key only when it has no value, and answers its revision; a key that has one is
|
|
// refused, naming it, so two writers racing to make it are one winner and one refusal.
|
|
func (c *Conn) StateCreate(module, name, key string, value json.RawMessage) (uint64, error) {
|
|
return c.stateWrite(module, name, key, value, func(ctx context.Context, kv jetstream.KeyValue) (uint64, error) {
|
|
return kv.Create(ctx, key, value)
|
|
})
|
|
}
|
|
|
|
// StateUpdate writes one key only when its revision is still the one given, and answers the new revision;
|
|
// a key changed meanwhile is refused, so of two writers that read one value only the first writes.
|
|
func (c *Conn) StateUpdate(module, name, key string, value json.RawMessage, revision uint64) (uint64, error) {
|
|
if revision == 0 {
|
|
// Revision zero would be read by the server as "the key is new": an update names the value it read.
|
|
return 0, fmt.Errorf("an update names the revision it read; zero is none — create the key instead")
|
|
}
|
|
return c.stateWrite(module, name, key, value, func(ctx context.Context, kv jetstream.KeyValue) (uint64, error) {
|
|
return kv.Update(ctx, key, value, revision)
|
|
})
|
|
}
|
|
|
|
// ErrStateChanged is a compare-and-set that lost: the key was made or changed by another writer first.
|
|
var ErrStateChanged = errors.New("the key was written by another writer first")
|
|
|
|
// changed is ErrStateChanged naming the key, with the code a bundle is answered it by.
|
|
type changed struct{ key string }
|
|
|
|
func (e changed) Error() string { return e.key + ": " + ErrStateChanged.Error() }
|
|
func (e changed) Is(target error) bool { return target == ErrStateChanged }
|
|
func (e changed) ErrorCode() int { return -32010 }
|
|
|
|
func (c *Conn) stateWrite(module, name, key string, value json.RawMessage,
|
|
write func(context.Context, jetstream.KeyValue) (uint64, error)) (uint64, error) {
|
|
s, err := c.writable(module, name)
|
|
if err != nil {
|
|
return 0, err
|
|
}
|
|
if err := checkKey(key); err != nil {
|
|
return 0, err
|
|
}
|
|
if len(value) == 0 || !json.Valid(value) {
|
|
return 0, fmt.Errorf("a state value is JSON")
|
|
}
|
|
if field := credentialField(value); field != "" {
|
|
return 0, fmt.Errorf("%s's %s.%s carries a field %q, which names a credential: no secret is kept in state "+
|
|
"(novox/hq ADR 0201)", module, name, key, field)
|
|
}
|
|
ctx, cancel := context.WithTimeout(context.Background(), StateTimeout)
|
|
defer cancel()
|
|
kv, err := c.bucket(ctx, s)
|
|
if err != nil {
|
|
return 0, err
|
|
}
|
|
rev, err := write(ctx, kv)
|
|
if errors.Is(err, jetstream.ErrKeyExists) || isWrongSequence(err) {
|
|
return 0, changed{name + "." + key}
|
|
}
|
|
return rev, err
|
|
}
|
|
|
|
// isWrongSequence is the server's refusal of an update whose revision is not the key's last.
|
|
func isWrongSequence(err error) bool {
|
|
var api *jetstream.APIError
|
|
return errors.As(err, &api) && api.ErrorCode == jetstream.JSErrCodeStreamWrongLastSequence
|
|
}
|