From babd7b2f4789cfc91b833a6e8f45155c02b7c865 Mon Sep 17 00:00:00 2001 From: jochen Date: Sun, 4 Oct 2026 11:13:34 +0200 Subject: [PATCH] Assert every declared state's bucket on each push, before the memberships that name it (novox/hq ADR 0201) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The raise at start was the only place buckets were asserted, so a module registered and assigned since had none until the control plane restarted — found on the first module to declare state. --- cmd/mesh-controller/push.go | 10 ++++++++++ internal/broker/jetstream.go | 7 +++++++ internal/broker/state_test.go | 14 ++++++++++++++ 3 files changed, 31 insertions(+) diff --git a/cmd/mesh-controller/push.go b/cmd/mesh-controller/push.go index e32721d..007807d 100644 --- a/cmd/mesh-controller/push.go +++ b/cmd/mesh-controller/push.go @@ -747,6 +747,16 @@ func issueMemberships(ctx context.Context, open *stores, server *link.Server, se if !ok { return nil } + // **Every declared state's bucket, before the memberships that name it** (novox/hq ADR 0201). The + // raise at start asserts them too, but a module registered and assigned since would otherwise have + // its bucket only after the control plane next restarts — found the first time a module declared + // state: its bundle asked for a bucket that did not exist. Idempotent and cheap; a failure is said + // and the push stands, as a membership's is. + if buckets, err := open.inventory.DeclaredBuckets(ctx); err != nil { + fmt.Printf(" the modules' state could not be read, so no bucket was asserted: %v\n", err) + } else if _, err := broker.RaiseBuckets(broker.OnConn(bus.Conn), buckets); err != nil { + fmt.Printf(" the modules' state could not be asserted on the bus: %v — the next push tries again\n", err) + } // The declarations are sent and recorded by now; a membership that cannot be issued is said // and does not unsay them. Every runtime without one serves the shape it derives (ADR 0160), so // the push stands, the first failure is named once, and the next push tries again. diff --git a/internal/broker/jetstream.go b/internal/broker/jetstream.go index aed8381..06150c5 100644 --- a/internal/broker/jetstream.go +++ b/internal/broker/jetstream.go @@ -86,6 +86,13 @@ func pinnedTo(path string) (*tls.Config, error) { return PinnedToFingerprint(want), nil } +// OnConn is the JetStream handle over a connection the caller already holds — the control plane's +// link — for asserting what the bus holds without dialling a second time. +func OnConn(conn *nats.Conn) *JetStream { + js, _ := conn.JetStream() + return &JetStream{conn: conn, js: js} +} + // DialPinned is Dial with the server's certificate pinned by a fingerprint the caller already holds // — a module or a build machine that was handed one beside its credential, and has no file. func DialPinned(url, fingerprint string, opts ...nats.Option) (*JetStream, error) { diff --git a/internal/broker/state_test.go b/internal/broker/state_test.go index ba78ebb..a6f7580 100644 --- a/internal/broker/state_test.go +++ b/internal/broker/state_test.go @@ -164,3 +164,17 @@ func TestABucketIsAssertedInPlace(t *testing.T) { } } } + +// Against a real server: the handle over a connection the control plane already holds asserts a +// bucket as Dial's does — what a push uses, so a module registered since the last start has its +// bucket before its membership names it. +func TestABucketIsAssertedOverAHeldConnection(t *testing.T) { + js := aLiveBus(t) + held := OnConn(js.Conn()) + if _, err := RaiseBuckets(held, []Bucket{{Module: "statetest", Name: "held"}}); err != nil { + t.Fatalf("asserting over a held connection failed: %v", err) + } + if _, err := js.Context().KeyValue("statetest_held"); err != nil { + t.Fatalf("the bucket is not there: %v", err) + } +}