Take the bus's snapshot through its machine's backup holder before the planned step (hq ADR 0235, 0236)
This commit is contained in:
@@ -9,6 +9,8 @@ import (
|
||||
"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"
|
||||
@@ -38,10 +40,11 @@ const (
|
||||
// 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)
|
||||
// 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
|
||||
@@ -174,15 +177,14 @@ func busCommand(ctx context.Context, args []string) error {
|
||||
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)
|
||||
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)
|
||||
}
|
||||
}
|
||||
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 {
|
||||
@@ -237,9 +239,7 @@ func busStatus(ctx context.Context) error {
|
||||
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)")
|
||||
}
|
||||
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
|
||||
@@ -325,3 +325,49 @@ func sortedKeys(set map[string]bool) []string {
|
||||
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
|
||||
}
|
||||
|
||||
@@ -403,22 +403,30 @@ func TestTheBusIsNeverRolledOutAutomatically(t *testing.T) {
|
||||
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.
|
||||
// The planned step refuses to start without its word on reversibility, and without a snapshot taken
|
||||
// first by the bus machine's backup holder.
|
||||
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)
|
||||
}
|
||||
wasSnapshot := takeBusSnapshot
|
||||
t.Cleanup(func() { takeBusSnapshot = wasSnapshot })
|
||||
takeBusSnapshot = func(context.Context, string, string) (string, error) { return "", errors.New("no holder answers") }
|
||||
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)
|
||||
!strings.Contains(err.Error(), "snapshotted") || len(sent) != 0 {
|
||||
t.Fatalf("a bus upgrade started without its snapshot: %v, sent %v", err, sent)
|
||||
}
|
||||
if err := busCommand(ctx, []string{"upgrade", "--why", "2.11", "--reversible", "--snapshot-taken", "nightly"}); err != nil {
|
||||
takeBusSnapshot = func(_ context.Context, module, node string) (string, error) {
|
||||
return node + "'s restore point of " + module, nil
|
||||
}
|
||||
if err := busCommand(ctx, []string{"upgrade", "--why", "2.11", "--reversible"}); 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"}) {
|
||||
if err != nil || !found || s.Snapshot != "anchor's restore point of nats" || s.Ended != nil ||
|
||||
!reflect.DeepEqual(s.Machines, []string{"anchor"}) {
|
||||
t.Fatalf("the step is %+v %v %v", s, found, err)
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user