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 } // 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. func stateGrants(module string, keeps []string, reads []string) []string { var out []string read := func(bucket string) { stream := "KV_" + bucket out = append(out, "$JS.API.STREAM.INFO."+stream, "$JS.API.DIRECT.GET."+stream+".>", "$JS.API.CONSUMER.CREATE."+stream+".>", "$JS.API.CONSUMER.DELETE."+stream+".>", "$JS.FC."+stream+".>") } for _, name := range keeps { if !safeSubject.MatchString(name) { continue } bucket := BucketName(module, name) read(bucket) out = append(out, "$KV."+bucket+".>") } for _, r := range reads { if bucket, ok := bucketOfRead(r); ok { read(bucket) } } 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). type StateIssued struct { Name string `json:"name"` Bucket string `json:"bucket"` Writes bool `json:"writes,omitempty"` } // stateIssuedFor is every bucket a module's code may reach, as its membership lists them. func stateIssuedFor(d Declared) []StateIssued { var out []StateIssued for _, b := range d.State { out = append(out, StateIssued{Name: b.Name, Bucket: BucketName(d.Module, b.Name), Writes: true}) } for _, r := range d.Reads { if bucket, ok := bucketOfRead(r); ok { out = append(out, StateIssued{Name: r, Bucket: bucket}) } } 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 { if !declared[n] { undeclared = append(undeclared, n) } } sort.Strings(undeclared) return undeclared, nil }