Files
mesh-controller/internal/broker/seattraffic.go
T

267 lines
9.8 KiB
Go

package broker
import "sort"
// What a module may say and take on the seats it holds, uses and watches, beyond tools (novox/hq ADR
// 0259 §3, to-be 46 §10).
//
// Three rules are added to the ones a seat always had, each a subject whose last token names who may
// publish it, granted to that publisher alone — the way a node seat's event about a machine carries the
// machine (ADR 0219):
//
// - **A verb named by its caller** (`by-caller`). A user of the seat submits that accept, and hears that
// event, under its own module's name and no other: `accept.ask.<module>`, `event.decided.<module>`. The
// holder takes every caller's accept and says the event to any caller. So an ask's asker is a fact the
// server enforces, and a warrant reaches only the asker it is for.
// - **A kinded bench** (ADR 0234 §2). A holder claims one kind and takes its own kind's accepts, says its
// own kind's events and asks its own kind's proofs, and nothing of another kind; a user submits to any
// kind. Each kind has its own worker on the seat's queue, so a holder that is away keeps its work and
// holds up no other kind.
// - **A proof** (`proofs`): core request and reply on `mesh.seat.<seat>.proof.<verb>.<kind>`. No stream's
// subjects cover it, so what travels there — a code typed by the operator — is never persisted. A kinded
// holder asks with its own kind; the modules that watch the seat answer.
//
// And one read: **a holder's records**, a bucket the seat names, read by each user under its own name only
// (`$KV.<bucket>.<module>.>`), so an asker reads the state of its own asks and no other asker's.
// Worker is one durable consumer a holder pulls a seat's work from.
type Worker struct {
Stream string
Consumer string
Filter string
}
// SeatTraffic is one module's seat traffic beyond tools. Publish and Subscribe are subject patterns, in the
// server's wildcards; the runtime that carries the module checks a bundle's request against them, since
// the runtime's own principal holds the union of every module it carries.
type SeatTraffic struct {
// Publish is what it submits (accepts of seats it uses), says (events of seats it holds) and asks
// (proofs of seats it holds a kind of).
Publish []string `json:"publish,omitempty"`
// Subscribe is what it hears: events of seats it uses that are named by caller, events of seats it
// watches, and accepts of seats it holds.
Subscribe []string `json:"subscribe,omitempty"`
// Answers is the proof subjects it answers, as a watcher of a kinded seat.
Answers []string `json:"answers,omitempty"`
// Workers are the work queues it takes from, as a holder.
Workers []Worker `json:"workers,omitempty"`
// Records is the direct-get subjects of the records it reads under its own name.
Records []string `json:"records,omitempty"`
// Kinds are the holders of every kinded bench it uses or watches, with what each promises: the one
// account of which channel is which, and what it can carry, that the router judges an answer by. The
// controller's, from the claims, never a channel's word (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"`
}
// KindedBenches are the seats that may be kinded (ADR 0234 §2): making another is a decision, recorded.
var KindedBenches = map[string]bool{"channel": true, "intake": true}
// WorkerName is the worker a seat's holders pull from: one for the seat, or one per kind on a kinded bench.
func WorkerName(seat, kind string) string {
if kind == "" {
return "SEAT_" + upperSnake(seat) + "_worker"
}
return "SEAT_" + upperSnake(seat) + "_" + upperSnake(kind) + "_worker"
}
func namesVerb(list []string, s string) bool {
for _, x := range list {
if x == s {
return true
}
}
return false
}
// SeatTrafficOf derives one module's seat traffic from the seats it holds, uses and watches. Only seats
// carrying one of the rules above are read: every other seat is composed as it always was.
func SeatTrafficOf(module string, holds, uses, watches []Seat) SeatTraffic {
var t SeatTraffic
for _, s := range holds {
if !s.isNewTraffic() {
continue
}
kind := ""
if s.Kinded {
kind = s.Kind
if kind == "" || !safeSubject.MatchString(kind) {
// A kinded claim without a usable kind is refused at registration; here it is granted
// nothing, which is the same answer at the last place it could be asked.
continue
}
}
if len(s.Accepts) > 0 {
stream := seatStreamName(s.Name)
filter := "mesh.seat." + s.Name + ".accept.>"
if kind != "" {
filter = "mesh.seat." + s.Name + ".accept.*." + kind
}
t.Workers = append(t.Workers, Worker{Stream: stream, Consumer: WorkerName(s.Name, kind), Filter: filter})
}
for _, a := range s.Accepts {
switch {
case kind != "":
t.Subscribe = append(t.Subscribe, seatSubject(s, "accept", a+"."+kind))
case namesVerb(s.ByCaller, a):
t.Subscribe = append(t.Subscribe, seatSubject(s, "accept", a+".*"))
}
}
for _, e := range s.Emits {
switch {
case kind != "":
t.Publish = append(t.Publish, seatSubject(s, "event", e+"."+kind))
case namesVerb(s.ByCaller, e):
t.Publish = append(t.Publish, seatSubject(s, "event", e+".*"))
}
}
if kind != "" {
for _, v := range s.Proofs {
t.Publish = append(t.Publish, seatSubject(s, "proof", v+"."+kind))
}
}
}
for _, s := range uses {
if !s.isNewTraffic() {
continue
}
for _, a := range s.Accepts {
switch {
case namesVerb(s.ByCaller, a):
t.Publish = append(t.Publish, seatSubject(s, "accept", a+"."+module))
case s.Kinded:
t.Publish = append(t.Publish, seatSubject(s, "accept", a+".*"))
}
}
for _, e := range s.Emits {
if namesVerb(s.ByCaller, e) {
t.Subscribe = append(t.Subscribe, seatSubject(s, "event", e+"."+module))
}
}
for _, b := range s.Records {
if !safeSubject.MatchString(b) {
continue
}
t.Records = append(t.Records, "$JS.API.DIRECT.GET.KV_"+b+".$KV."+b+"."+module+".>")
}
}
for _, w := range watches {
if w.Kinded {
for _, e := range w.Emits {
t.Subscribe = append(t.Subscribe, seatSubject(w, "event", e+".*"))
}
for _, v := range w.Proofs {
t.Answers = append(t.Answers, seatSubject(w, "proof", v+".*"))
}
}
}
t.Publish = unique(t.Publish)
t.Subscribe = unique(t.Subscribe)
t.Answers = unique(t.Answers)
t.Records = unique(t.Records)
sort.Slice(t.Workers, func(i, j int) bool { return t.Workers[i].Consumer < t.Workers[j].Consumer })
return t
}
// grants is the bus permissions seat traffic needs: the subjects themselves, and the JetStream API a
// worker is pulled and acknowledged through and a record is read through.
func (t SeatTraffic) grants() (pub, sub []string) {
pub = append(pub, t.Publish...)
sub = append(sub, t.Subscribe...)
sub = append(sub, t.Answers...)
for _, w := range t.Workers {
pub = append(pub,
"$JS.API.CONSUMER.INFO."+w.Stream+"."+w.Consumer,
"$JS.API.CONSUMER.MSG.NEXT."+w.Stream+"."+w.Consumer,
"$JS.ACK."+w.Stream+"."+w.Consumer+".>")
}
pub = append(pub, t.Records...)
return pub, sub
}
// isNewTraffic says whether a seat carries any of the rules above, so a seat that carries none is
// composed exactly as before them.
func (s Seat) isNewTraffic() bool {
return s.Kinded || len(s.ByCaller) > 0 || len(s.Proofs) > 0 || len(s.Records) > 0
}
// plainVerbs is a seat's accepts or emits with those the rules above compose taken out: a verb named by
// its caller and every verb of a kinded bench are composed by SeatTrafficOf and nowhere else.
func plainVerbs(s Seat, verbs []string) []string {
if s.Kinded {
return nil
}
var out []string
for _, v := range verbs {
if !namesVerb(s.ByCaller, v) {
out = append(out, v)
}
}
return out
}
// SeatTrafficObjects is the work queues and workers the seat traffic of every composed user implies
// (novox/hq ADR 0259 §3): a queue for each seat a holder takes work from, and each holder's worker on it —
// one per kind on a kinded bench, filtered to that kind, so the kinds never take each other's work. Only
// seats carrying the rules above; the mesh's own seats' queues are RaiseSeats'.
func SeatTrafficObjects(users []Principal) ([]Stream, []Consumer) {
streams := map[string]Stream{}
consumers := map[string]Consumer{}
add := func(module string, holds []Seat) {
for _, w := range SeatTrafficOf(module, holds, nil, nil).Workers {
seat := ""
for _, s := range holds {
if seatStreamName(s.Name) == w.Stream {
seat = s.Name
}
}
streams[w.Stream] = Stream{
Name: w.Stream,
Subjects: []string{"mesh.seat." + seat + ".accept.>"},
Retention: RetentionWorkQueue,
MaxAge: 7 * 24 * 60 * 60,
Why: "work submitted to the " + seat + " seat; its holders take it, each kind its own, " +
"and it queues while nobody does",
}
consumers[w.Consumer] = Consumer{
Name: w.Consumer,
Stream: w.Stream,
Filters: []string{w.Filter},
AckWaitSeconds: 60,
MaxDeliver: 5,
Why: module + " holds " + seat + "; it pulls one ask at a time and acknowledges once it has " +
"recorded it, so a crash redelivers rather than loses",
}
}
}
for _, p := range users {
switch p.Kind {
case KindModule:
add(p.Module, p.Holds)
case KindNodeTools:
for _, d := range p.Carries {
add(d.Module, d.Holds)
}
}
}
var ss []Stream
for _, s := range streams {
ss = append(ss, s)
}
sort.Slice(ss, func(i, j int) bool { return ss[i].Name < ss[j].Name })
var cs []Consumer
for _, c := range consumers {
cs = append(cs, c)
}
sort.Slice(cs, func(i, j int) bool { return cs[i].Name < cs[j].Name })
return ss, cs
}