plans said a build queued for minutes had been building for a few seconds: the line counted from the last save, and every advance saved the plan whether or not it moved, bumping its revision and saying plan-moved on the bus. The line, status and LATE now count from when the plan entered its tier with the stalled condition's bound, and an advance that changes nothing writes nothing.
697 lines
23 KiB
Go
697 lines
23 KiB
Go
package main
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"os"
|
|
"slices"
|
|
"sort"
|
|
"sync"
|
|
"time"
|
|
|
|
"github.com/nats-io/nats.go"
|
|
|
|
"github.com/novox/mesh-controller/internal/broker"
|
|
"github.com/novox/mesh-controller/internal/conditions"
|
|
"github.com/novox/mesh-controller/internal/inventory"
|
|
"github.com/novox/mesh-controller/internal/link"
|
|
)
|
|
|
|
// The watchdogs (novox/hq to-be 45 §3): every row of the signals table, run over what the serving
|
|
// controller knows, every half minute.
|
|
//
|
|
// **One loop, one gathering, every row.** The facts are gathered once a tick — the store's machines,
|
|
// plans and durations, the bus's queue and consumers, what this process heard — and each row's watch
|
|
// is a pure reading of them, which is what lets the test generated from the table suppress a signal by
|
|
// changing a fact. A part that cannot be gathered makes the rows that read it blind: each says so as
|
|
// a condition of its own (`probe.<row>.failed`) and keeps what it raised before, because a watchdog
|
|
// that cannot see must not read as one that sees nothing wrong (ADR 0227 rule 4).
|
|
//
|
|
// **The watchdogs are themselves watched.** Each tick is recorded; the self-check (doctor.go) fails
|
|
// its probe DW when the ticks stop, and the watchdogs raise S10 when the self-check stops — two loops
|
|
// watching each other, and mesh-watcher on a second machine watching the self-check's heartbeat for
|
|
// the case both stop with the process.
|
|
|
|
// watchEvery is how often the watchdogs run.
|
|
var watchEvery = 30 * time.Second
|
|
|
|
// signalFacts is what one tick of the watchdogs reads.
|
|
type signalFacts struct {
|
|
now time.Time
|
|
started time.Time
|
|
// host is the machine this controller runs on, as the mesh names it, for a condition about itself.
|
|
host string
|
|
|
|
machines []machineFacts
|
|
machinesErr error
|
|
// toolsHeardFrom is when this process began hearing node tools: one not heard since is silent
|
|
// since then at the most.
|
|
toolsHeardFrom time.Time
|
|
|
|
plans []planFacts
|
|
plansErr error
|
|
// waits are the walks waiting for their delivery's word (S16, novox/hq ADR 0239).
|
|
waits []waitFacts
|
|
|
|
loop loopFacts
|
|
loopErr error
|
|
|
|
merges []missedMerge
|
|
mergesPassed time.Time
|
|
mergesErr error
|
|
|
|
asks []askFacts
|
|
asksErr error
|
|
|
|
calls []link.Call
|
|
|
|
standings []conditions.Condition
|
|
standingsErr error
|
|
// providerWords are every open condition a provider's word raised — failing, waiting for a person
|
|
// to approve a retirement, kept by a rejection, cleanup waiting (ADR 0224, 0230) — so a provider no
|
|
// longer assigned has them all cleared.
|
|
providerWords []conditions.Condition
|
|
|
|
advisories []link.Advisory
|
|
lostConsumers map[string]bool
|
|
advisoriesErr error
|
|
|
|
selfCheck selfCheckFacts
|
|
|
|
// staleRefusals are the writers refused as older lately, and epochs the mesh's record of each epoch
|
|
// they name (novox/hq to-be 45 §6, S13).
|
|
staleRefusals []link.WriterRefusals
|
|
epochs map[int64]inventory.Epoch
|
|
|
|
// handActs are the acts done by hand within the fortnight S15 counts.
|
|
handActs []link.HandAct
|
|
handActsErr error
|
|
|
|
// lease is this controller's standing to the lease, and the epochs that ended lately (S12).
|
|
lease leaseFacts
|
|
leaseErr error
|
|
|
|
// facts is the snapshot this controller keeps for merge checks (S14).
|
|
facts factsFacts
|
|
}
|
|
|
|
// factsFacts is when the newest facts snapshot was taken, when this controller began keeping it, and
|
|
// the last attempt's error.
|
|
type factsFacts struct {
|
|
taken, began time.Time
|
|
err error
|
|
}
|
|
|
|
type leaseFacts struct {
|
|
held bool
|
|
epoch uint64
|
|
renewed time.Time
|
|
unleased string
|
|
ended []inventory.Epoch
|
|
// reset is when the lease bucket was found raised again from nothing; resetSaid what of it.
|
|
reset time.Time
|
|
resetSaid string
|
|
}
|
|
|
|
type machineFacts struct {
|
|
name string
|
|
control bool
|
|
// lastHeard is the store's last word from it, any word; zero when never.
|
|
lastHeard time.Time
|
|
// every is the heartbeat interval it said; zero when it said none.
|
|
every time.Duration
|
|
power link.PowerState
|
|
// sentAt is when it was last sent a declaration; reportedCurrent whether it has reported that one.
|
|
sentAt time.Time
|
|
reportedCurrent bool
|
|
reportedAt time.Time
|
|
// lastApply is its newest measured apply.
|
|
lastApply time.Duration
|
|
// tools is whether node-tools is assigned there; toolsHeard and toolsEvery its heartbeat.
|
|
tools bool
|
|
toolsHeard time.Time
|
|
toolsEvery time.Duration
|
|
}
|
|
|
|
// asleep says the machine said it would be away and has not said it is back (ADR 0211).
|
|
func (m machineFacts) asleep() bool { return m.power.Away() }
|
|
|
|
type planFacts struct {
|
|
id, repository, commit string
|
|
tier, tiers int
|
|
entered time.Time
|
|
bound time.Duration
|
|
waiting string
|
|
paused bool
|
|
}
|
|
|
|
// waitFacts is one walk waiting for its delivery's word: since its merge opened it.
|
|
type waitFacts struct {
|
|
id, repository, commit, awaits string
|
|
since time.Time
|
|
}
|
|
|
|
type loopFacts struct {
|
|
took time.Time
|
|
pending uint64
|
|
where string
|
|
}
|
|
|
|
type askFacts struct {
|
|
id, seat, what, state, on string
|
|
since time.Time
|
|
bound time.Duration
|
|
}
|
|
|
|
type selfCheckFacts struct {
|
|
last time.Time
|
|
every time.Duration
|
|
}
|
|
|
|
// watchdogs is the loop and what it needs.
|
|
type watchdogs struct {
|
|
open *stores
|
|
server *link.Server
|
|
js *broker.JetStream
|
|
keeper *conditions.Keeper
|
|
doctor *doctor
|
|
|
|
started time.Time
|
|
// acting says this controller is the one acting, not one standing by (link.Holding): a controller
|
|
// standing by hears no heartbeat and would call every machine silent.
|
|
acting func() bool
|
|
|
|
// confirm holds back a row blind for one tick: a read that timed out once on a loaded store.
|
|
confirm confirming
|
|
|
|
mu sync.Mutex
|
|
ticked time.Time
|
|
last *signalFacts
|
|
failed string
|
|
}
|
|
|
|
// lastTick is when the watchdogs last finished a tick.
|
|
func (w *watchdogs) lastTick() time.Time {
|
|
w.mu.Lock()
|
|
defer w.mu.Unlock()
|
|
return w.ticked
|
|
}
|
|
|
|
// lastFacts is the facts of the newest tick, for `doctor signals`; nil before the first.
|
|
func (w *watchdogs) lastFacts() *signalFacts {
|
|
w.mu.Lock()
|
|
defer w.mu.Unlock()
|
|
return w.last
|
|
}
|
|
|
|
// keep runs the watchdogs until ctx ends.
|
|
func (w *watchdogs) keep(ctx context.Context) {
|
|
tick := time.NewTicker(watchEvery)
|
|
defer tick.Stop()
|
|
for {
|
|
w.tick(ctx)
|
|
select {
|
|
case <-ctx.Done():
|
|
return
|
|
case <-tick.C:
|
|
}
|
|
}
|
|
}
|
|
|
|
// tick is one run of every row.
|
|
func (w *watchdogs) tick(ctx context.Context) {
|
|
if w.acting != nil && !w.acting() {
|
|
// Standing by: nothing seen, nothing said, and the tick counted — this process's self-check
|
|
// is not the one that matters while another acts.
|
|
w.mu.Lock()
|
|
w.ticked = time.Now()
|
|
w.mu.Unlock()
|
|
return
|
|
}
|
|
running, cancel := context.WithTimeout(ctx, watchEvery)
|
|
defer cancel()
|
|
w.see(running, w.gather(running))
|
|
}
|
|
|
|
// see runs every watched row over one gathering, keeps what each found, and records the tick.
|
|
func (w *watchdogs) see(running context.Context, f *signalFacts) {
|
|
var problems []string
|
|
var blind []conditions.Observation
|
|
for _, row := range watchedRows() {
|
|
if err := row.needs(f); err != nil {
|
|
blind = append(blind, blindRow(row, err))
|
|
continue
|
|
}
|
|
if err := w.keeper.Reconcile(running, row.Row, row.watch(f)); err != nil {
|
|
problems = append(problems, row.Row+": "+err.Error())
|
|
}
|
|
}
|
|
// The rows that could not see, said once they could not two ticks in a row; the ones that see
|
|
// again, cleared.
|
|
blind, _ = w.confirm.pass(running, w.keeper, sourceWatchdogs, blind)
|
|
if err := w.keeper.Reconcile(running, sourceWatchdogs, blind); err != nil {
|
|
problems = append(problems, err.Error())
|
|
}
|
|
// A provider no longer assigned where it ran: its standing is resolved by the assignment.
|
|
if f.standingsErr == nil && w.open != nil {
|
|
if err := unassignedProviders(running, w.open.inventory, w.keeper, f.providerWords); err != nil {
|
|
problems = append(problems, "S8: "+err.Error())
|
|
}
|
|
}
|
|
if err := w.keeper.EndSilences(running); err != nil {
|
|
problems = append(problems, err.Error())
|
|
}
|
|
failed := ""
|
|
if len(problems) > 0 {
|
|
failed = fmt.Sprintf("%v", problems)
|
|
}
|
|
w.mu.Lock()
|
|
said := w.failed
|
|
w.ticked, w.last, w.failed = time.Now(), f, failed
|
|
w.mu.Unlock()
|
|
if failed != said {
|
|
if failed != "" {
|
|
fmt.Printf("the watchdogs could not keep what they saw: %s\n", failed)
|
|
} else if said != "" {
|
|
fmt.Println("the watchdogs keep what they see again")
|
|
}
|
|
}
|
|
}
|
|
|
|
// sourceWatchdogs is what raises a blind row's condition.
|
|
const sourceWatchdogs = "watchdogs"
|
|
|
|
// blindRow is a row whose facts could not be gathered, as a condition of its own: raised when the next
|
|
// tick cannot gather them either (confirm.go), since one read that did not answer in time is not a
|
|
// watchdog gone blind.
|
|
func blindRow(row signalRow, err error) conditions.Observation {
|
|
return conditions.Observation{Scope: conditions.ScopeProbe, ID: row.Row, Kind: "probe-failed", Token: "failed",
|
|
Severity: conditions.Warning, Confirm: true,
|
|
Summary: fmt.Sprintf("the watchdog of %s (%s) cannot see: what it reads could not be read, so nothing "+
|
|
"it would raise can be — and nothing it raised before is cleared", row.Row, row.Signal),
|
|
Said: firstLine(err.Error())}
|
|
}
|
|
|
|
// gather reads every fact a tick needs. Each part's failure is kept beside it, never an empty part.
|
|
func (w *watchdogs) gather(ctx context.Context) *signalFacts {
|
|
now := time.Now()
|
|
f := &signalFacts{now: now, started: w.started, toolsHeardFrom: link.ToolsBeats.Started(), calls: link.Calls.Running(),
|
|
staleRefusals: link.StaleRefusals.Within(now.Add(-staleRefusalsWithin)), lostConsumers: map[string]bool{},
|
|
epochs: map[int64]inventory.Epoch{}}
|
|
if w.doctor != nil {
|
|
f.selfCheck = selfCheckFacts{last: w.doctor.lastRunEnded(), every: doctorEvery}
|
|
}
|
|
inv := w.open.inventory
|
|
f.host = controlHost(ctx, inv)
|
|
f.lease, f.leaseErr = gatherLease(ctx, inv, now)
|
|
for _, r := range f.staleRefusals {
|
|
if r.Epoch <= 0 {
|
|
continue
|
|
}
|
|
// Named where the record has it; a writer the record cannot name is still said by its epoch.
|
|
if e, found, err := inv.EpochOf(ctx, uint64(r.Epoch)); err == nil && found {
|
|
f.epochs[r.Epoch] = e
|
|
}
|
|
}
|
|
f.machines, f.machinesErr = w.gatherMachines(ctx, inv, now)
|
|
f.plans, f.waits, f.plansErr = gatherPlans(ctx, inv, now)
|
|
f.loop, f.loopErr = w.gatherLoop()
|
|
f.mergesPassed, f.merges, f.mergesErr = watchedMerges.last()
|
|
if f.mergesErr == nil && !f.mergesPassed.IsZero() && now.Sub(f.mergesPassed) > 3*mergeCatchUpEvery {
|
|
f.mergesErr = fmt.Errorf("the catch-up of merges has not passed since %s", f.mergesPassed.UTC().Format(time.RFC3339))
|
|
}
|
|
f.asks, f.asksErr = w.gatherAsks(ctx, inv)
|
|
if open, err := w.keeper.Open(ctx); err != nil {
|
|
f.standingsErr = err
|
|
} else {
|
|
f.standings = providerStandings(open)
|
|
f.providerWords = providerConditions(open)
|
|
}
|
|
f.advisories = link.Advisories.Since(now.Add(-advisoryQuiet))
|
|
f.lostConsumers, f.advisoriesErr = w.lostConsumers(ctx, f.advisories)
|
|
f.handActs, f.handActsErr = w.gatherHandActs(ctx, now)
|
|
f.facts.taken, _, f.facts.began, f.facts.err = exportedFacts.last()
|
|
return f
|
|
}
|
|
|
|
// gatherLease is this controller's standing to the lease and the epochs that ended within the hour.
|
|
func gatherLease(ctx context.Context, inv *inventory.Inventory, now time.Time) (leaseFacts, error) {
|
|
st := theLease.standing()
|
|
f := leaseFacts{held: st.Held, epoch: st.Epoch, renewed: st.Renewed, unleased: st.Unleased, reset: st.Reset,
|
|
resetSaid: st.ResetSaid}
|
|
ended, err := inv.EpochsSince(ctx, now.Add(-advisoryQuiet))
|
|
if err != nil {
|
|
return f, fmt.Errorf("the epochs the mesh issued cannot be read: %w", err)
|
|
}
|
|
f.ended = ended
|
|
return f, nil
|
|
}
|
|
|
|
// controlHost is the machine running the controller, as the mesh names it: the one the controller
|
|
// module is assigned to, or this process's host name where that is not one machine.
|
|
func controlHost(ctx context.Context, inv *inventory.Inventory) string {
|
|
if on, err := inv.Running(ctx, "mesh-controller"); err == nil && len(on) == 1 {
|
|
return on[0]
|
|
}
|
|
host, _ := os.Hostname()
|
|
return host
|
|
}
|
|
|
|
func (w *watchdogs) gatherMachines(ctx context.Context, inv *inventory.Inventory, now time.Time) ([]machineFacts, error) {
|
|
nodes, err := inv.Nodes(ctx)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("the machines cannot be read: %w", err)
|
|
}
|
|
reports, err := inv.LastReports(ctx)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("what each machine last reported cannot be read: %w", err)
|
|
}
|
|
reported := map[string]inventory.Reported{}
|
|
for _, r := range reports {
|
|
reported[r.Node] = r
|
|
}
|
|
control, err := inv.Running(ctx, "mesh-controller")
|
|
if err != nil {
|
|
return nil, fmt.Errorf("where the controller runs cannot be read: %w", err)
|
|
}
|
|
tooled, err := inv.Running(ctx, broker.RuntimeModule)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("where the node tools run cannot be read: %w", err)
|
|
}
|
|
applies, err := inv.Durations(ctx, inventory.DurationApply, now.Add(-7*24*time.Hour))
|
|
if err != nil {
|
|
return nil, fmt.Errorf("the measured applies cannot be read: %w", err)
|
|
}
|
|
lastApply := map[string]time.Duration{}
|
|
for _, d := range applies { // oldest first: the last one read is the newest
|
|
lastApply[d.Node] = d.Took
|
|
}
|
|
var out []machineFacts
|
|
anyPast := false
|
|
for _, n := range nodes {
|
|
m := machineFacts{name: n.Name, control: slices.Contains(control, n.Name), lastHeard: n.LastSeen,
|
|
lastApply: lastApply[n.Name], tools: slices.Contains(tooled, n.Name)}
|
|
if beat, ok := link.HostBeats.Of(n.Name); ok {
|
|
m.every = beat.Every
|
|
}
|
|
if beat, ok := link.ToolsBeats.Of(n.Name); ok {
|
|
m.toolsHeard, m.toolsEvery = beat.At, beat.Every
|
|
}
|
|
if r, ok := reported[n.Name]; ok {
|
|
if r.Sent != nil {
|
|
m.sentAt = *r.Sent
|
|
}
|
|
if r.At != nil {
|
|
m.reportedAt = *r.At
|
|
}
|
|
m.reportedCurrent = r.Current
|
|
}
|
|
if !m.lastHeard.IsZero() && now.Sub(m.lastHeard) > heartbeatBound(m.every) {
|
|
anyPast = true
|
|
}
|
|
if m.tools && now.Sub(later(m.toolsHeard, link.ToolsBeats.Started())) > heartbeatBound(m.toolsEvery) {
|
|
anyPast = true
|
|
}
|
|
out = append(out, m)
|
|
}
|
|
// What a machine past its bound last said of its power, read only then: a machine that said it
|
|
// is asleep is not lost (ADR 0211). Unreadable is said: a sleeping laptop is then called silent,
|
|
// which is the louder mistake and the right one to make.
|
|
if anyPast && w.server != nil {
|
|
states, err := w.server.PowerStates(ctx, now.Add(-7*24*time.Hour))
|
|
if err != nil {
|
|
fmt.Printf("what machines said of their power cannot be read, so a sleeping one is called silent: %v\n", err)
|
|
}
|
|
for i := range out {
|
|
out[i].power = states[out[i].name]
|
|
}
|
|
}
|
|
return out, nil
|
|
}
|
|
|
|
// gatherPlans is every open plan, its tier's bound from what was measured, and what it waits on.
|
|
func gatherPlans(ctx context.Context, inv *inventory.Inventory, now time.Time) ([]planFacts, []waitFacts, error) {
|
|
plans, err := inv.OpenPlans(ctx)
|
|
if err != nil {
|
|
return nil, nil, fmt.Errorf("the open plans cannot be read: %w", err)
|
|
}
|
|
if len(plans) == 0 {
|
|
return nil, nil, nil
|
|
}
|
|
bounds, err := measuredTierBounds(ctx, inv, now)
|
|
if err != nil {
|
|
return nil, nil, err
|
|
}
|
|
pause := buildSeatPause(ctx, inv, plans)
|
|
var out []planFacts
|
|
var waits []waitFacts
|
|
for _, p := range plans {
|
|
// A walk waiting for its delivery's word is no tier late (novox/hq ADR 0239); how long it has waited is
|
|
// S16's, whatever mesh-delivery says or does not say.
|
|
if p.Waiting() {
|
|
waits = append(waits, waitFacts{id: p.ID, repository: p.Repository, commit: p.Commit,
|
|
awaits: p.Delivery.Awaits, since: p.Created})
|
|
continue
|
|
}
|
|
_, paused := pausedWaiting(p, pause, now)
|
|
bound := bounds.of(p.Repository)
|
|
out = append(out, planFacts{id: p.ID, repository: p.Repository, commit: p.Commit, tier: p.Tier,
|
|
tiers: len(p.Tiers), entered: inTierSince(p), bound: bound,
|
|
waiting: planLineWith(p, now, pause, bound), paused: paused})
|
|
}
|
|
return out, waits, nil
|
|
}
|
|
|
|
// tierBounds is how long a plan's tier may take before it is stalled (S3) and `plans` and `status` call
|
|
// it LATE, by repository: three times the ninetieth percentile of the tiers measured for it, and never
|
|
// less than tierAtLeast. One bound for the watchdog and the lines a person reads, so the two cannot
|
|
// disagree about one plan (novox/hq issue 296). A repository nothing was measured for, or a nil
|
|
// tierBounds, is given tierAtLeast.
|
|
type tierBounds map[string]time.Duration
|
|
|
|
func (b tierBounds) of(repository string) time.Duration {
|
|
if d, ok := b[repository]; ok {
|
|
return d
|
|
}
|
|
return tierAtLeast
|
|
}
|
|
|
|
// measuredTierBounds is the tier bounds from the last fourteen days of measured plan tiers.
|
|
func measuredTierBounds(ctx context.Context, inv *inventory.Inventory, now time.Time) (tierBounds, error) {
|
|
tiers, err := inv.Durations(ctx, inventory.DurationPlanTier, now.Add(-14*24*time.Hour))
|
|
if err != nil {
|
|
return nil, fmt.Errorf("the measured plan tiers cannot be read: %w", err)
|
|
}
|
|
return boundsOf(tiers), nil
|
|
}
|
|
|
|
// readTierBounds is measuredTierBounds for a line a person reads: what cannot be read is said, and
|
|
// every plan is then given the least bound, as the watchdog gives a repository nothing was measured for.
|
|
func readTierBounds(ctx context.Context, inv *inventory.Inventory, now time.Time) tierBounds {
|
|
bounds, err := measuredTierBounds(ctx, inv, now)
|
|
if err != nil {
|
|
fmt.Fprintf(os.Stderr, "%v; a plan is called late past %s\n", err, tierAtLeast)
|
|
}
|
|
return bounds
|
|
}
|
|
|
|
func boundsOf(tiers []inventory.Duration) tierBounds {
|
|
measured := map[string][]time.Duration{}
|
|
for _, d := range tiers {
|
|
measured[d.Subject] = append(measured[d.Subject], d.Took)
|
|
}
|
|
out := tierBounds{}
|
|
for repository, took := range measured {
|
|
out[repository] = max(tierAtLeast, 3*p90(took))
|
|
}
|
|
return out
|
|
}
|
|
|
|
// p90 is the ninetieth percentile of measurements; zero for none.
|
|
func p90(took []time.Duration) time.Duration {
|
|
if len(took) == 0 {
|
|
return 0
|
|
}
|
|
sorted := append([]time.Duration(nil), took...)
|
|
sort.Slice(sorted, func(i, j int) bool { return sorted[i] < sorted[j] })
|
|
return sorted[int(0.9*float64(len(sorted)-1))]
|
|
}
|
|
|
|
// gatherLoop is when the event loop last took a message and what its consumers hold.
|
|
func (w *watchdogs) gatherLoop() (loopFacts, error) {
|
|
f := loopFacts{took: link.Loop.Last()}
|
|
if w.js == nil {
|
|
return f, errors.New("this controller is not on the bus")
|
|
}
|
|
var where []string
|
|
for _, c := range broker.MeshConsumers() {
|
|
info, err := w.js.Context().ConsumerInfo(c.Stream, c.Name)
|
|
if err != nil {
|
|
return f, fmt.Errorf("the controller's consumer on %s cannot be read: %w", c.Stream, err)
|
|
}
|
|
if held := info.NumPending + uint64(info.NumAckPending); held > 0 {
|
|
f.pending += held
|
|
where = append(where, fmt.Sprintf("%d on %s", held, c.Stream))
|
|
}
|
|
}
|
|
f.where = fmt.Sprint(where)
|
|
return f, nil
|
|
}
|
|
|
|
// gatherAsks is every ask in the build seat's queue that is in flight or dead, with its bound.
|
|
func (w *watchdogs) gatherAsks(ctx context.Context, inv *inventory.Inventory) ([]askFacts, error) {
|
|
if w.js == nil {
|
|
return nil, errors.New("this controller is not on the bus")
|
|
}
|
|
entries, err := inv.Catalogued(ctx)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("the catalogue cannot be read: %w", err)
|
|
}
|
|
seat := buildSeatAmong(entries)
|
|
q, err := link.ReadQueue(ctx, w.js, seat)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
builds, err := inv.Durations(ctx, inventory.DurationBuild, time.Now().Add(-14*24*time.Hour))
|
|
if err != nil {
|
|
return nil, fmt.Errorf("the measured builds cannot be read: %w", err)
|
|
}
|
|
var took []time.Duration
|
|
for _, d := range builds {
|
|
took = append(took, d.Took)
|
|
}
|
|
bound := askDefault
|
|
if len(took) > 0 {
|
|
bound = max(askAtLeast, 3*p90(took))
|
|
}
|
|
var out []askFacts
|
|
for _, a := range q.Asks {
|
|
if a.State != link.AskInFlight && a.State != link.AskDead {
|
|
continue
|
|
}
|
|
since := a.Started
|
|
if since.IsZero() {
|
|
since = a.AskedAt
|
|
}
|
|
what := a.Repository
|
|
if a.Path != "" {
|
|
what += " " + a.Path
|
|
}
|
|
out = append(out, askFacts{id: a.ID, seat: seat, what: what, state: a.State, on: a.On, since: since, bound: bound})
|
|
}
|
|
return out, nil
|
|
}
|
|
|
|
// lostConsumers is, of the consumers the bus said were deleted, each that is still missing and that
|
|
// the mesh expects: a consumer removed because its module was unassigned is not lost.
|
|
func (w *watchdogs) lostConsumers(ctx context.Context, heard []link.Advisory) (map[string]bool, error) {
|
|
out := map[string]bool{}
|
|
var deleted []link.Advisory
|
|
for _, a := range heard {
|
|
if a.Kind == link.AdvisoryConsumerLost {
|
|
deleted = append(deleted, a)
|
|
}
|
|
}
|
|
if len(deleted) == 0 {
|
|
return out, nil
|
|
}
|
|
if w.js == nil {
|
|
return nil, errors.New("this controller is not on the bus")
|
|
}
|
|
_, expected, err := expectedBusObjects(ctx, w.open.inventory)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
want := map[string]bool{}
|
|
for _, c := range expected {
|
|
want[c.Stream+"."+c.Name] = true
|
|
}
|
|
for _, a := range deleted {
|
|
_, err := w.js.Context().ConsumerInfo(a.Stream, a.Consumer)
|
|
switch {
|
|
case errors.Is(err, nats.ErrConsumerNotFound):
|
|
out[a.Stream+"."+a.Consumer] = want[a.Stream+"."+a.Consumer]
|
|
case err != nil:
|
|
return nil, fmt.Errorf("whether %s exists cannot be read: %w", link.ConsumerInWords(a.Stream, a.Consumer), err)
|
|
}
|
|
}
|
|
return out, nil
|
|
}
|
|
|
|
// gatherHandActs is the hand-act log's fortnight (S15), read from the bus.
|
|
func (w *watchdogs) gatherHandActs(ctx context.Context, now time.Time) ([]link.HandAct, error) {
|
|
if w.js == nil {
|
|
return nil, errors.New("this controller is not on the bus")
|
|
}
|
|
reading, cancel := context.WithTimeout(ctx, 10*time.Second)
|
|
defer cancel()
|
|
acts, err := link.HandActs(reading, w.js.Conn(), now.Add(-handActsWithin))
|
|
if err != nil {
|
|
return nil, fmt.Errorf("the hand-act log cannot be read: %w", err)
|
|
}
|
|
return acts, nil
|
|
}
|
|
|
|
// later is the later of two moments.
|
|
func later(a, b time.Time) time.Time {
|
|
if a.After(b) {
|
|
return a
|
|
}
|
|
return b
|
|
}
|
|
|
|
// watchTheMesh opens the condition store and starts the watchdogs, the bus's advisories and the
|
|
// self-check, for the serving controller; the returned function stops them. Nil when the store could
|
|
// not be opened, which is said.
|
|
func watchTheMesh(ctx context.Context, open *stores, server *link.Server, bus link.OverNATS) func() {
|
|
keeper, err := keeperOn(ctx, bus.Conn)
|
|
if err != nil {
|
|
fmt.Printf("the condition store could NOT be opened, so nothing that goes wrong is kept or said, and "+
|
|
"status says the conditions cannot be read: %v\n", err)
|
|
return nil
|
|
}
|
|
conditionsFrom = keeper
|
|
logf := func(format string, args ...any) { fmt.Printf(format+"\n", args...) }
|
|
stopHearing, err := server.HearAdvisories(logf)
|
|
if err != nil {
|
|
fmt.Printf("what the bus says about itself cannot be heard (S9 is blind): %v\n", err)
|
|
stopHearing = func() {}
|
|
}
|
|
host := controlHost(ctx, open.inventory)
|
|
w := &watchdogs{open: open, server: server, js: server.JetStream(), keeper: keeper, started: time.Now(),
|
|
acting: link.Holding}
|
|
d := &doctor{open: open, js: server.JetStream(), keeper: keeper, teller: bus, watchdogs: w, host: host}
|
|
w.doctor = d
|
|
doctorFrom = d
|
|
watching, stop := context.WithCancel(ctx)
|
|
go w.keep(watching)
|
|
go d.keep(watching)
|
|
// And the healers (novox/hq to-be 45 §7): what is open that a registered healer answers, repaired
|
|
// under the lease and the brake, every act said.
|
|
healers := newHealing(open, keeper, bus, server.JetStream())
|
|
go healers.keep(watching)
|
|
go forgettingOldHeals(watching, open.inventory)
|
|
fmt.Printf("watching the mesh: %d signal(s) every %s, %d probe(s) every %s; what is wrong is kept in %s "+
|
|
"and said as %s events\n", len(watchedRows()), watchEvery, len(runnableProbes()), doctorEvery,
|
|
broker.ConditionsBucket, conditions.Seat)
|
|
return func() {
|
|
stop()
|
|
stopHearing()
|
|
flushing, cancel := context.WithTimeout(context.Background(), 10*time.Second)
|
|
defer cancel()
|
|
keeper.Close(flushing)
|
|
}
|
|
}
|
|
|
|
// runnableProbes are the probes the registry runs.
|
|
func runnableProbes() []probe {
|
|
var out []probe
|
|
for _, p := range probeRegistry {
|
|
if p.run != nil {
|
|
out = append(out, p)
|
|
}
|
|
}
|
|
return out
|
|
}
|