Carry a bundle's seat traffic under its own name and kind, and start a module as its own account (hq ADR 0259)

This commit is contained in:
jochen
2026-10-09 10:12:39 +02:00
committed by jschoubben
parent 66de6774e0
commit 6d9906b6c9
9 changed files with 922 additions and 4 deletions
+3
View File
@@ -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.
+341
View File
@@ -0,0 +1,341 @@
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"`
}
// 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
// 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) {
if strings.Contains(subject, ".proof.") {
return 0, 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, err
}
if len(body) == 0 || !json.Valid(body) {
return 0, fmt.Errorf("what is said on %s is JSON", subject)
}
if id == "" {
id = newID()
}
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, err
}
return ack.Sequence, 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 {
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)
}
if err := deliver(Work{Subject: msg.Subject(), Body: json.RawMessage(msg.Data()), Headers: headers}); err != nil {
_ = msg.NakWithDelay(NakDelay)
return
}
_ = msg.Ack()
})
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
}
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) {
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")
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, fmt.Errorf("%s.%s: %w", name, key, ErrStateChanged)
}
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
}
+239
View File
@@ -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()
+39 -1
View File
@@ -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()
@@ -166,6 +181,12 @@ func Start(module, entry string, env []string, mesh Bus, logf func(string, ...an
start := func() (*child, error) {
cmd := exec.Command(entry)
cmd.Env = env
// **A module that names an account of its own runs as it** (novox/hq ADR 0259 §8): its secrets and
// state are that account's, and the operator's account — which every agent runs as — reaches
// neither. Never root, never the operator's.
if err := runAs(cmd, env); err != nil {
return nil, err
}
stdin, err := cmd.StdinPipe()
if err != nil {
return nil, err
@@ -237,7 +258,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":
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
@@ -329,6 +364,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
+68
View File
@@ -0,0 +1,68 @@
package launch
import (
"fmt"
"os/exec"
"os/user"
"strconv"
"strings"
"syscall"
)
// RunAs is the word in a bundle's environment naming the account it runs as (novox/hq ADR 0259 §8), and
// OperatorAccount the word naming the operator's account on this machine, which the runtime is given.
const (
RunAs = "MESH_RUN_AS"
OperatorAccount = "MESH_OPERATOR_ACCOUNT"
)
// lookupUser is user.Lookup, a variable so a test can name accounts this machine does not have.
var lookupUser = user.Lookup
func word(env []string, name string) string {
for _, kv := range env {
if k, v, ok := strings.Cut(kv, "="); ok && k == name {
return v
}
}
return ""
}
// runAs starts cmd as the account its environment names, with that account's home, or leaves it as the
// runtime's own when none is named. It refuses root and the operator's account: a module asks for an
// account of its own to keep what it holds from the agents, and every agent runs as the operator.
func runAs(cmd *exec.Cmd, env []string) error {
name := word(env, RunAs)
if name == "" {
return nil
}
if operator := word(env, OperatorAccount); operator != "" && name == operator {
return fmt.Errorf("%s names the operator's account %s; a module runs as an account of its own, never "+
"the one every agent runs as (novox/hq ADR 0259)", RunAs, name)
}
u, err := lookupUser(name)
if err != nil {
return fmt.Errorf("%s names the account %s, which this machine does not have: the module makes it with "+
"a user resource (novox/hq ADR 0259): %w", RunAs, name, err)
}
uid, err := strconv.ParseUint(u.Uid, 10, 32)
if err != nil {
return fmt.Errorf("the account %s has no usable uid %q", name, u.Uid)
}
gid, err := strconv.ParseUint(u.Gid, 10, 32)
if err != nil {
return fmt.Errorf("the account %s has no usable gid %q", name, u.Gid)
}
if uid == 0 || name == "root" {
return fmt.Errorf("%s names root; a module that asks for an account of its own is not given root", RunAs)
}
cmd.SysProcAttr = &syscall.SysProcAttr{Credential: &syscall.Credential{Uid: uint32(uid), Gid: uint32(gid)}}
var kept []string
for _, kv := range cmd.Env {
if !strings.HasPrefix(kv, "HOME=") && !strings.HasPrefix(kv, "USER=") && !strings.HasPrefix(kv, "LOGNAME=") {
kept = append(kept, kv)
}
}
cmd.Env = append(kept, "HOME="+u.HomeDir, "USER="+name, "LOGNAME="+name)
return nil
}
+48
View File
@@ -0,0 +1,48 @@
package launch
import (
"errors"
"os/exec"
"os/user"
"strings"
"testing"
)
// novox/hq ADR 0259 §8: a module naming an account of its own runs as it, never as root or the operator.
func TestABundleRunsAsTheAccountItNamesAndNeverRootOrTheOperator(t *testing.T) {
was := lookupUser
t.Cleanup(func() { lookupUser = was })
lookupUser = func(name string) (*user.User, error) {
switch name {
case "telegram":
return &user.User{Username: "telegram", Uid: "961", Gid: "961", HomeDir: "/var/lib/telegram"}, nil
case "root":
return &user.User{Username: "root", Uid: "0", Gid: "0", HomeDir: "/root"}, nil
case "toor":
return &user.User{Username: "toor", Uid: "0", Gid: "0", HomeDir: "/root"}, nil
}
return nil, errors.New("unknown user")
}
cmd := exec.Command("/bin/true")
cmd.Env = []string{"HOME=/root", "PATH=/usr/bin"}
if err := runAs(cmd, []string{RunAs + "=telegram", OperatorAccount + "=jo"}); err != nil {
t.Fatal(err)
}
if c := cmd.SysProcAttr.Credential; c == nil || c.Uid != 961 || c.Gid != 961 {
t.Fatalf("not started as telegram: %+v", cmd.SysProcAttr)
}
if env := strings.Join(cmd.Env, " "); !strings.Contains(env, "HOME=/var/lib/telegram") || strings.Contains(env, "HOME=/root") {
t.Fatalf("the account's home is not its own: %s", env)
}
for name, want := range map[string]string{"jo": "operator's account", "root": "names root", "toor": "names root",
"nobody-here": "does not have"} {
err := runAs(exec.Command("/bin/true"), []string{RunAs + "=" + name, OperatorAccount + "=jo"})
if err == nil || !strings.Contains(err.Error(), want) {
t.Errorf("%s: %v, want a refusal saying %q", name, err, want)
}
}
plain := exec.Command("/bin/true")
if err := runAs(plain, nil); err != nil || plain.SysProcAttr != nil {
t.Error("a bundle naming no account was changed")
}
}
+71
View File
@@ -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)
}
}
+110 -3
View File
@@ -458,6 +458,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 +477,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 +610,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 +652,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 +687,89 @@ 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, err := conn.SeatPublish(b.module, asked.Subject, asked.Body, asked.ID)
if err != nil {
return nil, err
}
return json.Marshal(map[string]any{"sequence": seq})
case "prove":
return conn.SeatProve(b.module, asked.Subject, asked.Body)
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
}