Gate a release plan's first machine and roll a failed build back there (hq ADR 0236)
A build that reported applied was sent everywhere; one that then did nothing, served no tools or broke its machine's word reached every machine. Now the first machine is judged by the component's health (the core's definitions, as doctor probes H-*, or a module's own) three times over two minutes within ten; a failing gate puts the previous build back there once, marks the build, and says it as a condition and an event. Upgrades roll out by default; the bus is a planned step; a module deleted at its source is not built (the public-acme plan failure).
This commit is contained in:
@@ -206,7 +206,7 @@ func (a *actor) serveUnderTheLease(ctx context.Context, inv *inventory.Inventory
|
||||
}
|
||||
say := func(format string, args ...any) { fmt.Printf(format+"\n", args...) }
|
||||
l, err := lease.Open(ctx, api, broker.LeaseBucket, lease.Options{Holder: holderOf(instance),
|
||||
Floor: inv.HighestEpoch, Say: say, Moved: func(was, floor uint64) {
|
||||
Floor: inv.HighestEpoch, Say: say, Health: controllerHealth, Moved: func(was, floor uint64) {
|
||||
a.mu.Lock()
|
||||
defer a.mu.Unlock()
|
||||
a.reset = time.Now()
|
||||
|
||||
@@ -570,6 +570,16 @@ func takeIn(ctx context.Context, inv *inventory.Inventory, result link.BuildResu
|
||||
return manifest, kept, fmt.Errorf("%s built %s (%s), and the mesh does not register it: %w",
|
||||
result.On, result.Repository, short(result.Commit), err)
|
||||
}
|
||||
// **A build that failed its gate is never registered again** (novox/hq ADR 0235): an outcome heard
|
||||
// twice, or replayed, would otherwise make the build a rollback put back what the module is again,
|
||||
// and the next push would send it.
|
||||
if failed, err := inv.GateFailed(ctx, kept.ID); err != nil {
|
||||
return manifest, kept, err
|
||||
} else if failed {
|
||||
return manifest, kept, fmt.Errorf("%s built %s (%s), which failed its gate on its first machine and was put "+
|
||||
"back: it is recorded and not registered again — a newer build is", result.On, manifest.Module,
|
||||
short(result.Commit))
|
||||
}
|
||||
if err := inv.RegisterModule(ctx, manifest, recorded); err != nil {
|
||||
if errors.Is(err, inventory.ErrSuperseded) {
|
||||
return manifest, kept, fmt.Errorf("%s built %s (%s), recorded and not registered: %w",
|
||||
|
||||
@@ -0,0 +1,327 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"flag"
|
||||
"fmt"
|
||||
"sort"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/novox/mesh-controller/internal/catalogue"
|
||||
"github.com/novox/mesh-controller/internal/conditions"
|
||||
"github.com/novox/mesh-controller/internal/inventory"
|
||||
"github.com/novox/mesh-controller/internal/link"
|
||||
)
|
||||
|
||||
// The bus as a planned step (novox/hq to-be 45 §8, ADR 0227 rule 8, ADR 0235).
|
||||
//
|
||||
// **A bus upgrade is never rolled out.** The bus carries every declaration, every report and the
|
||||
// controller's own lease; a new bus build that does not come up is a mesh nobody can tell anything,
|
||||
// and a version whose data format moved (2.10 → 2.11) cannot be undone by putting the old one back.
|
||||
// So nothing sends a new bus build on its own: its policy is `record` whatever anyone says (the
|
||||
// catalogue's DerivedUpgrade, and `upgrade` refuses a roll-out), a plan builds it and sends nothing, a
|
||||
// cascade holds the machine (ADR 0221), and a push naming that machine is refused while a new bus build
|
||||
// waits for it (busHeld). The one way is this verb: a person, with why, the streams snapshotted first,
|
||||
// whether it can be reverted said before it starts, the step said as `bus-maintenance` while it runs,
|
||||
// and every stream, durable consumer and a round trip checked after (H-bus) — or the step said failed,
|
||||
// with its snapshot as the way back.
|
||||
|
||||
// busStepProbe is the registry's row for the step; its kinds.
|
||||
const (
|
||||
busStepProbe = "DB"
|
||||
kindBusMaintenance = "bus-maintenance"
|
||||
kindBusUpgradeFailed = "bus-upgrade-failed"
|
||||
)
|
||||
|
||||
// busStepBound is how long after its start a bus upgrade must be followed by a healthy bus.
|
||||
var busStepBound = 15 * time.Minute
|
||||
|
||||
// takeBusSnapshot snapshots every stream before a bus upgrade and answers where the snapshot is. Nil
|
||||
// until the controller's JetStream snapshot is built (to-be 45 §8: "the streams are snapshotted"); until
|
||||
// then a person takes it by hand and says where with --snapshot-taken.
|
||||
var takeBusSnapshot func(ctx context.Context) (string, error)
|
||||
|
||||
// busPending is what a bus upgrade would do: the bus's module, the machines running it, and, per
|
||||
// machine, the build it was last sent against the build the mesh holds. Empty machines: the mesh holds
|
||||
// no bus module.
|
||||
type busPending struct {
|
||||
module string
|
||||
machines []string
|
||||
from map[string]string
|
||||
to string
|
||||
}
|
||||
|
||||
// moves is whether sending the machine would replace its bus.
|
||||
func (b busPending) moves(machine string) bool {
|
||||
from, known := b.from[machine]
|
||||
return b.module != "" && b.to != "" && (!known || !sameCommit(from, b.to))
|
||||
}
|
||||
|
||||
// pendingBus reads what a bus upgrade would do.
|
||||
func pendingBus(ctx context.Context, inv *inventory.Inventory) (busPending, error) {
|
||||
var b busPending
|
||||
shelf, err := inv.Catalogue(ctx)
|
||||
if err != nil {
|
||||
return b, err
|
||||
}
|
||||
for name, m := range shelf {
|
||||
if catalogue.ProvidesBus(m) {
|
||||
b.module = name
|
||||
}
|
||||
}
|
||||
if b.module == "" {
|
||||
return b, nil
|
||||
}
|
||||
current, err := inv.CurrentBuilds(ctx)
|
||||
if err != nil {
|
||||
return b, err
|
||||
}
|
||||
b.to = current[b.module].Commit
|
||||
if b.machines, err = inv.Running(ctx, b.module); err != nil {
|
||||
return b, err
|
||||
}
|
||||
b.from = map[string]string{}
|
||||
for _, n := range b.machines {
|
||||
sent, known, err := inv.SentBuilds(ctx, n)
|
||||
if err != nil {
|
||||
return b, err
|
||||
}
|
||||
if known {
|
||||
if c, carried := sent[b.module]; carried {
|
||||
b.from[n] = c
|
||||
}
|
||||
}
|
||||
}
|
||||
return b, nil
|
||||
}
|
||||
|
||||
// busHeld names the machines a push may not send because sending them would replace the bus: the
|
||||
// planned step's, not a push's (ADR 0235). Said with the remedy.
|
||||
func busHeld(ctx context.Context, inv *inventory.Inventory, machines []string) (map[string]string, error) {
|
||||
b, err := pendingBus(ctx, inv)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
out := map[string]string{}
|
||||
for _, n := range machines {
|
||||
for _, holder := range b.machines {
|
||||
if n == holder && b.moves(n) {
|
||||
out[n] = fmt.Sprintf("sending %s would replace the bus (%s %s → %s), which is a planned step: "+
|
||||
"`bus upgrade --why …` snapshots its streams first and checks them after (novox/hq ADR 0235)",
|
||||
n, b.module, short(orNotKnown(b.from[n])), short(b.to))
|
||||
}
|
||||
}
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
// busCommand is `bus` — what a bus upgrade would do and how the last went — and `bus upgrade`.
|
||||
func busCommand(ctx context.Context, args []string) error {
|
||||
sub := ""
|
||||
if len(args) > 0 && !strings.HasPrefix(args[0], "-") {
|
||||
sub, args = args[0], args[1:]
|
||||
}
|
||||
set := flag.NewFlagSet("bus", flag.ContinueOnError)
|
||||
snapshot := set.String("snapshot-taken", "", "where the streams' snapshot a person took is, while the mesh takes none itself")
|
||||
reversible := set.Bool("reversible", false, "the new version can be undone by putting the old one back")
|
||||
irreversible := set.Bool("irreversible", false, "the new version cannot be undone by putting the old one back, "+
|
||||
"and this is the person's explicit word that it runs anyway")
|
||||
why := addHandActFlags(set)
|
||||
if rest, err := parseAround(set, args); err != nil {
|
||||
return err
|
||||
} else if len(rest) > 0 {
|
||||
return errors.New("bus [upgrade --why … --reversible|--irreversible [--snapshot-taken <where>]]")
|
||||
}
|
||||
switch sub {
|
||||
case "":
|
||||
return busStatus(ctx)
|
||||
case "upgrade":
|
||||
default:
|
||||
return fmt.Errorf("bus says what a bus upgrade would do, or `bus upgrade` — not %q", sub)
|
||||
}
|
||||
// Everything refused before anything is done.
|
||||
if err := why.require("bus upgrade"); err != nil {
|
||||
return err
|
||||
}
|
||||
if *reversible == *irreversible {
|
||||
return errors.New("bus upgrade says, before it starts, whether the new version can be undone by putting the " +
|
||||
"old one back: --reversible, or --irreversible as your explicit word that it runs anyway (to-be 45 §8). " +
|
||||
"Nothing was done")
|
||||
}
|
||||
open, err := openStores(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer open.Close()
|
||||
inv := open.inventory
|
||||
b, err := pendingBus(ctx, inv)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if b.module == "" {
|
||||
return errors.New("the mesh holds no module that provides its bus: there is nothing to upgrade")
|
||||
}
|
||||
var moving []string
|
||||
for _, n := range b.machines {
|
||||
if b.moves(n) {
|
||||
moving = append(moving, n)
|
||||
}
|
||||
}
|
||||
if len(moving) == 0 {
|
||||
fmt.Printf("every machine running %s runs the build the mesh holds (%s): nothing to upgrade\n", b.module, short(b.to))
|
||||
return nil
|
||||
}
|
||||
where := strings.TrimSpace(*snapshot)
|
||||
switch {
|
||||
case takeBusSnapshot != nil:
|
||||
if where, err = takeBusSnapshot(ctx); err != nil {
|
||||
return fmt.Errorf("the streams could not be snapshotted, so the bus is not replaced: %w", err)
|
||||
}
|
||||
case where == "":
|
||||
return errors.New("the bus is replaced only after its streams are snapshotted. The mesh does not take the " +
|
||||
"snapshot itself yet (to-be 45 §8: the controller's JetStream snapshot); take one by hand and say where " +
|
||||
"with --snapshot-taken <where>. Nothing was done")
|
||||
}
|
||||
from := map[string]bool{}
|
||||
for _, n := range moving {
|
||||
from[orNotKnown(b.from[n])] = true
|
||||
}
|
||||
step, err := inv.StartBusStep(ctx, inventory.BusStep{Module: b.module, Machines: moving,
|
||||
From: strings.Join(sortedKeys(from), ", "), To: b.to, Snapshot: where, Reversible: *reversible,
|
||||
By: link.Caller(), Why: strings.TrimSpace(*why.why)})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
cause := "bus-upgrade"
|
||||
if strings.TrimSpace(*why.cause) == "" {
|
||||
why.cause = &cause
|
||||
}
|
||||
why.record(ctx, "bus upgrade", append([]string{b.module}, moving...))
|
||||
fmt.Printf("bus upgrade %d: %s %s → %s on %s; streams snapshotted at %s; %s\n", step.ID, b.module, step.From,
|
||||
short(b.to), strings.Join(moving, ", "), where, map[bool]string{true: "reversible: putting the old build back undoes it",
|
||||
false: "NOT reversible: the snapshot is the only way back"}[*reversible])
|
||||
sent, err := sendRollout(ctx, open, moving)
|
||||
if err != nil {
|
||||
_ = inv.EndBusStep(ctx, step.ID, "failed", "the send was refused: "+err.Error())
|
||||
return fmt.Errorf("the bus's machine could not be sent its new build: %w — nothing was replaced", err)
|
||||
}
|
||||
fmt.Printf("sent %s; `bus-maintenance` is open until the bus answers healthy again — every stream, every durable "+
|
||||
"consumer, a round trip to the machines (H-bus) — within %s, or the step is said failed with its snapshot "+
|
||||
"as the way back. `bus` says how it went\n", strings.Join(sent, ", "), busStepBound)
|
||||
return nil
|
||||
}
|
||||
|
||||
// busStatus is `bus`: what an upgrade would do, and the last step.
|
||||
func busStatus(ctx context.Context) error {
|
||||
open, err := openStores(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer open.Close()
|
||||
b, err := pendingBus(ctx, open.inventory)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if b.module == "" {
|
||||
fmt.Println("the mesh holds no module that provides its bus")
|
||||
} else {
|
||||
fmt.Printf("the bus is %s, built %s, on %s; never rolled out — a planned step (`bus upgrade`)\n", b.module,
|
||||
short(b.to), orNone(strings.Join(b.machines, ", ")))
|
||||
for _, n := range b.machines {
|
||||
state := "runs it"
|
||||
if b.moves(n) {
|
||||
state = "runs " + short(orNotKnown(b.from[n])) + ": `bus upgrade` replaces it"
|
||||
}
|
||||
fmt.Printf(" %-10s %s\n", n, state)
|
||||
}
|
||||
}
|
||||
if takeBusSnapshot == nil {
|
||||
fmt.Println(" the mesh takes no snapshot of the streams itself yet: `bus upgrade` asks where yours is (--snapshot-taken)")
|
||||
}
|
||||
s, found, err := open.inventory.LatestBusStep(ctx)
|
||||
if err != nil || !found {
|
||||
return err
|
||||
}
|
||||
state := "running since " + s.Started.Local().Format("2006-01-02 15:04")
|
||||
if s.Ended != nil {
|
||||
state = s.Outcome + " at " + s.Ended.Local().Format("2006-01-02 15:04")
|
||||
}
|
||||
fmt.Printf("last step %d: %s → %s on %s by %s (%s): %s; snapshot %s\n", s.ID, s.From, short(s.To),
|
||||
strings.Join(s.Machines, ", "), orNone(s.By), s.Why, state, s.Snapshot)
|
||||
if s.Found != "" {
|
||||
fmt.Printf(" %s\n", s.Found)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// probeBusStep is DB: a bus upgrade running is said as `bus-maintenance`; the bus healthy again after
|
||||
// the machines reported the new build ends it done; past its bound, unhealthy, it ends failed and is
|
||||
// said — urgent, with its snapshot — while the bus is still not healthy.
|
||||
func probeBusStep(ctx context.Context, d *doctor) ([]conditions.Observation, error) {
|
||||
inv := d.open.inventory
|
||||
s, found, err := inv.LatestBusStep(ctx)
|
||||
if err != nil || !found {
|
||||
return nil, err
|
||||
}
|
||||
if s.Ended != nil && s.Outcome != "failed" {
|
||||
return nil, nil
|
||||
}
|
||||
problems, err := busHealth(ctx, d)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
reports, err := inv.LastReports(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
applied := true
|
||||
for _, r := range reports {
|
||||
for _, n := range s.Machines {
|
||||
if r.Node == n && (!r.Current || r.Outcome != inventory.OutcomeApplied || r.At == nil || r.At.Before(s.Started)) {
|
||||
applied = false
|
||||
problems = append(problems, n+" has not reported the new bus applied")
|
||||
}
|
||||
}
|
||||
}
|
||||
id := s.Module
|
||||
if s.Ended == nil {
|
||||
switch {
|
||||
case applied && len(problems) == 0:
|
||||
return nil, inv.EndBusStep(ctx, s.ID, "done", "the bus answered healthy after the upgrade")
|
||||
case time.Since(s.Started) > busStepBound:
|
||||
found := strings.Join(problems, "; ")
|
||||
if err := inv.EndBusStep(ctx, s.ID, "failed", found); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
s.Found = found
|
||||
default:
|
||||
return []conditions.Observation{{Scope: conditions.ScopeBus, ID: id, Token: "maintenance",
|
||||
Kind: kindBusMaintenance, Severity: conditions.Warning,
|
||||
Summary: fmt.Sprintf("the bus is being upgraded (step %d, %s → %s on %s, by %s: %s); its snapshot is %s",
|
||||
s.ID, s.From, short(s.To), strings.Join(s.Machines, ", "), orNone(s.By), s.Why, s.Snapshot),
|
||||
Said: orNone(strings.Join(problems, "; "))}}, nil
|
||||
}
|
||||
}
|
||||
if len(problems) == 0 {
|
||||
return nil, nil // failed, and healthy since: nothing wrong now
|
||||
}
|
||||
way := "put the old build back"
|
||||
if !s.Reversible {
|
||||
way = "restore the snapshot"
|
||||
}
|
||||
return []conditions.Observation{{Scope: conditions.ScopeBus, ID: id, Token: "upgrade-failed",
|
||||
Kind: kindBusUpgradeFailed, Severity: conditions.Urgent, Resolver: conditions.ResolverOperator,
|
||||
Summary: fmt.Sprintf("the bus upgrade (step %d, %s → %s) did not end healthy within %s: %s — the way back is to %s (%s)",
|
||||
s.ID, s.From, short(s.To), busStepBound, strings.Join(problems, "; "), way, s.Snapshot)}}, nil
|
||||
}
|
||||
|
||||
func sortedKeys(set map[string]bool) []string {
|
||||
out := make([]string, 0, len(set))
|
||||
for k := range set {
|
||||
out = append(out, k)
|
||||
}
|
||||
sort.Strings(out)
|
||||
return out
|
||||
}
|
||||
@@ -113,6 +113,27 @@ var probeRegistry = []probe{
|
||||
Asks: []broker.SeatVerb{{Seat: "node-backup", Verb: "backed-up"}}},
|
||||
{ID: "DW", Asserts: "the watchdogs of the signals table ran within three of their intervals",
|
||||
From: "ADR 0227 rule 6: the watchers are watched", Kind: "watchdogs-silent", Phase: 1, run: probeWatchdogs},
|
||||
// The core's health definitions (novox/hq to-be 45 §8, ADR 0235): what a core component's new build is
|
||||
// judged by on its first machine, run against every machine between upgrades too.
|
||||
{ID: "H-controller", Asserts: "the controller lease is held, renewed in time, by a controller that says it is " +
|
||||
"ready: its self-check ran and status answered in full within ten seconds", From: "ADR 0235, to-be 45 §8",
|
||||
Kind: kindCoreUnhealthy, Phase: 4, run: probeControllerHealth},
|
||||
{ID: "H-engine", Asserts: "every machine heard from has reported its current declaration, under a node-engine " +
|
||||
"build it names", From: "ADR 0235, to-be 45 §8", Kind: kindCoreUnhealthy, Phase: 4, run: probeEngineHealth},
|
||||
{ID: "H-tools", Asserts: "every machine heard from that runs the node tools has them answering the bus",
|
||||
From: "ADR 0235, to-be 45 §8", Kind: kindCoreUnhealthy, Phase: 4, run: probeToolsHealth},
|
||||
{ID: "H-bus", Asserts: "every stream and durable consumer the mesh defines is on the bus, and a request crosses " +
|
||||
"it to the machines' node tools and back", From: "ADR 0235, to-be 45 §8", Kind: kindCoreUnhealthy, Phase: 4,
|
||||
run: probeBusHealth},
|
||||
// The gate's verdicts and the witnesses' rollbacks (ADR 0235): each build that failed its gate keeps its
|
||||
// condition until a newer build passes; each rollback a witness stands by is said.
|
||||
{ID: gateProbe, Asserts: "no build that failed its gate, and no core component a witness put back, goes unsaid; " +
|
||||
"a newer build that passes its gate clears it", From: "ADR 0235, to-be 45 §8", Kind: kindRolledBack,
|
||||
Raises: []string{kindRollbackFailed}, Phase: 4, run: probeGates},
|
||||
// The bus's planned step (ADR 0235): open while a person's bus upgrade runs, then checked by H-bus.
|
||||
{ID: busStepProbe, Asserts: "a bus upgrade a person started is said while it runs, and is followed by the bus's " +
|
||||
"health within its bound — or is said failed, with its snapshot as the way back", From: "ADR 0235, to-be 45 §8",
|
||||
Kind: kindBusMaintenance, Raises: []string{kindBusUpgradeFailed}, Phase: 4, run: probeBusStep},
|
||||
}
|
||||
|
||||
// probeVerdict is one probe's outcome in a run.
|
||||
|
||||
@@ -0,0 +1,692 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"slices"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/nats-io/nats.go"
|
||||
"github.com/nats-io/nats.go/micro"
|
||||
|
||||
"github.com/novox/mesh-controller/internal/broker"
|
||||
"github.com/novox/mesh-controller/internal/catalogue"
|
||||
"github.com/novox/mesh-controller/internal/conditions"
|
||||
"github.com/novox/mesh-controller/internal/inventory"
|
||||
"github.com/novox/mesh-controller/internal/lease"
|
||||
"github.com/novox/mesh-controller/internal/link"
|
||||
)
|
||||
|
||||
// The gate on a release plan's first machine, and the rollback after it (novox/hq ADR 0235, to-be 45
|
||||
// §8, ADR 0227 rule 8).
|
||||
//
|
||||
// **"Reported applied" is not enough.** ADR 0218 sent a module to one machine first and the rest once
|
||||
// that machine reported it applied. A build that applies and then does nothing, crashes, serves no
|
||||
// tools, or breaks the machine's word to the mesh passed that test. Now the first machine is judged by
|
||||
// the component's health — the core's health definitions, or a module's own — passing three times over
|
||||
// at least two minutes, within ten minutes of the apply. Only then are the rest sent.
|
||||
//
|
||||
// **A failing gate stops the plan and puts the previous build back there.** The build is marked failed
|
||||
// at its gate — registration refuses it and nothing sends it again on its own — the module's registered
|
||||
// build goes back to the one the first machine ran before, and that machine is sent it: the ordinary
|
||||
// path, again. Once per build: the verdict is written before the send, and a verdict once written is
|
||||
// not written over. Said as a condition, urgent for the core and when it could not be put back, and as
|
||||
// the event `rolled-back`.
|
||||
|
||||
// The gate's bounds (to-be 45 §8). Variables so a test can judge in a second, not in minutes.
|
||||
var (
|
||||
// gateSettle is how long the first machine must stay healthy, at the least.
|
||||
gateSettle = 2 * time.Minute
|
||||
// gateBound is how long after the apply the build may take to become healthy.
|
||||
gateBound = 10 * time.Minute
|
||||
// gatePasses is how many consecutive judgings must find it healthy, gateEvery apart at the least.
|
||||
gatePasses = 3
|
||||
gateEvery = 40 * time.Second
|
||||
)
|
||||
|
||||
// gateProbe is the registry's row for the gate's verdicts and the witnesses' rollbacks: its conditions
|
||||
// are raised by a plan as it judges, and kept or cleared by the probe on every run.
|
||||
const gateProbe = "DG"
|
||||
|
||||
// The kinds a gate raises.
|
||||
const (
|
||||
kindRolledBack = "rolled-back"
|
||||
kindRollbackFailed = "rollback-failed"
|
||||
)
|
||||
|
||||
// coreComponent is the core component a module is, as a witness names it — controller, node-engine,
|
||||
// node-tools — or empty: the core is judged by its health definitions, anything else by its own health.
|
||||
func coreComponent(module string) string {
|
||||
switch module {
|
||||
case catalogue.ControllerSeatName:
|
||||
return lease.ComponentController
|
||||
case hostModule:
|
||||
return lease.ComponentEngine
|
||||
case broker.RuntimeModule:
|
||||
return lease.ComponentNodeTools
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
// health is a judging's word on one machine.
|
||||
type health int
|
||||
|
||||
const (
|
||||
healthGood health = iota
|
||||
healthNotYet
|
||||
healthBroken
|
||||
)
|
||||
|
||||
// served is what one machine's node tools answered the bus's discovery with.
|
||||
type served struct {
|
||||
runtime bool
|
||||
tools map[string]bool
|
||||
}
|
||||
|
||||
// gateFacts is what one judging reads, gathered once for every machine it judges.
|
||||
type gateFacts struct {
|
||||
now time.Time
|
||||
reports map[string]inventory.Reported
|
||||
// engines is each machine's node-engine build as it last reported it.
|
||||
engines map[string]string
|
||||
// served is what the bus's discovery answered; servedErr why it could not be asked.
|
||||
served map[string]served
|
||||
servedErr error
|
||||
// open are the open conditions; openErr why they could not be read (nil keeper: not judged).
|
||||
open []conditions.Condition
|
||||
openErr error
|
||||
judged bool
|
||||
// rolledBack is what each machine's witnesses say they decided.
|
||||
rolledBack map[string][]lease.Rollback
|
||||
// holder is who holds the controller lease, for judging the controller.
|
||||
holder *lease.Holder
|
||||
holderErr error
|
||||
}
|
||||
|
||||
// gatherGateFacts reads what a judging needs, from the store, the bus and this controller's memory. A
|
||||
// variable so a test can hand a judging its facts.
|
||||
var gatherGateFacts = func(ctx context.Context, open *stores, component string) (gateFacts, error) {
|
||||
inv := open.inventory
|
||||
f := gateFacts{now: time.Now(), reports: map[string]inventory.Reported{}, engines: map[string]string{},
|
||||
rolledBack: witnessed.all()}
|
||||
reports, err := inv.LastReports(ctx)
|
||||
if err != nil {
|
||||
return f, err
|
||||
}
|
||||
for _, r := range reports {
|
||||
f.reports[r.Node] = r
|
||||
}
|
||||
nodes, err := inv.Nodes(ctx)
|
||||
if err != nil {
|
||||
return f, err
|
||||
}
|
||||
for _, n := range nodes {
|
||||
f.engines[n.Name] = n.HostVersion
|
||||
}
|
||||
if d := doctorFrom; d != nil {
|
||||
if d.keeper != nil {
|
||||
f.judged = true
|
||||
f.open, f.openErr = d.keeper.Open(ctx)
|
||||
}
|
||||
if d.js != nil {
|
||||
f.served, f.servedErr = servedOnTheBus(ctx, d.js.Conn())
|
||||
} else {
|
||||
f.servedErr = errors.New("this controller has no bus to ask")
|
||||
}
|
||||
} else {
|
||||
f.servedErr = errors.New("this process does not serve the mesh, so it cannot ask the bus who serves what")
|
||||
}
|
||||
if component == lease.ComponentController {
|
||||
h, found, err := theLease.holder(ctx)
|
||||
switch {
|
||||
case err != nil:
|
||||
f.holderErr = err
|
||||
case found:
|
||||
f.holder = &h
|
||||
}
|
||||
}
|
||||
return f, nil
|
||||
}
|
||||
|
||||
// judgeHealth is one machine's health for a module's new build, from what one judging read: healthy,
|
||||
// not yet (with what is wanting), or healthBroken — a witness put it back, or the machine refused or failed
|
||||
// what it was sent. Pure.
|
||||
func judgeHealth(module, component string, m catalogue.Manifest, machine string, since time.Time, f gateFacts) (health, string) {
|
||||
// A witness's verdict made since the build was sent: the build failed its health there. One that
|
||||
// could not judge at all (unwitnessed) is not a verdict on the build.
|
||||
for _, r := range f.rolledBack[machine] {
|
||||
if component == "" || r.Component != component || r.Outcome == lease.OutcomeUnwitnessed || r.At.Before(since) {
|
||||
continue
|
||||
}
|
||||
return healthBroken, fmt.Sprintf("the witness on %s judged the %s %s and %s: %s", machine, r.Component,
|
||||
short(r.From), r.Outcome, r.Why)
|
||||
}
|
||||
r, said := f.reports[machine]
|
||||
switch {
|
||||
case !said || r.At == nil || !r.Current:
|
||||
return healthNotYet, fmt.Sprintf("%s has not reported on what it was sent", machine)
|
||||
case r.Outcome == inventory.OutcomeFailed || r.Outcome == inventory.OutcomeRefused:
|
||||
return healthBroken, fmt.Sprintf("%s %s what it was sent", machine, r.Outcome)
|
||||
case r.Outcome != inventory.OutcomeApplied:
|
||||
return healthNotYet, fmt.Sprintf("%s reported %q", machine, r.Outcome)
|
||||
}
|
||||
// **No new condition about it**: about the machine itself, or naming the module on that machine,
|
||||
// raised since the judging began. The gate's own are not evidence about the build.
|
||||
if f.judged {
|
||||
if f.openErr != nil {
|
||||
return healthNotYet, "what is wrong cannot be read, so whether the build made anything wrong is not known: " +
|
||||
firstLine(f.openErr.Error())
|
||||
}
|
||||
for _, c := range f.open {
|
||||
if c.Source == gateProbe || c.Raised.Before(since) {
|
||||
continue
|
||||
}
|
||||
onIt := c.Subject.Machine == machine || slices.Contains(c.Subject.Also, machine) ||
|
||||
(c.Subject.Scope == conditions.ScopeMachine && c.Subject.ID == machine)
|
||||
if !onIt {
|
||||
continue
|
||||
}
|
||||
if c.Subject.Scope == conditions.ScopeMachine || slices.Contains(strings.Split(c.Subject.ID, "."), module) {
|
||||
return healthNotYet, fmt.Sprintf("raised since it was sent: %s — %s", c.Key, c.Summary)
|
||||
}
|
||||
}
|
||||
}
|
||||
switch component {
|
||||
case lease.ComponentEngine:
|
||||
// The node-engine has reported its current declaration under its own build.
|
||||
want := deliveredVersions(m)
|
||||
if len(want) > 0 && !slices.Contains(want, f.engines[machine]) {
|
||||
return healthNotYet, fmt.Sprintf("%s's node-engine reports build %s, not the new %s", machine,
|
||||
orNotKnown(f.engines[machine]), strings.Join(want, " or "))
|
||||
}
|
||||
case lease.ComponentNodeTools:
|
||||
// The node tools are announced and answer.
|
||||
if f.servedErr != nil {
|
||||
return healthNotYet, "whether the node tools answer cannot be asked: " + firstLine(f.servedErr.Error())
|
||||
}
|
||||
if !f.served[machine].runtime {
|
||||
return healthNotYet, fmt.Sprintf("the node tools on %s do not answer the bus", machine)
|
||||
}
|
||||
case lease.ComponentController:
|
||||
// The new controller holds the lease and says it is ready.
|
||||
switch {
|
||||
case f.holderErr != nil:
|
||||
return healthNotYet, "who holds the controller lease cannot be read: " + firstLine(f.holderErr.Error())
|
||||
case f.holder == nil:
|
||||
return healthNotYet, "no controller holds the lease"
|
||||
case f.holder.Taken.Before(since.Add(-time.Minute)):
|
||||
return healthNotYet, fmt.Sprintf("the lease is held since %s, by a controller older than the new build",
|
||||
f.holder.Taken.UTC().Format(time.RFC3339))
|
||||
case f.holder.Health == nil || !f.holder.Health.Ready:
|
||||
why := "the controller holding the lease does not say it is ready"
|
||||
if f.holder.Health != nil && f.holder.Health.Why != "" {
|
||||
why += ": " + f.holder.Health.Why
|
||||
}
|
||||
return healthNotYet, why
|
||||
}
|
||||
default:
|
||||
// A module's tools answer, where it has any and the machine runs the node tools that serve them.
|
||||
if len(m.Tools) > 0 {
|
||||
if f.servedErr != nil {
|
||||
return healthNotYet, "whether its tools are served cannot be asked: " + firstLine(f.servedErr.Error())
|
||||
}
|
||||
if s := f.served[machine]; s.runtime && !s.tools[module] {
|
||||
return healthNotYet, fmt.Sprintf("the node tools on %s do not serve %s's tools", machine, module)
|
||||
}
|
||||
}
|
||||
}
|
||||
return healthGood, ""
|
||||
}
|
||||
|
||||
// judgeGate takes one judging of a module's first machines and records it in the plan's gate: a pass
|
||||
// counted, a pass missed (and why), or the verdict. Answers the verdict once there is one.
|
||||
func judgeGate(ctx context.Context, open *stores, p *inventory.Plan, module string, state *inventory.PlanModule,
|
||||
running []string, now time.Time) (string, error) {
|
||||
g := state.Gate
|
||||
if g == nil {
|
||||
// Judged from the send: a witness's verdict, a condition, the bound — all counted from when the
|
||||
// first machine was sent the build.
|
||||
start := now
|
||||
if state.FirstAt != nil {
|
||||
start = *state.FirstAt
|
||||
}
|
||||
var machines []string
|
||||
for _, n := range state.First {
|
||||
if slices.Contains(running, n) {
|
||||
machines = append(machines, n)
|
||||
}
|
||||
}
|
||||
g = &inventory.PlanGate{Component: coreComponent(module), Machines: machines, From: state.Previous,
|
||||
To: state.Commit, Since: &start}
|
||||
state.Gate = g
|
||||
}
|
||||
if g.Verdict != "" {
|
||||
return g.Verdict, nil
|
||||
}
|
||||
if g.LastPass != nil && now.Sub(*g.LastPass) < gateEvery {
|
||||
return "", nil
|
||||
}
|
||||
shelf, err := open.inventory.Catalogue(ctx)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
facts, err := gatherGateFacts(ctx, open, g.Component)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
worst, why := healthGood, ""
|
||||
for _, n := range g.Machines {
|
||||
h, said := judgeHealth(module, g.Component, shelf[module], n, *g.Since, facts)
|
||||
if h > worst {
|
||||
worst, why = h, said
|
||||
} else if h == worst && h != healthGood && why == "" {
|
||||
why = said
|
||||
}
|
||||
}
|
||||
switch {
|
||||
case worst == healthBroken:
|
||||
decide(g, inventory.GateFailed, why, now)
|
||||
case worst == healthNotYet:
|
||||
g.Passes, g.LastPass, g.Last = 0, nil, why
|
||||
if now.Sub(*g.Since) > gateBound {
|
||||
decide(g, inventory.GateFailed, fmt.Sprintf("not healthy within %s of its apply: %s", gateBound, why), now)
|
||||
}
|
||||
default:
|
||||
g.Passes++
|
||||
g.LastPass, g.Last = &now, ""
|
||||
if g.Passes >= gatePasses && now.Sub(*g.Since) >= gateSettle {
|
||||
decide(g, inventory.GatePassed, fmt.Sprintf("healthy %d times over %s", g.Passes,
|
||||
now.Sub(*g.Since).Round(time.Second)), now)
|
||||
}
|
||||
}
|
||||
return g.Verdict, nil
|
||||
}
|
||||
|
||||
// decide sets a gate's verdict.
|
||||
func decide(g *inventory.PlanGate, verdict, why string, now time.Time) {
|
||||
g.Verdict, g.Why, g.JudgedAt = verdict, why, &now
|
||||
g.Took = now.Sub(*g.Since).Round(time.Second).String()
|
||||
}
|
||||
|
||||
// gatePassed keeps a passing build's verdict, so `plans` and the gate's probe can read it.
|
||||
func gatePassed(ctx context.Context, open *stores, p *inventory.Plan, module string, state *inventory.PlanModule) {
|
||||
g := state.Gate
|
||||
err := open.inventory.RecordGate(ctx, inventory.GateVerdict{Build: state.Build, Module: module,
|
||||
Commit: state.Commit, Previous: state.Previous, Plan: p.ID, Machines: g.Machines,
|
||||
Verdict: inventory.GatePassed, Why: g.Why, Component: g.Component, JudgingFrom: g.Since})
|
||||
if err != nil && state.Build != "" {
|
||||
fmt.Printf("%s: %s passed its gate, and the verdict could not be kept: %v\n", p.ID, module, err)
|
||||
}
|
||||
fmt.Printf("%s: %s passed its gate on %s (%s); the rest are sent\n", p.ID, module,
|
||||
strings.Join(g.Machines, ", "), g.Why)
|
||||
if module == catalogue.ControllerSeatName {
|
||||
carryUserList(ctx, open, p)
|
||||
}
|
||||
}
|
||||
|
||||
// carryUserList sends the machine holding the bus the user list a new controller composes, once that
|
||||
// controller passed its gate (ADR 0235). The controller's own grants travel in that list, and the old
|
||||
// controller composed the list the plan sent; on 2026-10-06 eight pushes by hand carried a new
|
||||
// controller's grant into it. Not when a build its policy or a plan holds back would go with it (ADR
|
||||
// 0221): then it is said, as a push would say it.
|
||||
func carryUserList(ctx context.Context, open *stores, p *inventory.Plan) {
|
||||
holder, behind, err := brokerBehind(ctx, open, nil)
|
||||
if err != nil || holder == "" || !behind {
|
||||
if err != nil {
|
||||
fmt.Printf("%s: whether the bus's user list is behind the new controller cannot be read: %v\n", p.ID, err)
|
||||
}
|
||||
return
|
||||
}
|
||||
held, err := heldMachines(ctx, open, []string{holder})
|
||||
if err != nil {
|
||||
fmt.Printf("%s: whether %s may be sent the new user list cannot be read: %v\n", p.ID, holder, err)
|
||||
return
|
||||
}
|
||||
if why, isHeld := held[holder]; isHeld {
|
||||
fmt.Printf("%s: the new controller's user list is not carried to %s, which holds the bus: %s — `push %s` "+
|
||||
"carries it\n", p.ID, holder, strings.Join(why, "; "), holder)
|
||||
return
|
||||
}
|
||||
if _, err := sendRollout(ctx, open, []string{holder}); err != nil {
|
||||
fmt.Printf("%s: the new controller's user list could not be carried to %s: %v\n", p.ID, holder, err)
|
||||
return
|
||||
}
|
||||
fmt.Printf("%s: carried the new controller's user list to %s, which holds the bus\n", p.ID, holder)
|
||||
}
|
||||
|
||||
// gateFailed stops the plan at a build that failed its gate and puts the previous build back on the
|
||||
// machines it was judged on — once per build, said as a condition and an event. The plan is saved
|
||||
// before the send: a controller that is itself the build being put back does not outlive it.
|
||||
func gateFailed(ctx context.Context, open *stores, p *inventory.Plan, module string, state *inventory.PlanModule,
|
||||
machines []string, why string) {
|
||||
inv := open.inventory
|
||||
now := time.Now().UTC()
|
||||
if state.Gate == nil {
|
||||
state.Gate = &inventory.PlanGate{Component: coreComponent(module), Machines: machines,
|
||||
From: state.Previous, To: state.Commit, Since: state.FirstAt}
|
||||
}
|
||||
g := state.Gate
|
||||
if g.Verdict == "" {
|
||||
if g.Since == nil {
|
||||
g.Since = &now
|
||||
}
|
||||
decide(g, inventory.GateFailed, why, now)
|
||||
}
|
||||
if len(g.Machines) == 0 {
|
||||
g.Machines = machines
|
||||
}
|
||||
state.Why = "failed its gate: " + g.Why
|
||||
p.State = inventory.PlanFailed
|
||||
p.Note = fmt.Sprintf("%s failed its gate on %s in tier %d: %s", module, strings.Join(g.Machines, ", "), p.Tier, g.Why)
|
||||
|
||||
verdict := inventory.GateVerdict{Build: state.Build, Module: module, Commit: state.Commit, Previous: state.Previous,
|
||||
Plan: p.ID, Machines: g.Machines, Verdict: inventory.GateFailed, Rollback: inventory.RollingBack, Why: g.Why,
|
||||
Component: g.Component, JudgingFrom: g.Since}
|
||||
if state.Build == "" {
|
||||
// A plan from before builds were asked by id: nothing to mark, so nothing is put back by the
|
||||
// mesh — said, for a person.
|
||||
g.Rollback = inventory.NotRolledBack
|
||||
p.Note += "; not put back: the plan does not know which build it sent"
|
||||
sayRollback(ctx, open, module, g, "")
|
||||
return
|
||||
}
|
||||
if err := inv.RecordGate(ctx, verdict); err != nil {
|
||||
if errors.Is(err, inventory.ErrGateKept) {
|
||||
// Already judged and acted on, by this controller before a restart or by another: never twice.
|
||||
if kept, found, _ := inv.GateOf(ctx, state.Build); found {
|
||||
g.Rollback = kept.Rollback
|
||||
}
|
||||
p.Note += "; its rollback was already made once and is not made again"
|
||||
return
|
||||
}
|
||||
g.Rollback = inventory.NotRolledBack
|
||||
p.Note += "; not put back: its verdict could not be kept, and a rollback that cannot be counted is not made — " + err.Error()
|
||||
sayRollback(ctx, open, module, g, "")
|
||||
return
|
||||
}
|
||||
// The previous build: the one the first machine ran, from the build records.
|
||||
notBack := func(why string) {
|
||||
g.Rollback = inventory.NotRolledBack
|
||||
p.Note += "; NOT put back: " + why
|
||||
if err := inv.SetRollback(ctx, state.Build, inventory.NotRolledBack, g.Why+"; not put back: "+why); err != nil {
|
||||
fmt.Printf("%s: how %s's rollback went could not be kept: %v\n", p.ID, module, err)
|
||||
}
|
||||
sayRollback(ctx, open, module, g, why)
|
||||
}
|
||||
if state.Previous == "" {
|
||||
notBack(fmt.Sprintf("%s had not been sent %s before, or what it was sent is not known — `unassign` takes it "+
|
||||
"off, or a newer merge replaces it", strings.Join(g.Machines, ", "), module))
|
||||
return
|
||||
}
|
||||
failed, _, err := inv.BuildByID(ctx, state.Build)
|
||||
if err != nil {
|
||||
notBack("the failed build's record cannot be read: " + err.Error())
|
||||
return
|
||||
}
|
||||
previous, found, err := inv.PreviousBuild(ctx, module, state.Previous, failed)
|
||||
if err != nil {
|
||||
notBack("the build records cannot be read: " + err.Error())
|
||||
return
|
||||
}
|
||||
if !found {
|
||||
notBack(fmt.Sprintf("no build of %s from %s is still kept to put back", module, short(state.Previous)))
|
||||
return
|
||||
}
|
||||
if err := inv.RestoreModule(ctx, previous); err != nil {
|
||||
notBack(err.Error())
|
||||
return
|
||||
}
|
||||
g.Rollback = inventory.RollingBack
|
||||
if err := inv.SavePlan(ctx, p); err != nil {
|
||||
fmt.Printf("%s: the plan could not be kept before %s is put back: %v\n", p.ID, module, err)
|
||||
}
|
||||
sent, err := sendRollout(ctx, open, g.Machines)
|
||||
if err != nil {
|
||||
notBack(fmt.Sprintf("its registered build is back at %s, and sending it to %s was refused: %v — `push %s` "+
|
||||
"sends it", short(previous.Commit), strings.Join(g.Machines, ", "), err, g.Machines[0]))
|
||||
return
|
||||
}
|
||||
g.Rollback = inventory.RolledBack
|
||||
p.Note += fmt.Sprintf("; put back to %s on %s", short(previous.Commit), strings.Join(sent, ", "))
|
||||
if err := inv.SetRollback(ctx, state.Build, inventory.RolledBack, g.Why); err != nil {
|
||||
fmt.Printf("%s: how %s's rollback went could not be kept: %v\n", p.ID, module, err)
|
||||
}
|
||||
sayRollback(ctx, open, module, g, "")
|
||||
}
|
||||
|
||||
// rolledBackEvent is the body of `rolled-back` (ADR 0235): a contract, like a condition's events.
|
||||
type rolledBackEvent struct {
|
||||
Event string `json:"event"`
|
||||
At time.Time `json:"at"`
|
||||
Module string `json:"module"`
|
||||
Component string `json:"component,omitempty"`
|
||||
Machines []string `json:"machines"`
|
||||
From string `json:"from,omitempty"`
|
||||
To string `json:"to,omitempty"`
|
||||
Why string `json:"why"`
|
||||
// Rollback is rolled-back, or not-rolled-back with NotWhy.
|
||||
Rollback string `json:"rollback"`
|
||||
NotWhy string `json:"not_why,omitempty"`
|
||||
Show string `json:"show"`
|
||||
}
|
||||
|
||||
// sayRollback raises the gate's condition at once — the probe keeps it from then — and says the event.
|
||||
func sayRollback(ctx context.Context, open *stores, module string, g *inventory.PlanGate, notWhy string) {
|
||||
o := gateObservation(module, g.Component, g.Machines, g.Rollback, g.Why, notWhy, g.To, g.From)
|
||||
fmt.Println(o.Summary)
|
||||
d := doctorFrom
|
||||
if d == nil {
|
||||
return
|
||||
}
|
||||
if d.keeper != nil {
|
||||
o.Source = gateProbe
|
||||
if _, err := d.keeper.Observe(ctx, o); err != nil {
|
||||
fmt.Printf("the gate's condition %s could not be kept: %v\n", o.Key(), err)
|
||||
}
|
||||
}
|
||||
if d.teller == nil {
|
||||
return
|
||||
}
|
||||
body, err := json.Marshal(rolledBackEvent{Event: link.KeyRolledBack, At: time.Now().UTC(), Module: module,
|
||||
Component: g.Component, Machines: g.Machines, From: g.To, To: g.From, Why: g.Why, Rollback: g.Rollback,
|
||||
NotWhy: notWhy, Show: conditions.Condition{Key: o.Key()}.Show()})
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
saying, cancel := context.WithTimeout(ctx, 10*time.Second)
|
||||
defer cancel()
|
||||
if err := d.teller.PublishSeatEvent(saying, conditions.Seat, link.KeyRolledBack, body); err != nil {
|
||||
fmt.Printf("%s's rollback could NOT be said on the bus: %v\n", module, err)
|
||||
}
|
||||
}
|
||||
|
||||
// gateObservation is a failed gate as a condition: `core.<component>.<machine>.rolled-back` for the
|
||||
// core (to-be 45 §2), `build.<module>.<machine>.rolled-back` for any other module; `rollback-failed`
|
||||
// when it could not be put back. Urgent for the core and for a build left in place; a warning for a
|
||||
// module put back, which runs what it ran before.
|
||||
func gateObservation(module, component string, machines []string, rollback, why, notWhy, failed, previous string) conditions.Observation {
|
||||
machine := strings.Join(machines, ",")
|
||||
o := conditions.Observation{Scope: conditions.ScopeBuild, ID: module + "." + machine, Kind: kindRolledBack,
|
||||
Token: kindRolledBack, Machine: firstOf(machines), Severity: conditions.Warning, Source: gateProbe,
|
||||
Summary: fmt.Sprintf("%s's build %s failed its gate on %s and was put back to %s: %s", module, short(failed),
|
||||
machine, short(previous), why)}
|
||||
if component != "" {
|
||||
o.Scope, o.ID, o.Severity = conditions.ScopeCore, component+"."+machine, conditions.Urgent
|
||||
}
|
||||
if len(machines) > 1 {
|
||||
o.Also = machines[1:]
|
||||
}
|
||||
if rollback != inventory.RolledBack && rollback != inventory.RollingBack {
|
||||
o.Kind, o.Token, o.Severity = kindRollbackFailed, kindRollbackFailed, conditions.Urgent
|
||||
o.Resolver = conditions.ResolverOperator
|
||||
o.Summary = fmt.Sprintf("%s's build %s failed its gate on %s and was NOT put back: %s — %s", module, short(failed),
|
||||
machine, why, orNone(notWhy))
|
||||
}
|
||||
return o
|
||||
}
|
||||
|
||||
// probeGates is DG: no build that failed its gate is left without its condition, and no witness's
|
||||
// rollback goes unsaid. Each module whose newest verdict is a failure keeps its condition; a newer build
|
||||
// that passes clears it. Each core component a machine's witness says it put back is `core.<component>.
|
||||
// <machine>.rolled-back`, urgent, until the machine's reports stop saying it.
|
||||
func probeGates(ctx context.Context, d *doctor) ([]conditions.Observation, error) {
|
||||
latest, err := d.open.inventory.LatestGates(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
var out []conditions.Observation
|
||||
for _, v := range latest {
|
||||
if v.Verdict != inventory.GateFailed {
|
||||
continue
|
||||
}
|
||||
notWhy := ""
|
||||
if v.Rollback == inventory.NotRolledBack {
|
||||
notWhy = v.Why
|
||||
}
|
||||
out = append(out, gateObservation(v.Module, v.Component, v.Machines, v.Rollback, v.Why, notWhy, v.Commit, v.Previous))
|
||||
}
|
||||
for node, list := range witnessed.all() {
|
||||
for _, r := range list {
|
||||
out = append(out, witnessObservation(node, r))
|
||||
}
|
||||
}
|
||||
return sortedFound(dedupeObservations(out)), nil
|
||||
}
|
||||
|
||||
// dedupeObservations keeps one observation per key, the first.
|
||||
func dedupeObservations(list []conditions.Observation) []conditions.Observation {
|
||||
seen := map[string]bool{}
|
||||
var out []conditions.Observation
|
||||
for _, o := range list {
|
||||
if seen[o.Key()] {
|
||||
continue
|
||||
}
|
||||
seen[o.Key()] = true
|
||||
out = append(out, o)
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
// servedOnTheBus asks the bus's discovery what every machine's node tools answer and which modules'
|
||||
// tools are served where. A variable so a test needs no runtime.
|
||||
var servedOnTheBus = func(ctx context.Context, conn *nats.Conn) (map[string]served, error) {
|
||||
inbox := conn.NewRespInbox()
|
||||
sub, err := conn.SubscribeSync(inbox)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer func() { _ = sub.Unsubscribe() }()
|
||||
if err := conn.PublishRequest("$SRV.INFO", inbox, nil); err != nil {
|
||||
return nil, fmt.Errorf("asking the bus who serves what: %w", err)
|
||||
}
|
||||
out := map[string]served{}
|
||||
add := func(node string) served {
|
||||
s, ok := out[node]
|
||||
if !ok {
|
||||
s = served{tools: map[string]bool{}}
|
||||
}
|
||||
return s
|
||||
}
|
||||
deadline := time.Now().Add(discoveryPatience)
|
||||
for time.Now().Before(deadline) {
|
||||
wait, cancel := context.WithTimeout(ctx, discoveryQuiet)
|
||||
msg, err := sub.NextMsgWithContext(wait)
|
||||
cancel()
|
||||
if err != nil {
|
||||
if ctx.Err() != nil {
|
||||
return nil, ctx.Err()
|
||||
}
|
||||
break
|
||||
}
|
||||
var info micro.Info
|
||||
if json.Unmarshal(msg.Data, &info) != nil {
|
||||
continue
|
||||
}
|
||||
if info.Name == broker.RuntimeModule {
|
||||
node := info.Metadata["node"]
|
||||
if node == "" {
|
||||
node = info.ID
|
||||
}
|
||||
s := add(node)
|
||||
s.runtime = true
|
||||
out[node] = s
|
||||
}
|
||||
for _, e := range info.Endpoints {
|
||||
if e.Metadata["kind"] != "tool" || e.Metadata["module"] == "" {
|
||||
continue
|
||||
}
|
||||
node := e.Metadata["node"]
|
||||
if node == "" {
|
||||
node = info.ID
|
||||
}
|
||||
s := add(node)
|
||||
s.tools[e.Metadata["module"]] = true
|
||||
out[node] = s
|
||||
}
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
// gateLines is a plan's gates as `plans <id>` says them: the rollout's record.
|
||||
func gateLine(g *inventory.PlanGate) string {
|
||||
if g == nil {
|
||||
return ""
|
||||
}
|
||||
what := "judging"
|
||||
switch g.Verdict {
|
||||
case inventory.GatePassed:
|
||||
what = "passed"
|
||||
case inventory.GateFailed:
|
||||
what = "FAILED"
|
||||
}
|
||||
line := fmt.Sprintf("gate on %s: %s", strings.Join(g.Machines, ", "), what)
|
||||
if g.Component != "" {
|
||||
line += " (core: " + g.Component + ")"
|
||||
}
|
||||
if g.From != "" || g.To != "" {
|
||||
line += fmt.Sprintf(", %s → %s", short(orNone(g.From)), short(g.To))
|
||||
}
|
||||
if g.Took != "" {
|
||||
line += ", after " + g.Took
|
||||
}
|
||||
if g.Why != "" {
|
||||
line += ": " + g.Why
|
||||
} else if g.Last != "" {
|
||||
line += fmt.Sprintf(" (%d of %d passes; wanting: %s)", g.Passes, gatePasses, g.Last)
|
||||
} else if g.Verdict == "" {
|
||||
line += fmt.Sprintf(" (%d of %d passes)", g.Passes, gatePasses)
|
||||
}
|
||||
if g.Rollback != "" {
|
||||
line += "; " + g.Rollback
|
||||
}
|
||||
return line
|
||||
}
|
||||
|
||||
// witnessObservation is a witness's verdict as a condition (ADR 0236, the host's contract):
|
||||
// `core.<component>.<node>.<outcome>`, urgent for rolled-back, not-reversible, restore-failed and
|
||||
// halted, a warning for nothing-to-restore and unwitnessed.
|
||||
func witnessObservation(node string, r lease.Rollback) conditions.Observation {
|
||||
severity := conditions.Warning
|
||||
if lease.Urgent(r.Outcome) {
|
||||
severity = conditions.Urgent
|
||||
}
|
||||
outcome := r.Outcome
|
||||
if outcome == "" {
|
||||
outcome = lease.OutcomeRolledBack
|
||||
}
|
||||
what := map[string]string{
|
||||
lease.OutcomeRolledBack: "and put back " + short(orNone(r.To)),
|
||||
lease.OutcomeNotReversible: "and left it: it is not reversible",
|
||||
lease.OutcomeNothingToRestore: "and had nothing to put back",
|
||||
lease.OutcomeRestoreFailed: "and could not put the previous build back",
|
||||
lease.OutcomeUnwitnessed: "and could not judge it at all",
|
||||
lease.OutcomeHalted: "and gave up: the build before it fails too",
|
||||
}[outcome]
|
||||
return conditions.Observation{Scope: conditions.ScopeCore, ID: r.Component + "." + node, Kind: kindRolledBack,
|
||||
Token: outcome, Machine: node, Severity: severity,
|
||||
Summary: fmt.Sprintf("the witness on %s judged the %s %s not healthy %s at %s: %s", node, r.Component,
|
||||
short(r.From), what, r.At.UTC().Format(time.RFC3339), r.Why)}
|
||||
}
|
||||
@@ -0,0 +1,523 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"reflect"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/novox/mesh-controller/internal/catalogue"
|
||||
"github.com/novox/mesh-controller/internal/conditions"
|
||||
"github.com/novox/mesh-controller/internal/inventory"
|
||||
"github.com/novox/mesh-controller/internal/lease"
|
||||
"github.com/novox/mesh-controller/internal/link"
|
||||
)
|
||||
|
||||
// The gate on a plan's first machine and the rollback after it (novox/hq ADR 0235, to-be 45 §8).
|
||||
|
||||
// gateMesh is a mesh with `app` running on anchor and laptop at build c1, a newer build c2 registered,
|
||||
// a plan whose tier built c2, and every send recorded and answered — the machine applies what it is
|
||||
// sent and reports it — with the gate's bounds shortened to judge in three steps.
|
||||
type gateMesh struct {
|
||||
open *stores
|
||||
sent [][]string
|
||||
health map[string]health // per machine, what a judging finds; healthy when unsaid
|
||||
keeper *conditions.Keeper
|
||||
told *conditions.Told
|
||||
}
|
||||
|
||||
func aGateMesh(t *testing.T) *gateMesh {
|
||||
t.Helper()
|
||||
open := aMesh(t)
|
||||
ctx := t.Context()
|
||||
inv := open.inventory
|
||||
g := &gateMesh{open: open, health: map[string]health{}}
|
||||
g.keeper, _ = withConditionsInMemory(t)
|
||||
g.told = &conditions.Told{}
|
||||
was := doctorFrom
|
||||
doctorFrom = &doctor{open: open, keeper: g.keeper, teller: g.told}
|
||||
t.Cleanup(func() { doctorFrom = was })
|
||||
|
||||
for _, b := range []inventory.Build{
|
||||
{ID: "build-1", Module: "app", Commit: "c1", Repository: "novox/mesh-catalog", Path: "modules/app",
|
||||
Asked: time.Now().Add(-2 * time.Hour), At: time.Now().Add(-2 * time.Hour)},
|
||||
{ID: "build-2", Module: "app", Commit: "c2", Repository: "novox/mesh-catalog", Path: "modules/app",
|
||||
Asked: time.Now().Add(-time.Minute), At: time.Now().Add(-time.Minute)},
|
||||
} {
|
||||
manifest, _ := json.Marshal(catalogue.Manifest{Module: "app", Version: b.Commit})
|
||||
b.Manifest = manifest
|
||||
if err := inv.RecordBuild(ctx, b); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
register := func(commit string, asked time.Time) {
|
||||
if err := inv.RegisterModule(ctx, catalogue.Manifest{Module: "app", Version: commit},
|
||||
inventory.Source{Repository: "novox/mesh-catalog", Seat: "git", Path: "modules/app", BuiltFrom: commit,
|
||||
Head: commit, Asked: asked}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
register("c1", time.Now().Add(-2*time.Hour))
|
||||
for _, n := range []string{"anchor", "laptop"} {
|
||||
if _, err := inv.Assign(ctx, n, "app"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := inv.RecordSent(ctx, nodeID(t, open, n), "d-"+n+"-c1", map[string]string{"app": "c1"}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
register("c2", time.Now().Add(-time.Minute))
|
||||
|
||||
// Every send: recorded with the build the module is at, applied and reported by the machine.
|
||||
n := 0
|
||||
wasSend := sendRollout
|
||||
sendRollout = func(ctx context.Context, open *stores, names []string) ([]string, error) {
|
||||
g.sent = append(g.sent, append([]string(nil), names...))
|
||||
current, err := open.inventory.CurrentBuilds(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
for _, node := range names {
|
||||
n++
|
||||
digest := fmt.Sprintf("d-%s-%d", node, n)
|
||||
if err := open.inventory.RecordSent(ctx, nodeID(t, open, node), digest, map[string]string{"app": current["app"].Commit}); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if _, err := open.inventory.RecordDoing(ctx, nodeID(t, open, node), inventory.Doing{Node: node,
|
||||
Outcome: inventory.OutcomeApplied, Declared: digest, Applied: 1, At: time.Now()}); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
}
|
||||
return names, nil
|
||||
}
|
||||
t.Cleanup(func() { sendRollout = wasSend })
|
||||
|
||||
// What a judging reads: the store's reports, and a health the test says per machine.
|
||||
wasGather := gatherGateFacts
|
||||
gatherGateFacts = func(ctx context.Context, open *stores, component string) (gateFacts, error) {
|
||||
f := gateFacts{now: time.Now(), reports: map[string]inventory.Reported{}, engines: map[string]string{},
|
||||
rolledBack: map[string][]lease.Rollback{}, served: map[string]served{}}
|
||||
reports, err := open.inventory.LastReports(ctx)
|
||||
if err != nil {
|
||||
return f, err
|
||||
}
|
||||
for _, r := range reports {
|
||||
if g.health[r.Node] == healthBroken {
|
||||
r.Outcome = inventory.OutcomeFailed
|
||||
}
|
||||
f.reports[r.Node] = r
|
||||
}
|
||||
for node, h := range g.health {
|
||||
if h == healthNotYet {
|
||||
f.rolledBack[node] = nil
|
||||
r := f.reports[node]
|
||||
r.Current = false
|
||||
f.reports[node] = r
|
||||
}
|
||||
}
|
||||
return f, nil
|
||||
}
|
||||
t.Cleanup(func() { gatherGateFacts = wasGather })
|
||||
|
||||
wasSettle, wasEvery, wasBound := gateSettle, gateEvery, gateBound
|
||||
// One judging per advance until a test says otherwise: a pass waits gateEvery for the next.
|
||||
gateSettle, gateEvery = 0, time.Hour
|
||||
t.Cleanup(func() { gateSettle, gateEvery, gateBound = wasSettle, wasEvery, wasBound })
|
||||
|
||||
built := time.Now().UTC()
|
||||
plan := inventory.Plan{ID: "plan-gate", Repository: "novox/mesh-catalog", Branch: "main", Commit: "c2",
|
||||
Created: built, State: inventory.PlanBuilding, Tiers: [][]string{{"app"}},
|
||||
Modules: map[string]*inventory.PlanModule{"app": {State: "built", BuiltAt: &built, Commit: "c2",
|
||||
Build: "build-2"}}}
|
||||
if err := inv.SavePlan(ctx, &plan); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return g
|
||||
}
|
||||
|
||||
func (g *gateMesh) plan(t *testing.T) inventory.Plan {
|
||||
t.Helper()
|
||||
p, err := g.open.inventory.PlanByID(t.Context(), "plan-gate")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return p
|
||||
}
|
||||
|
||||
// A module's build that fails its gate on the first machine is put back there — the previous build
|
||||
// registered and sent to it again — and never reaches the second machine; it is marked, said as a
|
||||
// condition and an event, never registered or rolled back again.
|
||||
func TestABuildThatFailsItsGateIsRolledBackOnItsFirstMachineAndGoesNoFurther(t *testing.T) {
|
||||
g := aGateMesh(t)
|
||||
ctx := t.Context()
|
||||
inv := g.open.inventory
|
||||
|
||||
advancePlans(ctx, g.open) // the first machine is sent the new build
|
||||
if !reflect.DeepEqual(g.sent, [][]string{{"anchor"}}) {
|
||||
t.Fatalf("sent %v, not the first machine alone", g.sent)
|
||||
}
|
||||
if p := g.plan(t); p.Modules["app"].Previous != "c1" {
|
||||
t.Fatalf("the build the first machine ran before was not kept: %+v", p.Modules["app"])
|
||||
}
|
||||
|
||||
g.health["anchor"] = healthBroken // the new build breaks its first machine, after one good judging
|
||||
gateEvery = 0
|
||||
advancePlans(ctx, g.open)
|
||||
|
||||
p := g.plan(t)
|
||||
gate := p.Modules["app"].Gate
|
||||
if p.State != inventory.PlanFailed || gate == nil || gate.Verdict != inventory.GateFailed ||
|
||||
gate.Rollback != inventory.RolledBack {
|
||||
t.Fatalf("the plan is %s (%s), its gate %+v", p.State, p.Note, gate)
|
||||
}
|
||||
if !reflect.DeepEqual(g.sent, [][]string{{"anchor"}, {"anchor"}}) {
|
||||
t.Fatalf("sent %v: the rollback goes to the first machine, and nothing reaches laptop", g.sent)
|
||||
}
|
||||
if current, _ := inv.CurrentBuilds(ctx); current["app"].Commit != "c1" {
|
||||
t.Fatalf("the module is registered at %s, not put back to c1", current["app"].Commit)
|
||||
}
|
||||
if sent, _, _ := inv.SentBuilds(ctx, "anchor"); sent["app"] != "c1" {
|
||||
t.Fatalf("anchor was last sent %v, not the previous build", sent)
|
||||
}
|
||||
if sent, _, _ := inv.SentBuilds(ctx, "laptop"); sent["app"] != "c1" {
|
||||
t.Fatalf("laptop was sent the failed build: %v", sent)
|
||||
}
|
||||
if failed, err := inv.GateFailed(ctx, "build-2"); err != nil || !failed {
|
||||
t.Fatalf("the build is not marked failed at its gate: %v %v", failed, err)
|
||||
}
|
||||
|
||||
// Said: a condition for the operator, and the event.
|
||||
open, err := g.keeper.Open(ctx)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
var key string
|
||||
for _, c := range open {
|
||||
if c.Kind == kindRolledBack {
|
||||
key = c.Key
|
||||
}
|
||||
}
|
||||
if key != "build.app.anchor.rolled-back" {
|
||||
t.Fatalf("no rolled-back condition for app on anchor: %+v", open)
|
||||
}
|
||||
saidIt := false
|
||||
for _, e := range g.told.Said() {
|
||||
saidIt = saidIt || e.Event == link.KeyRolledBack
|
||||
}
|
||||
if !saidIt {
|
||||
t.Fatalf("the rollback was not said as an event: %+v", g.told.Said())
|
||||
}
|
||||
|
||||
// Never again: more passes send nothing, a second judging rolls nothing back, and the failed build
|
||||
// heard again is not registered.
|
||||
advancePlans(ctx, g.open)
|
||||
state := p.Modules["app"]
|
||||
gateFailed(ctx, g.open, &p, "app", state, []string{"anchor"}, "again")
|
||||
if len(g.sent) != 2 {
|
||||
t.Fatalf("sent again after the rollback: %v", g.sent)
|
||||
}
|
||||
_, _, err = takeIn(ctx, inv, link.BuildResult{ID: "build-2", Repository: "novox/mesh-catalog", Path: "modules/app",
|
||||
Commit: "c2", Manifest: mustJSON(t, catalogue.Manifest{Module: "app", Version: "c2"})})
|
||||
if err == nil || !strings.Contains(err.Error(), "failed its gate") {
|
||||
t.Fatalf("the failed build was registered again: %v", err)
|
||||
}
|
||||
if current, _ := inv.CurrentBuilds(ctx); current["app"].Commit != "c1" {
|
||||
t.Fatalf("the module moved to %s", current["app"].Commit)
|
||||
}
|
||||
if _, err := retryPlan(ctx, g.open, "plan-gate"); err == nil || !strings.Contains(err.Error(), "failed its gate") {
|
||||
t.Fatalf("a plan stopped at a failed gate was retried: %v", err)
|
||||
}
|
||||
|
||||
// The probe keeps the condition while the newest verdict is the failure.
|
||||
obs, err := probeGates(ctx, doctorFrom)
|
||||
if err != nil || len(obs) != 1 || obs[0].Key() != key {
|
||||
t.Fatalf("the gate's probe found %+v %v", obs, err)
|
||||
}
|
||||
}
|
||||
|
||||
// A passing build rolls everywhere with no hand: judged healthy on its first machine three times, then
|
||||
// the rest are sent, the verdict kept, and the plan done.
|
||||
func TestABuildThatPassesItsGateRollsEverywhereUnattended(t *testing.T) {
|
||||
g := aGateMesh(t)
|
||||
ctx := t.Context()
|
||||
gateEvery = 0
|
||||
for i := 0; i < 6; i++ {
|
||||
advancePlans(ctx, g.open)
|
||||
}
|
||||
p := g.plan(t)
|
||||
if !reflect.DeepEqual(g.sent, [][]string{{"anchor"}, {"laptop"}}) {
|
||||
t.Fatalf("sent %v", g.sent)
|
||||
}
|
||||
gate := p.Modules["app"].Gate
|
||||
if p.State != inventory.PlanDone || gate == nil || gate.Verdict != inventory.GatePassed || gate.Passes < gatePasses {
|
||||
t.Fatalf("the plan is %s (%s), its gate %+v", p.State, p.Note, gate)
|
||||
}
|
||||
if v, found, err := g.open.inventory.GateOf(ctx, "build-2"); err != nil || !found || v.Verdict != inventory.GatePassed {
|
||||
t.Fatalf("the verdict was not kept: %+v %v %v", v, found, err)
|
||||
}
|
||||
if obs, err := probeGates(ctx, doctorFrom); err != nil || len(obs) != 0 {
|
||||
t.Fatalf("a passing build left a condition: %+v %v", obs, err)
|
||||
}
|
||||
}
|
||||
|
||||
// A first machine that never becomes healthy fails its gate at the bound, and is put back.
|
||||
func TestABuildNeverHealthyFailsAtTheBound(t *testing.T) {
|
||||
g := aGateMesh(t)
|
||||
ctx := t.Context()
|
||||
advancePlans(ctx, g.open)
|
||||
g.health["anchor"] = healthNotYet
|
||||
gateEvery = 0
|
||||
advancePlans(ctx, g.open)
|
||||
if p := g.plan(t); !p.Open() || len(g.sent) != 1 {
|
||||
t.Fatalf("judged before the bound: %s %v", p.State, g.sent)
|
||||
}
|
||||
gateBound = -time.Second
|
||||
advancePlans(ctx, g.open)
|
||||
p := g.plan(t)
|
||||
if gate := p.Modules["app"].Gate; p.State != inventory.PlanFailed || gate.Verdict != inventory.GateFailed ||
|
||||
!strings.Contains(gate.Why, "not healthy within") || gate.Rollback != inventory.RolledBack {
|
||||
t.Fatalf("the plan is %s, its gate %+v", p.State, gate)
|
||||
}
|
||||
}
|
||||
|
||||
// The health a judging finds, per core component and for a module.
|
||||
func TestTheHealthDefinitions(t *testing.T) {
|
||||
now := time.Now()
|
||||
since := now.Add(-time.Minute)
|
||||
applied := func(node string) gateFacts {
|
||||
at := now
|
||||
return gateFacts{now: now, reports: map[string]inventory.Reported{node: {Node: node,
|
||||
Outcome: inventory.OutcomeApplied, At: &at, Current: true}}, engines: map[string]string{},
|
||||
served: map[string]served{}, rolledBack: map[string][]lease.Rollback{}}
|
||||
}
|
||||
m := catalogue.Manifest{Module: "app", Tools: []string{"app_list"}}
|
||||
|
||||
f := applied("anchor")
|
||||
f.served["anchor"] = served{runtime: true, tools: map[string]bool{}}
|
||||
if h, why := judgeHealth("app", "", m, "anchor", since, f); h != healthNotYet || !strings.Contains(why, "tools") {
|
||||
t.Errorf("a module whose tools are not served is %v (%s)", h, why)
|
||||
}
|
||||
f.served["anchor"].tools["app"] = true
|
||||
if h, why := judgeHealth("app", "", m, "anchor", since, f); h != healthGood {
|
||||
t.Errorf("a module applied with its tools served is %v (%s)", h, why)
|
||||
}
|
||||
// A condition raised about it on that machine since it was sent: not yet healthy.
|
||||
f.judged = true
|
||||
f.open = []conditions.Condition{{Key: "provider.app.anchor.x.failing", Subject: conditions.Subject{
|
||||
Scope: conditions.ScopeProvider, ID: "app.anchor.x", Machine: "anchor"}, Raised: now, Summary: "failing"}}
|
||||
if h, _ := judgeHealth("app", "", m, "anchor", since, f); h != healthNotYet {
|
||||
t.Errorf("a module with a new condition about it is %v", h)
|
||||
}
|
||||
f.open[0].Raised = since.Add(-time.Hour)
|
||||
if h, _ := judgeHealth("app", "", m, "anchor", since, f); h != healthGood {
|
||||
t.Errorf("a condition older than the send counted against it: %v", h)
|
||||
}
|
||||
// A witness that put it back: broken.
|
||||
f.rolledBack["anchor"] = []lease.Rollback{{Component: lease.ComponentNodeTools, From: "sha256:aa", To: "sha256:bb",
|
||||
Outcome: lease.OutcomeRolledBack, Why: "x", At: now}}
|
||||
if h, _ := judgeHealth("node-tools", lease.ComponentNodeTools, catalogue.Manifest{}, "anchor", since, f); h != healthBroken {
|
||||
t.Errorf("a component a witness put back is %v", h)
|
||||
}
|
||||
// The node tools answer, or not.
|
||||
f = applied("anchor")
|
||||
if h, _ := judgeHealth("node-tools", lease.ComponentNodeTools, catalogue.Manifest{}, "anchor", since, f); h != healthNotYet {
|
||||
t.Errorf("node tools not answering are %v", h)
|
||||
}
|
||||
f.served["anchor"] = served{runtime: true}
|
||||
if h, _ := judgeHealth("node-tools", lease.ComponentNodeTools, catalogue.Manifest{}, "anchor", since, f); h != healthGood {
|
||||
t.Errorf("node tools answering are %v", h)
|
||||
}
|
||||
// The controller: the lease held since the send, by a controller that says it is ready.
|
||||
f = applied("control")
|
||||
f.holder = &lease.Holder{Taken: since.Add(-time.Hour), Health: &lease.Health{Ready: true}}
|
||||
if h, _ := judgeHealth("mesh-controller", lease.ComponentController, catalogue.Manifest{}, "control", since, f); h != healthNotYet {
|
||||
t.Errorf("a lease held by a controller older than the build is %v", h)
|
||||
}
|
||||
f.holder.Taken = now
|
||||
if h, _ := judgeHealth("mesh-controller", lease.ComponentController, catalogue.Manifest{}, "control", since, f); h != healthGood {
|
||||
t.Errorf("a new controller holding the lease and ready is %v", h)
|
||||
}
|
||||
f.holder.Health = &lease.Health{Why: "the self-check has not finished its first run"}
|
||||
if h, why := judgeHealth("mesh-controller", lease.ComponentController, catalogue.Manifest{}, "control", since, f); h != healthNotYet ||
|
||||
!strings.Contains(why, "first run") {
|
||||
t.Errorf("a controller not ready is %v (%s)", h, why)
|
||||
}
|
||||
// A machine that refused what it was sent: broken.
|
||||
f = applied("anchor")
|
||||
r := f.reports["anchor"]
|
||||
r.Outcome = inventory.OutcomeRefused
|
||||
f.reports["anchor"] = r
|
||||
if h, _ := judgeHealth("app", "", catalogue.Manifest{}, "anchor", since, f); h != healthBroken {
|
||||
t.Errorf("a refusal is %v", h)
|
||||
}
|
||||
}
|
||||
|
||||
// The bus is never rolled out: its policy records whatever is said, a person's roll-out is refused,
|
||||
// a plan builds it and sends nothing, and a push naming its machine is refused while a new build waits.
|
||||
func TestTheBusIsNeverRolledOutAutomatically(t *testing.T) {
|
||||
open := aMesh(t)
|
||||
ctx := t.Context()
|
||||
inv := open.inventory
|
||||
bus := catalogue.Manifest{Module: "nats", Version: "2", Provides: []catalogue.Offer{{Name: "mesh-bus"}},
|
||||
Upgrade: &catalogue.UpgradePolicy{Policy: catalogue.PolicyRoll}}
|
||||
if err := inv.RegisterModule(ctx, bus, inventory.Source{Repository: "novox/mesh-catalog", Seat: "git",
|
||||
Path: "modules/nats", BuiltFrom: "n2", Head: "n2"}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := inv.Assign(ctx, "anchor", "nats"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := inv.RecordSent(ctx, nodeID(t, open, "anchor"), "d-anchor", map[string]string{"nats": "n1"}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
u, err := inv.UpgradeOf(ctx, "nats")
|
||||
if err != nil || u.RollOut || u.From != catalogue.FromBus {
|
||||
t.Fatalf("the bus's policy is %+v %v", u, err)
|
||||
}
|
||||
if err := inv.SetUpgradeOf(ctx, "nats", inventory.Upgrade{RollOut: true}); !errors.Is(err, inventory.ErrBusIsPlanned) {
|
||||
t.Fatalf("a person rolled the bus out: %v", err)
|
||||
}
|
||||
var sent [][]string
|
||||
was := sendRollout
|
||||
sendRollout = func(_ context.Context, _ *stores, names []string) ([]string, error) {
|
||||
sent = append(sent, names)
|
||||
return names, nil
|
||||
}
|
||||
t.Cleanup(func() { sendRollout = was })
|
||||
now := time.Now().UTC()
|
||||
plan := inventory.Plan{ID: "plan-bus", Repository: "novox/mesh-catalog", Commit: "n2", Created: now,
|
||||
State: inventory.PlanBuilding, Tiers: [][]string{{"nats"}},
|
||||
Modules: map[string]*inventory.PlanModule{"nats": {State: "built", BuiltAt: &now, Commit: "n2", Build: "b"}}}
|
||||
if err := inv.SavePlan(ctx, &plan); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
advancePlans(ctx, open)
|
||||
if p, _ := inv.PlanByID(ctx, "plan-bus"); p.State != inventory.PlanDone || len(sent) != 0 {
|
||||
t.Fatalf("a plan sent the bus: %s %v", p.State, sent)
|
||||
}
|
||||
held, err := busHeld(ctx, inv, []string{"anchor", "laptop"})
|
||||
if err != nil || !strings.Contains(held["anchor"], "planned step") || held["laptop"] != "" {
|
||||
t.Fatalf("a push may send the bus's machine: %v %v", held, err)
|
||||
}
|
||||
// The planned step refuses to start without its word on reversibility, and without a snapshot.
|
||||
if err := busCommand(ctx, []string{"upgrade", "--why", "2.11"}); err == nil || !strings.Contains(err.Error(), "reversible") {
|
||||
t.Fatalf("a bus upgrade started without saying whether it can be reverted: %v", err)
|
||||
}
|
||||
if err := busCommand(ctx, []string{"upgrade", "--why", "2.11", "--reversible"}); err == nil ||
|
||||
!strings.Contains(err.Error(), "snapshot") {
|
||||
t.Fatalf("a bus upgrade started without a snapshot: %v", err)
|
||||
}
|
||||
if err := busCommand(ctx, []string{"upgrade", "--why", "2.11", "--reversible", "--snapshot-taken", "nightly"}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if !reflect.DeepEqual(sent, [][]string{{"anchor"}}) {
|
||||
t.Fatalf("the step sent %v, not the bus's machine", sent)
|
||||
}
|
||||
s, found, err := inv.LatestBusStep(ctx)
|
||||
if err != nil || !found || s.Snapshot != "nightly" || s.Ended != nil || !reflect.DeepEqual(s.Machines, []string{"anchor"}) {
|
||||
t.Fatalf("the step is %+v %v %v", s, found, err)
|
||||
}
|
||||
}
|
||||
|
||||
// A merge that deletes a module's directory builds nothing for it: the module is forgotten where
|
||||
// nothing holds it, and a plan whose build finds no manifest goes on instead of failing.
|
||||
func TestAMergeThatDeletesAModulePlansNothingToBuildForIt(t *testing.T) {
|
||||
open := aMesh(t)
|
||||
ctx := t.Context()
|
||||
inv := open.inventory
|
||||
asked := asksRecorded(t)
|
||||
for _, name := range []string{"gone", "kept"} {
|
||||
if err := inv.RegisterModule(ctx, catalogue.Manifest{Module: name, Version: "1"}, inventory.Source{
|
||||
Repository: "novox/mesh-catalog", Seat: "git", Path: "modules/" + name, BuiltFrom: "c0", Head: "c0"}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
m := link.SourceMoved{Owner: "novox", Repo: "mesh-catalog", Base: "main", Commit: "c1deadbeef",
|
||||
Paths: []string{"modules/gone/module.json", "modules/gone/index.ts"},
|
||||
Removed: []string{"modules/gone/module.json", "modules/gone/index.ts"}}
|
||||
if err := (following{open: open}).SourceMoved(ctx, m); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if len(*asked) != 0 {
|
||||
t.Fatalf("a deleted module was asked to build: %v", *asked)
|
||||
}
|
||||
if plans, _ := inv.OpenPlans(ctx); len(plans) != 0 {
|
||||
t.Fatalf("a merge that only deleted a module made a plan: %+v", plans)
|
||||
}
|
||||
shelf, err := inv.Catalogue(ctx)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, still := shelf["gone"]; still {
|
||||
t.Fatal("a deleted module nothing holds was not forgotten")
|
||||
}
|
||||
|
||||
// Without the announcer saying what went: the build finds no manifest, and the plan goes on.
|
||||
asked0 := time.Now().UTC().Add(-time.Minute)
|
||||
plan := inventory.Plan{ID: "plan-deleted", Repository: "novox/mesh-catalog", Commit: "c2", Created: asked0,
|
||||
State: inventory.PlanBuilding, Tiers: [][]string{{"kept"}},
|
||||
Modules: map[string]*inventory.PlanModule{"kept": {State: "asked", AskedAt: &asked0, Build: "build-k"}}}
|
||||
if err := inv.SavePlan(ctx, &plan); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
planBuilt(ctx, open, "kept", "", "http://forge/novox/mesh-catalog.git has no module.json at modules/kept, so "+
|
||||
"there is nothing saying what it is: open …/module.json: no such file or directory", asked0, "build-k")
|
||||
p, _ := inv.PlanByID(ctx, "plan-deleted")
|
||||
if p.State == inventory.PlanFailed || p.Modules["kept"].State != planDeleted {
|
||||
t.Fatalf("a module deleted at its source failed the plan: %s %+v", p.State, p.Modules["kept"])
|
||||
}
|
||||
}
|
||||
|
||||
func mustJSON(t *testing.T, v any) []byte {
|
||||
t.Helper()
|
||||
b, err := json.Marshal(v)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return b
|
||||
}
|
||||
|
||||
func nodeID(t *testing.T, open *stores, name string) string {
|
||||
t.Helper()
|
||||
n, err := open.inventory.NodeByName(context.Background(), name)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return n.ID
|
||||
}
|
||||
|
||||
// What a machine's witness says in its reports is a condition while it says it, by the host's words:
|
||||
// `core.<component>.<node>.<outcome>`, urgent for a rollback, a warning when it could not judge — and
|
||||
// gone with the first report that no longer carries it.
|
||||
func TestAWitnessesVerdictIsAConditionWhileItsReportsSayIt(t *testing.T) {
|
||||
open := aMesh(t)
|
||||
d := &doctor{open: open}
|
||||
at := time.Now().UTC()
|
||||
witnessed.heard("control", []lease.Rollback{
|
||||
{Component: lease.ComponentController, From: "sha256:new", To: "sha256:old", Outcome: lease.OutcomeRolledBack,
|
||||
Why: "the new controller did not take the lease within 60s", At: at},
|
||||
{Component: lease.ComponentNodeTools, From: "sha256:t2", Outcome: lease.OutcomeUnwitnessed, Why: "no grant", At: at},
|
||||
}, at)
|
||||
t.Cleanup(func() { witnessed.heard("control", nil, time.Now()) })
|
||||
obs, err := probeGates(t.Context(), d)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
got := map[string]conditions.Severity{}
|
||||
for _, o := range obs {
|
||||
got[o.Key()] = o.Severity
|
||||
}
|
||||
want := map[string]conditions.Severity{"core.controller.control.rolled-back": conditions.Urgent,
|
||||
"core.node-tools.control.unwitnessed": conditions.Warning}
|
||||
if !reflect.DeepEqual(got, want) {
|
||||
t.Fatalf("the witness's verdicts are %v", got)
|
||||
}
|
||||
witnessed.heard("control", nil, time.Now())
|
||||
if obs, err := probeGates(t.Context(), d); err != nil || len(obs) != 0 {
|
||||
t.Fatalf("a verdict the reports no longer carry is still said: %+v %v", obs, err)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,226 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"slices"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/nats-io/nats.go"
|
||||
|
||||
"github.com/novox/mesh-controller/internal/broker"
|
||||
"github.com/novox/mesh-controller/internal/catalogue"
|
||||
"github.com/novox/mesh-controller/internal/conditions"
|
||||
"github.com/novox/mesh-controller/internal/inventory"
|
||||
"github.com/novox/mesh-controller/internal/lease"
|
||||
)
|
||||
|
||||
// The core components' health definitions, as probes in the self-check's registry (novox/hq to-be 45
|
||||
// §8, §4 `H-*`). The same definitions judge a core component's new build on its first machine (the
|
||||
// gate, gate.go); here they are run against every machine on the self-check's schedule, so a core
|
||||
// component that stops being healthy between upgrades is said too.
|
||||
//
|
||||
// controller holds the lease, and says it is ready: its self-check ran, `status` answered in bound
|
||||
// node-engine has reported its current declaration, under a build the mesh delivered
|
||||
// node tools announced, and answer the bus's discovery
|
||||
// bus every stream and durable consumer the mesh defines is there, and a request crosses the
|
||||
// bus to every machine's node tools and back
|
||||
const kindCoreUnhealthy = "core-unhealthy"
|
||||
|
||||
// probeControllerHealth is H-controller: the lease is held, fresh, by a controller that says it is
|
||||
// ready — or one that started less than the witness's bound ago.
|
||||
func probeControllerHealth(ctx context.Context, d *doctor) ([]conditions.Observation, error) {
|
||||
h, found, err := theLease.holder(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
unhealthy := func(why string) []conditions.Observation {
|
||||
node := d.host
|
||||
if found && h.Host != "" {
|
||||
node = h.Host
|
||||
}
|
||||
return []conditions.Observation{{Scope: conditions.ScopeCore, ID: lease.ComponentController + "." + node,
|
||||
Token: "unhealthy", Kind: kindCoreUnhealthy, Machine: node, Severity: conditions.Urgent,
|
||||
Summary: "the controller is not healthy: " + why}}
|
||||
}
|
||||
switch {
|
||||
case !found:
|
||||
return unhealthy("nobody holds the controller lease"), nil
|
||||
case time.Since(h.Renewed) > lease.FreshWithin:
|
||||
return unhealthy(fmt.Sprintf("the lease was last renewed %s ago", time.Since(h.Renewed).Round(time.Second))), nil
|
||||
case (h.Health == nil || !h.Health.Ready) && time.Since(h.Taken) > lease.ReadyWithin:
|
||||
why := "the controller holding the lease does not say it is ready"
|
||||
if h.Health != nil && h.Health.Why != "" {
|
||||
why += ": " + h.Health.Why
|
||||
}
|
||||
return unhealthy(why), nil
|
||||
}
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
// probeEngineHealth is H-engine: every machine heard from has reported its current declaration — the
|
||||
// last one it was sent — under a node-engine build the mesh delivered, or is inside the bound of a send.
|
||||
func probeEngineHealth(ctx context.Context, d *doctor) ([]conditions.Observation, error) {
|
||||
return machinesHealth(ctx, d, hostModule)
|
||||
}
|
||||
|
||||
// probeToolsHealth is H-tools: every machine heard from that is assigned the node tools has them
|
||||
// announced, answering the bus's discovery.
|
||||
func probeToolsHealth(ctx context.Context, d *doctor) ([]conditions.Observation, error) {
|
||||
return machinesHealth(ctx, d, broker.RuntimeModule)
|
||||
}
|
||||
|
||||
// machinesHealth judges a component on every machine running it and heard from, by its health
|
||||
// definition, as the gate does — with no condition taken as the build's doing.
|
||||
func machinesHealth(ctx context.Context, d *doctor, component string) ([]conditions.Observation, error) {
|
||||
inv := d.open.inventory
|
||||
running, err := inv.Running(ctx, component)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
shelf, err := inv.Catalogue(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
f, err := gatherGateFacts(ctx, d.open, coreComponent(component))
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
f.judged = false
|
||||
heard := heardMachines(d)
|
||||
var out []conditions.Observation
|
||||
for _, node := range running {
|
||||
if !heard[node] {
|
||||
continue // a machine not heard from is S1's
|
||||
}
|
||||
if r, said := f.reports[node]; said && !r.Current && r.Sent != nil && time.Since(*r.Sent) < gateBound {
|
||||
continue // inside the bound of a send: S2's, and the gate's
|
||||
}
|
||||
// What the machine did with what it was sent is not the component's health: the node-engine's is
|
||||
// that it reported its current declaration at all, the node tools' that they answer.
|
||||
if r := f.reports[node]; r.Current || component == broker.RuntimeModule {
|
||||
now := time.Now()
|
||||
r.Current, r.Outcome = true, inventory.OutcomeApplied
|
||||
if r.At == nil {
|
||||
r.At = &now
|
||||
}
|
||||
f.reports[node] = r
|
||||
}
|
||||
if component == hostModule {
|
||||
// A node-engine reporting a build the mesh delivered, current or previous, is the version
|
||||
// split's (D10), not ill health; a build the mesh did not deliver is.
|
||||
f.engines[node] = deliveredOr(shelf[component], f.engines[node])
|
||||
}
|
||||
h, why := judgeHealth(component, coreComponent(component), shelf[component], node, time.Time{}, f)
|
||||
if h == healthGood {
|
||||
continue
|
||||
}
|
||||
out = append(out, conditions.Observation{Scope: conditions.ScopeCore, ID: coreComponent(component) + "." + node,
|
||||
Token: "unhealthy", Kind: kindCoreUnhealthy, Machine: node, Severity: conditions.Warning,
|
||||
Summary: fmt.Sprintf("the %s on %s is not healthy: %s", componentWords(component), node, why)})
|
||||
}
|
||||
return sortedFound(out), nil
|
||||
}
|
||||
|
||||
// deliveredOr is the reported engine build, or the current delivered one when the report names any
|
||||
// build at all: what H-engine judges is that it reports, not which.
|
||||
func deliveredOr(m catalogue.Manifest, reported string) string {
|
||||
if reported == "" {
|
||||
return reported
|
||||
}
|
||||
if v := deliveredVersions(m); len(v) > 0 {
|
||||
return v[0]
|
||||
}
|
||||
return reported
|
||||
}
|
||||
|
||||
// componentWords is a core component as a sentence names it.
|
||||
func componentWords(component string) string {
|
||||
switch component {
|
||||
case hostModule:
|
||||
return "node-engine"
|
||||
case broker.RuntimeModule:
|
||||
return "node tools"
|
||||
case catalogue.ControllerSeatName:
|
||||
return "controller"
|
||||
}
|
||||
return component
|
||||
}
|
||||
|
||||
// probeBusHealth is H-bus: every stream and durable consumer the mesh defines is on the bus, and a
|
||||
// request crosses it to every machine's node tools and back.
|
||||
func probeBusHealth(ctx context.Context, d *doctor) ([]conditions.Observation, error) {
|
||||
problems, err := busHealth(ctx, d)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if len(problems) == 0 {
|
||||
return nil, nil
|
||||
}
|
||||
return []conditions.Observation{{Scope: conditions.ScopeBus, ID: "mesh", Token: "unhealthy",
|
||||
Kind: kindCoreUnhealthy, Severity: conditions.Urgent,
|
||||
Summary: fmt.Sprintf("the bus is not healthy: %s", strings.Join(problems, "; ")),
|
||||
Said: strings.Join(problems, "; ")}}, nil
|
||||
}
|
||||
|
||||
// busHealth is the bus's health definition, as what is wanting: nothing when healthy. What the bus's
|
||||
// planned step checks after the bus is replaced, too.
|
||||
var busHealth = func(ctx context.Context, d *doctor) ([]string, error) {
|
||||
if d.js == nil {
|
||||
return nil, errors.New("this controller has no bus to ask")
|
||||
}
|
||||
streams, consumers, err := expectedBusObjects(ctx, d.open.inventory)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
js := d.js.Context()
|
||||
var problems []string
|
||||
for _, s := range streams {
|
||||
if _, err := js.StreamInfo(s.Name, nats.Context(ctx)); errors.Is(err, nats.ErrStreamNotFound) {
|
||||
problems = append(problems, "the stream "+s.Name+" is missing")
|
||||
} else if err != nil {
|
||||
return nil, fmt.Errorf("the stream %s cannot be read: %w", s.Name, err)
|
||||
}
|
||||
}
|
||||
for _, c := range consumers {
|
||||
_, err := js.ConsumerInfo(c.Stream, c.Name, nats.Context(ctx))
|
||||
if errors.Is(err, nats.ErrConsumerNotFound) || errors.Is(err, nats.ErrStreamNotFound) {
|
||||
problems = append(problems, consumerWords(c)+" is missing")
|
||||
} else if err != nil {
|
||||
return nil, fmt.Errorf("%s cannot be read: %w", consumerWords(c), err)
|
||||
}
|
||||
}
|
||||
// The round trip: every machine heard from that runs the node tools answers across the bus.
|
||||
running, err := d.open.inventory.Running(ctx, broker.RuntimeModule)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
heard := heardMachines(d)
|
||||
var expected []string
|
||||
for _, n := range running {
|
||||
if heard[n] {
|
||||
expected = append(expected, n)
|
||||
}
|
||||
}
|
||||
if len(expected) > 0 {
|
||||
answered, err := servedOnTheBus(ctx, d.js.Conn())
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
var silent []string
|
||||
for _, n := range expected {
|
||||
if !answered[n].runtime {
|
||||
silent = append(silent, n)
|
||||
}
|
||||
}
|
||||
// One machine's tools silent is that machine's (H-tools); none answering is the bus.
|
||||
if len(silent) == len(expected) {
|
||||
slices.Sort(silent)
|
||||
problems = append(problems, "no request crossed the bus and came back: no machine's node tools answered ("+
|
||||
strings.Join(silent, ", ")+")")
|
||||
}
|
||||
}
|
||||
return problems, nil
|
||||
}
|
||||
@@ -193,6 +193,10 @@ func TestANamedPushLeavesAMachineAPolicyHoldsBack(t *testing.T) {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
// Recorded by a person's choice: the default rolls out since novox/hq ADR 0235.
|
||||
if err := inv.SetUpgradeOf(ctx, "resolver", inventory.Upgrade{Why: "each machine checked by hand"}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
gens, err := generators(ctx, open)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
@@ -214,7 +218,7 @@ func TestANamedPushLeavesAMachineAPolicyHoldsBack(t *testing.T) {
|
||||
}
|
||||
before := digestOfLaptop()
|
||||
|
||||
// The change merges; the policy is the default, record. `push anchor` sends the anchor...
|
||||
// The change merges; the policy is record. `push anchor` sends the anchor...
|
||||
aResolver(t, open, "c2c2c2c2c2", asked.Add(time.Minute))
|
||||
d.declared = nil
|
||||
if _, err := sendRound(ctx, open, []string{"anchor"}, compose, d, ""); err != nil {
|
||||
|
||||
@@ -116,6 +116,9 @@ func run() error {
|
||||
return serve(ctx)
|
||||
case "upgrade":
|
||||
return upgradeCommand(ctx, args[1:])
|
||||
// The bus as a planned step (novox/hq to-be 45 §8, ADR 0235).
|
||||
case "bus":
|
||||
return busCommand(ctx, args[1:])
|
||||
case "declare":
|
||||
return declare(ctx, args[1:])
|
||||
case "overlay":
|
||||
|
||||
@@ -178,6 +178,15 @@ func retryRefusal(p inventory.Plan, plans []inventory.Plan) error {
|
||||
return fmt.Errorf("nothing in tier %d of %s failed to build or stopped rolling out — it stopped at: %s",
|
||||
p.Tier, p.ID, p.Note)
|
||||
}
|
||||
// **A build that failed its gate is not sent again** (novox/hq ADR 0235): it was put back on its first
|
||||
// machine, and retrying would judge the build the mesh put back, or send the failed one by hand.
|
||||
for _, m := range stopped {
|
||||
if g := p.Modules[m].Gate; g != nil && g.Verdict == inventory.GateFailed {
|
||||
return fmt.Errorf("%s failed its gate on %s (%s) and was put back: a build that failed its gate is not "+
|
||||
"sent again — a newer merge, or `rebuild %s`, makes a new build, judged at the gate again",
|
||||
m, strings.Join(g.Machines, ", "), g.Why, m)
|
||||
}
|
||||
}
|
||||
// **A rollout is retried unless the module has moved on**: a newer plan holding it sends — or
|
||||
// sent — a newer build, and sending this one again would put the older build back on its machines.
|
||||
for _, m := range stopped {
|
||||
|
||||
@@ -466,6 +466,27 @@ func pushCommand(ctx context.Context, args []string) error {
|
||||
which, len(asked), strings.Join(asked, ", "))
|
||||
}
|
||||
|
||||
// **The bus is replaced only as a planned step** (novox/hq ADR 0235, to-be 45 §8): a machine whose bus
|
||||
// would move is not sent by a push — named, it is refused; otherwise it is left and said.
|
||||
busKept, err := busHeld(ctx, inv, asked)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if len(busKept) > 0 {
|
||||
if len(args) == 1 {
|
||||
return fmt.Errorf("%s. Nothing was sent", busKept[args[0]])
|
||||
}
|
||||
var kept []string
|
||||
for _, n := range asked {
|
||||
if why, h := busKept[n]; h {
|
||||
fmt.Printf(" %s is left: %s\n", n, why)
|
||||
continue
|
||||
}
|
||||
kept = append(kept, n)
|
||||
}
|
||||
asked = kept
|
||||
}
|
||||
|
||||
// **The machine holding the bus first** (novox/hq issue 249): its declaration carries the bus's
|
||||
// user list, and a module's new grants are refused by the bus until that list says them. Among
|
||||
// the machines asked it goes first; not among them and behind, it is added — a named push whose
|
||||
|
||||
@@ -625,6 +625,10 @@ func TestAPlanStoppedAtItsFirstMachineIsRetried(t *testing.T) {
|
||||
}
|
||||
t.Cleanup(func() { sendRollout = was })
|
||||
|
||||
// a records: a person's choice, since the default rolls out (novox/hq ADR 0235).
|
||||
if err := open.inventory.SetUpgradeOf(ctx, "a", inventory.Upgrade{Why: "test"}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
long := time.Now().UTC().Add(-2 * time.Hour)
|
||||
stopped := inventory.Plan{ID: "plan-rollout", Repository: "novox/a", Branch: "main", Commit: "c0ffee",
|
||||
Created: long, State: inventory.PlanFailed, Tiers: [][]string{{"a"}, {"b"}},
|
||||
|
||||
@@ -5,6 +5,7 @@ import (
|
||||
"errors"
|
||||
"flag"
|
||||
"fmt"
|
||||
"slices"
|
||||
"sort"
|
||||
"strings"
|
||||
"time"
|
||||
@@ -223,6 +224,9 @@ func supersededBy(newer inventory.Plan, open []inventory.Plan, rollsOut func(str
|
||||
}
|
||||
var took []string
|
||||
for name, s := range old.Modules {
|
||||
if s != nil && s.State == planDeleted {
|
||||
continue // deleted at its source: nothing to plan again
|
||||
}
|
||||
if s == nil || s.State != "built" || (s.SentAt == nil && rollsOut(name)) {
|
||||
folded[name] = true
|
||||
took = append(took, name)
|
||||
@@ -427,7 +431,11 @@ func planBuilt(ctx context.Context, open *stores, module, commit, failed string,
|
||||
} else if askedBefore(asked, state.AskedAt) {
|
||||
continue
|
||||
}
|
||||
if failed != "" {
|
||||
if failed != "" && deletedAtSource(failed) {
|
||||
// Deleted at its source by the merge, not broken (novox/hq ADR 0235): the plan goes on.
|
||||
state.State, state.Why = planDeleted, "deleted at its source: "+firstLine(failed)
|
||||
forgetDeleted(ctx, inv, module, p)
|
||||
} else if failed != "" {
|
||||
state.State = "failed"
|
||||
state.Why = failed
|
||||
p.State = inventory.PlanFailed
|
||||
@@ -575,7 +583,7 @@ func advanceOnce(ctx context.Context, open *stores, p *inventory.Plan,
|
||||
var latest time.Time
|
||||
for _, m := range tier {
|
||||
s := p.Modules[m]
|
||||
if s == nil || s.State != "built" {
|
||||
if s == nil || (s.State != "built" && s.State != planDeleted) {
|
||||
return false, nil
|
||||
}
|
||||
if s.BuiltAt != nil && s.BuiltAt.After(latest) {
|
||||
@@ -598,7 +606,7 @@ func advanceOnce(ctx context.Context, open *stores, p *inventory.Plan,
|
||||
var pending []string
|
||||
for _, m := range tier {
|
||||
state := p.Modules[m]
|
||||
if state == nil || state.SentAt != nil || !rollsOut(m) {
|
||||
if state == nil || state.SentAt != nil || state.State == planDeleted || !rollsOut(m) {
|
||||
continue
|
||||
}
|
||||
running, err := inv.Running(ctx, m)
|
||||
@@ -620,29 +628,66 @@ func advanceOnce(ctx context.Context, open *stores, p *inventory.Plan,
|
||||
step := nextRollout(*state, running, policy.Together, reports, now, planWaitBound)
|
||||
switch {
|
||||
case step.failed != "":
|
||||
state.Why = step.failed
|
||||
p.State = inventory.PlanFailed
|
||||
p.Note = fmt.Sprintf("%s stopped at its first machine in tier %d: %s; %s left as it was",
|
||||
m, p.Tier, step.failed, orNone(strings.Join(step.rest, ", ")))
|
||||
// The first machine refused or failed what it was sent, or never said: the gate failed, and
|
||||
// the build is put back there (novox/hq ADR 0235); the rest are left as they were.
|
||||
gateFailed(ctx, open, p, m, state, firstRunning(state.First, running), step.failed)
|
||||
p.Note += fmt.Sprintf("; %s left as it was", orNone(strings.Join(step.rest, ", ")))
|
||||
fmt.Printf("%s: %s\n", p.ID, p.Note)
|
||||
return true, nil
|
||||
case step.waiting != "":
|
||||
pending = append(pending, fmt.Sprintf("%s on %s, sent first at %s", m, step.waiting, state.FirstAt.Local().Format("15:04")))
|
||||
continue
|
||||
}
|
||||
// **The gate** (novox/hq ADR 0235, to-be 45 §8): the first machine reported the build applied;
|
||||
// it is judged by its health before anything else is sent — the rest, or, where it is the only
|
||||
// machine, the plan's next step.
|
||||
if !policy.Together && state.FirstAt != nil && len(state.First) > 0 {
|
||||
verdict, err := judgeGate(ctx, open, p, m, state, running, now)
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
switch verdict {
|
||||
case "":
|
||||
pending = append(pending, fmt.Sprintf("%s judged on %s: %s", m,
|
||||
strings.Join(state.Gate.Machines, ", "), gateLine(state.Gate)))
|
||||
continue
|
||||
case inventory.GateFailed:
|
||||
gateFailed(ctx, open, p, m, state, state.Gate.Machines, state.Gate.Why)
|
||||
p.Note += fmt.Sprintf("; %s left as it was", orNone(strings.Join(step.rest, ", ")))
|
||||
fmt.Printf("%s: %s\n", p.ID, p.Note)
|
||||
return true, nil
|
||||
case inventory.GatePassed:
|
||||
if !state.Gate.Kept {
|
||||
gatePassed(ctx, open, p, m, state)
|
||||
state.Gate.Kept = true
|
||||
}
|
||||
}
|
||||
}
|
||||
switch {
|
||||
case len(step.send) == 0:
|
||||
// No machine runs it: nothing to send, and nothing to wait for.
|
||||
state.SentAt = &now
|
||||
continue
|
||||
}
|
||||
// Read before the send, which records what it carries.
|
||||
var before map[string]string
|
||||
var beforeKnown bool
|
||||
if step.first {
|
||||
before, beforeKnown, _ = inv.SentBuilds(ctx, step.send[0])
|
||||
}
|
||||
// What sendToEach answers, not what was asked: the machine holding the bus is sent before
|
||||
// the first when its user list must change (issue 249), and the plan waits for it too.
|
||||
sent, err := sendToEach(ctx, open, step.send)
|
||||
sent, err := sendRollout(ctx, open, step.send)
|
||||
if err != nil {
|
||||
// Not marked sent, so the next step tries again (issue 249): a grant that could not be
|
||||
// issued is a send that did not happen.
|
||||
return false, fmt.Errorf("sending %s to %s after tier %d: %w", m, strings.Join(step.send, ", "), p.Tier, err)
|
||||
}
|
||||
if step.first {
|
||||
// What the first machine ran before: what a failed gate puts back (ADR 0235).
|
||||
if beforeKnown {
|
||||
state.Previous = before[m]
|
||||
}
|
||||
state.First = sent
|
||||
state.FirstAt = &now
|
||||
p.State = inventory.PlanRolling
|
||||
@@ -681,6 +726,9 @@ func advanceOnce(ctx context.Context, open *stores, p *inventory.Plan,
|
||||
state = &inventory.PlanModule{}
|
||||
p.Modules[m] = state
|
||||
}
|
||||
if state.State == planDeleted {
|
||||
continue
|
||||
}
|
||||
// The plan sends what it waits for. A rebuild from the same source commit is not a
|
||||
// move the catalogue announces — the build machine rebuilt for a controller change
|
||||
// is one — so the roll-out that opens this gate is the plan's to make, once, and
|
||||
@@ -1033,6 +1081,11 @@ func plansCommand(ctx context.Context, args []string) error {
|
||||
}
|
||||
}
|
||||
fmt.Printf(" %-22s %s\n", m, state)
|
||||
// The rollout's record at its gate (novox/hq to-be 45 §8): first machine, from and to,
|
||||
// verdict, time to it, rolled back or not.
|
||||
if s != nil && s.Gate != nil {
|
||||
fmt.Printf(" %-22s %s\n", "", gateLine(s.Gate))
|
||||
}
|
||||
}
|
||||
}
|
||||
return nil
|
||||
@@ -1231,6 +1284,12 @@ func settleFromRecords(p *inventory.Plan, tier []string, recorded map[string][]i
|
||||
continue
|
||||
}
|
||||
at := outcome.At
|
||||
if !outcome.Worked() && deletedAtSource(outcome.Failed) {
|
||||
s.State, s.Why = planDeleted, "deleted at its source: "+firstLine(outcome.Failed)
|
||||
fmt.Printf("%s: %s was deleted at its source (%s): not built, and the plan goes on\n", p.ID, m, outcome.ID)
|
||||
changed = true
|
||||
continue
|
||||
}
|
||||
if outcome.Worked() {
|
||||
s.State = "built"
|
||||
s.BuiltAt = &at
|
||||
@@ -1252,3 +1311,34 @@ func settleFromRecords(p *inventory.Plan, tier []string, recorded map[string][]i
|
||||
func askedBefore(asked time.Time, planAsked *time.Time) bool {
|
||||
return !asked.IsZero() && planAsked != nil && asked.Before(*planAsked)
|
||||
}
|
||||
|
||||
// firstRunning is the machines sent first that run the module — not the machine holding the bus, sent
|
||||
// with them only for its user list.
|
||||
func firstRunning(first, running []string) []string {
|
||||
var out []string
|
||||
for _, n := range first {
|
||||
if slices.Contains(running, n) {
|
||||
out = append(out, n)
|
||||
}
|
||||
}
|
||||
if len(out) == 0 {
|
||||
return first
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
// planDeleted is a plan's module that the merge deleted at its source: not built, not sent, and no
|
||||
// failure of the plan (novox/hq ADR 0235).
|
||||
const planDeleted = "deleted"
|
||||
|
||||
// forgetDeleted forgets a module a plan found deleted at its source, where nothing holds it, and says
|
||||
// it otherwise.
|
||||
func forgetDeleted(ctx context.Context, inv *inventory.Inventory, module string, p *inventory.Plan) {
|
||||
switch err := inv.ForgetModule(ctx, module); {
|
||||
case err == nil:
|
||||
fmt.Printf("%s: %s was deleted at its source: forgotten, and the plan goes on\n", p.ID, module)
|
||||
default:
|
||||
fmt.Printf("%s: %s was deleted at its source and is not built; it is kept until a person unassigns it and "+
|
||||
"`module forget %s`: %v\n", p.ID, module, module, firstLine(err.Error()))
|
||||
}
|
||||
}
|
||||
|
||||
@@ -388,14 +388,7 @@ func rolloutMint(ctx context.Context, again bool) error {
|
||||
}
|
||||
|
||||
// providesBus is whether a manifest provides the mesh's bus.
|
||||
func providesBus(m catalogue.Manifest) bool {
|
||||
for _, o := range m.Provides {
|
||||
if o.Name == "mesh-bus" {
|
||||
return true
|
||||
}
|
||||
}
|
||||
return false
|
||||
}
|
||||
func providesBus(m catalogue.Manifest) bool { return catalogue.ProvidesBus(m) }
|
||||
|
||||
// rolloutHand mints a machine its credential for the new bus afresh and prints its membership
|
||||
// once, for an operator to carry by hand — the rescue for a machine that cannot be reached over
|
||||
|
||||
@@ -508,6 +508,43 @@ func (a *verbArguments) commandLine() ([]string, error) {
|
||||
argv = append(argv, "--retired")
|
||||
}
|
||||
return argv, nil
|
||||
case "upgrade":
|
||||
m, policy := str("module"), str("policy")
|
||||
if m == "" {
|
||||
if policy != "" {
|
||||
return nil, errors.New("upgrade takes a policy only for a module")
|
||||
}
|
||||
return []string{"upgrade"}, nil
|
||||
}
|
||||
argv := []string{"upgrade", m}
|
||||
if policy != "" {
|
||||
argv = append(argv, policy)
|
||||
if on("together") {
|
||||
argv = append(argv, "--together")
|
||||
}
|
||||
if w := str("why"); w != "" {
|
||||
argv = append(argv, "--why", w)
|
||||
}
|
||||
}
|
||||
return argv, nil
|
||||
case "bus":
|
||||
if !on("upgrade") {
|
||||
return []string{"bus"}, nil
|
||||
}
|
||||
argv := []string{"bus", "upgrade", "--why", str("why")}
|
||||
if c := str("cause"); c != "" {
|
||||
argv = append(argv, "--cause", c)
|
||||
}
|
||||
if on("reversible") {
|
||||
argv = append(argv, "--reversible")
|
||||
}
|
||||
if on("irreversible") {
|
||||
argv = append(argv, "--irreversible")
|
||||
}
|
||||
if w := str("snapshot-taken"); w != "" {
|
||||
argv = append(argv, "--snapshot-taken", w)
|
||||
}
|
||||
return argv, nil
|
||||
case "doctor":
|
||||
which := 0
|
||||
argv := []string{"doctor"}
|
||||
|
||||
@@ -99,7 +99,7 @@ func TestAVerbMissingWhatItNeedsIsRefused(t *testing.T) {
|
||||
if _, err := argvFor("node", map[string]any{}); err == nil || !strings.Contains(err.Error(), `node needs "node"`) {
|
||||
t.Fatalf("node without a machine was accepted: %v", err)
|
||||
}
|
||||
if _, err := argvFor("upgrade", map[string]any{}); err == nil {
|
||||
if _, err := argvFor("no-such-verb", map[string]any{}); err == nil {
|
||||
t.Fatal("a verb the seat does not serve was accepted")
|
||||
}
|
||||
}
|
||||
|
||||
@@ -195,6 +195,11 @@ func (l nudgingListener) Heard(ctx context.Context, report link.Report) (bool, e
|
||||
if report.Ordered() {
|
||||
link.StaleRefusals.Lifetime(report.Node, report.RefusedOlder, now)
|
||||
}
|
||||
// What the machine's witnesses put back and stand by (novox/hq ADR 0235): read by the gate and its
|
||||
// probe. Only from an account of the machine — not a word that a declaration was set aside, nor a rekey.
|
||||
if report.Superseded == "" && report.Rekey == nil && report.Node != "" {
|
||||
witnessed.heard(report.Node, report.Rollbacks, now)
|
||||
}
|
||||
news, err := l.Enrolment.Heard(ctx, report)
|
||||
if news {
|
||||
l.summary.nudge()
|
||||
|
||||
+125
-17
@@ -7,6 +7,7 @@ import (
|
||||
"fmt"
|
||||
"path"
|
||||
"regexp"
|
||||
"sort"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
@@ -143,20 +144,17 @@ func isAre(n int) string {
|
||||
return "are"
|
||||
}
|
||||
|
||||
// upgradeCommand says what should happen when a module's current version moves.
|
||||
// upgradeCommand says what should happen when a module's current version moves (ADR 0162 §3, ADR
|
||||
// 0235): every module's policy and where it comes from; one module's; or a person's choice for one.
|
||||
func upgradeCommand(ctx context.Context, args []string) error {
|
||||
set := flag.NewFlagSet("upgrade", flag.ContinueOnError)
|
||||
together := set.Bool("together", false,
|
||||
"send every machine running it at once, instead of one after another")
|
||||
"send every machine running it at once, instead of one machine first")
|
||||
why := set.String("why", "", "why this choice — required for record; kept and said with the policy")
|
||||
positionals, err := parseAround(set, args)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if len(positionals) == 0 {
|
||||
return errors.New("upgrade <module> [roll-out|record] [--together]")
|
||||
}
|
||||
module := positionals[0]
|
||||
|
||||
open, err := openStores(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
@@ -164,6 +162,27 @@ func upgradeCommand(ctx context.Context, args []string) error {
|
||||
defer open.Close()
|
||||
inv := open.inventory
|
||||
|
||||
if len(positionals) == 0 {
|
||||
all, err := inv.Upgrades(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
names := make([]string, 0, len(all))
|
||||
for n := range all {
|
||||
names = append(names, n)
|
||||
}
|
||||
sort.Strings(names)
|
||||
for _, n := range names {
|
||||
u := all[n]
|
||||
line := fmt.Sprintf("%-28s %-9s %s", n, u.Policy(), u.From)
|
||||
if u.Why != "" {
|
||||
line += ": " + u.Why
|
||||
}
|
||||
fmt.Println(line)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
module := positionals[0]
|
||||
if len(positionals) == 1 {
|
||||
decision, err := inv.UpgradeOf(ctx, module)
|
||||
if err != nil {
|
||||
@@ -173,10 +192,10 @@ func upgradeCommand(ctx context.Context, args []string) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
var decision inventory.Upgrade
|
||||
decision := inventory.Upgrade{Why: strings.TrimSpace(*why), By: link.Caller()}
|
||||
switch positionals[1] {
|
||||
case "roll-out":
|
||||
decision = inventory.Upgrade{RollOut: true, Together: *together}
|
||||
case "roll-out", "roll":
|
||||
decision.RollOut, decision.Together = true, *together
|
||||
case "record":
|
||||
if *together {
|
||||
// Refused rather than ignored: --together only means anything for a roll-out, and
|
||||
@@ -184,28 +203,53 @@ func upgradeCommand(ctx context.Context, args []string) error {
|
||||
return errors.New("`--together` says how to roll out, so it cannot be given with " +
|
||||
"`record`, which is the choice not to")
|
||||
}
|
||||
decision = inventory.Upgrade{}
|
||||
if decision.Why == "" {
|
||||
return fmt.Errorf("holding %s back from every merge is a choice a person reads later: --why <text> "+
|
||||
"(novox/hq ADR 0235). Nothing was changed", module)
|
||||
}
|
||||
case "default":
|
||||
if err := inv.ClearUpgradeOf(ctx, module); err != nil {
|
||||
return err
|
||||
}
|
||||
decision, err := inv.UpgradeOf(ctx, module)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
fmt.Println(sayUpgrade(module, decision))
|
||||
return nil
|
||||
default:
|
||||
return fmt.Errorf("upgrade <module> roll-out|record — not %q", positionals[1])
|
||||
return fmt.Errorf("upgrade <module> roll-out|record|default — not %q", positionals[1])
|
||||
}
|
||||
if err := inv.SetUpgradeOf(ctx, module, decision); err != nil {
|
||||
return err
|
||||
}
|
||||
decision, err = inv.UpgradeOf(ctx, module)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
fmt.Println(sayUpgrade(module, decision))
|
||||
return nil
|
||||
}
|
||||
|
||||
func sayUpgrade(module string, u inventory.Upgrade) string {
|
||||
from := " (" + u.From
|
||||
if u.By != "" {
|
||||
from += ", " + u.By
|
||||
}
|
||||
if u.Why != "" {
|
||||
from += ": " + u.Why
|
||||
}
|
||||
from += ")"
|
||||
if !u.RollOut {
|
||||
return fmt.Sprintf("when %s moves, the mesh records it and the machines running it are "+
|
||||
"reported as behind", module)
|
||||
"reported as behind until a person pushes them%s", module, from)
|
||||
}
|
||||
if u.Together {
|
||||
return fmt.Sprintf("when %s moves, every machine running it is sent the new version "+
|
||||
"together", module)
|
||||
"together%s", module, from)
|
||||
}
|
||||
return fmt.Sprintf("when %s moves, the machines running it are sent the new version one at "+
|
||||
"a time, stopping at the first that fails", module)
|
||||
return fmt.Sprintf("when %s moves, one machine running it is sent the new version first and judged at "+
|
||||
"the gate; the rest follow once it passes, and a build that fails is put back there%s", module, from)
|
||||
}
|
||||
|
||||
// Announceable is every build this mesh recorded, in the shape the builder announces one.
|
||||
@@ -310,6 +354,13 @@ func (f following) SourceMoved(ctx context.Context, m link.SourceMoved) error {
|
||||
}
|
||||
}
|
||||
touched := whatTheMergeTouched(from, entries, m)
|
||||
// **A module the merge deleted is not built** (novox/hq ADR 0235): its manifest is gone, so the build
|
||||
// seat finds nothing saying what it is, and the plan failed on it (`has no module.json at …`) with
|
||||
// every other module of its tier left unsent. It is forgotten where nothing holds it, said otherwise.
|
||||
touched, deleted := splitDeleted(touched, m)
|
||||
for _, e := range deleted {
|
||||
retireDeleted(ctx, inv, e, m)
|
||||
}
|
||||
for _, e := range touched {
|
||||
if err := inv.SourceMoved(ctx, e.Manifest.Module, m.Commit); err != nil {
|
||||
return notNow(err)
|
||||
@@ -465,7 +516,52 @@ func mergeCandidates(m link.SourceMoved, entries []inventory.Entry,
|
||||
func wouldMove(m link.SourceMoved, entries []inventory.Entry,
|
||||
read map[string][]inventory.ReadRepository) []inventory.Entry {
|
||||
from, _, _ := mergeCandidates(m, entries, read)
|
||||
return whatTheMergeTouched(from, entries, m)
|
||||
touched, _ := splitDeleted(whatTheMergeTouched(from, entries, m), m)
|
||||
return touched
|
||||
}
|
||||
|
||||
// splitDeleted parts the modules a merge touched into those it changed and those whose manifest it
|
||||
// removed — deleted at their source (ADR 0235).
|
||||
func splitDeleted(touched []inventory.Entry, m link.SourceMoved) (kept, deleted []inventory.Entry) {
|
||||
removed := map[string]bool{}
|
||||
for _, p := range m.Removed {
|
||||
removed[strings.Trim(p, "/")] = true
|
||||
}
|
||||
for _, e := range touched {
|
||||
dir := strings.Trim(e.Source.Path, "/")
|
||||
manifest := moduleManifestFile
|
||||
if dir != "" {
|
||||
manifest = dir + "/" + manifest
|
||||
}
|
||||
if len(removed) > 0 && removed[manifest] {
|
||||
deleted = append(deleted, e)
|
||||
continue
|
||||
}
|
||||
kept = append(kept, e)
|
||||
}
|
||||
return kept, deleted
|
||||
}
|
||||
|
||||
// retireDeleted is what the mesh does with a module deleted at its source: its record says the merge
|
||||
// was looked at, so it is not acted on again; it is forgotten where nothing holds it; where a machine
|
||||
// runs it or the mesh holds something for it, that is said, and nothing is built.
|
||||
func retireDeleted(ctx context.Context, inv *inventory.Inventory, e inventory.Entry, m link.SourceMoved) {
|
||||
name := e.Manifest.Module
|
||||
if err := inv.SourceMoved(ctx, name, m.Commit); err != nil {
|
||||
fmt.Printf(" %s was deleted from %s/%s at %.8s, and that could not be recorded: %v\n", name, m.Owner, m.Repo, m.Commit, err)
|
||||
}
|
||||
err := inv.ForgetModule(ctx, name)
|
||||
switch {
|
||||
case err == nil:
|
||||
fmt.Printf(" %s was deleted from %s/%s at %.8s: forgotten, nothing built\n", name, m.Owner, m.Repo, m.Commit)
|
||||
case errors.Is(err, inventory.ErrStillAssigned) || errors.Is(err, inventory.ErrStillHolds):
|
||||
fmt.Printf(" %s was deleted from %s/%s at %.8s and is not built; the mesh still runs or holds it, so it is "+
|
||||
"kept until a person unassigns it and `module forget %s`: %v\n", name, m.Owner, m.Repo, m.Commit, name,
|
||||
firstLine(err.Error()))
|
||||
default:
|
||||
fmt.Printf(" %s was deleted from %s/%s at %.8s and is not built; forgetting it failed: %v\n", name, m.Owner,
|
||||
m.Repo, m.Commit, err)
|
||||
}
|
||||
}
|
||||
|
||||
// sourceIs is whether a recorded source is the repository and branch a merge announced. A source on
|
||||
@@ -749,3 +845,15 @@ func dependentsOf(moved, entries []inventory.Entry, against map[string][]string)
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
// moduleManifestFile is the file that says what a module is, in its directory (the build seat's
|
||||
// ManifestName).
|
||||
const moduleManifestFile = "module.json"
|
||||
|
||||
// deletedAtSource is whether a build failed because its module's manifest is not at its source any
|
||||
// more — the build seat's own words (internal/builder) — which a merge that deleted the module causes
|
||||
// when its announcer did not say which files went (ADR 0235). Such a module is not a failure of the plan.
|
||||
func deletedAtSource(failed string) bool {
|
||||
return strings.Contains(failed, "has no "+moduleManifestFile+" at ") &&
|
||||
strings.Contains(failed, "so there is nothing saying what it is")
|
||||
}
|
||||
|
||||
@@ -0,0 +1,119 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/novox/mesh-controller/internal/lease"
|
||||
)
|
||||
|
||||
// The controller's side of its own rollback witness (novox/hq to-be 45 §8, ADR 0235; the contract is
|
||||
// internal/lease/witness.go): what the serving controller says of itself in every write of the lease's
|
||||
// key, for the node-engine on the control node to judge a new controller build by.
|
||||
|
||||
// processStarted is when this process started.
|
||||
var processStarted = time.Now().UTC()
|
||||
|
||||
// readiness remembers when this process first became ready.
|
||||
var readiness struct {
|
||||
sync.Mutex
|
||||
at time.Time
|
||||
}
|
||||
|
||||
// controllerHealth is the controller's health definition (to-be 45 §8) as this process meets it now:
|
||||
// it holds the lease — it is asked only while writing the key — its self-check has finished a run, and
|
||||
// in that run `status` answered in full within ten seconds (D9). Anything short of that is said, never
|
||||
// read as ready.
|
||||
func controllerHealth() *lease.Health {
|
||||
h := &lease.Health{Started: processStarted}
|
||||
d := doctorFrom
|
||||
if d == nil {
|
||||
h.Why = "the self-check is not running in this process yet"
|
||||
return h
|
||||
}
|
||||
run := d.lastRun()
|
||||
if run == nil {
|
||||
h.Why = "the self-check has not finished its first run"
|
||||
return h
|
||||
}
|
||||
h.DoctorRan = run.At
|
||||
ran := false
|
||||
for _, p := range run.Probes {
|
||||
if p.ID != "D9" {
|
||||
continue
|
||||
}
|
||||
ran = true
|
||||
if p.Verdict != verdictPass {
|
||||
h.Why = "in the last self-check, status did not answer in full within ten seconds (D9 " + p.Verdict + ")"
|
||||
return h
|
||||
}
|
||||
}
|
||||
if !ran {
|
||||
h.Why = "the last self-check did not judge status (D9)"
|
||||
return h
|
||||
}
|
||||
readiness.Lock()
|
||||
if readiness.at.IsZero() {
|
||||
readiness.at = time.Now().UTC()
|
||||
}
|
||||
h.Ready, h.ReadyAt = true, readiness.at
|
||||
readiness.Unlock()
|
||||
return h
|
||||
}
|
||||
|
||||
// witnessed is every rollback the machines' witnesses say they made and still stand by, as their last
|
||||
// report carried it (Report.RolledBack). Kept in the serving controller's memory only: a witness says it
|
||||
// on every report until it is resolved, so a controller started afresh hears it again with the next
|
||||
// report — the host's reconcile is minutes apart, and a restored controller hears the report that
|
||||
// follows its own restoring.
|
||||
var witnessed = &rollbacksHeard{byNode: map[string]heardRollbacks{}}
|
||||
|
||||
type rollbacksHeard struct {
|
||||
mu sync.Mutex
|
||||
byNode map[string]heardRollbacks
|
||||
}
|
||||
|
||||
type heardRollbacks struct {
|
||||
at time.Time
|
||||
list []lease.Rollback
|
||||
}
|
||||
|
||||
// heard keeps what one machine's account said, replacing what it said before: an account without any
|
||||
// says the machine's witnesses stand by none.
|
||||
func (w *rollbacksHeard) heard(node string, list []lease.Rollback, at time.Time) {
|
||||
w.mu.Lock()
|
||||
defer w.mu.Unlock()
|
||||
w.byNode[node] = heardRollbacks{at: at, list: append([]lease.Rollback(nil), list...)}
|
||||
}
|
||||
|
||||
// all is what every machine heard from said, by machine.
|
||||
func (w *rollbacksHeard) all() map[string][]lease.Rollback {
|
||||
w.mu.Lock()
|
||||
defer w.mu.Unlock()
|
||||
out := map[string][]lease.Rollback{}
|
||||
for node, h := range w.byNode {
|
||||
if len(h.list) > 0 {
|
||||
out[node] = append([]lease.Rollback(nil), h.list...)
|
||||
}
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
// holder is who holds the controller lease now, read from the bucket: through the serving controller's
|
||||
// own lease, or the bucket a command opened.
|
||||
func (a *actor) holder(ctx context.Context) (lease.Holder, bool, error) {
|
||||
a.mu.Lock()
|
||||
held, kv := a.held, a.kv
|
||||
a.mu.Unlock()
|
||||
reading, cancel := context.WithTimeout(ctx, 5*time.Second)
|
||||
defer cancel()
|
||||
switch {
|
||||
case held != nil:
|
||||
return held.Current(reading)
|
||||
case kv != nil:
|
||||
return lease.Current(reading, kv)
|
||||
}
|
||||
return lease.Holder{}, false, errors.New("this process has no view of the controller lease")
|
||||
}
|
||||
Reference in New Issue
Block a user