mesh/merge-gate error: the check could not run: a throwaway postgres:17-alpine could not be raised: docker run --label mesh.build=build-1791317509716888018…
mesh/delivery delivered
An image is not byte-reproducible, so ADR 0236's 'same artifacts is no move' never held for one: a catalogue merge that did not touch the bus rebuilt it, and every send to the control node waited for a planned bus upgrade. The builder now records a source fingerprint per build (module tree, context trees, bases and toolchains by digest). A rebuild with the fingerprint of the build it repeats is registered with that build's artifacts, handed to modules standing on it, holds no push, demands no bus step, and a plan sends and gates nothing for it. Identical artifacts remain a second way to be no move.
411 lines
15 KiB
Go
411 lines
15 KiB
Go
package main
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"flag"
|
|
"fmt"
|
|
"sort"
|
|
"strings"
|
|
"time"
|
|
|
|
"github.com/nats-io/nats.go"
|
|
|
|
"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 0236).
|
|
//
|
|
// **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: the bus
|
|
// machine's backup holder backs the bus module up now, whose dump is the streams' snapshot (novox/hq ADR
|
|
// 0235). A variable so a test takes none. A person who took one by hand says where with --snapshot-taken,
|
|
// and then none is taken — for a bus whose module does not yet carry the snapshot program.
|
|
var takeBusSnapshot = snapshotTheBusNow
|
|
|
|
// 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
|
|
// same are the commits whose build was made from the same source as the build the mesh holds, or
|
|
// made the same artifacts and manifest (novox/hq issue 280).
|
|
same map[string]bool
|
|
}
|
|
|
|
// moves is whether sending the machine would replace its bus: a build it was not last sent, unless the
|
|
// two builds were made from the same source, or made the same artifacts from the same manifest — a
|
|
// rebuild of the same source for another module's merge changes nothing the machine runs, whatever
|
|
// image digest it made (novox/hq issue 280).
|
|
func (b busPending) moves(machine string) bool {
|
|
from, known := b.from[machine]
|
|
if b.module == "" || b.to == "" || (known && sameCommit(from, b.to)) {
|
|
return false
|
|
}
|
|
return !known || !b.same[from]
|
|
}
|
|
|
|
// 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, b.same = map[string]string{}, map[string]bool{}
|
|
made, err := inv.BuildFingerprints(ctx, b.module)
|
|
if err != nil {
|
|
return b, err
|
|
}
|
|
for commit, refs := range made {
|
|
b.same[commit] = refs != "" && refs == made[b.to]
|
|
}
|
|
sources, err := inv.BuildSourceFingerprints(ctx, b.module)
|
|
if err != nil {
|
|
return b, err
|
|
}
|
|
for commit, src := range sources {
|
|
if src != "" && src == sources[b.to] {
|
|
b.same[commit] = true
|
|
}
|
|
}
|
|
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 0236). 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 0236)",
|
|
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)
|
|
if where == "" {
|
|
for _, n := range moving {
|
|
fmt.Printf("snapshotting the bus's streams on %s first (its backup holder, ADR 0235)…\n", n)
|
|
if where, err = takeBusSnapshot(ctx, b.module, n); err != nil {
|
|
return fmt.Errorf("the streams could not be snapshotted, so the bus is not replaced: %w — a snapshot "+
|
|
"taken by hand is said with --snapshot-taken <where>", err)
|
|
}
|
|
}
|
|
}
|
|
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(withBusStep(withScope(ctx, sendScope{person: true})), 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)
|
|
}
|
|
}
|
|
fmt.Println(" `bus upgrade` has the bus machine's backup holder snapshot the streams first (ADR 0235)")
|
|
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
|
|
}
|
|
|
|
// busSnapshotWithin is how long the bus machine's backup holder is given to take the bus's snapshot.
|
|
var busSnapshotWithin = 15 * time.Minute
|
|
|
|
// snapshotTheBusNow asks the bus machine's backup holder to back the bus module up now — its dump is
|
|
// the streams' snapshot (novox/hq ADR 0235) — and waits until it says a backup newer than the ask:
|
|
// where the snapshot is, as a person reads it. The serving controller's connection, or one of its own.
|
|
func snapshotTheBusNow(ctx context.Context, module, node string) (string, error) {
|
|
var where string
|
|
err := onTheBus(func(conn *nats.Conn) error {
|
|
asked := time.Now()
|
|
answer, err := link.AskSeatTool(ctx, conn, catalogue.BackupSeat, "now", node,
|
|
map[string]any{"module": module}, 30*time.Second)
|
|
if err != nil {
|
|
return fmt.Errorf("%s's backup holder was not asked to take the bus's snapshot: %w", node, err)
|
|
}
|
|
if answer.Error != "" {
|
|
return fmt.Errorf("%s's backup holder would not take the bus's snapshot: %s", node, answer.Error)
|
|
}
|
|
deadline := time.Now().Add(busSnapshotWithin)
|
|
for {
|
|
answer, err := link.AskSeatTool(ctx, conn, catalogue.BackupSeat, "backed-up", node, map[string]any{}, 10*time.Second)
|
|
if err == nil && answer.Error == "" {
|
|
if measured, err := readHolder(answer.Result); err == nil {
|
|
for item, m := range measured[module] {
|
|
if m.LastBackup != nil && m.LastBackup.After(asked) && m.Error == "" {
|
|
where = fmt.Sprintf("%s's restore point of %s (%s) taken %s", node, module, item,
|
|
m.LastBackup.UTC().Format(time.RFC3339))
|
|
return nil
|
|
}
|
|
}
|
|
}
|
|
}
|
|
if time.Now().After(deadline) {
|
|
return fmt.Errorf("%s's backup holder did not say the bus's snapshot was taken within %s; the bus is "+
|
|
"not replaced", node, busSnapshotWithin)
|
|
}
|
|
select {
|
|
case <-ctx.Done():
|
|
return ctx.Err()
|
|
case <-time.After(10 * time.Second):
|
|
}
|
|
}
|
|
})
|
|
return where, err
|
|
}
|
|
|
|
// errBusWaits is a send refused because it would replace the bus outside its planned step.
|
|
var errBusWaits = errors.New("a new bus build waits for its planned step")
|
|
|
|
type busStepKey struct{}
|
|
|
|
// withBusStep marks a send as the bus's planned step: the one send that may replace the bus.
|
|
func withBusStep(ctx context.Context) context.Context {
|
|
return context.WithValue(ctx, busStepKey{}, true)
|
|
}
|
|
|
|
func busStepSending(ctx context.Context) bool { on, _ := ctx.Value(busStepKey{}).(bool); return on }
|