// Package lease is the controller's lease (novox/hq to-be 45 §6, ADR 0227 rule 1): the one key that // says which controller instance may act, and the epoch every act of it carries. // // **A controller instance acts only while it holds the key `holder`** in its lease bucket — sends a // declaration, writes a plan, a condition or a call. The key is written by compare-and-set, lives a // bucket's age (fifteen seconds) unless renewed, and is renewed every five. The **epoch** is the // bucket's revision at which the instance took it: a later holder's is higher, so a machine that has // heard from epoch 57 can refuse anything still arriving from 41 (issue 204). // // controller A (epoch 41) ──renew──renew──╳ (renewal refused)──► stops acting, exits // controller B ──wait──────────────take (epoch 57)──► acts // // **Three ways an instance stops acting, all at once.** A renewal refused — somebody else wrote the key // — is a loss, said and final: Lost closes and the process exits, so its service manager restarts it as // a candidate. A renewal that cannot be made in time is the same, because past its time the key may be // somebody else's. And Epoch, the gate every act passes, answers an error from the moment the last // renewal is older than the key's age less a margin — **before** another instance could have taken // it, whatever the renewing goroutine is doing — so a send in flight at the moment of loss is refused // by the clock, not by a goroutine that may be stuck. package lease import ( "context" "encoding/json" "errors" "fmt" "sync" "time" "github.com/nats-io/nats.go/jetstream" ) // Key is the one key of the lease bucket. const Key = "holder" // Holder is who holds the lease, as the key says it. type Holder struct { // Instance names one controller process: its machine, its process and when it started. Two // instances on one machine — the handover of issue 213 — are two names. Instance string `json:"instance"` Host string `json:"host,omitempty"` Build string `json:"build,omitempty"` // Epoch is the revision the key was taken at. Zero in the value written to take it, because the // revision is known only once written; every renewal says it, and Current reads the revision of a // key never renewed. Epoch uint64 `json:"epoch,omitempty"` Taken time.Time `json:"taken"` Renewed time.Time `json:"renewed"` } // Defaults, as to-be 45 §6 sets them. The age is the bucket's, read from it (Open), so the bucket // and the lease cannot disagree about it. const ( RenewEvery = 5 * time.Second // Margin is how long before the key could expire this instance stops acting on it: the gap // between its own clock and the bus's, and a send already on its way. Margin = 3 * time.Second // Poll is how often a candidate looks again at a key somebody else holds. Poll = time.Second ) // ErrNotHeld is an act asked of an instance that does not hold the lease. var ErrNotHeld = errors.New("this controller does not hold the lease") // ErrUnwritable is a lease nobody holds that the bus would not let this instance write: the one case a // caller may serve without it, since nothing else holds it either. Every other failure to take it — // the holder unreadable, the floor unreadable — says nothing about whether another instance acts. var ErrUnwritable = errors.New("nobody holds the lease, and the bus would not let this controller write it") // Options are what a Lease is made with. type Options struct { Holder Holder // RenewEvery, Margin and Poll default to the package's. RenewEvery, Margin, Poll time.Duration // Floor is the highest epoch this mesh has issued, from the controller's own store. **An epoch is // never issued twice**: a lease bucket raised again from nothing — a bus whose data was replaced — // starts its revisions over, and every machine that heard epoch 57 would refuse the next controller // for ever. Below the floor, the bucket's revisions are moved past it before the key is taken. // Nil is no floor. Floor func(ctx context.Context) (uint64, error) // Say is where what the lease does is said: waiting, taking, losing. Say func(format string, args ...any) // Moved is told when the bucket's revisions were moved past the floor: the bucket was raised again // from nothing, which the caller says as a condition. Nil tells nobody. Moved func(was, floor uint64) // Now is the clock; nil is time.Now. Now func() time.Time } // Lease is one instance's hold, or its wait for one. type Lease struct { kv jetstream.KeyValue stream jetstream.Stream ttl time.Duration o Options mu sync.Mutex // held is whether this instance holds the key; epoch the revision it took it at; revision the // revision of its last write, which the next renewal must find. held bool epoch uint64 revision uint64 // validUntil is when this instance stops acting unless it renews: the moment before its last // renewal was sent, plus the key's age, less the margin. validUntil time.Time // renewed is when this instance last took or renewed the key. renewed time.Time lostWhy error lost chan struct{} taken time.Time } // Open is a lease over the bucket, read for its age. The bucket is the caller's to assert. func Open(ctx context.Context, js jetstream.JetStream, bucket string, o Options) (*Lease, error) { kv, err := js.KeyValue(ctx, bucket) if err != nil { return nil, fmt.Errorf("the lease bucket %s is not on the bus: %w", bucket, err) } status, err := kv.Status(ctx) if err != nil { return nil, fmt.Errorf("the lease bucket %s cannot be read: %w", bucket, err) } if status.TTL() <= 0 { // A key that never expires is a lease a dead controller holds for ever. return nil, fmt.Errorf("the lease bucket %s keeps its keys for ever, so a controller that died "+ "holding it would hold it for ever: it must have an age", bucket) } stream, err := js.Stream(ctx, "KV_"+bucket) if err != nil { return nil, fmt.Errorf("the stream under the lease bucket %s cannot be read: %w", bucket, err) } if o.RenewEvery <= 0 { o.RenewEvery = RenewEvery } if o.Margin <= 0 { o.Margin = Margin } if o.Poll <= 0 { o.Poll = Poll } if o.Say == nil { o.Say = func(string, ...any) {} } if o.Now == nil { o.Now = time.Now } if o.RenewEvery+o.Margin >= status.TTL() { return nil, fmt.Errorf("a lease renewed every %s with a margin of %s does not fit a key that lives %s", o.RenewEvery, o.Margin, status.TTL()) } return &Lease{kv: kv, stream: stream, ttl: status.TTL(), o: o, lost: make(chan struct{})}, nil } // TTL is how long the key lives unrenewed. func (l *Lease) TTL() time.Duration { return l.ttl } // Current is who holds the lease now, and false when nobody does. An error is an error: never "nobody". func Current(ctx context.Context, kv jetstream.KeyValue) (Holder, bool, error) { entry, err := kv.Get(ctx, Key) if errors.Is(err, jetstream.ErrKeyNotFound) { return Holder{}, false, nil } if err != nil { return Holder{}, false, err } var h Holder if err := json.Unmarshal(entry.Value(), &h); err != nil { return Holder{}, false, fmt.Errorf("the lease's key holds something that is not a holder: %w", err) } if h.Epoch == 0 { // Taken and not yet renewed: the revision it was taken at is the one it has. h.Epoch = entry.Revision() } return h, true, nil } // Current is who holds this lease now. func (l *Lease) Current(ctx context.Context) (Holder, bool, error) { return Current(ctx, l.kv) } // ErrTaken is a take that found the key held by somebody else; Take waits on it, TryTake answers it. var ErrTaken = errors.New("another controller holds the lease") // Take waits until the key is absent or expired and takes it, answering the epoch. A key somebody // else holds is waited on and said once; an error reading or writing the bucket is answered, because // a candidate that cannot see the key cannot tell whether it may act. func (l *Lease) Take(ctx context.Context) (uint64, error) { said := "" for { epoch, err := l.TryTake(ctx) if err == nil { return epoch, nil } var held *heldBy if !errors.As(err, &held) { return 0, err } if held.h.Instance != said { l.o.Say("another controller holds the lease (%s, epoch %d, taken %s); waiting until it lets go "+ "or stops renewing", held.h.Instance, held.h.Epoch, held.h.Taken.UTC().Format(time.RFC3339)) said = held.h.Instance } select { case <-ctx.Done(): return 0, ctx.Err() case <-time.After(l.o.Poll): } } } // heldBy is ErrTaken naming the holder. type heldBy struct{ h Holder } func (e *heldBy) Error() string { return fmt.Sprintf("%v: %s, epoch %d", ErrTaken, e.h.Instance, e.h.Epoch) } func (e *heldBy) Unwrap() error { return ErrTaken } // TryTake takes the key if nobody holds it, once, without waiting: the epoch, or ErrTaken naming the // holder, or what the bus said. func (l *Lease) TryTake(ctx context.Context) (uint64, error) { l.mu.Lock() switch { case l.lostWhy != nil: // A lease lost is not taken again by the instance that lost it: that instance exits, and its // successor is a new candidate with a new name. defer l.mu.Unlock() return 0, fmt.Errorf("%w: it lost it — %v", ErrNotHeld, l.lostWhy) case l.held: defer l.mu.Unlock() return l.epoch, nil } l.mu.Unlock() if h, found, err := l.Current(ctx); err != nil { return 0, fmt.Errorf("who holds the lease cannot be read: %w", err) } else if found { return 0, &heldBy{h} } if err := l.aboveTheFloor(ctx); err != nil { return 0, err } now := l.o.Now() h := l.o.Holder h.Taken, h.Renewed, h.Epoch = now, now, 0 value, err := json.Marshal(h) if err != nil { return 0, err } anchor := l.o.Now() revision, err := l.kv.Create(ctx, Key, value) if errors.Is(err, jetstream.ErrKeyExists) { // Somebody took it between the read and the write: theirs. if h, found, rerr := l.Current(ctx); rerr == nil && found { return 0, &heldBy{h} } return 0, ErrTaken } if err != nil { return 0, fmt.Errorf("%w: %w", ErrUnwritable, err) } if floor, err := l.floor(ctx); err != nil { _ = l.kv.Delete(ctx, Key, jetstream.LastRevision(revision)) return 0, err } else if revision <= floor { // The bucket moved under the floor between the check and the write (another candidate raising // it again): never an epoch already issued. Given back, and taken again on the next try. _ = l.kv.Delete(ctx, Key, jetstream.LastRevision(revision)) return 0, fmt.Errorf("the lease was taken at revision %d, not above the highest epoch issued (%d); given "+ "back to be taken again", revision, floor) } l.mu.Lock() l.held, l.epoch, l.revision, l.taken, l.renewed = true, revision, revision, now, anchor l.validUntil = anchor.Add(l.ttl - l.o.Margin) l.mu.Unlock() l.o.Say("took the controller lease at epoch %d (%s)", revision, h.Instance) return revision, nil } // floor is the highest epoch issued, zero without a floor. func (l *Lease) floor(ctx context.Context) (uint64, error) { if l.o.Floor == nil { return 0, nil } floor, err := l.o.Floor(ctx) if err != nil { return 0, fmt.Errorf("the highest epoch this mesh has issued cannot be read, so no epoch can be "+ "issued that is surely higher: %w", err) } return floor, nil } // aboveTheFloor moves the bucket's revisions past the highest epoch issued, when a bucket raised again // from nothing is behind it. Only while nobody holds the key: the stream is compacted to the floor, which // sets the next revision above it. func (l *Lease) aboveTheFloor(ctx context.Context) error { floor, err := l.floor(ctx) if err != nil || floor == 0 { return err } info, err := l.stream.Info(ctx) if err != nil { return fmt.Errorf("the lease bucket's revisions cannot be read: %w", err) } if info.State.LastSeq > floor { return nil } l.o.Say("the lease bucket is at revision %d and this mesh has issued epoch %d: the bucket was raised "+ "again from nothing, so its revisions are moved past %d before the lease is taken", info.State.LastSeq, floor, floor) if err := l.stream.Purge(ctx, jetstream.WithPurgeSequence(floor+1)); err != nil { return fmt.Errorf("the lease bucket's revisions could not be moved past epoch %d: %w", floor, err) } if l.o.Moved != nil { l.o.Moved(info.State.LastSeq, floor) } return nil } // Epoch is the gate every act passes: this instance's epoch while it holds the lease and its last // renewal is in time, and an error otherwise — not held, lost, or a renewal overdue. func (l *Lease) Epoch() (uint64, error) { l.mu.Lock() defer l.mu.Unlock() switch { case l.lostWhy != nil: return 0, fmt.Errorf("%w: it lost it — %v", ErrNotHeld, l.lostWhy) case !l.held: return 0, ErrNotHeld case l.o.Now().After(l.validUntil): return 0, fmt.Errorf("%w in time: its last renewal was due by %s", ErrNotHeld, l.validUntil.UTC().Format(time.RFC3339)) } return l.epoch, nil } // Renewed is when this instance last took or renewed the key; zero before it took it. func (l *Lease) Renewed() time.Time { l.mu.Lock() defer l.mu.Unlock() return l.renewed } // Held says this instance holds the lease now, as Epoch's gate does. func (l *Lease) Held() bool { _, err := l.Epoch() return err == nil } // Lost is closed when this instance loses the lease it held. A caller exits on it. func (l *Lease) Lost() <-chan struct{} { return l.lost } // LostWhy is why it was lost, nil while it was not. func (l *Lease) LostWhy() error { l.mu.Lock() defer l.mu.Unlock() return l.lostWhy } // Keep renews the lease until ctx ends or a renewal fails; on a failure it is lost. func (l *Lease) Keep(ctx context.Context) { tick := time.NewTicker(l.o.RenewEvery) defer tick.Stop() for { select { case <-ctx.Done(): return case <-tick.C: } if err := l.Renew(ctx); err != nil { return // lost, said by lose; or stopping } } } // Renew writes the key once more at the revision this instance last wrote. A refusal or a failure is // the lease lost: past this renewal's time the key may be somebody else's, and an instance that is not // sure it holds the lease does not act on it. func (l *Lease) Renew(ctx context.Context) error { l.mu.Lock() if !l.held || l.lostWhy != nil { defer l.mu.Unlock() if l.lostWhy != nil { return l.lostWhy } return ErrNotHeld } revision, epoch, taken := l.revision, l.epoch, l.taken l.mu.Unlock() h := l.o.Holder h.Epoch, h.Taken = epoch, taken anchor := l.o.Now() h.Renewed = anchor value, err := json.Marshal(h) if err != nil { return l.lose(err) } renewing, cancel := context.WithTimeout(ctx, l.o.RenewEvery) defer cancel() next, err := l.kv.Update(renewing, Key, value, revision) if err != nil { if ctx.Err() != nil { // Stopping, not losing: Release gives it back. return ctx.Err() } why := fmt.Errorf("the renewal at revision %d was refused: %w", revision, err) if !errors.Is(err, jetstream.ErrKeyExists) && !isWrongLast(err) { why = fmt.Errorf("the renewal could not be made: %w", err) } return l.lose(why) } l.mu.Lock() l.revision, l.renewed = next, anchor l.validUntil = anchor.Add(l.ttl - l.o.Margin) l.mu.Unlock() return nil } // isWrongLast is the bus refusing a compare-and-set because the key moved. func isWrongLast(err error) bool { var apiErr *jetstream.APIError return errors.As(err, &apiErr) && apiErr.ErrorCode == jetstream.JSErrCodeStreamWrongLastSequence } // lose marks the lease lost, once, and says why. func (l *Lease) lose(why error) error { l.mu.Lock() defer l.mu.Unlock() if l.lostWhy == nil { l.lostWhy = why l.held = false close(l.lost) l.o.Say("LOST the controller lease (epoch %d): %v — this controller stops acting at once", l.epoch, why) } return l.lostWhy } // Release gives the lease back, so the next candidate takes it at once rather than after its age. // Only the key this instance last wrote is deleted: one that moved is somebody else's. func (l *Lease) Release(ctx context.Context) { l.mu.Lock() held, revision, epoch := l.held && l.lostWhy == nil, l.revision, l.epoch l.held = false if l.lostWhy == nil { l.lostWhy = errors.New("given back") close(l.lost) } l.mu.Unlock() if !held { return } releasing, cancel := context.WithTimeout(context.WithoutCancel(ctx), 5*time.Second) defer cancel() err := l.kv.Delete(releasing, Key, jetstream.LastRevision(revision)) if err != nil && isWrongLast(err) { // A renewal was in flight as this began and landed first: the key moved, and is still this // instance's if it still names it at this epoch. Deleted at the revision it has now. if entry, gerr := l.kv.Get(releasing, Key); gerr == nil { var h Holder if json.Unmarshal(entry.Value(), &h) == nil && h.Instance == l.o.Holder.Instance && (h.Epoch == epoch || h.Epoch == 0 && entry.Revision() == epoch) { err = l.kv.Delete(releasing, Key, jetstream.LastRevision(entry.Revision())) } } } if err != nil { l.o.Say("the controller lease (epoch %d) could not be given back, so the next controller waits for "+ "it to expire: %v", epoch, err) return } l.o.Say("gave the controller lease back (epoch %d)", epoch) }