Compare commits
4
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
127047edd6 | ||
|
|
c6087129ce | ||
|
|
7b4428c510 | ||
|
|
6d9906b6c9 |
@@ -33,6 +33,9 @@ if ! command -v node >/dev/null 2>&1 || [ ! -d node_modules/@novox/mesh-sdk ]; t
|
||||
# What of the console launches no bundle runs anyway: how an address resolves (novox/hq issue 287),
|
||||
# and how a search ranks (novox/hq ADR 0245).
|
||||
bundleless='^(TestASeatAndAModuleOfOneNameAreEachReached|TestSearch.*|TestASearchByCommandSaysWhatItReplaces)$'
|
||||
# And what of the runtime launches no bundle: how a bundle's seat traffic reaches the bus, what its
|
||||
# environment holds, and what a runtime on a module's own account serves (novox/hq ADR 0259).
|
||||
runtimeless='^(TestARuntimeOnAModulesOwnAccountServesThatModuleAlone|TestNoBusWordReachesABundle|TestAToolCallNeverReachesASeatsEventAcceptOrProof|TestABundlesSeatTrafficGoesThroughItsModulesBus|TestALostCompareAndSetCarriesItsCode)$'
|
||||
fi
|
||||
if command -v gcc >/dev/null 2>&1; then
|
||||
CGO_ENABLED=1 go test -race -count=1 $packages
|
||||
@@ -43,4 +46,7 @@ fi
|
||||
if [ -n "$bundleless" ]; then
|
||||
CGO_ENABLED=0 go test -count=1 -run "$bundleless" ./internal/console
|
||||
fi
|
||||
if [ -n "${runtimeless:-}" ]; then
|
||||
CGO_ENABLED=0 go test -count=1 -run "$runtimeless" ./internal/runtime
|
||||
fi
|
||||
echo "NOT TESTED HERE: node-tools/src (TypeScript) needs @novox/mesh-sdk from the package registry"
|
||||
|
||||
@@ -64,6 +64,9 @@ type Membership struct {
|
||||
Tools string `json:"tools"`
|
||||
// State is every bucket the module's code may reach, by the name it uses (novox/hq ADR 0201).
|
||||
State []StateIssued `json:"state,omitempty"`
|
||||
// SeatTraffic is what the module's code may submit, say, take, ask, answer and read on seats that name
|
||||
// their caller or their kind (novox/hq ADR 0259): the runtime does for it only what is listed here.
|
||||
SeatTraffic *SeatTraffic `json:"seat-traffic,omitempty"`
|
||||
}
|
||||
|
||||
// Served is an address a tool is answered on; `{tool}` stands for the tool's name.
|
||||
|
||||
@@ -0,0 +1,439 @@
|
||||
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
|
||||
}
|
||||
@@ -0,0 +1,239 @@
|
||||
package bus_test
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"strings"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/nats-io/nats.go"
|
||||
|
||||
"github.com/novox/mesh-tools/node-tools/internal/bus"
|
||||
mt "github.com/novox/mesh-tools/node-tools/internal/meshtest"
|
||||
)
|
||||
|
||||
// The seats of novox/hq ADR 0259 §3 as the controller issues them: an asker, the router, and a channel of
|
||||
// kind telegram, each with the seat traffic its membership lists.
|
||||
func seatMesh(t *testing.T) (*mt.Mesh, *bus.Conn, nats.JetStreamContext) {
|
||||
t.Helper()
|
||||
mesh := mt.New(t)
|
||||
nc, err := nats.Connect(mt.URL(t))
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
t.Cleanup(nc.Close)
|
||||
js, _ := nc.JetStream()
|
||||
for _, s := range []*nats.StreamConfig{
|
||||
{Name: "SEAT_OPERATOR_CHANNEL", Subjects: []string{"mesh.seat.operator-channel.accept.>"}, Retention: nats.WorkQueuePolicy},
|
||||
{Name: "SEAT_CHANNEL", Subjects: []string{"mesh.seat.channel.accept.>"}, Retention: nats.WorkQueuePolicy},
|
||||
{Name: "SEAT_EVENTS", Subjects: []string{"mesh.seat.*.event.>"}},
|
||||
} {
|
||||
if _, err := js.AddStream(s); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
for stream, c := range map[string]*nats.ConsumerConfig{
|
||||
"SEAT_OPERATOR_CHANNEL": {Durable: "SEAT_OPERATOR_CHANNEL_worker", AckPolicy: nats.AckExplicitPolicy,
|
||||
FilterSubject: "mesh.seat.operator-channel.accept.>", AckWait: 2 * time.Second},
|
||||
"SEAT_CHANNEL": {Durable: "SEAT_CHANNEL_TELEGRAM_worker", AckPolicy: nats.AckExplicitPolicy,
|
||||
FilterSubject: "mesh.seat.channel.accept.*.telegram", AckWait: 2 * time.Second},
|
||||
} {
|
||||
if _, err := js.AddConsumer(stream, c); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
if _, err := js.CreateKeyValue(&nats.KeyValueConfig{Bucket: "messenger_asks"}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
asker := mt.MembershipOf("mesh-delivery", "anchor", false, nil)
|
||||
asker.SeatTraffic = &bus.SeatTraffic{
|
||||
Publish: []string{"mesh.seat.operator-channel.accept.ask.mesh-delivery"},
|
||||
Subscribe: []string{"mesh.seat.operator-channel.event.decided.mesh-delivery"},
|
||||
Records: []string{"$JS.API.DIRECT.GET.KV_messenger_asks.$KV.messenger_asks.mesh-delivery.>"},
|
||||
}
|
||||
router := mt.MembershipOf("messenger", "anchor", false, nil)
|
||||
router.SeatTraffic = &bus.SeatTraffic{
|
||||
Publish: []string{"mesh.seat.operator-channel.event.decided.*", "mesh.seat.channel.accept.show.*"},
|
||||
Answers: []string{"mesh.seat.intake.proof.code.*"},
|
||||
Workers: []bus.SeatWorker{{Stream: "SEAT_OPERATOR_CHANNEL", Consumer: "SEAT_OPERATOR_CHANNEL_worker",
|
||||
Filter: "mesh.seat.operator-channel.accept.>"}},
|
||||
}
|
||||
router.State = []bus.StateIssued{{Name: "asks", Bucket: "messenger_asks", Writes: true}}
|
||||
telegram := mt.MembershipOf("telegram", "anchor", false, nil)
|
||||
telegram.SeatTraffic = &bus.SeatTraffic{
|
||||
Publish: []string{"mesh.seat.intake.event.choice.telegram", "mesh.seat.intake.proof.code.telegram"},
|
||||
Workers: []bus.SeatWorker{{Stream: "SEAT_CHANNEL", Consumer: "SEAT_CHANNEL_TELEGRAM_worker",
|
||||
Filter: "mesh.seat.channel.accept.*.telegram"}},
|
||||
}
|
||||
for _, m := range []bus.Membership{asker, router, telegram} {
|
||||
mesh.Issue(t, m)
|
||||
}
|
||||
conn, err := bus.Connect(bus.Credential{URL: mt.URL(t), Module: "node-tools", Node: "anchor"})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
conn.Logf = func(string, ...any) {}
|
||||
t.Cleanup(conn.Close)
|
||||
for _, m := range []string{"mesh-delivery", "messenger", "telegram"} {
|
||||
conn.Follow(m)
|
||||
}
|
||||
mt.Until(t, func() error {
|
||||
for _, m := range []string{"mesh-delivery", "messenger", "telegram"} {
|
||||
if conn.Membership(m) == nil {
|
||||
return errors.New(m + " not issued yet")
|
||||
}
|
||||
}
|
||||
return nil
|
||||
})
|
||||
return mesh, conn, js
|
||||
}
|
||||
|
||||
func TestAModuleSubmitsAndSaysOnlyUnderItsOwnNameAndKind(t *testing.T) {
|
||||
_, conn, _ := seatMesh(t)
|
||||
body := json.RawMessage(`{"id":"a1"}`)
|
||||
if _, err := conn.SeatPublish("mesh-delivery", "mesh.seat.operator-channel.accept.ask.mesh-delivery", body, "a1"); err != nil {
|
||||
t.Fatalf("an asker could not ask under its own name: %v", err)
|
||||
}
|
||||
for module, subject := range map[string]string{
|
||||
"mesh-delivery": "mesh.seat.operator-channel.accept.ask.mesh-controller",
|
||||
"telegram": "mesh.seat.intake.event.choice.desktop",
|
||||
"messenger": "mesh.seat.intake.event.choice.telegram",
|
||||
} {
|
||||
if _, err := conn.SeatPublish(module, subject, body, ""); err == nil || !strings.Contains(err.Error(), "may not publish") {
|
||||
t.Errorf("%s published %s: %v", module, subject, err)
|
||||
}
|
||||
}
|
||||
if _, err := conn.SeatPublish("messenger", "mesh.seat.operator-channel.event.decided.*", body, ""); err == nil {
|
||||
t.Error("a wildcard was published")
|
||||
}
|
||||
if _, err := conn.SeatPublish("telegram", "mesh.seat.intake.proof.code.telegram", body, ""); err == nil ||
|
||||
!strings.Contains(err.Error(), "never kept") {
|
||||
t.Errorf("a proof was published onto a stream: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestTheHolderTakesItsWorkOnceAndAFailureComesBack(t *testing.T) {
|
||||
_, conn, _ := seatMesh(t)
|
||||
if _, err := conn.SeatPublish("mesh-delivery", "mesh.seat.operator-channel.accept.ask.mesh-delivery",
|
||||
json.RawMessage(`{"id":"a1"}`), "a1"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
// The same id again is one ask: a retry after a lost answer does not ask twice.
|
||||
if _, err := conn.SeatPublish("mesh-delivery", "mesh.seat.operator-channel.accept.ask.mesh-delivery",
|
||||
json.RawMessage(`{"id":"a1"}`), "a1"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
var mu sync.Mutex
|
||||
var got []string
|
||||
fail := true
|
||||
stop, err := conn.SeatTake("messenger", "SEAT_OPERATOR_CHANNEL_worker", func(w bus.Work) error {
|
||||
mu.Lock()
|
||||
defer mu.Unlock()
|
||||
got = append(got, w.Subject+" from "+w.Headers["x-source"])
|
||||
if fail {
|
||||
fail = false
|
||||
return errors.New("not recorded")
|
||||
}
|
||||
return nil
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer stop()
|
||||
deadline := time.Now().Add(15 * time.Second)
|
||||
for {
|
||||
mu.Lock()
|
||||
n := len(got)
|
||||
mu.Unlock()
|
||||
if n >= 2 || time.Now().After(deadline) {
|
||||
break
|
||||
}
|
||||
time.Sleep(50 * time.Millisecond)
|
||||
}
|
||||
time.Sleep(500 * time.Millisecond)
|
||||
mu.Lock()
|
||||
defer mu.Unlock()
|
||||
if len(got) != 2 || got[0] != "mesh.seat.operator-channel.accept.ask.mesh-delivery from mesh-delivery" || got[0] != got[1] {
|
||||
t.Fatalf("want the one ask, offered again after the failure: %v", got)
|
||||
}
|
||||
if _, err := conn.SeatTake("telegram", "SEAT_OPERATOR_CHANNEL_worker", func(bus.Work) error { return nil }); err == nil {
|
||||
t.Error("a channel took the router's work")
|
||||
}
|
||||
}
|
||||
|
||||
func TestAProofIsAskedAndAnsweredAndNeverKept(t *testing.T) {
|
||||
_, conn, js := seatMesh(t)
|
||||
stop, err := conn.SeatAnswer("messenger", "mesh.seat.intake.proof.code.*", func(subject string, body json.RawMessage) (json.RawMessage, error) {
|
||||
if !strings.HasSuffix(subject, ".telegram") {
|
||||
return nil, errors.New("asked by another kind")
|
||||
}
|
||||
return json.RawMessage(`{"linked":true}`), nil
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer stop()
|
||||
var answer json.RawMessage
|
||||
mt.Until(t, func() error {
|
||||
answer, err = conn.SeatProve("telegram", "mesh.seat.intake.proof.code.telegram", json.RawMessage(`{"code":"123456"}`))
|
||||
return err
|
||||
})
|
||||
if string(answer) != `{"linked":true}` {
|
||||
t.Fatalf("the proof was answered %s", answer)
|
||||
}
|
||||
if _, err := conn.SeatProve("telegram", "mesh.seat.intake.proof.code.desktop", json.RawMessage(`{}`)); err == nil {
|
||||
t.Error("a channel asked a proof as another kind")
|
||||
}
|
||||
if _, err := conn.SeatAnswer("telegram", "mesh.seat.intake.proof.code.*", nil); err == nil {
|
||||
t.Error("a channel answers proofs")
|
||||
}
|
||||
for name := range js.StreamNames() {
|
||||
info, _ := js.StreamInfo(name)
|
||||
for _, s := range info.Config.Subjects {
|
||||
if bus.SubjectMatches(s, "mesh.seat.intake.proof.code.telegram") {
|
||||
t.Errorf("%s keeps proofs", name)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestAnAskerReadsItsOwnRecordAndNoOther(t *testing.T) {
|
||||
_, conn, js := seatMesh(t)
|
||||
kv, _ := js.KeyValue("messenger_asks")
|
||||
if _, err := kv.Put("mesh-delivery.a1", []byte(`{"state":"open"}`)); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := kv.Put("mesh-controller.a1", []byte(`{"state":"open"}`)); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
got, err := conn.SeatRecord("mesh-delivery", "messenger_asks", "a1")
|
||||
if err != nil || got == nil || string(got.Value) != `{"state":"open"}` {
|
||||
t.Fatalf("the asker could not read its own ask: %v %v", got, err)
|
||||
}
|
||||
if missing, err := conn.SeatRecord("mesh-delivery", "messenger_asks", "a2"); err != nil || missing != nil {
|
||||
t.Fatalf("a missing ask read as %v, %v", missing, err)
|
||||
}
|
||||
if _, err := conn.SeatRecord("telegram", "messenger_asks", "a1"); err == nil {
|
||||
t.Error("a module read records it is not given")
|
||||
}
|
||||
}
|
||||
|
||||
func TestOfTwoAnswersToOneAskOnlyTheFirstIsWritten(t *testing.T) {
|
||||
_, conn, _ := seatMesh(t)
|
||||
rev, err := conn.StateCreate("messenger", "asks", "mesh-delivery.a1", json.RawMessage(`{"state":"open"}`))
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := conn.StateCreate("messenger", "asks", "mesh-delivery.a1", json.RawMessage(`{"state":"open"}`)); !errors.Is(err, bus.ErrStateChanged) {
|
||||
t.Errorf("a second create stood: %v", err)
|
||||
}
|
||||
if _, err := conn.StateUpdate("messenger", "asks", "mesh-delivery.a1", json.RawMessage(`{"state":"answered"}`), rev); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := conn.StateUpdate("messenger", "asks", "mesh-delivery.a1", json.RawMessage(`{"state":"answered"}`), rev); !errors.Is(err, bus.ErrStateChanged) {
|
||||
t.Errorf("a second answer from the same revision stood: %v", err)
|
||||
}
|
||||
}
|
||||
@@ -40,6 +40,9 @@ func (b *subscribing) State(string, json.RawMessage) (json.RawMessage, error) {
|
||||
func (b *subscribing) Watch(json.RawMessage, func(json.RawMessage) error) (func(), error) {
|
||||
return func() {}, nil
|
||||
}
|
||||
func (b *subscribing) Seat(string, json.RawMessage, func(string, any) (json.RawMessage, error)) (json.RawMessage, error) {
|
||||
return nil, nil
|
||||
}
|
||||
func (b *subscribing) Subscribe(d func(json.RawMessage) error) error {
|
||||
b.mu.Lock()
|
||||
defer b.mu.Unlock()
|
||||
|
||||
@@ -60,6 +60,10 @@ type Bus interface {
|
||||
Subscribe(deliver func(envelope json.RawMessage) error) error
|
||||
State(verb string, params json.RawMessage) (json.RawMessage, error)
|
||||
Watch(params json.RawMessage, deliver func(change json.RawMessage) error) (stop func(), err error)
|
||||
// Seat answers a bundle's seat traffic (novox/hq ADR 0259 §3): `publish`, `prove` and `record` at once;
|
||||
// `take` and `answer` bind once and hand each message or proof on through handOn, to whichever child
|
||||
// of the module asked last — a restarted child asks again as its code runs again.
|
||||
Seat(verb string, params json.RawMessage, handOn func(method string, params any) (json.RawMessage, error)) (json.RawMessage, error)
|
||||
}
|
||||
|
||||
// EventTimeout bounds how long a bundle has to handle one event before it is offered again.
|
||||
@@ -147,6 +151,17 @@ func Start(module, entry string, env []string, mesh Bus, logf func(string, ...an
|
||||
// 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
|
||||
// The child that last asked to take a seat's work or answer its proofs is the one they are handed to.
|
||||
var seatTaker *child
|
||||
handOn := func(method string, params any) (json.RawMessage, error) {
|
||||
mu.Lock()
|
||||
c := seatTaker
|
||||
mu.Unlock()
|
||||
if c == nil {
|
||||
return nil, errors.New(module + "'s bundle is not running to take it")
|
||||
}
|
||||
return c.ask(module, method, params, EventTimeout)
|
||||
}
|
||||
stopped := false
|
||||
deliver := func(envelope json.RawMessage) error {
|
||||
mu.Lock()
|
||||
@@ -237,7 +252,21 @@ func Start(module, entry string, env []string, mesh Bus, logf func(string, ...an
|
||||
subscriber = c
|
||||
mu.Unlock()
|
||||
err = mesh.Subscribe(deliver)
|
||||
case "mesh/state.get", "mesh/state.put", "mesh/state.delete", "mesh/state.keys":
|
||||
case "mesh/seat.publish", "mesh/seat.prove", "mesh/seat.record", "mesh/seat.kinds":
|
||||
var answered json.RawMessage
|
||||
if answered, err = mesh.Seat(strings.TrimPrefix(m.Method, "mesh/seat."), m.Params, nil); err == nil {
|
||||
result = answered
|
||||
}
|
||||
case "mesh/seat.take", "mesh/seat.answer":
|
||||
mu.Lock()
|
||||
seatTaker = c
|
||||
mu.Unlock()
|
||||
var answered json.RawMessage
|
||||
if answered, err = mesh.Seat(strings.TrimPrefix(m.Method, "mesh/seat."), m.Params, handOn); err == nil {
|
||||
result = answered
|
||||
}
|
||||
case "mesh/state.get", "mesh/state.put", "mesh/state.delete", "mesh/state.keys",
|
||||
"mesh/state.create", "mesh/state.update":
|
||||
var answered json.RawMessage
|
||||
if answered, err = mesh.State(strings.TrimPrefix(m.Method, "mesh/state."), m.Params); err == nil {
|
||||
result = answered
|
||||
@@ -275,7 +304,7 @@ func Start(module, entry string, env []string, mesh Bus, logf func(string, ...an
|
||||
}
|
||||
if err != nil {
|
||||
_ = c.write(map[string]any{"jsonrpc": "2.0", "id": m.ID,
|
||||
"error": map[string]any{"code": -32000, "message": err.Error()}})
|
||||
"error": map[string]any{"code": errorCode(err), "message": err.Error()}})
|
||||
return
|
||||
}
|
||||
_ = c.write(map[string]any{"jsonrpc": "2.0", "id": m.ID, "result": result})
|
||||
@@ -329,6 +358,9 @@ func Start(module, entry string, env []string, mesh Bus, logf func(string, ...an
|
||||
if subscriber == c {
|
||||
subscriber = nil
|
||||
}
|
||||
if seatTaker == c {
|
||||
seatTaker = nil
|
||||
}
|
||||
wasStopped := stopped
|
||||
// Started again at once if it ran a while; after a growing pause while it keeps exiting.
|
||||
wait := RestartFirst
|
||||
@@ -561,3 +593,17 @@ func (c *child) deathReason() error {
|
||||
}
|
||||
return errors.New("the bundle exited")
|
||||
}
|
||||
|
||||
// CodeChanged is the error code a bundle is answered a lost compare-and-set with (mesh-sdk's
|
||||
// stdio.CodeChanged): told by its code, never by its words (novox/hq ADR 0259).
|
||||
const CodeChanged = -32010
|
||||
|
||||
// errorCode is the code an error is answered to a bundle with: a lost compare-and-set's own, else the
|
||||
// general one.
|
||||
func errorCode(err error) int {
|
||||
var coded interface{ ErrorCode() int }
|
||||
if errors.As(err, &coded) {
|
||||
return coded.ErrorCode()
|
||||
}
|
||||
return -32000
|
||||
}
|
||||
|
||||
@@ -0,0 +1,71 @@
|
||||
package launch
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
// A bundle that asks to take a worker's work and answers each piece it is handed.
|
||||
const takingBundle = `#!/bin/sh
|
||||
while IFS= read -r line; do
|
||||
id=$(printf '%s' "$line" | sed -n 's/.*"id":\([0-9]*\)[,}].*/\1/p')
|
||||
case "$line" in
|
||||
*'"method":"initialize"'*) printf '{"jsonrpc":"2.0","id":%s,"result":{}}\n' "$id" ;;
|
||||
*'"method":"tools/list"'*)
|
||||
printf '{"jsonrpc":"2.0","id":%s,"result":{"tools":[]}}\n' "$id"
|
||||
printf '{"jsonrpc":"2.0","id":"s1","method":"mesh/seat.take","params":{"worker":"SEAT_CHANNEL_TELEGRAM_worker"}}\n' ;;
|
||||
*'"method":"mesh/work"'*) printf '{"jsonrpc":"2.0","id":%s,"result":{"taken":true}}\n' "$id" ;;
|
||||
esac
|
||||
done
|
||||
`
|
||||
|
||||
type seating struct {
|
||||
subscribing
|
||||
mu sync.Mutex
|
||||
verb string
|
||||
params string
|
||||
handOn func(string, any) (json.RawMessage, error)
|
||||
}
|
||||
|
||||
func (b *seating) Seat(verb string, params json.RawMessage, handOn func(string, any) (json.RawMessage, error)) (json.RawMessage, error) {
|
||||
b.mu.Lock()
|
||||
defer b.mu.Unlock()
|
||||
b.verb, b.params, b.handOn = verb, string(params), handOn
|
||||
return json.RawMessage(`{}`), nil
|
||||
}
|
||||
|
||||
// novox/hq ADR 0259 §3: a bundle's seat traffic reaches the runtime's bus from its own channel, and the work
|
||||
// it takes is handed back to it.
|
||||
func TestABundleTakesASeatsWorkThroughItsOwnChannel(t *testing.T) {
|
||||
entry := filepath.Join(t.TempDir(), "bundle")
|
||||
if err := os.WriteFile(entry, []byte(takingBundle), 0o755); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
b := &seating{}
|
||||
l, err := Start("telegram", entry, os.Environ(), b, t.Logf)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
t.Cleanup(l.Stop)
|
||||
var handOn func(string, any) (json.RawMessage, error)
|
||||
for i := 0; handOn == nil; i++ {
|
||||
if i > 100 {
|
||||
t.Fatal("the bundle never asked to take its work")
|
||||
}
|
||||
time.Sleep(20 * time.Millisecond)
|
||||
b.mu.Lock()
|
||||
handOn = b.handOn
|
||||
b.mu.Unlock()
|
||||
}
|
||||
if b.verb != "take" || b.params != `{"worker":"SEAT_CHANNEL_TELEGRAM_worker"}` {
|
||||
t.Fatalf("asked %s %s", b.verb, b.params)
|
||||
}
|
||||
got, err := handOn("mesh/work", map[string]any{"work": map[string]any{"subject": "mesh.seat.channel.accept.show.telegram"}})
|
||||
if err != nil || string(got) != `{"taken":true}` {
|
||||
t.Fatalf("the work was not handed to the bundle: %s %v", got, err)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,56 @@
|
||||
package runtime
|
||||
|
||||
import (
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"github.com/novox/mesh-tools/node-tools/internal/bus"
|
||||
)
|
||||
|
||||
// novox/hq ADR 0259 §8: a runtime on a module's own account serves that module alone, and the machine's
|
||||
// runtime never serves itself.
|
||||
func TestARuntimeOnAModulesOwnAccountServesThatModuleAlone(t *testing.T) {
|
||||
if got, err := ServedModulesFrom("telegram=/b/telegram/tools/telegram", "telegram"); err != nil || len(got) != 1 {
|
||||
t.Fatalf("a module's own runtime could not serve it: %v", err)
|
||||
}
|
||||
if _, err := ServedModulesFrom("telegram=/b/t,dunst=/b/d", "telegram"); err == nil || !strings.Contains(err.Error(), "serves telegram alone") {
|
||||
t.Errorf("a module's own runtime served another: %v", err)
|
||||
}
|
||||
if _, err := ServedModulesFrom("node-tools=/b/n", "node-tools"); err == nil {
|
||||
t.Error("the machine's runtime served itself")
|
||||
}
|
||||
}
|
||||
|
||||
// No bus word reaches a bundle (security review of 2026-10-08), not even one its module's words name.
|
||||
func TestNoBusWordReachesABundle(t *testing.T) {
|
||||
base := []string{"MESH_BROKER_FILE=/etc/mesh/broker", "MESH_BROKER_URL=nats://x", "MESH_TOOL_ENV={}",
|
||||
"MESH_TOOL_MODULES=a=/b", "MESH_CONSOLE_LISTEN=127.0.0.1:1", "PATH=/usr/bin"}
|
||||
env := strings.Join(BundleEnv(base, map[string]string{"MESH_BROKER_FILE": "/again", "OWN": "1"}, "telegram", "anchor"), "\n")
|
||||
for _, never := range []string{"MESH_BROKER_FILE", "MESH_BROKER_URL", "MESH_TOOL_ENV", "MESH_TOOL_MODULES", "MESH_CONSOLE_LISTEN"} {
|
||||
if strings.Contains(env, never+"=") {
|
||||
t.Errorf("%s reached the bundle", never)
|
||||
}
|
||||
}
|
||||
for _, want := range []string{"PATH=/usr/bin", "OWN=1", "MESH_MODULE=telegram", "MESH_NODE=anchor"} {
|
||||
if !strings.Contains(env, want) {
|
||||
t.Errorf("%s did not reach the bundle", want)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// A tool call reaches only a tool's subject: whatever key a caller names, it is never a seat's event, accept
|
||||
// or proof — so no tool call says a choice, a link or a code, or submits an ask, in anybody's name.
|
||||
func TestAToolCallNeverReachesASeatsEventAcceptOrProof(t *testing.T) {
|
||||
for _, key := range []string{"seat:intake.event.choice.telegram", "seat:intake.proof.code.telegram",
|
||||
"seat:operator-channel.accept.ask.mesh-delivery", "intake.event", "seat:operator-channel.event.decided.x@anchor",
|
||||
"telegram.anything", "x"} {
|
||||
subject, err := bus.ToolSubject(key, "node-tools")
|
||||
if err != nil {
|
||||
continue
|
||||
}
|
||||
tokens := strings.Split(subject, ".")
|
||||
if len(tokens) < 4 || tokens[3] != "tool" {
|
||||
t.Errorf("%q reaches %s, which is not a tool's subject", key, subject)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -53,6 +53,13 @@ type ListedTool struct {
|
||||
Subjects []string `json:"subjects,omitempty"`
|
||||
}
|
||||
|
||||
// NodeRuntimeModule is the module that is a machine's runtime, serving every module on it.
|
||||
const NodeRuntimeModule = "node-tools"
|
||||
|
||||
// runtimeOnly are the words of the runtime's own environment no bundle is given.
|
||||
var runtimeOnly = map[string]bool{ToolEnv: true, ToolModules: true, "MESH_BROKER_FILE": true,
|
||||
"MESH_BROKER_URL": true, "MESH_CONSOLE_LISTEN": true}
|
||||
|
||||
// ServedModulesFrom reads MESH_TOOL_MODULES: `<module>=<entrypoint>` entries, comma-separated, several
|
||||
// per module. The one-module form — a bare path, or the runtime's own module — is the per-module
|
||||
// containers' (to-be 38 WP4c) and refused here: the node's runtime imports nothing.
|
||||
@@ -66,9 +73,20 @@ func ServedModulesFrom(spec, own string) ([]Served, error) {
|
||||
}
|
||||
module, path, ok := strings.Cut(entry, "=")
|
||||
module, path = strings.TrimSpace(module), strings.TrimSpace(path)
|
||||
if !ok || module == "" || path == "" || module == own {
|
||||
return nil, fmt.Errorf("%s: %q is not <module>=<entrypoint> of another module; the node's "+
|
||||
"runtime launches the bundles it is given and imports nothing (novox/hq ADR 0193)", ToolModules, entry)
|
||||
if !ok || module == "" || path == "" {
|
||||
return nil, fmt.Errorf("%s: %q is not <module>=<entrypoint>; the runtime launches the bundles it "+
|
||||
"is given and imports nothing (novox/hq ADR 0193)", ToolModules, entry)
|
||||
}
|
||||
// The node's runtime serves the machine's modules and never itself. **A runtime on a module's own
|
||||
// account serves that module and nothing else** (novox/hq ADR 0259 §8): the router and a channel that
|
||||
// proves its sender reach the bus on an account of their own, never the machine's runtime.
|
||||
switch {
|
||||
case own == NodeRuntimeModule && module == own:
|
||||
return nil, fmt.Errorf("%s: %q is the runtime itself; it launches the bundles it is given and "+
|
||||
"imports nothing (novox/hq ADR 0193)", ToolModules, entry)
|
||||
case own != NodeRuntimeModule && module != own:
|
||||
return nil, fmt.Errorf("%s: %q is another module's; a runtime on %s's own account serves %s alone "+
|
||||
"(novox/hq ADR 0259)", ToolModules, entry, own, own)
|
||||
}
|
||||
if _, seen := by[module]; !seen {
|
||||
order = append(order, module)
|
||||
@@ -135,28 +153,7 @@ func Run(conn *bus.Conn, served []Served, envs map[string]map[string]string, log
|
||||
node := conn.Node()
|
||||
|
||||
base := os.Environ()
|
||||
envFor := func(module string) []string {
|
||||
words := map[string]string{}
|
||||
for _, kv := range base {
|
||||
if k, v, ok := strings.Cut(kv, "="); ok && k != ToolEnv {
|
||||
words[k] = v
|
||||
}
|
||||
}
|
||||
for k, v := range envs[module] {
|
||||
words[k] = v
|
||||
}
|
||||
words["MESH_SERVED_MODULE"] = module
|
||||
words["MESH_MODULE"] = module
|
||||
if node != "" {
|
||||
words["MESH_NODE"] = node
|
||||
}
|
||||
out := make([]string, 0, len(words))
|
||||
for k, v := range words {
|
||||
out = append(out, k+"="+v)
|
||||
}
|
||||
sort.Strings(out)
|
||||
return out
|
||||
}
|
||||
envFor := func(module string) []string { return BundleEnv(base, envs[module], module, node) }
|
||||
|
||||
// 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.
|
||||
@@ -458,6 +455,8 @@ type consumers struct {
|
||||
logf func(string, ...any)
|
||||
mu sync.Mutex
|
||||
of map[string]*moduleEvents
|
||||
// seatBound is each module's seat traffic bound once: a worker taken, a proof subject answered.
|
||||
seatBound map[string]func()
|
||||
}
|
||||
|
||||
type moduleEvents struct {
|
||||
@@ -475,6 +474,10 @@ func (c *consumers) stopAll() {
|
||||
}
|
||||
}
|
||||
c.of = map[string]*moduleEvents{}
|
||||
for _, stop := range c.seatBound {
|
||||
stop()
|
||||
}
|
||||
c.seatBound = map[string]func(){}
|
||||
}
|
||||
|
||||
// forModule is the bus one launched bundle of a module reaches the mesh through.
|
||||
@@ -604,9 +607,10 @@ func eventRef(env bus.Envelope) string {
|
||||
// stateAsked is what a bundle names when it reaches its state (ADR 0201): the state by the name its
|
||||
// module uses, a key, and for a put the value.
|
||||
type stateAsked struct {
|
||||
State string `json:"state"`
|
||||
Key string `json:"key"`
|
||||
Value json.RawMessage `json:"value"`
|
||||
State string `json:"state"`
|
||||
Key string `json:"key"`
|
||||
Value json.RawMessage `json:"value"`
|
||||
Revision uint64 `json:"revision"`
|
||||
}
|
||||
|
||||
// State answers a bundle's `get`, `put`, `delete` and `keys` on its module's state (ADR 0201). The
|
||||
@@ -645,6 +649,20 @@ func (b *moduleBus) State(verb string, params json.RawMessage) (json.RawMessage,
|
||||
return nil, err
|
||||
}
|
||||
answer = keys
|
||||
case "create":
|
||||
// Only when the key has no value: of two writers making it, one wins (novox/hq ADR 0259).
|
||||
revision, err := conn.StateCreate(b.module, asked.State, asked.Key, asked.Value)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
answer = map[string]any{"revision": revision}
|
||||
case "update":
|
||||
// Only when the key is still at the revision read: a second answer to one ask loses (ADR 0259).
|
||||
revision, err := conn.StateUpdate(b.module, asked.State, asked.Key, asked.Value, asked.Revision)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
answer = map[string]any{"revision": revision}
|
||||
default:
|
||||
return nil, fmt.Errorf("the runtime answers no mesh/state.%s", verb)
|
||||
}
|
||||
@@ -666,3 +684,124 @@ func (b *moduleBus) Watch(params json.RawMessage, deliver func(json.RawMessage)
|
||||
return deliver(raw)
|
||||
})
|
||||
}
|
||||
|
||||
// seatAsked is what a bundle names for its seat traffic (novox/hq ADR 0259 §3).
|
||||
type seatAsked struct {
|
||||
Subject string `json:"subject"`
|
||||
Body json.RawMessage `json:"body"`
|
||||
ID string `json:"id"`
|
||||
Worker string `json:"worker"`
|
||||
Bucket string `json:"bucket"`
|
||||
Key string `json:"key"`
|
||||
}
|
||||
|
||||
// Seat answers a bundle's seat traffic, each as the module and only as far as its membership lists:
|
||||
// `publish` an accept or an event, `prove` (ask a proof), `record` (read a record kept for it), `take` a
|
||||
// worker's work and `answer` a proof subject — the last two bound once per module and subject, each
|
||||
// message or proof handed to the module's child that asked last (novox/hq ADR 0259 §3).
|
||||
func (b *moduleBus) Seat(verb string, params json.RawMessage, handOn func(string, any) (json.RawMessage, error)) (json.RawMessage, error) {
|
||||
var asked seatAsked
|
||||
if err := json.Unmarshal(params, &asked); err != nil {
|
||||
return nil, fmt.Errorf("mesh/seat.%s: %w", verb, err)
|
||||
}
|
||||
conn := b.all.conn
|
||||
switch verb {
|
||||
case "publish":
|
||||
seq, duplicate, err := conn.SeatPublishSaid(b.module, asked.Subject, asked.Body, asked.ID)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return json.Marshal(map[string]any{"sequence": seq, "duplicate": duplicate})
|
||||
case "prove":
|
||||
return conn.SeatProve(b.module, asked.Subject, asked.Body)
|
||||
case "kinds":
|
||||
kinds, err := conn.SeatKinds(b.module)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return json.Marshal(kinds)
|
||||
case "record":
|
||||
entry, err := conn.SeatRecord(b.module, asked.Bucket, asked.Key)
|
||||
if err != nil || entry == nil {
|
||||
return json.RawMessage("null"), err
|
||||
}
|
||||
return json.Marshal(entry)
|
||||
case "take":
|
||||
key := "take " + asked.Worker
|
||||
if b.bound(key) {
|
||||
return json.RawMessage(`{}`), nil
|
||||
}
|
||||
stop, err := conn.SeatTake(b.module, asked.Worker, func(w bus.Work) error {
|
||||
_, err := handOn("mesh/work", map[string]any{"work": w})
|
||||
return err
|
||||
})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
b.bind(key, stop)
|
||||
b.all.logf("[mesh-tools] %s takes its work from %s", b.module, asked.Worker)
|
||||
return json.RawMessage(`{}`), nil
|
||||
case "answer":
|
||||
key := "answer " + asked.Subject
|
||||
if b.bound(key) {
|
||||
return json.RawMessage(`{}`), nil
|
||||
}
|
||||
stop, err := conn.SeatAnswer(b.module, asked.Subject, func(subject string, body json.RawMessage) (json.RawMessage, error) {
|
||||
return handOn("mesh/proof", map[string]any{"subject": subject, "body": body})
|
||||
})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
b.bind(key, stop)
|
||||
b.all.logf("[mesh-tools] %s answers the proofs on %s", b.module, asked.Subject)
|
||||
return json.RawMessage(`{}`), nil
|
||||
}
|
||||
return nil, fmt.Errorf("the runtime answers no mesh/seat.%s", verb)
|
||||
}
|
||||
|
||||
// bound and bind keep, per module, what its seat traffic is bound to, so a child that starts again and
|
||||
// asks again is handed what is already bound instead of a second reader of it.
|
||||
func (b *moduleBus) bound(key string) bool {
|
||||
b.all.mu.Lock()
|
||||
defer b.all.mu.Unlock()
|
||||
_, ok := b.all.seatBound[b.module+" "+key]
|
||||
return ok
|
||||
}
|
||||
|
||||
func (b *moduleBus) bind(key string, stop func()) {
|
||||
b.all.mu.Lock()
|
||||
defer b.all.mu.Unlock()
|
||||
if b.all.seatBound == nil {
|
||||
b.all.seatBound = map[string]func(){}
|
||||
}
|
||||
b.all.seatBound[b.module+" "+key] = stop
|
||||
}
|
||||
|
||||
// BundleEnv is what one module's bundle is started with: the runtime's environment without its own words,
|
||||
// the module's composed words over it, and the module's and machine's names. **No bus word reaches a
|
||||
// bundle** (security review of 2026-10-08): the credential and the bus's address are the runtime's, and a
|
||||
// bundle reaches the bus only through the runtime.
|
||||
func BundleEnv(base []string, given map[string]string, module, node string) []string {
|
||||
words := map[string]string{}
|
||||
for _, kv := range base {
|
||||
if k, v, ok := strings.Cut(kv, "="); ok && !runtimeOnly[k] {
|
||||
words[k] = v
|
||||
}
|
||||
}
|
||||
for k, v := range given {
|
||||
if !runtimeOnly[k] {
|
||||
words[k] = v
|
||||
}
|
||||
}
|
||||
words["MESH_SERVED_MODULE"] = module
|
||||
words["MESH_MODULE"] = module
|
||||
if node != "" {
|
||||
words["MESH_NODE"] = node
|
||||
}
|
||||
out := make([]string, 0, len(words))
|
||||
for k, v := range words {
|
||||
out = append(out, k+"="+v)
|
||||
}
|
||||
sort.Strings(out)
|
||||
return out
|
||||
}
|
||||
|
||||
@@ -0,0 +1,149 @@
|
||||
package runtime
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"strings"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/nats-io/nats.go"
|
||||
|
||||
"github.com/novox/mesh-tools/node-tools/internal/bus"
|
||||
mt "github.com/novox/mesh-tools/node-tools/internal/meshtest"
|
||||
)
|
||||
|
||||
// A bundle's seat traffic through moduleBus.Seat (novox/hq ADR 0259 §3), against a real bus: publish under
|
||||
// its own name with the id prefixed and a duplicate said; the record read without naming the bucket; a
|
||||
// worker and a proof subject bound once however often a restarted child asks; and nothing beyond what the
|
||||
// membership lists.
|
||||
func TestABundlesSeatTrafficGoesThroughItsModulesBus(t *testing.T) {
|
||||
mesh := mt.New(t)
|
||||
nc, err := nats.Connect(mt.URL(t))
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
t.Cleanup(nc.Close)
|
||||
js, _ := nc.JetStream()
|
||||
if _, err := js.AddStream(&nats.StreamConfig{Name: "SEAT_OPERATOR_CHANNEL", Retention: nats.WorkQueuePolicy,
|
||||
Subjects: []string{"mesh.seat.operator-channel.accept.>"}}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := js.AddConsumer("SEAT_OPERATOR_CHANNEL", &nats.ConsumerConfig{Durable: "SEAT_OPERATOR_CHANNEL_worker",
|
||||
AckPolicy: nats.AckExplicitPolicy, FilterSubject: "mesh.seat.operator-channel.accept.>"}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
kv, err := js.CreateKeyValue(&nats.KeyValueConfig{Bucket: "messenger_asks"})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
_, _ = kv.Put("mesh-delivery.d-1", []byte(`{"state":"open"}`))
|
||||
asker := mt.MembershipOf("mesh-delivery", "anchor", false, nil)
|
||||
asker.SeatTraffic = &bus.SeatTraffic{Publish: []string{"mesh.seat.operator-channel.accept.ask.mesh-delivery"},
|
||||
Records: []string{"$JS.API.DIRECT.GET.KV_messenger_asks.$KV.messenger_asks.mesh-delivery.>"}}
|
||||
router := mt.MembershipOf("messenger", "anchor", false, nil)
|
||||
router.SeatTraffic = &bus.SeatTraffic{Answers: []string{"mesh.seat.intake.proof.code.*"},
|
||||
Workers: []bus.SeatWorker{{Stream: "SEAT_OPERATOR_CHANNEL", Consumer: "SEAT_OPERATOR_CHANNEL_worker",
|
||||
Filter: "mesh.seat.operator-channel.accept.>"}}}
|
||||
mesh.Issue(t, asker)
|
||||
mesh.Issue(t, router)
|
||||
conn := connect(t, "node-tools", "anchor")
|
||||
conn.Follow("mesh-delivery")
|
||||
conn.Follow("messenger")
|
||||
mt.Until(t, func() error {
|
||||
if conn.Membership("mesh-delivery") == nil || conn.Membership("messenger") == nil {
|
||||
return errors.New("not issued yet")
|
||||
}
|
||||
return nil
|
||||
})
|
||||
all := &consumers{conn: conn, logf: t.Logf, of: map[string]*moduleEvents{}}
|
||||
asking := all.forModule("mesh-delivery").(*moduleBus)
|
||||
routing := all.forModule("messenger").(*moduleBus)
|
||||
|
||||
publish := json.RawMessage(`{"subject":"mesh.seat.operator-channel.accept.ask.mesh-delivery","body":{"id":"d-1"},"id":"d-1"}`)
|
||||
first, err := asking.Seat("publish", publish, nil)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
again, _ := asking.Seat("publish", publish, nil)
|
||||
if !strings.Contains(string(first), `"duplicate":false`) || !strings.Contains(string(again), `"duplicate":true`) {
|
||||
t.Errorf("the same id twice: %s, then %s", first, again)
|
||||
}
|
||||
if _, err := asking.Seat("publish", json.RawMessage(`{"subject":"mesh.seat.operator-channel.accept.ask.mesh-controller","body":{}}`), nil); err == nil {
|
||||
t.Error("a module submitted under another's name")
|
||||
}
|
||||
got, err := asking.Seat("record", json.RawMessage(`{"key":"d-1"}`), nil)
|
||||
if err != nil || !strings.Contains(string(got), `"state":"open"`) {
|
||||
t.Errorf("the record read without its bucket: %s %v", got, err)
|
||||
}
|
||||
|
||||
var mu sync.Mutex
|
||||
var taken []string
|
||||
handOn := func(method string, params any) (json.RawMessage, error) {
|
||||
mu.Lock()
|
||||
defer mu.Unlock()
|
||||
raw, _ := json.Marshal(params)
|
||||
taken = append(taken, method+" "+string(raw))
|
||||
return json.RawMessage(`{}`), nil
|
||||
}
|
||||
for i := 0; i < 2; i++ { // a child that started again asks again
|
||||
if _, err := routing.Seat("take", json.RawMessage(`{"worker":"SEAT_OPERATOR_CHANNEL_worker"}`), handOn); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
deadline := time.Now().Add(10 * time.Second)
|
||||
for {
|
||||
mu.Lock()
|
||||
n := len(taken)
|
||||
mu.Unlock()
|
||||
if n >= 1 || time.Now().After(deadline) {
|
||||
break
|
||||
}
|
||||
time.Sleep(50 * time.Millisecond)
|
||||
}
|
||||
time.Sleep(300 * time.Millisecond)
|
||||
mu.Lock()
|
||||
if len(taken) != 1 || !strings.Contains(taken[0], "mesh/work") || !strings.Contains(taken[0], `"x-event-id":"mesh-delivery.d-1"`) {
|
||||
t.Errorf("taken %v", taken)
|
||||
}
|
||||
mu.Unlock()
|
||||
if _, err := asking.Seat("take", json.RawMessage(`{"worker":"SEAT_OPERATOR_CHANNEL_worker"}`), handOn); err == nil {
|
||||
t.Error("an asker took the router's work")
|
||||
}
|
||||
if _, err := asking.Seat("answer", json.RawMessage(`{"subject":"mesh.seat.intake.proof.code.*"}`), handOn); err == nil {
|
||||
t.Error("an asker answered proofs")
|
||||
}
|
||||
if _, err := asking.Seat("anything", json.RawMessage(`{}`), nil); err == nil {
|
||||
t.Error("an unknown verb was answered")
|
||||
}
|
||||
}
|
||||
|
||||
// A lost compare-and-set is answered with its code, so the SDK knows it without reading words.
|
||||
func TestALostCompareAndSetCarriesItsCode(t *testing.T) {
|
||||
mesh := mt.New(t)
|
||||
mesh.Bucket(t, "messenger_asks")
|
||||
m := mt.MembershipOf("messenger", "anchor", false, nil)
|
||||
m.State = []bus.StateIssued{{Name: "asks", Bucket: "messenger_asks", Writes: true}}
|
||||
mesh.Issue(t, m)
|
||||
conn := connect(t, "node-tools", "anchor")
|
||||
conn.Follow("messenger")
|
||||
mt.Until(t, func() error {
|
||||
if conn.Membership("messenger") == nil {
|
||||
return errors.New("not issued yet")
|
||||
}
|
||||
return nil
|
||||
})
|
||||
b := (&consumers{conn: conn, logf: t.Logf, of: map[string]*moduleEvents{}}).forModule("messenger")
|
||||
if _, err := b.State("create", json.RawMessage(`{"state":"asks","key":"a","value":{}}`)); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
_, err := b.State("create", json.RawMessage(`{"state":"asks","key":"a","value":{}}`))
|
||||
var coded interface{ ErrorCode() int }
|
||||
if !errors.As(err, &coded) || coded.ErrorCode() != -32010 || !errors.Is(err, bus.ErrStateChanged) {
|
||||
t.Errorf("a second create: %v", err)
|
||||
}
|
||||
if _, err := b.State("update", json.RawMessage(`{"state":"asks","key":"a","value":{},"revision":0}`)); err == nil {
|
||||
t.Error("an update of revision zero was taken")
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user