A module's state on the bus: buckets from the catalogue, grants, membership (novox/hq ADR 0201)
A manifest names the state it keeps (state) and reads (reads); the controller asserts a key-value bucket per name on every raise, grants owners write and readers read (measured against a running server), issues each assignment its buckets in the membership, and reports buckets nothing declares without removing them.
This commit is contained in:
@@ -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()
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user