diff --git a/cmd/mesh-controller/bus_step.go b/cmd/mesh-controller/bus_step.go index 415230a..c45ee8f 100644 --- a/cmd/mesh-controller/bus_step.go +++ b/cmd/mesh-controller/bus_step.go @@ -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 ", 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 . 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 +} diff --git a/cmd/mesh-controller/gate_test.go b/cmd/mesh-controller/gate_test.go index ff2135b..4c8bcb9 100644 --- a/cmd/mesh-controller/gate_test.go +++ b/cmd/mesh-controller/gate_test.go @@ -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) } } diff --git a/internal/broker/nats.go b/internal/broker/nats.go index 84284ab..5b88417 100644 --- a/internal/broker/nats.go +++ b/internal/broker/nats.go @@ -152,6 +152,12 @@ type SeatVerb struct{ Seat, Verb string } var VerbsTheSelfCheckAsks = []SeatVerb{{Seat: "node-intrusion-prevention", Verb: "banned"}, {Seat: "node-backup", Verb: "backed-up"}} +// VerbsTheBusStepAsks are the seat verbs the bus's planned step calls (novox/hq to-be 45 §8, ADR 0236): +// the bus machine's backup holder takes the bus's snapshot now, before the bus is replaced (ADR 0235's +// dump), and says when it is taken (`backed-up`, granted with the self-check's). The one verb of +// the controller's grant that acts, and only through the step a person starts. +var VerbsTheBusStepAsks = []SeatVerb{{Seat: "node-backup", Verb: "now"}} + // perMachineEvents are a node-scoped seat's events about the holder itself, whose last token is the // holder's machine (novox/hq ADR 0219): `paused.`, the build agent saying whether it takes work. var perMachineEvents = map[string]bool{"paused.*": true} @@ -342,6 +348,10 @@ func PermissionsFor(p Principal) (Permissions, error) { for _, v := range VerbsTheSelfCheckAsks { pub = append(pub, "mesh.seat."+v.Seat+".tool."+v.Verb+".*") } + // And the bus's planned step: a snapshot taken now, before the bus is replaced (ADR 0236). + for _, v := range VerbsTheBusStepAsks { + pub = append(pub, "mesh.seat."+v.Seat+".tool."+v.Verb+".*") + } // And asks who answers (novox/hq to-be 45 §4, D3): the self-check finds every seat's holder by // the same discovery the console reads. The question only; the answers come to its own inbox. pub = append(pub, "$SRV.INFO") diff --git a/internal/broker/testdata/composed.conf b/internal/broker/testdata/composed.conf index e53cc8a..c4cdabe 100644 --- a/internal/broker/testdata/composed.conf +++ b/internal/broker/testdata/composed.conf @@ -24,7 +24,7 @@ accounts { jetstream: enabled users = [ { user: "controller", password: "$2a$11$cccccccccccccccccccccc", permissions: { - publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.EVENTS.controller.>", "$JS.API.>", "$KV.SEAT_MESH_BUILD_MACHINE_cancelled.>", "$KV.SEAT_NODE_BUILD_AGENT_cancelled.>", "$KV.mesh-controller_calls.>", "$KV.mesh-controller_condition-history.>", "$KV.mesh-controller_conditions.>", "$KV.mesh-controller_hand-acts.>", "$KV.mesh-controller_lease.>", "$SRV.INFO", "_INBOX.enrol.>", "mesh.assignment.>", "mesh.mod.*.tool.>", "mesh.node.>", "mesh.seat.mesh-build-machine.accept.>", "mesh.seat.mesh-build-machine.tool.>", "mesh.seat.mesh-controller.event.applied", "mesh.seat.mesh-controller.event.built-before", "mesh.seat.mesh-controller.event.condition-changed", "mesh.seat.mesh-controller.event.condition-cleared", "mesh.seat.mesh-controller.event.condition-raised", "mesh.seat.mesh-controller.event.doctor-heartbeat", "mesh.seat.mesh-controller.event.healer-acted", "mesh.seat.mesh-controller.event.refused", "mesh.seat.mesh-controller.event.rolled-back", "mesh.seat.mesh-controller.event.secret-replaced", "mesh.seat.node-backup.tool.backed-up.*", "mesh.seat.node-build-agent.accept.>", "mesh.seat.node-build-agent.tool.>", "mesh.seat.node-intrusion-prevention.tool.banned.*"] } + publish: { allow: ["$JS.ACK.CONTROL.controller.>", "$JS.ACK.EVENTS.controller.>", "$JS.API.>", "$KV.SEAT_MESH_BUILD_MACHINE_cancelled.>", "$KV.SEAT_NODE_BUILD_AGENT_cancelled.>", "$KV.mesh-controller_calls.>", "$KV.mesh-controller_condition-history.>", "$KV.mesh-controller_conditions.>", "$KV.mesh-controller_hand-acts.>", "$KV.mesh-controller_lease.>", "$SRV.INFO", "_INBOX.enrol.>", "mesh.assignment.>", "mesh.mod.*.tool.>", "mesh.node.>", "mesh.seat.mesh-build-machine.accept.>", "mesh.seat.mesh-build-machine.tool.>", "mesh.seat.mesh-controller.event.applied", "mesh.seat.mesh-controller.event.built-before", "mesh.seat.mesh-controller.event.condition-changed", "mesh.seat.mesh-controller.event.condition-cleared", "mesh.seat.mesh-controller.event.condition-raised", "mesh.seat.mesh-controller.event.doctor-heartbeat", "mesh.seat.mesh-controller.event.healer-acted", "mesh.seat.mesh-controller.event.refused", "mesh.seat.mesh-controller.event.rolled-back", "mesh.seat.mesh-controller.event.secret-replaced", "mesh.seat.node-backup.tool.backed-up.*", "mesh.seat.node-backup.tool.now.*", "mesh.seat.node-build-agent.accept.>", "mesh.seat.node-build-agent.tool.>", "mesh.seat.node-intrusion-prevention.tool.banned.*"] } subscribe: { allow: ["$JS.API.>", "$JS.EVENT.ADVISORY.CONSUMER.DELETED.>", "$JS.EVENT.ADVISORY.CONSUMER.MAX_DELIVERIES.>", "$SRV.INFO", "$SRV.INFO.mesh-controller", "$SRV.INFO.mesh-controller.>", "$SRV.PING", "$SRV.PING.mesh-controller", "$SRV.PING.mesh-controller.>", "$SRV.STATS", "$SRV.STATS.mesh-controller", "$SRV.STATS.mesh-controller.>", "_DELIVER.controller", "_DELIVER.controller.>", "_INBOX.controller.>", "mesh.control.>", "mesh.mod.*.event.provisioner.failing", "mesh.mod.*.event.provisioner.recovered", "mesh.mod.*.event.provisioner.retirement", "mesh.mod.gitea.event.pull.merged", "mesh.mod.mesh-catalog.event.catching-up", "mesh.mod.mesh-catalog.event.upgraded", "mesh.seat.mesh-build-machine.event.built", "mesh.seat.mesh-controller.tool.>", "mesh.seat.node-build-agent.event.built"] } allow_responses: { max: 1, ttl: "1m" } } }