diff --git a/cmd/mesh-controller/push.go b/cmd/mesh-controller/push.go index 3f5f73f..13c587e 100644 --- a/cmd/mesh-controller/push.go +++ b/cmd/mesh-controller/push.go @@ -898,6 +898,21 @@ func raiseTheBus(ctx context.Context, inv *inventory.Inventory, address string) if err := broker.RaiseSeats(js, inventory.MeshSeats(), holders); err != nil { return err } + // Every module's state (novox/hq ADR 0201), from the catalogue: a bucket exists from + // registration, so a module reading one may watch it before its owner runs anywhere. One that + // nothing declares any more is said and kept — what it holds is data. + buckets, err := inv.DeclaredBuckets(ctx) + if err != nil { + return err + } + undeclared, err := broker.RaiseBuckets(js, buckets) + if err != nil { + return err + } + if len(undeclared) > 0 { + fmt.Printf("the bus holds state nothing declares any more, kept because it is data: %s — "+ + "removing it is a person's act\n", strings.Join(undeclared, ", ")) + } // And how every module hears what it consumes. Derived from the same records the user list is // composed from, so a module the mesh grants a consumer's subjects has that consumer waiting. // Done on every raise, not only when a credential is issued: every module moved onto this bus @@ -921,8 +936,8 @@ func raiseTheBus(ctx context.Context, inv *inventory.Inventory, address string) } hearing++ } - fmt.Printf("the bus at %s has its streams, %d machine(s) can hear a declaration, and %d module(s) "+ - "can hear what they consume\n", broker.BareAddress(address), len(names), hearing) + fmt.Printf("the bus at %s has its streams, %d machine(s) can hear a declaration, %d module(s) "+ + "can hear what they consume, and %d bucket(s) of state\n", broker.BareAddress(address), len(names), hearing, len(buckets)) return nil } diff --git a/internal/broker/jetstream.go b/internal/broker/jetstream.go index cf86519..20b5479 100644 --- a/internal/broker/jetstream.go +++ b/internal/broker/jetstream.go @@ -1,6 +1,7 @@ package broker import ( + "context" "crypto/sha256" "crypto/tls" "crypto/x509" @@ -12,6 +13,7 @@ import ( "time" "github.com/nats-io/nats.go" + "github.com/nats-io/nats.go/jetstream" ) // The JetStream side of the controller: the one place the mesh's streams and consumers are @@ -316,3 +318,50 @@ func retentionOf(r Retention) nats.RetentionPolicy { return nats.LimitsPolicy } } + +// EnsureBucket creates a module's bucket if it is absent and brings its options to match if it is +// present (novox/hq ADR 0201). +// +// **An update, never a delete and recreate**, for the reason a stream is updated: recreating +// discards what the bucket holds, and what a module's state holds is data. The mesh's caps are +// asserted with the owner's options, so a bucket made by hand converges to them. +func (j *JetStream) EnsureBucket(b Bucket) error { + history := b.History + if history == 0 { + history = 1 + } + js, err := jetstream.New(j.conn) + if err != nil { + return err + } + ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + if _, err := js.CreateOrUpdateKeyValue(ctx, jetstream.KeyValueConfig{ + Bucket: b.Bucket(), + Description: b.Why(), + History: uint8(history), + TTL: time.Duration(b.TTLSeconds) * time.Second, + MaxValueSize: StateMaxValueBytes, + MaxBytes: StateMaxBytes, + Storage: jetstream.FileStorage, + }); err != nil { + return fmt.Errorf("asserting bucket %s: %w", b.Bucket(), err) + } + return nil +} + +// BucketNames is every key-value bucket on the server, the mesh's and anybody else's. +func (j *JetStream) BucketNames() ([]string, error) { + js, err := jetstream.New(j.conn) + if err != nil { + return nil, err + } + ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + lister := js.KeyValueStoreNames(ctx) + var out []string + for name := range lister.Name() { + out = append(out, name) + } + return out, lister.Error() +} diff --git a/internal/broker/membership.go b/internal/broker/membership.go index 8cf9b98..2fca555 100644 --- a/internal/broker/membership.go +++ b/internal/broker/membership.go @@ -45,6 +45,11 @@ type Membership struct { // module that must tell the mesh from the world, the route proxy serving an internal name, reads // it here rather than keeping a definition of its own. Mesh []string `json:"mesh,omitempty"` + // State is every bucket this module's code may reach, by the name it uses for each, and whether + // it may write it (novox/hq ADR 0201): the runtime answers a bundle's state verbs from this list + // and refuses, with the reason, what is not on it — the bus enforces only the union over every + // module on the machine. + State []StateIssued `json:"state,omitempty"` } // Served is one address a tool is answered on. @@ -103,6 +108,7 @@ func MembershipFor(node string, d Declared, where Placements) Membership { m.Seats = append(m.Seats, SeatServed{Seat: s.Name, Verb: verb, Subject: seatToolSubject(s, verb, node)}) } } + m.State = stateIssuedFor(d) if len(d.Invokes) > 0 { m.Reaches = map[string][]string{} for _, t := range d.Invokes { diff --git a/internal/broker/nats.go b/internal/broker/nats.go index 86ecace..556f8e3 100644 --- a/internal/broker/nats.go +++ b/internal/broker/nats.go @@ -103,6 +103,12 @@ type Principal struct { // permission and nothing beside it. Invokes []string + // State is the local names of the state this principal's module keeps, and Reads the state of + // others it reads as `.` (novox/hq ADR 0201): a bucket each, kept by the owner's + // instances and read by whoever declares it. + State []string + Reads []string + // PasswordHash is the bcrypt hash the mesh minted. The plaintext is sealed to the principal // and never appears here: this file is written to a node's disk and read by a server, and a // secret that can be read from a configuration file is a secret with a wider blast radius @@ -421,6 +427,10 @@ func PermissionsFor(p Principal) (Permissions, error) { } } + // 5. Its state, and the state of others it reads (novox/hq ADR 0201): every one read and + // watched, its own written too. + pub = append(pub, stateGrants(p.Module, p.State, p.Reads)...) + case KindNodeTools: // **One process serves what every module on the machine would have served for itself** // (novox/hq ADR 0175). Each carried module's whole tool namespace — the same grant that @@ -487,6 +497,13 @@ func PermissionsFor(p Principal) (Permissions, error) { "$JS.API.CONSUMER.MSG.NEXT."+stream+"."+durable, "$JS.ACK."+stream+"."+durable+".>") } + // **And it keeps and reads state for the modules it carries** (novox/hq ADR 0201): the union + // of what each may do with a bucket — an owner's write, a reader's read. That one module's code + // does not write another's bucket through it is the runtime's to keep, from the membership + // each assignment is issued, as it keeps each module's events under that module's own name. + for _, d := range p.Carries { + pub = append(pub, stateGrants(d.Module, stateNames(d.State), d.Reads)...) + } sub = unique(sub) pub = unique(pub) } diff --git a/internal/broker/state.go b/internal/broker/state.go new file mode 100644 index 0000000..0324543 --- /dev/null +++ b/internal/broker/state.go @@ -0,0 +1,169 @@ +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 +} diff --git a/internal/broker/state_test.go b/internal/broker/state_test.go new file mode 100644 index 0000000..ba78ebb --- /dev/null +++ b/internal/broker/state_test.go @@ -0,0 +1,166 @@ +package broker + +import ( + "slices" + "strings" + "testing" + + "github.com/nats-io/nats.go" +) + +// The grants measured against a running server (novox/hq research 024): an owner reads and writes +// its bucket, a reader only reads, and neither reaches any other bucket. +func TestAnOwnerWritesItsStateAndAReaderOnlyReads(t *testing.T) { + owner, err := PermissionsFor(Principal{Kind: KindModule, Node: "one", Module: "claude-code", + State: []string{"servers"}, PasswordHash: "x"}) + if err != nil { + t.Fatal(err) + } + for _, s := range []string{ + "$KV.claude-code_servers.>", + "$JS.API.STREAM.INFO.KV_claude-code_servers", + "$JS.API.DIRECT.GET.KV_claude-code_servers.>", + "$JS.API.CONSUMER.CREATE.KV_claude-code_servers.>", + "$JS.API.CONSUMER.DELETE.KV_claude-code_servers.>", + "$JS.FC.KV_claude-code_servers.>", + } { + has(t, owner.Publish, s) + } + hasNot(t, owner.Publish, "$KV.>") + hasNot(t, owner.Publish, "$JS.API.>") + + reader, err := PermissionsFor(Principal{Kind: KindModule, Node: "one", Module: "console", + Reads: []string{"claude-code.servers"}, PasswordHash: "x"}) + if err != nil { + t.Fatal(err) + } + has(t, reader.Publish, "$JS.API.DIRECT.GET.KV_claude-code_servers.>") + has(t, reader.Publish, "$JS.API.CONSUMER.CREATE.KV_claude-code_servers.>") + hasNot(t, reader.Publish, "$KV.claude-code_servers.>") + for _, s := range reader.Subscribe { + if s == "$KV.claude-code_servers.>" { + t.Fatalf("a reader subscribes the bucket's subjects directly: %v", reader.Subscribe) + } + } +} + +// One runtime carries every module on its machine, so its grant is the union: the owner's write +// where an owner is carried, a read where only a reader is. +func TestTheRuntimeKeepsAndReadsStateForItsModules(t *testing.T) { + perms, err := PermissionsFor(Principal{Kind: KindNodeTools, Node: "one", Module: RuntimeModule, + Carries: []Declared{ + {Module: "claude-code", State: []Bucket{{Module: "claude-code", Name: "servers"}}, + Reads: []string{"licence-manager.bindings"}}, + {Module: "audit"}, + }, PasswordHash: "x"}) + if err != nil { + t.Fatal(err) + } + has(t, perms.Publish, "$KV.claude-code_servers.>") + has(t, perms.Publish, "$JS.API.DIRECT.GET.KV_licence-manager_bindings.>") + hasNot(t, perms.Publish, "$KV.licence-manager_bindings.>") +} + +// A module with no state is granted nothing of any bucket — the composition of every module that +// existed before this is unchanged. +func TestAModuleWithNoStateReachesNoBucket(t *testing.T) { + perms, err := PermissionsFor(Principal{Kind: KindModule, Node: "one", Module: "billing", + Emits: []string{"order.placed"}, PasswordHash: "x"}) + if err != nil { + t.Fatal(err) + } + for _, s := range perms.Publish { + if strings.HasPrefix(s, "$KV.") || strings.HasPrefix(s, "$JS.FC.") || strings.Contains(s, ".KV_") { + t.Fatalf("granted %q without declaring state", s) + } + } +} + +// A read that names no bucket grants nothing rather than something that happens to parse. +func TestAReadThatNamesNoBucketGrantsNothing(t *testing.T) { + if got := stateGrants("a", nil, []string{"nodot", "x.", ".y", "a.b>"}); len(got) != 0 { + t.Fatalf("granted %v for reads that name no bucket", got) + } +} + +// The membership lists every bucket the module's code may reach, by the name the module uses for +// it, and whether it may write it — the list the runtime refuses from. +func TestAMembershipListsTheStateItsModuleMayReach(t *testing.T) { + m := MembershipFor("one", Declared{Module: "claude-code", + State: []Bucket{{Module: "claude-code", Name: "servers"}}, + Reads: []string{"licence-manager.bindings"}}, Placements{}) + want := []StateIssued{ + {Name: "servers", Bucket: "claude-code_servers", Writes: true}, + {Name: "licence-manager.bindings", Bucket: "licence-manager_bindings"}, + } + if !slices.Equal(m.State, want) { + t.Fatalf("issued %+v, want %+v", m.State, want) + } + if none := MembershipFor("one", Declared{Module: "audit"}, Placements{}); none.State != nil { + t.Fatalf("a module with no state was issued %+v", none.State) + } +} + +type buckets struct { + ensured []string + on []string +} + +func (b *buckets) EnsureBucket(x Bucket) error { + b.ensured = append(b.ensured, x.Bucket()) + return nil +} +func (b *buckets) BucketNames() ([]string, error) { return b.on, nil } + +// Every declared bucket is asserted; one on the server that nothing declares is said, not removed. +func TestRaisingStateReportsWhatNothingDeclares(t *testing.T) { + b := &buckets{on: []string{"claude-code_servers", "gone_old", "ours_by_hand"}} + undeclared, err := RaiseBuckets(b, []Bucket{{Module: "claude-code", Name: "servers"}, {Module: "a", Name: "b"}}) + if err != nil { + t.Fatal(err) + } + if !slices.Equal(b.ensured, []string{"a_b", "claude-code_servers"}) { + t.Fatalf("asserted %v", b.ensured) + } + if !slices.Equal(undeclared, []string{"gone_old", "ours_by_hand"}) { + t.Fatalf("reported %v", undeclared) + } +} + +// Against a real server: a bucket is created with the owner's options and the mesh's caps, +// asserting it again changes nothing and keeps what it holds, and a changed option is brought to +// match in place. +func TestABucketIsAssertedInPlace(t *testing.T) { + js := aLiveBus(t) + b := Bucket{Module: "statetest", Name: "servers"} + if _, err := RaiseBuckets(js, []Bucket{b}); err != nil { + t.Fatalf("a real server refused a module's bucket: %v", err) + } + kv, err := js.Context().KeyValue(b.Bucket()) + if err != nil { + t.Fatal(err) + } + if _, err := kv.Put("all.one", []byte(`{"kept":true}`)); err != nil { + t.Fatal(err) + } + b.History = 3 + if _, err := RaiseBuckets(js, []Bucket{b}); err != nil { + t.Fatalf("asserting the bucket again failed, so a restart would: %v", err) + } + got, err := kv.Get("all.one") + if err != nil || string(got.Value()) != `{"kept":true}` { + t.Fatalf("asserting again lost what the bucket held: %v %v", got, err) + } + status, err := kv.Status() + if err != nil { + t.Fatal(err) + } + if status.History() != 3 { + t.Fatalf("history is %d, the owner declared 3", status.History()) + } + if s, ok := status.(*nats.KeyValueBucketStatus); ok { + if c := s.StreamInfo().Config; c.MaxMsgSize != StateMaxValueBytes || c.MaxBytes != StateMaxBytes { + t.Fatalf("the mesh's caps are not on the bucket: value %d, bucket %d", c.MaxMsgSize, c.MaxBytes) + } + } +} diff --git a/internal/broker/users.go b/internal/broker/users.go index 1395467..d831b15 100644 --- a/internal/broker/users.go +++ b/internal/broker/users.go @@ -33,6 +33,10 @@ type Declared struct { Watches []Seat // Invokes are the tools it calls, `.` or `*` (novox/hq ADR 0152). Invokes []string + // State is the state it keeps, each a bucket its instances write (novox/hq ADR 0201). + State []Bucket + // Reads are other modules' state it reads, each `.` (novox/hq ADR 0201). + Reads []string } // Records is what composing a user list needs to know about the mesh, and nothing more. @@ -81,6 +85,7 @@ func Users(r Records) ([]Principal, error) { Kind: KindModule, Node: node, Module: d.Module, Emits: d.Emits, Consumes: d.Consumes, Serves: d.Serves, Holds: d.Holds, Uses: d.Uses, Watches: d.Watches, Invokes: d.Invokes, + State: stateNames(d.State), Reads: d.Reads, }) } if runtimeHere { diff --git a/internal/catalogue/manifest.go b/internal/catalogue/manifest.go index 8528c9f..6eb3b39 100644 --- a/internal/catalogue/manifest.go +++ b/internal/catalogue/manifest.go @@ -325,6 +325,15 @@ type Manifest struct { // person's account (design 25 §7) already had the same shape. Invokes []string `json:"invokes,omitempty"` + // State is the current state this module keeps on the bus, by local name: each a key-value + // bucket the controller creates, which every instance of the module writes and reads + // (novox/hq ADR 0201). Not history — that is an event — and never a secret, sealed or not. + State []StateDeclaration `json:"state,omitempty"` + + // Reads are other modules' state this module reads and watches, each `.` + // (novox/hq ADR 0201). Read-only: only the owner's instances write. + Reads []string `json:"reads,omitempty"` + // Capabilities the machine must have. A different field from Requires because the remedy // differs: a missing module can be assigned, and a missing capability means the wrong // machine. @@ -1316,6 +1325,8 @@ func ParseManifest(raw []byte) (Manifest, error) { // module whose event names are wrong installs, starts, connects and reacts to nothing, with // every log line saying it is fine (novox/hq 04-ISSUES/127). problems = append(problems, EventProblems(m)...) + // And what it may call its state, and whose it may read (state.go, novox/hq ADR 0201). + problems = append(problems, StateProblems(m)...) wellFormed := true for _, c := range m.Claims { if !name.MatchString(c.Name) { diff --git a/internal/catalogue/state.go b/internal/catalogue/state.go new file mode 100644 index 0000000..fa49f71 --- /dev/null +++ b/internal/catalogue/state.go @@ -0,0 +1,150 @@ +package catalogue + +import ( + "bytes" + "encoding/json" + "fmt" + "regexp" + "strings" +) + +// What a module may call its state, and whose state it may ask to read (novox/hq ADR 0201). +// +// A module names its state **locally** — `servers`, never a bucket or a subject — and another +// module's as `.`, the way a consumed event names its emitter (design 32 §1). The +// mesh derives the bucket from the two names, so the module and the local name must each be one +// token: the bucket joins them with an underscore, which neither may contain, so two modules can +// never derive one bucket. + +// stateName is one local name of a module's state: lower-case, no dot, no underscore. +var stateName = regexp.MustCompile(`^[a-z0-9][a-z0-9-]*$`) + +// The mesh's caps on what a module may ask of a bucket's history. +const ( + // StateMostHistory is the most past values a key may keep. The server's own limit. + StateMostHistory = 64 +) + +// StateDeclaration is one bucket a module owns: its local name, and the options that are the +// owner's to choose, as a seat chooses how long its backlog survives (design 32 §3). +type StateDeclaration struct { + Name string `json:"name"` + // History is how many values a key keeps, the current one included; zero is one. + History int `json:"history,omitempty"` + // TTLSeconds is how long a value lives once written; zero is until it is replaced or deleted. + TTLSeconds int `json:"ttl-seconds,omitempty"` +} + +// UnmarshalJSON reads a bucket as its bare name, or as {name, history, ttl-seconds}. +func (s *StateDeclaration) UnmarshalJSON(raw []byte) error { + trimmed := bytes.TrimSpace(raw) + if len(trimmed) > 0 && trimmed[0] == '"' { + return json.Unmarshal(trimmed, &s.Name) + } + type plain StateDeclaration + var full plain + dec := json.NewDecoder(bytes.NewReader(trimmed)) + dec.DisallowUnknownFields() + if err := dec.Decode(&full); err != nil { + return fmt.Errorf("a state is either a name or {name, history, ttl-seconds}: %w", err) + } + *s = StateDeclaration(full) + return nil +} + +// MarshalJSON writes back the short form when there is nothing else to say. +func (s StateDeclaration) MarshalJSON() ([]byte, error) { + if s.History == 0 && s.TTLSeconds == 0 { + return json.Marshal(s.Name) + } + type plain StateDeclaration + return json.Marshal(plain(s)) +} + +// ReadState splits a read into the owning module and the local name, or says why it is not one. +func ReadState(read string) (module, local string, err error) { + at := strings.LastIndex(read, ".") + if at <= 0 || at == len(read)-1 { + return "", "", fmt.Errorf("%q does not name a module and its state: a read is .", read) + } + module, local = read[:at], read[at+1:] + if !stateName.MatchString(module) { + return "", "", fmt.Errorf("%q cannot own state: a module whose state is read is one plain name", module) + } + if !stateName.MatchString(local) { + return "", "", fmt.Errorf("%q is not a state name: lower-case letters, digits and hyphens", local) + } + return module, local, nil +} + +// StateProblems is what is wrong with a manifest's state and reads. +// +// Refused at registration, because a bucket name the bus cannot hold is a module that installs, +// starts, and is refused on its first write with a reason about a bucket nobody named. +func StateProblems(m Manifest) []string { + var problems []string + if len(m.State) > 0 && !stateName.MatchString(m.Module) { + problems = append(problems, fmt.Sprintf( + "%s keeps state, and a module's name is part of its buckets' names, which take one plain "+ + "name — no dot (novox/hq ADR 0201)", m.Module)) + } + seen := map[string]bool{} + for _, s := range m.State { + switch { + case !stateName.MatchString(s.Name): + problems = append(problems, fmt.Sprintf( + "%s keeps state %q: a state is named locally — lower-case letters, digits and hyphens, "+ + "no dot and no underscore; the mesh derives the bucket (novox/hq ADR 0201)", m.Module, s.Name)) + case seen[s.Name]: + problems = append(problems, fmt.Sprintf("%s keeps state %q twice", m.Module, s.Name)) + } + seen[s.Name] = true + if s.History < 0 || s.History > StateMostHistory { + problems = append(problems, fmt.Sprintf( + "%s keeps %d values of %q; a key keeps between 1 and %d", m.Module, s.History, s.Name, StateMostHistory)) + } + if s.TTLSeconds < 0 { + problems = append(problems, fmt.Sprintf("%s gives %q a negative lifetime", m.Module, s.Name)) + } + } + for _, r := range m.Reads { + module, _, err := ReadState(r) + if err != nil { + problems = append(problems, fmt.Sprintf("%s reads %v", m.Module, err)) + continue + } + if module == m.Module { + problems = append(problems, fmt.Sprintf( + "%s reads %q, which is its own state: a module reads and writes what it keeps already", m.Module, r)) + } + } + return problems +} + +// StateReadsNothingDeclares is every read across a catalogue whose owner is present and declares no +// such state. An absent owner says nothing — a module may be installed long before the one whose +// state it reads, as a consumer may before its emitter (design 32 §1). +func StateReadsNothingDeclares(manifests []Manifest) []string { + declared := map[string]map[string]bool{} + for _, m := range manifests { + own := map[string]bool{} + for _, s := range m.State { + own[s.Name] = true + } + declared[m.Module] = own + } + var problems []string + for _, m := range manifests { + for _, r := range m.Reads { + module, local, err := ReadState(r) + if err != nil { + continue + } + if own, present := declared[module]; present && !own[local] { + problems = append(problems, fmt.Sprintf( + "%s reads %q, and %s keeps no state called %q", m.Module, r, module, local)) + } + } + } + return problems +} diff --git a/internal/catalogue/state_test.go b/internal/catalogue/state_test.go new file mode 100644 index 0000000..4cf6206 --- /dev/null +++ b/internal/catalogue/state_test.go @@ -0,0 +1,82 @@ +package catalogue + +import ( + "encoding/json" + "strings" + "testing" +) + +// A module declares the state it keeps and the state it reads (novox/hq ADR 0201), a bucket by its +// bare name or with the owner's options. +func TestAManifestMaySayWhatStateItKeepsAndReads(t *testing.T) { + m, err := ParseManifest([]byte(`{"module":"claude-code","version":"1",` + + `"state":["servers",{"name":"seen","history":5,"ttl-seconds":3600}],` + + `"reads":["licence-manager.bindings"]}`)) + if err != nil { + t.Fatal(err) + } + if len(m.State) != 2 || m.State[0].Name != "servers" || m.State[1].History != 5 || m.State[1].TTLSeconds != 3600 { + t.Fatalf("state not read: %+v", m.State) + } + if len(m.Reads) != 1 || m.Reads[0] != "licence-manager.bindings" { + t.Fatalf("reads not read: %v", m.Reads) + } + // Written back as it came in: the short form where nothing else is said. + out, _ := json.Marshal(m.State) + if string(out) != `["servers",{"name":"seen","history":5,"ttl-seconds":3600}]` { + t.Fatalf("written back as %s", out) + } +} + +// A name the bus could not hold, or that would let two modules derive one bucket, is refused at +// registration in the manifest's words. +func TestAStateNameIsLocalAndOneToken(t *testing.T) { + for _, c := range []struct{ manifest, says string }{ + {`{"module":"a","version":"1","state":["mesh.servers"]}`, `keeps state "mesh.servers": a state is named locally`}, + {`{"module":"a","version":"1","state":["my_servers"]}`, `keeps state "my_servers"`}, + {`{"module":"a","version":"1","state":["s","s"]}`, `keeps state "s" twice`}, + {`{"module":"a","version":"1","state":[{"name":"s","history":65}]}`, `a key keeps between 1 and 64`}, + {`{"module":"a.b","version":"1","state":["s"]}`, `no dot`}, + {`{"module":"a","version":"1","reads":["bindings"]}`, `a read is .`}, + {`{"module":"a","version":"1","reads":["a.s"]}`, `which is its own state`}, + {`{"module":"a","version":"1","state":[{"name":"s","shared":true}]}`, `{name, history, ttl-seconds}`}, + } { + _, err := ParseManifest([]byte(c.manifest)) + if err == nil { + t.Errorf("%s was accepted", c.manifest) + continue + } + if !strings.Contains(err.Error(), c.says) { + t.Errorf("%s refused for the wrong reason: %v", c.manifest, err) + } + } +} + +// A read whose owner is present must name a state that owner keeps; an absent owner says nothing, +// because a module may be installed before the one whose state it reads. +func TestAReadNamesStateItsOwnerKeeps(t *testing.T) { + owner := Manifest{Module: "licence-manager", State: []StateDeclaration{{Name: "bindings"}}} + good := Manifest{Module: "claude-code", Reads: []string{"licence-manager.bindings", "absent.anything"}} + bad := Manifest{Module: "other", Reads: []string{"licence-manager.tokens"}} + if p := StateReadsNothingDeclares([]Manifest{owner, good}); len(p) != 0 { + t.Fatalf("a read of declared state was refused: %v", p) + } + p := StateReadsNothingDeclares([]Manifest{owner, bad}) + if len(p) != 1 || !strings.Contains(p[0], `licence-manager keeps no state called "tokens"`) { + t.Fatalf("a read of state nobody keeps was not named: %v", p) + } +} + +// **Across the whole catalogue**: every state name is local, and every read whose owner is present +// names state that owner keeps. +func TestEveryManifestsStateIsLocalAndEveryReadIsKept(t *testing.T) { + manifests := theCatalogue(t) + var problems []string + for _, m := range manifests { + problems = append(problems, StateProblems(m)...) + } + problems = append(problems, StateReadsNothingDeclares(manifests)...) + if len(problems) > 0 { + t.Fatalf("the catalogue's state is not what ADR 0201 says:\n %s", strings.Join(problems, "\n ")) + } +} diff --git a/internal/inventory/busrecords.go b/internal/inventory/busrecords.go index 431e5cc..88fccf1 100644 --- a/internal/inventory/busrecords.go +++ b/internal/inventory/busrecords.go @@ -132,6 +132,9 @@ func declaredFor(m catalogue.Manifest, seats map[string]catalogue.SeatDeclaratio Serves: m.Tools, // And what it calls (novox/hq ADR 0152) — the console's `*`, nothing else's. Invokes: m.Invokes, + // And the state it keeps and reads (novox/hq ADR 0201). + State: bucketsOf(m), + Reads: m.Reads, } for _, c := range m.Claims { // Every seat with a protocol, the mesh's own included. One that says only who does a job is @@ -148,6 +151,30 @@ func declaredFor(m catalogue.Manifest, seats map[string]catalogue.SeatDeclaratio return d } +// bucketsOf is the state a module keeps, as the bus holds it. +func bucketsOf(m catalogue.Manifest) []broker.Bucket { + var out []broker.Bucket + for _, s := range m.State { + out = append(out, broker.Bucket{Module: m.Module, Name: s.Name, History: s.History, TTLSeconds: s.TTLSeconds}) + } + return out +} + +// DeclaredBuckets is every bucket the catalogue declares, registered modules assigned or not: a +// bucket exists from registration, like a seat's stream, so a module reading it may watch before its +// owner runs anywhere (novox/hq ADR 0201). +func (i *Inventory) DeclaredBuckets(ctx context.Context) ([]broker.Bucket, error) { + declared, err := i.Catalogue(ctx) + if err != nil { + return nil, fmt.Errorf("cannot read the catalogue: %w", err) + } + var out []broker.Bucket + for _, m := range declared { + out = append(out, bucketsOf(m)...) + } + return out, nil +} + func asSeat(s catalogue.SeatDeclaration) broker.Seat { return broker.Seat{Name: s.Name, Scope: s.Scope, Accepts: s.Accepts, Emits: s.Emits, Serves: catalogue.VerbNames(s.Serves)}