Compare commits
7
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
75b2676d32 | ||
|
|
1363a2fe27 | ||
|
|
17b8f14fe1 | ||
|
|
42c394acc2 | ||
|
|
b9ad7a2948 | ||
|
|
cde22ff627 | ||
|
|
aec55b7072 |
@@ -7,6 +7,7 @@ import (
|
||||
"os"
|
||||
"path/filepath"
|
||||
"sort"
|
||||
"strings"
|
||||
|
||||
"github.com/novox/mesh-controller/internal/catalogue"
|
||||
)
|
||||
@@ -90,6 +91,17 @@ func moduleCheck(paths []string, out io.Writer) error {
|
||||
if len(m.Invokes) > 0 {
|
||||
fmt.Fprintf(out, ", invokes %s", joinInvokes(m.Invokes))
|
||||
}
|
||||
// The state it keeps and reads (novox/hq ADR 0201), so a reviewer sees what lands on the bus.
|
||||
if len(m.State) > 0 {
|
||||
kept := make([]string, 0, len(m.State))
|
||||
for _, s := range m.State {
|
||||
kept = append(kept, s.Name)
|
||||
}
|
||||
fmt.Fprintf(out, ", keeps state %s", strings.Join(kept, ", "))
|
||||
}
|
||||
if len(m.Reads) > 0 {
|
||||
fmt.Fprintf(out, ", reads %s", strings.Join(m.Reads, ", "))
|
||||
}
|
||||
fmt.Fprintln(out)
|
||||
}
|
||||
if failed > 0 {
|
||||
|
||||
@@ -5,6 +5,7 @@ import (
|
||||
"errors"
|
||||
"fmt"
|
||||
"os"
|
||||
"time"
|
||||
|
||||
"github.com/novox/mesh-controller/internal/artifacts"
|
||||
"github.com/novox/mesh-controller/internal/inventory"
|
||||
@@ -49,28 +50,43 @@ func collect(ctx context.Context, inv *inventory.Inventory) {
|
||||
return
|
||||
}
|
||||
|
||||
// **Bounded, because this runs inside somebody's build.** The first sweep of a mesh that has
|
||||
// never collected has the whole history to get through, and a person waiting on `build` should
|
||||
// not pay for it. Two bounds, and what is left over is simply offered again next time —
|
||||
// builds are frequent, and the point is that the store stops growing, not that it empties
|
||||
// tonight.
|
||||
within, stop := context.WithTimeout(ctx, sweepBudget)
|
||||
defer stop()
|
||||
store := artifacts.Store{Address: address}
|
||||
|
||||
var done []string
|
||||
var refused int
|
||||
for _, reference := range references {
|
||||
switch err := store.LetGo(ctx, reference); {
|
||||
case err == nil, errors.Is(err, artifacts.Gone):
|
||||
var left int
|
||||
for i, reference := range references {
|
||||
if i >= mostPerSweep || within.Err() != nil {
|
||||
left = len(references) - i
|
||||
break
|
||||
}
|
||||
err := store.LetGo(within, reference)
|
||||
if err == nil || errors.Is(err, artifacts.Gone) {
|
||||
// Gone is the outcome wanted, already true. Recorded so the next sweep does not ask
|
||||
// again for ever.
|
||||
done = append(done, reference)
|
||||
default:
|
||||
refused++
|
||||
if refused == 1 {
|
||||
// Once per sweep. A store that refuses one refuses all of them, and a hundred
|
||||
// identical lines would bury the reason.
|
||||
fmt.Fprintf(os.Stderr, "the artifact store kept %s: %v\n", reference, err)
|
||||
}
|
||||
continue
|
||||
}
|
||||
// **Stopped at the first refusal, not pushed through.** A store that refuses one refuses
|
||||
// all of them — deletion disabled, the store down, the network gone — so going on would
|
||||
// be a hundred identical failures and a hundred identical log lines in front of whoever
|
||||
// was building something.
|
||||
fmt.Fprintf(os.Stderr, "the artifact store kept %s, so nothing more was asked of it: %v\n",
|
||||
reference, err)
|
||||
left = len(references) - i
|
||||
break
|
||||
}
|
||||
|
||||
if len(done) > 0 {
|
||||
// Recorded outside `within`: the deletions happened, and losing the record of them because
|
||||
// the sweep ran out of budget would mean asking about them again for ever.
|
||||
if err := inv.MarkCollected(ctx, done); err != nil {
|
||||
// Said, and that is all: the artifacts are gone either way, and the only cost of an
|
||||
// unrecorded collection is that the next sweep asks about them again.
|
||||
fmt.Fprintf(os.Stderr, "the store let go of %d artifact(s) and the record of it did not keep: %v\n",
|
||||
len(done), err)
|
||||
return
|
||||
@@ -78,7 +94,15 @@ func collect(ctx context.Context, inv *inventory.Inventory) {
|
||||
fmt.Fprintf(os.Stderr, "the artifact store let go of %d artifact(s) the mesh no longer keeps\n",
|
||||
len(done))
|
||||
}
|
||||
if refused > 0 {
|
||||
fmt.Fprintf(os.Stderr, "%d artifact(s) were not collected; the next build asks again\n", refused)
|
||||
if left > 0 {
|
||||
fmt.Fprintf(os.Stderr, "%d more to collect; the next build asks again\n", left)
|
||||
}
|
||||
}
|
||||
|
||||
// mostPerSweep is how many artifacts one sweep will ask about. Enough that a mesh building
|
||||
// several times a day converges within days of this landing; small enough that no single build
|
||||
// waits on the whole backlog.
|
||||
const mostPerSweep = 200
|
||||
|
||||
// sweepBudget is the longest a sweep will keep a build waiting.
|
||||
const sweepBudget = 60 * time.Second
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
|
||||
@@ -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()
|
||||
}
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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 `<module>.<name>` (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)
|
||||
}
|
||||
|
||||
@@ -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, `<module>.<name>`, 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 `<module>.<name>` (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
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -33,6 +33,10 @@ type Declared struct {
|
||||
Watches []Seat
|
||||
// Invokes are the tools it calls, `<module>.<tool>` 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 `<module>.<name>` (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 {
|
||||
|
||||
@@ -61,7 +61,7 @@ func knownFor(m Manifest, needs []Needed, node string) (map[string]map[string]st
|
||||
"as": as,
|
||||
}
|
||||
// What the provider derives for this consumer rather than for all of them
|
||||
// (novox/hq ADR 0201). Filled here, the one place a provision and the module
|
||||
// (novox/hq ADR 0202). Filled here, the one place a provision and the module
|
||||
// requiring it are both in hand.
|
||||
served, err := ServedTo(n.Serves, as)
|
||||
if err != nil {
|
||||
|
||||
@@ -8,7 +8,7 @@ import (
|
||||
)
|
||||
|
||||
// What a provider derives for one consumer, said once in the provider's definition and delivered
|
||||
// to both ends (novox/hq ADR 0201, issue 124).
|
||||
// to both ends (novox/hq ADR 0202, issue 124).
|
||||
//
|
||||
// A `serves` block is otherwise literal: the same values for every consumer. Where the provider
|
||||
// *names the resource* — a bucket, a database, a vhost — the name is derived from who is asking,
|
||||
@@ -165,7 +165,7 @@ func sortedAnyKeys(values map[string]any) []string {
|
||||
}
|
||||
|
||||
// derivedFor is what the provider on this machine derives for one consumer of one provision
|
||||
// (novox/hq ADR 0201).
|
||||
// (novox/hq ADR 0202).
|
||||
//
|
||||
// Settled first, then derived: an operator may set a prefix on what the provider serves and the
|
||||
// mesh still fills the consumer's half of it ([ADR 0174]). Only the keys that actually name the
|
||||
@@ -177,7 +177,7 @@ func sortedAnyKeys(values map[string]any) []string {
|
||||
// choice servedOnThisMachine makes for the consumer's half. Nothing serving it on this machine is
|
||||
// not an error: a contribution can reach a machine whose provider is a record or an adapter, and
|
||||
// then there is nothing derived to tell.
|
||||
func (r Resolution) derivedFor(provision, as string, settings SettingsBy) (map[string]any, error) {
|
||||
func (r Resolution) derivedFor(provision, as, consumer, local string, settings SettingsBy) (map[string]any, error) {
|
||||
for _, m := range r.Modules {
|
||||
serves, said := m.Serves[provision]
|
||||
if !said {
|
||||
@@ -195,6 +195,27 @@ func (r Resolution) derivedFor(provision, as string, settings SettingsBy) (map[s
|
||||
if names == nil {
|
||||
return nil, nil
|
||||
}
|
||||
// **A consumer that keeps several holders of this provision is refused** — this is issue
|
||||
// 124's own failure one case to the side, and it would be just as quiet.
|
||||
//
|
||||
// Each holder gets its own login, `…_<local>` (ADR 0094), and a provider derives from the
|
||||
// login, so it would make one resource per holder. The consumer's side has no such
|
||||
// dimension: one binding file per provision, one `${bound:<provision>:<key>}`, both
|
||||
// derived from the un-suffixed identity. So the provider would create the holder's
|
||||
// resource and the consumer would be configured against a name nothing made — it would
|
||||
// authenticate successfully and be refused on every object, which reads like a credential
|
||||
// fault and is not one.
|
||||
//
|
||||
// Lifting this means giving the consumer's side a local dimension. That is a decision,
|
||||
// not an omission, and until it is taken the mesh says so rather than guessing.
|
||||
if local != "" {
|
||||
return nil, fmt.Errorf(
|
||||
"%s keeps several holders of %s (this one is %q), and %s derives %s for each "+
|
||||
"consumer from the login the mesh minted. Each holder has its own login, and a "+
|
||||
"consumer is told one value per requirement — so the two ends would name "+
|
||||
"different things and nothing would compare them (novox/hq ADR 0202)",
|
||||
consumer, local, provision, m.Module, orNothing(sortedAnyKeys(names)))
|
||||
}
|
||||
settled, err := Settle(names, settings[m.Module])
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("%s serving %s: %w", m.Module, provision, err)
|
||||
@@ -209,7 +230,7 @@ func (r Resolution) derivedFor(provision, as string, settings SettingsBy) (map[s
|
||||
}
|
||||
|
||||
// notTranscribed refuses a consumer's file that writes out the value its provider derives for it,
|
||||
// instead of asking for it (novox/hq ADR 0201, issue 124).
|
||||
// instead of asking for it (novox/hq ADR 0202, issue 124).
|
||||
//
|
||||
// **What would have caught the one wrong instance.** The object store's three consumers each wrote
|
||||
// their bucket into their own configuration by hand. One of them named a predecessor's bucket, and
|
||||
|
||||
@@ -647,7 +647,7 @@ func (r Resolution) compose(with Rendering, owner map[string]string,
|
||||
}
|
||||
as := ConsumerIdentity(r.Node, IdentitySource(m.Slug, m.Module))
|
||||
// What the provider derives for THIS consumer, filled here where the consumer is
|
||||
// known (novox/hq ADR 0201). The same fill knownFor does below, so the binding file
|
||||
// known (novox/hq ADR 0202). The same fill knownFor does below, so the binding file
|
||||
// and the module's `${bound:…}` substitutions cannot say different things.
|
||||
told := *found
|
||||
told.Serves, err = ServedTo(told.Serves, as)
|
||||
@@ -777,7 +777,7 @@ func (r Resolution) compose(with Rendering, owner map[string]string,
|
||||
// And the machine underneath, which no binding of its own can tell it.
|
||||
thisMachine := machineFacts(r, with.Names, with.MeshRange)
|
||||
|
||||
// **A definition that already holds the answer transcribed it** (novox/hq ADR 0201).
|
||||
// **A definition that already holds the answer transcribed it** (novox/hq ADR 0202).
|
||||
// Judged over what the module itself declares, and before anything is substituted: the
|
||||
// mesh's own generated files — the binding, the contributions — legitimately carry the
|
||||
// derived value, and after substitution so does every consumer's file, so this is the one
|
||||
@@ -907,6 +907,15 @@ func (r Resolution) compose(with Rendering, owner map[string]string,
|
||||
if renamed := reflectsRenamed(m.Module, resource["reload-on"]); renamed != nil {
|
||||
copied["reload-on"] = renamed
|
||||
}
|
||||
// And which of its module's containers a scheduled step holds still (novox/hq ADR 0189).
|
||||
// **The loudest of the three when it is missed.** An unprefixed `restart-on` matches
|
||||
// nothing and a service quietly never restarts; an unprefixed `while-stopped` names a
|
||||
// container the declaration does not contain, and the host refuses the whole
|
||||
// declaration — so the machine takes nothing at all, for every push, until this is
|
||||
// right. That is what it did on the control node (2026-10-04).
|
||||
if renamed := reflectsRenamed(m.Module, resource[WhileStopped]); renamed != nil {
|
||||
copied[WhileStopped] = renamed
|
||||
}
|
||||
// And what a process replaces (novox/hq issue 213): a resource of this module's that it
|
||||
// no longer declares, named as the host recorded it, or the host hands nothing over and
|
||||
// removes it first.
|
||||
@@ -1203,7 +1212,7 @@ type Contribution struct {
|
||||
// is what makes swapping one for another cost nothing.
|
||||
Values map[string]any `json:"values"`
|
||||
// Derived is what this provider's own definition said it derives for this consumer, already
|
||||
// derived (novox/hq ADR 0201).
|
||||
// derived (novox/hq ADR 0202).
|
||||
//
|
||||
// **The provider is told, rather than recomputing it.** A served value may name the consumer's
|
||||
// identity — a bucket named for who is asking, a database prefixed with it — and before this
|
||||
@@ -1307,7 +1316,7 @@ func (r Resolution) contributions(settings SettingsBy, grants []Grant,
|
||||
continue
|
||||
}
|
||||
as := holderAs(ConsumerIdentity(g.Consumer, IdentitySource(g.Slug, g.From)), g.Local)
|
||||
derived, err := r.derivedFor(g.Provision, as, settings)
|
||||
derived, err := r.derivedFor(g.Provision, as, g.From, g.Local, settings)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
@@ -7,7 +7,7 @@ import (
|
||||
)
|
||||
|
||||
// What a provider derives for each consumer, said once and delivered to both ends
|
||||
// (novox/hq ADR 0201, issue 124).
|
||||
// (novox/hq ADR 0202, issue 124).
|
||||
//
|
||||
// The failure these are written against: the object store's provisioner derived each consumer's
|
||||
// bucket from the login the mesh minted, in its own code, and the mesh had no channel to tell the
|
||||
@@ -295,3 +295,49 @@ func storeGrants(t *testing.T, out []map[string]any) []Contribution {
|
||||
t.Fatalf("the provider was given no contributions file: %v", out)
|
||||
return nil
|
||||
}
|
||||
|
||||
// A consumer that keeps SEVERAL holders of one provision is refused, rather than told one thing
|
||||
// while its provider is told another.
|
||||
//
|
||||
// **This is issue 124's own failure, one case to the side.** The mesh gives each holder its own
|
||||
// login — `mesh_node_mod_<local>` (ADR 0094) — and the provider derives from the login, so it
|
||||
// would make one resource per holder. The consumer's side has no such dimension: there is one
|
||||
// binding file per provision and one `${bound:<provision>:<key>}`, both derived from the
|
||||
// un-suffixed identity. So the provider would create `…-mod-cold` and the consumer would be
|
||||
// configured against `…-mod`: it would authenticate successfully and be refused on every object,
|
||||
// which is exactly the fault this whole record exists to end.
|
||||
//
|
||||
// Refused, loudly, at the one place that can see both halves. Lifting it means giving the
|
||||
// consumer's side a local dimension, which is a decision and not an omission.
|
||||
func TestAConsumerWithSeveralHoldersOfADerivingProviderIsRefused(t *testing.T) {
|
||||
m := files()
|
||||
// Two holders of the one provision, the shape ADR 0094 gives a module that keeps several.
|
||||
m.Secrets = nil
|
||||
m.SecretsMany = map[string]map[string]string{"s3-bucket": {
|
||||
"hot": "/var/lib/files/hot.secret",
|
||||
"cold": "/var/lib/files/cold.secret",
|
||||
}}
|
||||
m.Resources = []map[string]any{{
|
||||
"id": "env", "type": "file", "path": "/var/lib/files/env", "mode": "0600",
|
||||
"content": "BUCKET=${bound:s3-bucket:bucket}\n",
|
||||
}}
|
||||
r, err := Resolve(shelf(store(), m), []string{"store", "files"}, reachable(), World{})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
_, err = r.Declaration(Rendering{Grants: []Grant{
|
||||
{Provision: "s3-bucket", Consumer: "workstation", From: "files", Slug: "files",
|
||||
Local: "hot", Values: map[string]any{}, Sealed: "c2VhbGVk"},
|
||||
{Provision: "s3-bucket", Consumer: "workstation", From: "files", Slug: "files",
|
||||
Local: "cold", Values: map[string]any{}, Sealed: "c2VhbGVk"},
|
||||
}})
|
||||
if err == nil {
|
||||
t.Fatal("a consumer with several holders of a deriving provider was accepted; " +
|
||||
"its two ends would have disagreed in silence")
|
||||
}
|
||||
for _, want := range []string{"files", "s3-bucket", "bucket"} {
|
||||
if !strings.Contains(err.Error(), want) {
|
||||
t.Errorf("the refusal does not name %q: %v", want, err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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 `<module>.<name>`
|
||||
// (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) {
|
||||
@@ -1428,7 +1439,7 @@ func ParseManifest(raw []byte) (Manifest, error) {
|
||||
"%s serves %q to whoever requires it, and does not provide it", m.Module, to))
|
||||
}
|
||||
}
|
||||
// A served value may be derived for the consumer it is served to (novox/hq ADR 0201). Read
|
||||
// A served value may be derived for the consumer it is served to (novox/hq ADR 0202). Read
|
||||
// here, where the definition is, rather than when somebody first requires it: a rule that
|
||||
// would be refused at the first consumer is wrong from the moment it is written.
|
||||
problems = append(problems, CheckServes(m)...)
|
||||
|
||||
@@ -215,6 +215,13 @@ func CatalogueProblems(shelf Shelf) []string {
|
||||
}
|
||||
}
|
||||
}
|
||||
// A read of a module's state that module does not keep (novox/hq ADR 0201) — said only where the
|
||||
// owner is on the shelf, as a consumer may be installed before its emitter.
|
||||
var manifests []Manifest
|
||||
for _, module := range shelfOrder(shelf) {
|
||||
manifests = append(manifests, shelf[module])
|
||||
}
|
||||
problems = append(problems, StateReadsNothingDeclares(manifests)...)
|
||||
sort.Strings(problems)
|
||||
return problems
|
||||
}
|
||||
|
||||
@@ -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 `<module>.<name>`, 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 <module>.<name>", 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
|
||||
}
|
||||
@@ -0,0 +1,94 @@
|
||||
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>.<name>`},
|
||||
{`{"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 "))
|
||||
}
|
||||
}
|
||||
|
||||
// `module check` says it too: the cross-catalogue pass names a read nothing on the shelf keeps.
|
||||
func TestTheCataloguePassNamesAReadItsOwnerDoesNotKeep(t *testing.T) {
|
||||
shelf := Shelf{
|
||||
"licence-manager": {Module: "licence-manager", State: []StateDeclaration{{Name: "bindings"}}},
|
||||
"claude-code": {Module: "claude-code", Reads: []string{"licence-manager.tokens"}},
|
||||
}
|
||||
problems := CatalogueProblems(shelf)
|
||||
if len(problems) != 1 || !strings.Contains(problems[0], `keeps no state called "tokens"`) {
|
||||
t.Fatalf("the catalogue pass said %v", problems)
|
||||
}
|
||||
}
|
||||
@@ -2,6 +2,7 @@ package catalogue
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"strings"
|
||||
"testing"
|
||||
)
|
||||
@@ -87,3 +88,58 @@ func TestAMaintenanceWindowIsRefusedWhereTheDefinitionShowsItCannotMean(t *testi
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// A composed declaration names the step's held containers the way the machine knows them.
|
||||
//
|
||||
// **The gap that let a bug through to the control node.** The manifest says `while-stopped:
|
||||
// ["store"]`, because a module names its own resources locally; the declaration a machine
|
||||
// receives calls that container `distribution.store`, because every resource is composed under
|
||||
// its module. `restart-on` and `reload-on` are rewritten for exactly this reason, and
|
||||
// `while-stopped` was not — so the host found no container by that id and refused the whole
|
||||
// declaration, every push, until it was fixed.
|
||||
//
|
||||
// It passed every test on both sides: the controller's tests read manifests, the host's read
|
||||
// hand-written declarations with bare ids. Only composing one and judging the result catches it.
|
||||
func TestAComposedWindowNamesTheContainerAsTheMachineKnowsIt(t *testing.T) {
|
||||
store := Manifest{
|
||||
Module: "distribution", Version: "1",
|
||||
Provides: FromAnywhere("artifact-store"),
|
||||
Listens: []Listening{{Port: 5000, Protocol: "tcp", From: FromMesh}},
|
||||
Serves: map[string]map[string]any{"artifact-store": {"port": 5000}},
|
||||
Resources: []map[string]any{
|
||||
{"id": "store", "type": "container", "name": "mesh-registry",
|
||||
"image": "registry@sha256:" + strings.Repeat("a", 64), "ports": []any{"5000"}},
|
||||
{"id": "collect", "type": "container", "name": "mesh-registry-collect",
|
||||
"image": "registry@sha256:" + strings.Repeat("a", 64),
|
||||
"schedule": "30 3 * * *", WhileStopped: []any{"store"}},
|
||||
},
|
||||
}
|
||||
r, err := Resolve(shelf(store), []string{"distribution"}, reachable(), World{})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
out, err := r.Declaration(Rendering{})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
collect := fileNamed(out, "distribution.collect")
|
||||
if collect == nil {
|
||||
for _, res := range out {
|
||||
if res["id"] == "distribution.collect" {
|
||||
collect = res
|
||||
}
|
||||
}
|
||||
}
|
||||
if collect == nil {
|
||||
t.Fatalf("the step was not composed at all: %v", out)
|
||||
}
|
||||
held, _ := collect[WhileStopped].([]any)
|
||||
if len(held) != 1 {
|
||||
t.Fatalf("the composed step holds %v still; want one container", collect[WhileStopped])
|
||||
}
|
||||
if got := fmt.Sprint(held[0]); got != "distribution.store" {
|
||||
t.Fatalf("the composed step says it holds %q still, and the machine's container is "+
|
||||
"called %q — the host refuses a declaration naming a container it does not have, "+
|
||||
"whole, so the machine would take nothing at all", got, "distribution.store")
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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)}
|
||||
|
||||
Reference in New Issue
Block a user