Merge pull request 'Assert every declared state's bucket on each push (hq ADR 0201)' (#262) from fix/buckets-on-push into main
This commit was merged in pull request #262.
This commit is contained in:
@@ -747,6 +747,16 @@ func issueMemberships(ctx context.Context, open *stores, server *link.Server, se
|
|||||||
if !ok {
|
if !ok {
|
||||||
return nil
|
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
|
// 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
|
// 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.
|
// the push stands, the first failure is named once, and the next push tries again.
|
||||||
|
|||||||
@@ -86,6 +86,13 @@ func pinnedTo(path string) (*tls.Config, error) {
|
|||||||
return PinnedToFingerprint(want), nil
|
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
|
// 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.
|
// — 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) {
|
func DialPinned(url, fingerprint string, opts ...nats.Option) (*JetStream, error) {
|
||||||
|
|||||||
@@ -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)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user