194 lines
9.2 KiB
Go
194 lines
9.2 KiB
Go
package broker
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"time"
|
|
|
|
"github.com/nats-io/nats.go/jetstream"
|
|
)
|
|
|
|
// The controller's own key-value buckets (novox/hq to-be 45 §1, §6, §7).
|
|
//
|
|
// **What the controller must remember across its own restart, it keeps on the bus.** A call's
|
|
// outcome lived in the memory of the process that served it (novox/hq issue 265), so a controller
|
|
// replaced while a push ran answered "no such call" for the one thing its caller had been told to
|
|
// ask about. The bus already outlives the controller and is the shape ADR 0201 gives a module's
|
|
// current state: one value per key, written by one owner, read by anybody granted it. These are the
|
|
// controller's, written by it alone — the writers table of to-be 45 §1 — and asserted on every start
|
|
// like the streams, so a bus raised from nothing has them before the first call is served.
|
|
|
|
// CallsBucket keeps every call of the mesh's own verbs and what came of it; HandActsBucket every act
|
|
// a person did by hand, with why; ConditionsBucket every condition open now (to-be 45 §2), one key
|
|
// each, and ConditionHistoryBucket every transition of one — raised, changed, silenced, cleared —
|
|
// for ninety days.
|
|
//
|
|
// **The history is a bucket of its own** because its keys expire and an open condition's must not:
|
|
// a bucket has one age for every key, and a condition open longer than the history is kept would
|
|
// otherwise vanish from the store while still true.
|
|
var (
|
|
CallsBucket = BucketName(ControllerSeat, "calls")
|
|
HandActsBucket = BucketName(ControllerSeat, "hand-acts")
|
|
ConditionsBucket = BucketName(ControllerSeat, "conditions")
|
|
ConditionHistoryBucket = BucketName(ControllerSeat, "condition-history")
|
|
// LeaseBucket holds the controller's lease (to-be 45 §6): one key, `holder`, which the instance
|
|
// allowed to act writes by compare-and-set and renews; its revision when taken is the epoch.
|
|
LeaseBucket = BucketName(ControllerSeat, "lease")
|
|
// AskedBucket keeps what the controller asked the operator about its conditions (novox/hq ADR 0259):
|
|
// each ask by its id, its options and the actions they stand for, how it ended and whether the
|
|
// controller acted on its warrant — so a restart neither asks twice nor acts twice.
|
|
AskedBucket = BucketName(ControllerSeat, "asked")
|
|
)
|
|
|
|
// AskedKeptFor is how long an ask is kept after it was made: a month, as the router keeps its own.
|
|
const AskedKeptFor = 30 * 24 * time.Hour
|
|
|
|
// LeaseTTL is how long the lease's key lives unrenewed (to-be 45 §6): fifteen seconds, renewed
|
|
// every five. The bucket's age, so the bus forgets a holder that stopped renewing.
|
|
const LeaseTTL = 15 * time.Second
|
|
|
|
// The bounds to-be 45 §6 sets for calls: the last thousand, or fourteen days, whichever is fewer.
|
|
// A call is two keys — its record, and its answer apart so a listing does not read every answer —
|
|
// so the stream holds twice as many messages as it keeps calls.
|
|
const (
|
|
KeptCallsDurably = 1000
|
|
CallsKeptFor = 14 * 24 * time.Hour
|
|
// CallAnswerBytes is the most of one answer kept: a whole declaration is far smaller, and an
|
|
// answer larger is cut and says so.
|
|
CallAnswerBytes = 64 << 10
|
|
// HandActsKeptFor is as long as a condition's history (to-be 45 §2): an act by hand is read
|
|
// back beside what it addressed.
|
|
HandActsKeptFor = 90 * 24 * time.Hour
|
|
// ConditionHistoryKeptFor is how long a condition's transitions are kept (to-be 45 §2).
|
|
ConditionHistoryKeptFor = 90 * 24 * time.Hour
|
|
)
|
|
|
|
// IsControllerBucket says a bucket is the controller's own, not a module's state nothing declares.
|
|
func IsControllerBucket(bucket string) bool {
|
|
return bucket == CallsBucket || bucket == HandActsBucket || bucket == ConditionsBucket ||
|
|
bucket == ConditionHistoryBucket || bucket == LeaseBucket || bucket == AskedBucket
|
|
}
|
|
|
|
// ControllerBuckets are the controller's own buckets, in the order they are asserted.
|
|
func ControllerBuckets() []string {
|
|
return []string{LeaseBucket, CallsBucket, HandActsBucket, ConditionsBucket, ConditionHistoryBucket, AskedBucket}
|
|
}
|
|
|
|
// ControllerBucketsAsserter is what raising the controller's buckets needs of a connection.
|
|
type ControllerBucketsAsserter interface {
|
|
EnsureControllerBuckets() error
|
|
}
|
|
|
|
// EnsureControllerBuckets creates the controller's buckets if absent and brings their options to
|
|
// match. An update, never a delete: what they hold is the record of what the mesh was asked.
|
|
func (j *JetStream) EnsureControllerBuckets() error {
|
|
js, err := jetstream.New(j.conn)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
|
|
defer cancel()
|
|
if err := EnsureLeaseBucket(ctx, js); err != nil {
|
|
return err
|
|
}
|
|
if _, err := js.CreateOrUpdateKeyValue(ctx, jetstream.KeyValueConfig{
|
|
Bucket: CallsBucket,
|
|
Description: "the calls of the mesh's own verbs and what came of each (novox/hq to-be 45 §6, issue " +
|
|
"265): written by the controller alone, read through `calls`; the last thousand, or fourteen days",
|
|
History: 1,
|
|
TTL: CallsKeptFor,
|
|
MaxValueSize: CallAnswerBytes + 4<<10,
|
|
MaxBytes: 2 * KeptCallsDurably * (CallAnswerBytes + 4<<10),
|
|
Storage: jetstream.FileStorage,
|
|
}); err != nil {
|
|
return fmt.Errorf("asserting bucket %s: %w", CallsBucket, err)
|
|
}
|
|
// **The count, on the stream under the bucket.** A bucket has an age and a size and no count;
|
|
// the stream it is made of does, and with one value per key the oldest message is the oldest
|
|
// call. Asserted after the bucket, every time, because asserting the bucket writes the stream's
|
|
// configuration whole and puts the count back to none.
|
|
stream, err := js.Stream(ctx, "KV_"+CallsBucket)
|
|
if err != nil {
|
|
return fmt.Errorf("reading the stream under %s: %w", CallsBucket, err)
|
|
}
|
|
cfg := stream.CachedInfo().Config
|
|
if cfg.MaxMsgs != 2*KeptCallsDurably {
|
|
cfg.MaxMsgs = 2 * KeptCallsDurably
|
|
cfg.Discard = jetstream.DiscardOld
|
|
if _, err := js.UpdateStream(ctx, cfg); err != nil {
|
|
return fmt.Errorf("bounding %s to the last %d calls: %w", CallsBucket, KeptCallsDurably, err)
|
|
}
|
|
}
|
|
if _, err := js.CreateOrUpdateKeyValue(ctx, jetstream.KeyValueConfig{
|
|
Bucket: HandActsBucket,
|
|
Description: "every act a person did by hand, with why (novox/hq to-be 45 §7): written by the " +
|
|
"controller's repairing verbs and `hand-act record`, read through `hand-acts`",
|
|
History: 1,
|
|
TTL: HandActsKeptFor,
|
|
MaxValueSize: 16 << 10,
|
|
MaxBytes: 64 << 20,
|
|
Storage: jetstream.FileStorage,
|
|
}); err != nil {
|
|
return fmt.Errorf("asserting bucket %s: %w", HandActsBucket, err)
|
|
}
|
|
// **No age on the open conditions.** A condition is removed when observation clears it and at no
|
|
// other moment: one that expired would be a fault the store forgot while it was still true.
|
|
if _, err := js.CreateOrUpdateKeyValue(ctx, jetstream.KeyValueConfig{
|
|
Bucket: ConditionsBucket,
|
|
Description: "every condition open now, one key each (novox/hq to-be 45 §2): written by the " +
|
|
"controller alone, raised and cleared by observation, read through `conditions`",
|
|
History: 1,
|
|
MaxValueSize: 64 << 10,
|
|
MaxBytes: 64 << 20,
|
|
Storage: jetstream.FileStorage,
|
|
}); err != nil {
|
|
return fmt.Errorf("asserting bucket %s: %w", ConditionsBucket, err)
|
|
}
|
|
if _, err := js.CreateOrUpdateKeyValue(ctx, jetstream.KeyValueConfig{
|
|
Bucket: ConditionHistoryBucket,
|
|
Description: "every transition of a condition — raised, changed, silenced, cleared — kept ninety " +
|
|
"days (novox/hq to-be 45 §2): written by the controller alone, read through `conditions history`",
|
|
History: 1,
|
|
TTL: ConditionHistoryKeptFor,
|
|
MaxValueSize: 64 << 10,
|
|
MaxBytes: 256 << 20,
|
|
Storage: jetstream.FileStorage,
|
|
}); err != nil {
|
|
return fmt.Errorf("asserting bucket %s: %w", ConditionHistoryBucket, err)
|
|
}
|
|
if _, err := js.CreateOrUpdateKeyValue(ctx, jetstream.KeyValueConfig{
|
|
Bucket: AskedBucket,
|
|
Description: "what the controller asked the operator about its conditions, and what came of each (novox/hq " +
|
|
"ADR 0259): written by the controller alone; an ask acted on is acted on once",
|
|
History: 1,
|
|
TTL: AskedKeptFor,
|
|
MaxValueSize: 32 << 10,
|
|
MaxBytes: 32 << 20,
|
|
Storage: jetstream.FileStorage,
|
|
}); err != nil {
|
|
return fmt.Errorf("asserting bucket %s: %w", AskedBucket, err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// EnsureLeaseBucket creates the lease's bucket if absent and brings its options to match (to-be 45
|
|
// §6). **Before the lease is taken, by any candidate**: it is the lease's own precondition, and asserting
|
|
// a bucket that exists changes nothing. File storage, so its revisions — the epochs — outlive a restart of
|
|
// the bus; one raised again from nothing is moved past the highest epoch issued when the lease is taken.
|
|
func EnsureLeaseBucket(ctx context.Context, js jetstream.JetStream) error {
|
|
if _, err := js.CreateOrUpdateKeyValue(ctx, jetstream.KeyValueConfig{
|
|
Bucket: LeaseBucket,
|
|
Description: "the controller's lease (novox/hq to-be 45 §6): the key `holder`, written by compare-and-set " +
|
|
"by the one controller instance that may act and renewed every five seconds; its revision when taken " +
|
|
"is the epoch every declaration and plan write carries",
|
|
History: 1,
|
|
TTL: LeaseTTL,
|
|
MaxValueSize: 4 << 10,
|
|
MaxBytes: 1 << 20,
|
|
Storage: jetstream.FileStorage,
|
|
}); err != nil {
|
|
return fmt.Errorf("asserting bucket %s: %w", LeaseBucket, err)
|
|
}
|
|
return nil
|
|
}
|