Research 031 counted the repairs people made by hand: a push to unstick a plan waiting on a report, a controller restarted to make an object again, a plan closed, a consumer re-made from now. Each was the ordinary path taken again by someone who noticed. The healer registry makes each a registered response to one condition kind, with a budget, a settle and its event: - H1 sent-not-reported: ask the machine's node-engine to report again (mesh.node.<n>.ask.report); if it does not report what it was sent, send it again, never moving a build a policy or a plan holds back - H2 stalled: close a plan whose wait is superseded or finished - H3 holder-silent / consumer-lost: the send's own assertion of the bus's objects (issue 208's note) - H4 consumer-behind: consumer-reset, only for a consumer the stream table marks resettable (the controller's own events consumer) - H5 is the identity provider's own repair (ADR 0224 §5), registered only Success is the observation clearing the condition, never the healer; a spent budget hands the condition to the operator, urgent, with what was tried, and no healer touches it again. Every act is begun in the store before it is made (migration 0070), kept in the condition's tried as "healer Hn" and said as the seat event healer-acted; a heal is never a hand act. More than twelve acts in an hour stop every healer until an hour after the last, said urgently. Only the lease holder heals. S15 is live: a cause repaired by hand twice in a fortnight raises healer-wanted, naming the healer that was not enough where one exists. D6's far-behind finding has its own kind, consumer-behind. Nodes are granted the question; the controller's grant gains healer-acted (genesis lock in mesh-host). `healers` lists the registry, the acts and the brake; status counts the week's heals.
630 lines
21 KiB
Go
630 lines
21 KiB
Go
package conditions
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"strings"
|
|
"sync"
|
|
"time"
|
|
)
|
|
|
|
// Backend is where the open conditions are kept: one value per key, written by compare-and-set.
|
|
type Backend interface {
|
|
// Get is one key's value and revision; false when it holds none.
|
|
Get(ctx context.Context, key string) (Entry, bool, error)
|
|
// Create writes a key that holds nothing, and fails with ErrMoved when it holds something.
|
|
Create(ctx context.Context, key string, value []byte) error
|
|
// Update writes a key at the revision it was read at, and fails with ErrMoved when it moved.
|
|
Update(ctx context.Context, key string, value []byte, revision uint64) error
|
|
// Delete removes a key at the revision it was read at, and fails with ErrMoved when it moved.
|
|
Delete(ctx context.Context, key string, revision uint64) error
|
|
// All is every key's value. An error is an error: never an empty store (ADR 0227 rule 4).
|
|
All(ctx context.Context) (map[string]Entry, error)
|
|
}
|
|
|
|
// Entry is one key's value, at a revision.
|
|
type Entry struct {
|
|
Value []byte
|
|
Revision uint64
|
|
}
|
|
|
|
// ErrMoved is a compare-and-set that lost: somebody wrote the key since it was read.
|
|
var ErrMoved = errors.New("the condition was written by somebody else since it was read")
|
|
|
|
// History keeps every transition (to-be 45 §2): appended, read back from a moment.
|
|
type History interface {
|
|
Append(ctx context.Context, e Event) error
|
|
// Since is every transition from a moment, oldest first.
|
|
Since(ctx context.Context, since time.Time) ([]Event, error)
|
|
}
|
|
|
|
// Teller says a transition on the bus, as the mesh-controller seat's event. The link's bus is one.
|
|
type Teller interface {
|
|
PublishSeatEvent(ctx context.Context, seat, event string, body []byte) error
|
|
}
|
|
|
|
// Seat is the role the events are said under (novox/hq ADR 0134): the control plane's.
|
|
const Seat = "mesh-controller"
|
|
|
|
// Keeper raises, observes, silences and clears conditions, and says each transition.
|
|
type Keeper struct {
|
|
store Backend
|
|
history History
|
|
teller Teller
|
|
now func() time.Time
|
|
say func(format string, args ...any)
|
|
changed func()
|
|
epoch func() (uint64, error)
|
|
|
|
mu sync.Mutex
|
|
// cleared is when each recently cleared condition cleared and how often it had been raised, so
|
|
// one raised again within ReopenWithin is the same one again.
|
|
cleared map[string]clearing
|
|
|
|
// out is the transitions still to be said and kept, in order: said by one goroutine, so a
|
|
// condition's events arrive in the order they happened, and offered again while the bus is away.
|
|
out chan Event
|
|
drained chan struct{}
|
|
closing sync.Once
|
|
// Unsaid counts the transitions given up on, for the self-check to say.
|
|
unsaid int
|
|
}
|
|
|
|
type clearing struct {
|
|
at time.Time
|
|
count int
|
|
silenced *Silence
|
|
// tried is what healers tried before it cleared: a reopening is the same fault, and what was
|
|
// tried on it is still what was tried.
|
|
tried []Attempt
|
|
}
|
|
|
|
// Options are what a Keeper is made with.
|
|
type Options struct {
|
|
Store Backend
|
|
History History
|
|
// Teller says the transitions; nil says nothing (a test, or a command run with no bus to say on).
|
|
Teller Teller
|
|
Now func() time.Time
|
|
// Say is where a transition that could not be said or kept is said instead.
|
|
Say func(format string, args ...any)
|
|
// Changed is told of every transition, at once — for `status`, which leads with what is open.
|
|
Changed func()
|
|
// Epoch is the controller lease this keeper writes under (to-be 45 §6): every write carries its
|
|
// epoch, and none is made while it is not held. Nil writes with no epoch (a test, the memory store).
|
|
Epoch func() (uint64, error)
|
|
}
|
|
|
|
// TellFor is how long one transition is offered to the bus before it is said lost.
|
|
var TellFor = 10 * time.Minute
|
|
|
|
// NewKeeper is a keeper over a store. It reads what cleared lately from the history, so a condition
|
|
// that cleared just before this controller started and is raised again now is a reopening.
|
|
func NewKeeper(ctx context.Context, o Options) *Keeper {
|
|
k := &Keeper{store: o.Store, history: o.History, teller: o.Teller, now: o.Now, say: o.Say, changed: o.Changed,
|
|
epoch: o.Epoch,
|
|
cleared: map[string]clearing{}, out: make(chan Event, 1024), drained: make(chan struct{})}
|
|
if k.now == nil {
|
|
k.now = time.Now
|
|
}
|
|
if k.say == nil {
|
|
k.say = func(string, ...any) {}
|
|
}
|
|
if k.history != nil {
|
|
if recent, err := k.history.Since(ctx, k.now().Add(-ReopenWithin)); err == nil {
|
|
for _, e := range recent {
|
|
if e.Change == ChangeCleared {
|
|
k.cleared[e.Key] = clearing{at: e.At, count: e.Condition.Count, silenced: e.Condition.Silenced,
|
|
tried: e.Condition.Tried}
|
|
}
|
|
}
|
|
} else {
|
|
k.say("what cleared lately could not be read from the condition history, so a condition "+
|
|
"raised again now is said as new rather than reopened: %v", err)
|
|
}
|
|
}
|
|
go k.telling()
|
|
return k
|
|
}
|
|
|
|
// Close says what is still to be said, waiting at most until ctx ends.
|
|
func (k *Keeper) Close(ctx context.Context) {
|
|
k.closing.Do(func() { close(k.out) })
|
|
select {
|
|
case <-k.drained:
|
|
case <-ctx.Done():
|
|
k.say("%d condition transition(s) were not yet said when this process ended", len(k.out))
|
|
}
|
|
}
|
|
|
|
// Unsaid is how many transitions were given up on since this keeper started.
|
|
func (k *Keeper) Unsaid() int {
|
|
k.mu.Lock()
|
|
defer k.mu.Unlock()
|
|
return k.unsaid
|
|
}
|
|
|
|
// stamp is the gate every write passes: the epoch it carries, set on the condition, or why this
|
|
// keeper may not write — a controller that lost the lease raises, observes and clears nothing.
|
|
func (k *Keeper) stamp(c *Condition) error {
|
|
if k.epoch == nil {
|
|
return nil
|
|
}
|
|
epoch, err := k.epoch()
|
|
if err != nil {
|
|
return fmt.Errorf("the condition %s is not written: %w", c.Key, err)
|
|
}
|
|
c.Epoch = epoch
|
|
return nil
|
|
}
|
|
|
|
// tries bounds one compare-and-set: two writers rarely race more than once.
|
|
const tries = 8
|
|
|
|
// Observe records one observation: raises the condition if it is not open, and otherwise adds the
|
|
// evidence. Says a raising, a reopening, and a change of severity or resolver; an observation that
|
|
// changes neither is written and said nowhere.
|
|
func (k *Keeper) Observe(ctx context.Context, o Observation) (Condition, error) {
|
|
if err := o.check(); err != nil {
|
|
return Condition{}, err
|
|
}
|
|
key := o.Key()
|
|
for i := 0; i < tries; i++ {
|
|
now := k.now().UTC()
|
|
said := o.Said
|
|
if said == "" {
|
|
said = o.Summary
|
|
}
|
|
entry, found, err := k.store.Get(ctx, key)
|
|
if err != nil {
|
|
return Condition{}, fmt.Errorf("reading the condition %s: %w", key, err)
|
|
}
|
|
if !found {
|
|
c := Condition{Key: key, Kind: o.Kind, Subject: Subject{Scope: o.Scope, ID: o.ID, Machine: o.Machine, Also: o.Also},
|
|
Severity: o.Severity, Summary: o.Summary, Evidence: []Evidence{{At: now, Said: said}},
|
|
Source: o.Source, Raised: now, LastObserved: now, Observations: 1, Count: 1,
|
|
Resolver: orSelf(o.Resolver)}
|
|
change := ChangeRaised
|
|
k.mu.Lock()
|
|
if before, ok := k.cleared[key]; ok && now.Sub(before.at) <= ReopenWithin {
|
|
c.Count, change = before.count+1, ChangeReopened
|
|
// A silence a person gave the condition before it cleared still holds: they said
|
|
// they knew, and the same fault again ten minutes later is what they knew about.
|
|
if before.silenced != nil && now.Before(before.silenced.Until) {
|
|
c.Silenced = before.silenced
|
|
}
|
|
c.Tried = before.tried
|
|
}
|
|
k.mu.Unlock()
|
|
if err := k.stamp(&c); err != nil {
|
|
return Condition{}, err
|
|
}
|
|
body, err := json.Marshal(c)
|
|
if err != nil {
|
|
return Condition{}, err
|
|
}
|
|
if err := k.store.Create(ctx, key, body); errors.Is(err, ErrMoved) {
|
|
continue
|
|
} else if err != nil {
|
|
return Condition{}, fmt.Errorf("raising the condition %s: %w", key, err)
|
|
}
|
|
k.mu.Lock()
|
|
delete(k.cleared, key)
|
|
k.mu.Unlock()
|
|
k.tell(Event{Condition: c, At: now, Change: change})
|
|
return c, nil
|
|
}
|
|
var c Condition
|
|
if err := json.Unmarshal(entry.Value, &c); err != nil {
|
|
return Condition{}, fmt.Errorf("the condition %s on the bus cannot be read: %w", key, err)
|
|
}
|
|
var changes []Event
|
|
// A healer's budget spent made it urgent; the watchdog seeing it again does not undo that.
|
|
if o.Severity != c.Severity && !(c.Escalated() && o.Severity != Urgent) {
|
|
changes = append(changes, Event{Change: ChangeSeverity, Was: string(c.Severity)})
|
|
c.Severity = o.Severity
|
|
}
|
|
if r := orSelf(o.Resolver); o.Resolver != "" && r != c.Resolver {
|
|
changes = append(changes, Event{Change: ChangeResolver, Was: c.Resolver})
|
|
c.Resolver = r
|
|
}
|
|
// The kind as the source says it now: a source that gave the same key a kind of its own since
|
|
// (a probe's finding split out for a healer) is read by that kind from its next observation.
|
|
c.Kind, c.Summary, c.Source, c.LastObserved = o.Kind, o.Summary, o.Source, now
|
|
if o.Machine != "" {
|
|
c.Subject.Machine = o.Machine
|
|
}
|
|
if len(o.Also) > 0 {
|
|
c.Subject.Also = o.Also
|
|
}
|
|
c.Observations++
|
|
c.Evidence = append([]Evidence{{At: now, Said: said}}, c.Evidence...)
|
|
if len(c.Evidence) > KeptEvidence {
|
|
c.Evidence = c.Evidence[:KeptEvidence]
|
|
}
|
|
if err := k.stamp(&c); err != nil {
|
|
return Condition{}, err
|
|
}
|
|
body, err := json.Marshal(c)
|
|
if err != nil {
|
|
return Condition{}, err
|
|
}
|
|
if err := k.store.Update(ctx, key, body, entry.Revision); errors.Is(err, ErrMoved) {
|
|
continue
|
|
} else if err != nil {
|
|
return Condition{}, fmt.Errorf("observing the condition %s: %w", key, err)
|
|
}
|
|
for _, e := range changes {
|
|
e.At, e.Condition = now, c
|
|
k.tell(e)
|
|
}
|
|
return c, nil
|
|
}
|
|
return Condition{}, fmt.Errorf("the condition %s kept moving under this write; %d tries", key, tries)
|
|
}
|
|
|
|
// Clear removes a condition an observation says is resolved, and says so. False when none was open.
|
|
func (k *Keeper) Clear(ctx context.Context, key, why string) (bool, error) {
|
|
for i := 0; i < tries; i++ {
|
|
entry, found, err := k.store.Get(ctx, key)
|
|
if err != nil {
|
|
return false, fmt.Errorf("reading the condition %s: %w", key, err)
|
|
}
|
|
if !found {
|
|
return false, nil
|
|
}
|
|
var c Condition
|
|
if err := json.Unmarshal(entry.Value, &c); err != nil {
|
|
// Unreadable is not resolved: kept, and said, rather than removed unread.
|
|
return false, fmt.Errorf("the condition %s on the bus cannot be read, so it is not cleared: %w", key, err)
|
|
}
|
|
if err := k.stamp(&c); err != nil {
|
|
return false, err
|
|
}
|
|
if err := k.store.Delete(ctx, key, entry.Revision); errors.Is(err, ErrMoved) {
|
|
continue
|
|
} else if err != nil {
|
|
return false, fmt.Errorf("clearing the condition %s: %w", key, err)
|
|
}
|
|
now := k.now().UTC()
|
|
k.mu.Lock()
|
|
k.cleared[key] = clearing{at: now, count: c.Count, silenced: c.Silenced, tried: c.Tried}
|
|
k.mu.Unlock()
|
|
k.tell(Event{Condition: c, At: now, Change: ChangeCleared, Why: why, Cleared: &now})
|
|
return true, nil
|
|
}
|
|
return false, fmt.Errorf("the condition %s kept moving under this clearing; %d tries", key, tries)
|
|
}
|
|
|
|
// Tried records a healer's attempt on an open condition (to-be 45 §7) and makes the healer its
|
|
// resolver: said as `condition-changed` when the resolver changes, kept in `tried` either way. **It
|
|
// never clears the condition**: a repair that worked is seen by the observation that raised it, which
|
|
// clears it — a healer marking its own work done would be a second opinion of the fact. False when no
|
|
// condition is open under the key: it cleared meanwhile, and there is nothing to record against.
|
|
func (k *Keeper) Tried(ctx context.Context, key string, a Attempt, resolver string) (Condition, bool, error) {
|
|
return k.amend(ctx, key, func(c *Condition, now time.Time) []Event {
|
|
c.Tried = appendAttempt(c.Tried, a, now)
|
|
var changes []Event
|
|
if resolver != "" && resolver != c.Resolver && !c.Escalated() {
|
|
changes = append(changes, Event{Change: ChangeResolver, Was: c.Resolver, Why: a.Outcome})
|
|
c.Resolver = resolver
|
|
}
|
|
return changes
|
|
})
|
|
}
|
|
|
|
// Escalate is a healer's budget spent, or a repair it may not make (to-be 45 §2, §7): the attempt that
|
|
// says so is kept in `tried`, the resolver becomes the operator and the severity urgent, each said as
|
|
// `condition-changed`. Observation still clears it when the fault goes; nothing else does.
|
|
func (k *Keeper) Escalate(ctx context.Context, key string, a Attempt) (Condition, bool, error) {
|
|
return k.amend(ctx, key, func(c *Condition, now time.Time) []Event {
|
|
c.Tried = appendAttempt(c.Tried, a, now)
|
|
var changes []Event
|
|
if c.Severity != Urgent {
|
|
changes = append(changes, Event{Change: ChangeSeverity, Was: string(c.Severity), Why: a.Outcome})
|
|
c.Severity = Urgent
|
|
}
|
|
if c.Resolver != ResolverOperator {
|
|
changes = append(changes, Event{Change: ChangeResolver, Was: c.Resolver, Why: a.Outcome})
|
|
c.Resolver = ResolverOperator
|
|
}
|
|
return changes
|
|
})
|
|
}
|
|
|
|
// appendAttempt adds an attempt, stamped when it has no time, keeping the newest KeptAttempts.
|
|
func appendAttempt(tried []Attempt, a Attempt, now time.Time) []Attempt {
|
|
if a.At.IsZero() {
|
|
a.At = now
|
|
}
|
|
tried = append(tried, a)
|
|
if len(tried) > KeptAttempts {
|
|
tried = tried[len(tried)-KeptAttempts:]
|
|
}
|
|
return tried
|
|
}
|
|
|
|
// amend changes one open condition by compare-and-set and says what changed. False when none is open.
|
|
func (k *Keeper) amend(ctx context.Context, key string, change func(*Condition, time.Time) []Event) (Condition, bool, error) {
|
|
for i := 0; i < tries; i++ {
|
|
entry, found, err := k.store.Get(ctx, key)
|
|
if err != nil {
|
|
return Condition{}, false, fmt.Errorf("reading the condition %s: %w", key, err)
|
|
}
|
|
if !found {
|
|
return Condition{}, false, nil
|
|
}
|
|
var c Condition
|
|
if err := json.Unmarshal(entry.Value, &c); err != nil {
|
|
return Condition{}, false, fmt.Errorf("the condition %s on the bus cannot be read: %w", key, err)
|
|
}
|
|
now := k.now().UTC()
|
|
changes := change(&c, now)
|
|
if err := k.stamp(&c); err != nil {
|
|
return Condition{}, false, err
|
|
}
|
|
body, err := json.Marshal(c)
|
|
if err != nil {
|
|
return Condition{}, false, err
|
|
}
|
|
if err := k.store.Update(ctx, key, body, entry.Revision); errors.Is(err, ErrMoved) {
|
|
continue
|
|
} else if err != nil {
|
|
return Condition{}, false, fmt.Errorf("writing the condition %s: %w", key, err)
|
|
}
|
|
for _, e := range changes {
|
|
e.At, e.Condition = now, c
|
|
k.tell(e)
|
|
}
|
|
return c, true, nil
|
|
}
|
|
return Condition{}, false, fmt.Errorf("the condition %s kept moving under this write; %d tries", key, tries)
|
|
}
|
|
|
|
// Reconcile is one source's whole observation: every condition it observes is observed, and every
|
|
// condition it raised before and no longer observes is cleared — the observation says it is
|
|
// resolved. A source that could not observe must not call this: an empty observation clears all it
|
|
// raised, which is exactly the fault of saying "none" for "I could not tell" (ADR 0227 rule 4).
|
|
func (k *Keeper) Reconcile(ctx context.Context, source string, observed []Observation) error {
|
|
all, err := k.Open(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
seen := map[string]bool{}
|
|
var problems []string
|
|
for _, o := range observed {
|
|
o.Source = source
|
|
seen[o.Key()] = true
|
|
if _, err := k.Observe(ctx, o); err != nil {
|
|
problems = append(problems, err.Error())
|
|
}
|
|
}
|
|
for _, c := range all {
|
|
if c.Source != source || seen[c.Key] {
|
|
continue
|
|
}
|
|
if _, err := k.Clear(ctx, c.Key, source+" no longer observes it"); err != nil {
|
|
problems = append(problems, err.Error())
|
|
}
|
|
}
|
|
if len(problems) > 0 {
|
|
return errors.New(strings.Join(problems, "; "))
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// Silence stops a condition's messages for a while, with a reason, by somebody (to-be 45 §2). The
|
|
// condition stays open and `status` still says it; recording the act in the hand-act log is the
|
|
// caller's, which knows who acted.
|
|
func (k *Keeper) Silence(ctx context.Context, key string, d time.Duration, by, why string) (Condition, error) {
|
|
if strings.TrimSpace(why) == "" {
|
|
return Condition{}, errors.New("a silence says why: --why <text>")
|
|
}
|
|
if d <= 0 || d > MaxSilence {
|
|
return Condition{}, fmt.Errorf("a condition is silenced for a while, at most %s — not %s", MaxSilence, d)
|
|
}
|
|
for i := 0; i < tries; i++ {
|
|
entry, found, err := k.store.Get(ctx, key)
|
|
if err != nil {
|
|
return Condition{}, fmt.Errorf("reading the condition %s: %w", key, err)
|
|
}
|
|
if !found {
|
|
return Condition{}, fmt.Errorf("no condition %s is open — `conditions` lists them", key)
|
|
}
|
|
var c Condition
|
|
if err := json.Unmarshal(entry.Value, &c); err != nil {
|
|
return Condition{}, fmt.Errorf("the condition %s on the bus cannot be read: %w", key, err)
|
|
}
|
|
now := k.now().UTC()
|
|
c.Silenced = &Silence{Until: now.Add(d), By: by, Why: strings.TrimSpace(why), Since: now}
|
|
if err := k.stamp(&c); err != nil {
|
|
return Condition{}, err
|
|
}
|
|
body, err := json.Marshal(c)
|
|
if err != nil {
|
|
return Condition{}, err
|
|
}
|
|
if err := k.store.Update(ctx, key, body, entry.Revision); errors.Is(err, ErrMoved) {
|
|
continue
|
|
} else if err != nil {
|
|
return Condition{}, fmt.Errorf("silencing the condition %s: %w", key, err)
|
|
}
|
|
k.tell(Event{Condition: c, At: now, Change: ChangeSilenced, Why: c.Silenced.Why})
|
|
return c, nil
|
|
}
|
|
return Condition{}, fmt.Errorf("the condition %s kept moving under this silence; %d tries", key, tries)
|
|
}
|
|
|
|
// EndSilences ends every silence that has run out, and says each: the condition is still open, and
|
|
// its messages start again.
|
|
func (k *Keeper) EndSilences(ctx context.Context) error {
|
|
all, err := k.Open(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
now := k.now().UTC()
|
|
for _, c := range all {
|
|
if c.Silenced == nil || now.Before(c.Silenced.Until) {
|
|
continue
|
|
}
|
|
for i := 0; i < tries; i++ {
|
|
entry, found, err := k.store.Get(ctx, c.Key)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if !found {
|
|
break
|
|
}
|
|
var held Condition
|
|
if err := json.Unmarshal(entry.Value, &held); err != nil {
|
|
return fmt.Errorf("the condition %s on the bus cannot be read: %w", c.Key, err)
|
|
}
|
|
if held.Silenced == nil || now.Before(held.Silenced.Until) {
|
|
break
|
|
}
|
|
was := held.Silenced.Why
|
|
held.Silenced = nil
|
|
if err := k.stamp(&held); err != nil {
|
|
return err
|
|
}
|
|
body, err := json.Marshal(held)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if err := k.store.Update(ctx, c.Key, body, entry.Revision); errors.Is(err, ErrMoved) {
|
|
continue
|
|
} else if err != nil {
|
|
return err
|
|
}
|
|
k.tell(Event{Condition: held, At: now, Change: ChangeUnsilenced, Why: was})
|
|
break
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// Open is every open condition, urgent first and then oldest first.
|
|
func (k *Keeper) Open(ctx context.Context) ([]Condition, error) {
|
|
return Read(ctx, k.store)
|
|
}
|
|
|
|
// Get is one open condition.
|
|
func (k *Keeper) Get(ctx context.Context, key string) (Condition, bool, error) {
|
|
return ReadOne(ctx, k.store, key)
|
|
}
|
|
|
|
// HistorySince is every transition from a moment, oldest first.
|
|
func (k *Keeper) HistorySince(ctx context.Context, since time.Time) ([]Event, error) {
|
|
if k.history == nil {
|
|
return nil, errors.New("this keeper has no history to read")
|
|
}
|
|
return k.history.Since(ctx, since)
|
|
}
|
|
|
|
// Read is every open condition in a store, in the order status says them. A value that cannot be
|
|
// read is an error naming its key, never a condition left out (ADR 0227 rule 4).
|
|
func Read(ctx context.Context, store Backend) ([]Condition, error) {
|
|
all, err := store.All(ctx)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("the open conditions cannot be read: %w", err)
|
|
}
|
|
out := make([]Condition, 0, len(all))
|
|
for key, e := range all {
|
|
var c Condition
|
|
if err := json.Unmarshal(e.Value, &c); err != nil {
|
|
return nil, fmt.Errorf("the condition %s cannot be read: %w", key, err)
|
|
}
|
|
out = append(out, c)
|
|
}
|
|
Order(out)
|
|
return out, nil
|
|
}
|
|
|
|
// ReadOne is one open condition from a store.
|
|
func ReadOne(ctx context.Context, store Backend, key string) (Condition, bool, error) {
|
|
e, found, err := store.Get(ctx, key)
|
|
if err != nil || !found {
|
|
return Condition{}, found, err
|
|
}
|
|
var c Condition
|
|
if err := json.Unmarshal(e.Value, &c); err != nil {
|
|
return Condition{}, false, fmt.Errorf("the condition %s cannot be read: %w", key, err)
|
|
}
|
|
return c, true, nil
|
|
}
|
|
|
|
// tell queues a transition to be kept and said. Never blocks the caller for long: a queue that is
|
|
// full is a bus away for a long time, and the transition is said lost rather than holding a watchdog.
|
|
func (k *Keeper) tell(e Event) {
|
|
e.Event = eventFor(e.Change)
|
|
e.Show = e.Condition.Show()
|
|
if k.changed != nil {
|
|
k.changed()
|
|
}
|
|
defer func() {
|
|
// A keeper closed while a write was in flight: said, not a panic.
|
|
if recover() != nil {
|
|
k.lost(e, errors.New("the keeper was closed"))
|
|
}
|
|
}()
|
|
select {
|
|
case k.out <- e:
|
|
default:
|
|
k.lost(e, errors.New("too many transitions are waiting to be said"))
|
|
}
|
|
}
|
|
|
|
func (k *Keeper) lost(e Event, err error) {
|
|
k.mu.Lock()
|
|
k.unsaid++
|
|
k.mu.Unlock()
|
|
k.say("the condition %s was %s and that could NOT be said or kept: %v", e.Key, e.Change, err)
|
|
}
|
|
|
|
// telling keeps and says every transition in order, offering each again while the bus is away.
|
|
func (k *Keeper) telling() {
|
|
defer close(k.drained)
|
|
for e := range k.out {
|
|
body, err := json.Marshal(e)
|
|
if err != nil {
|
|
k.lost(e, err)
|
|
continue
|
|
}
|
|
deadline := time.Now().Add(TellFor)
|
|
wait := 200 * time.Millisecond
|
|
kept, said := k.history == nil, k.teller == nil
|
|
for {
|
|
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
|
|
if !kept {
|
|
kept = k.history.Append(ctx, e) == nil
|
|
}
|
|
if !said {
|
|
said = k.teller.PublishSeatEvent(ctx, Seat, e.Event, body) == nil
|
|
}
|
|
cancel()
|
|
if kept && said {
|
|
break
|
|
}
|
|
if time.Now().After(deadline) {
|
|
what := "said"
|
|
if !kept {
|
|
what = "kept in the history"
|
|
}
|
|
k.lost(e, fmt.Errorf("not %s within %s", what, TellFor))
|
|
break
|
|
}
|
|
time.Sleep(wait)
|
|
wait = min(2*wait, 10*time.Second)
|
|
}
|
|
}
|
|
}
|
|
|
|
func orSelf(resolver string) string {
|
|
if resolver == "" {
|
|
return ResolverSelf
|
|
}
|
|
return resolver
|
|
}
|