package broker import ( "fmt" "sort" "strings" ) // A module's state on the bus (novox/hq ADR 0201, design 32 §4, design 25 §3). // // A module names the state it keeps (`state`) and the state of others it reads (`reads`), and each // is a key-value bucket: the server's own last-per-subject stream with direct reads, delete markers // and watches, which is the state relationship the mesh already uses for declarations, opened to // modules. The controller creates every bucket from the catalogue — from registration, like a // seat's stream, so a reader may watch before the owner runs anywhere — and no module can. // // Pure, like everything else in this package that decides what the bus holds; jetstream.go is the // part that asks a server. // The mesh's caps on a bucket, the same for every module: a value is a piece of state, not a file, // and a bucket that grew without bound would be one module filling the bus's disk for everyone. const ( StateMaxValueBytes = 256 * 1024 StateMaxBytes = 64 * 1024 * 1024 ) // A Bucket is one module's declared state as the bus holds it. type Bucket struct { Module string Name string // History is how many values a key keeps; zero is one. History int // TTLSeconds is how long a value lives; zero is until replaced or deleted. TTLSeconds int // PerMachine says each machine's instance writes only the key named for its machine (novox/hq ADR // 0260), and is granted that key and no other. PerMachine bool } // BucketName is the bucket a module's state lives in: the module and the local name joined by an // underscore, which neither may contain, so two modules can never derive one bucket. func BucketName(module, name string) string { return module + "_" + name } // Bucket is this bucket's name on the bus. func (b Bucket) Bucket() string { return BucketName(b.Module, b.Name) } // Why is carried into the server's description of the bucket, so somebody reading the server's // own state finds whose it is and why it is kept. func (b Bucket) Why() string { return fmt.Sprintf("%s's state %q (novox/hq ADR 0201): its current value per key, written by %s, "+ "read by whatever declares it reads it; kept when %s is unassigned, because it is data", b.Module, b.Name, b.Module, b.Module) } // bucketOfRead is the bucket a read names, `.`, or false when it names none. func bucketOfRead(read string) (string, bool) { at := strings.LastIndex(read, ".") if at <= 0 || at == len(read)-1 { return "", false } module, name := read[:at], read[at+1:] if !safeSubject.MatchString(module) || !safeSubject.MatchString(name) { return "", false } return BucketName(module, name), true } // stateGrants is what a principal publishes to reach the state its modules keep and read: for every // bucket, binding to it, reading a key directly, and an ordered consumer for listing and watching, // created and deleted on the bucket's own stream, with its flow control answered; for a bucket an // owner keeps, writing under the bucket's own subjects too. // // **Measured against a running server, 2026-10-04** (novox/hq research 024), and each one is there // because leaving it out failed: without STREAM.INFO nothing binds; without DIRECT.GET nothing is // read; without CONSUMER.CREATE no key is listed and nothing is watched; without CONSUMER.DELETE a // watch cannot be stopped and lingers on the server. A write outside these is refused by the server // — and reaches the writer as a timeout, not a refusal, which is why the runtime refuses first. // // **One key, where the state is per machine** (novox/hq ADR 0260). The owner's instance on a machine // writes and reads the key named for that machine, and a holder granted a read for its machine reads // that key: a direct get of the key's own subject, and a consumer whose filter is that subject — the // server's create-with-filter form, which a client uses for a watch of one key and which the server // holds to the filter it names. What cannot be narrowed by key: binding (STREAM.INFO), deleting a // consumer (named at random by the client) and flow control; none of them reads or writes a value. func stateGrants(a stateAccess) []string { var out []string read := func(bucket, key string) { stream := "KV_" + bucket out = append(out, "$JS.API.STREAM.INFO."+stream) if key == "" { out = append(out, "$JS.API.DIRECT.GET."+stream+".>", "$JS.API.CONSUMER.CREATE."+stream+".>") } else { out = append(out, "$JS.API.DIRECT.GET."+stream+".$KV."+bucket+"."+key, "$JS.API.CONSUMER.CREATE."+stream+".*.$KV."+bucket+"."+key) } out = append(out, "$JS.API.CONSUMER.DELETE."+stream+".>", "$JS.FC."+stream+".>") } machineKey := "" if safeSubject.MatchString(a.Node) { machineKey = a.Node } perMachine := map[string]bool{} for _, n := range a.PerMachine { perMachine[n] = true } for _, name := range a.Keeps { if !safeSubject.MatchString(name) { continue } bucket := BucketName(a.Module, name) if perMachine[name] { if machineKey == "" { continue // a machine with no usable name reaches none of it } read(bucket, machineKey) out = append(out, "$KV."+bucket+"."+machineKey) continue } read(bucket, "") out = append(out, "$KV."+bucket+".>") } for _, r := range a.Reads { if bucket, ok := bucketOfRead(r); ok { read(bucket, "") } } for _, k := range a.KeyedReads { if bucket, ok := bucketOfRead(k.Read); ok && safeSubject.MatchString(k.Key) { read(bucket, k.Key) } } return out } // KeyedRead is one key of another module's state a principal reads, and no other key of it. type KeyedRead struct { Read string `json:"read"` Key string `json:"key"` } // stateAccess is what one principal may do with state: the module's own buckets (some per machine), // what it reads whole, and what it reads for its machine's key alone. type stateAccess struct { Module, Node string Keeps, PerMachine []string Reads []string KeyedReads []KeyedRead } // perMachineNames is the local names of a module's buckets keyed by machine. func perMachineNames(buckets []Bucket) []string { var out []string for _, b := range buckets { if b.PerMachine { out = append(out, b.Name) } } return out } // StateIssued is one bucket an assignment may reach, by the name its module uses for it: its own // state by the local name, another's as `.` (novox/hq ADR 0201). Key, when set, is the // one key it may reach (novox/hq ADR 0260): this machine's. type StateIssued struct { Name string `json:"name"` Bucket string `json:"bucket"` Writes bool `json:"writes,omitempty"` Key string `json:"key,omitempty"` } // stateIssuedFor is every bucket a module's code may reach on a machine, as its membership lists them. func stateIssuedFor(d Declared, node string) []StateIssued { var out []StateIssued for _, b := range d.State { issued := StateIssued{Name: b.Name, Bucket: BucketName(d.Module, b.Name), Writes: true} if b.PerMachine { issued.Key = node } out = append(out, issued) } for _, r := range d.Reads { if bucket, ok := bucketOfRead(r); ok { out = append(out, StateIssued{Name: r, Bucket: bucket}) } } for _, k := range d.KeyedReads { if bucket, ok := bucketOfRead(k.Read); ok { out = append(out, StateIssued{Name: k.Read, Bucket: bucket, Key: k.Key}) } } return out } // stateNames is the local names of a module's own buckets. func stateNames(buckets []Bucket) []string { out := make([]string, 0, len(buckets)) for _, b := range buckets { out = append(out, b.Name) } return out } // A BucketAsserter is the part of a JetStream connection bucket assertion needs. type BucketAsserter interface { // EnsureBucket creates the bucket if absent and brings its options to match if present, never // discarding what it holds. EnsureBucket(b Bucket) error // BucketNames is every key-value bucket on the server. BucketNames() ([]string, error) } // RaiseBuckets asserts every declared bucket and answers the buckets on the server that nothing // declares any more. // // **Those are reported, never removed** (novox/hq ADR 0201, ADR 0030): what a module stored is // data, and a manifest edited, a module renamed or a catalogue entry dropped is an ordinary day's // work that must not take data with it. Removing one is a person's act. func RaiseBuckets(a BucketAsserter, buckets []Bucket) (undeclared []string, err error) { sorted := append([]Bucket(nil), buckets...) sort.Slice(sorted, func(i, j int) bool { return sorted[i].Bucket() < sorted[j].Bucket() }) declared := map[string]bool{} for _, b := range sorted { if err := a.EnsureBucket(b); err != nil { return nil, fmt.Errorf("asserting %s's state %q: %w", b.Module, b.Name, err) } declared[b.Bucket()] = true } names, err := a.BucketNames() if err != nil { return nil, fmt.Errorf("listing the bus's state: %w", err) } for _, n := range names { // A seat's cancelled set is the mesh's own (novox/hq ADR 0219), not a module's state; so are // the controller's own buckets (novox/hq to-be 45 §1). if !declared[n] && !IsCancelledSet(n) && !IsControllerBucket(n) { undeclared = append(undeclared, n) } } sort.Strings(undeclared) return undeclared, nil }