diff --git a/node-tools/internal/bus/bus.go b/node-tools/internal/bus/bus.go index dd39bcd..5e0a0a8 100644 --- a/node-tools/internal/bus/bus.go +++ b/node-tools/internal/bus/bus.go @@ -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. diff --git a/node-tools/internal/bus/seat.go b/node-tools/internal/bus/seat.go new file mode 100644 index 0000000..7d6704f --- /dev/null +++ b/node-tools/internal/bus/seat.go @@ -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": ""}`, 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 `.`, 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 +} diff --git a/node-tools/internal/bus/seat_test.go b/node-tools/internal/bus/seat_test.go new file mode 100644 index 0000000..6da91d3 --- /dev/null +++ b/node-tools/internal/bus/seat_test.go @@ -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) + } +} diff --git a/node-tools/internal/launch/event_answer_test.go b/node-tools/internal/launch/event_answer_test.go index a102b0f..2046ec0 100644 --- a/node-tools/internal/launch/event_answer_test.go +++ b/node-tools/internal/launch/event_answer_test.go @@ -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() diff --git a/node-tools/internal/launch/launch.go b/node-tools/internal/launch/launch.go index a691111..67f56d8 100644 --- a/node-tools/internal/launch/launch.go +++ b/node-tools/internal/launch/launch.go @@ -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 diff --git a/node-tools/internal/launch/runas.go b/node-tools/internal/launch/runas.go new file mode 100644 index 0000000..1f9725f --- /dev/null +++ b/node-tools/internal/launch/runas.go @@ -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 +} diff --git a/node-tools/internal/launch/runas_test.go b/node-tools/internal/launch/runas_test.go new file mode 100644 index 0000000..69629c9 --- /dev/null +++ b/node-tools/internal/launch/runas_test.go @@ -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") + } +} diff --git a/node-tools/internal/launch/seat_test.go b/node-tools/internal/launch/seat_test.go new file mode 100644 index 0000000..acd45a9 --- /dev/null +++ b/node-tools/internal/launch/seat_test.go @@ -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) + } +} diff --git a/node-tools/internal/runtime/runtime.go b/node-tools/internal/runtime/runtime.go index d778471..be9919f 100644 --- a/node-tools/internal/runtime/runtime.go +++ b/node-tools/internal/runtime/runtime.go @@ -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 +}