diff --git a/cmd/mesh-controller/release_plan.go b/cmd/mesh-controller/release_plan.go index 267d2c7..3efb023 100644 --- a/cmd/mesh-controller/release_plan.go +++ b/cmd/mesh-controller/release_plan.go @@ -470,6 +470,15 @@ func advanceOnce(ctx context.Context, open *stores, p *inventory.Plan, // a module that packages another repository's source, keeps its commit; the catalogue announces // no move for it and its machines would keep the old image until somebody pushed (novox/hq // issue 189). A module whose policy records is built and left, as its policy says. + // + // **One machine first, unless the module's policy says together** (novox/hq issue 249, ADR + // 0218). The plan sent every machine running the module at once, and the operator's policy — + // one at a time, stopping at the first that fails, which an announced upgrade honours — was not + // read here at all: a module whose new declarations broke it broke everywhere in the same + // minute. Now the first machine is sent, the plan records it and waits for that machine's report + // after the send to say it applied what it was sent; only then are the rest sent. A first machine + // that fails or refuses stops the module's rollout and the plan with it, the rest untouched. + var pending []string for _, m := range tier { state := p.Modules[m] if state == nil || state.SentAt != nil || !rollsOut(m) { @@ -479,17 +488,67 @@ func advanceOnce(ctx context.Context, open *stores, p *inventory.Plan, if err != nil { return false, err } + policy, err := inv.UpgradeOf(ctx, m) + if err != nil { + return false, err + } + var reports []inventory.Reported + if state.FirstAt != nil && !policy.Together { + if reports, err = inv.LastReports(ctx); err != nil { + return false, err + } + } + step := nextRollout(*state, running, policy.Together, reports) now := time.Now().UTC() - state.SentAt = &now - if len(running) == 0 { + 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, ", "))) + fmt.Printf("%s: %s\n", p.ID, p.Note) + return true, nil + case step.waiting != "": + wait := fmt.Sprintf("%s on %s, sent first at %s", m, step.waiting, state.FirstAt.Format("15:04")) + if now.Sub(*state.FirstAt) > planWaitBound { + wait += " — LATE" + } + pending = append(pending, wait) + continue + case len(step.send) == 0: + // No machine runs it: nothing to send, and nothing to wait for. + state.SentAt = &now continue } - if err := sendTo(ctx, open, running); err != nil { - return false, fmt.Errorf("sending %s to %s after tier %d: %w", m, strings.Join(running, ", "), p.Tier, err) + // 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) + 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) } - fmt.Printf("%s: tier %d built; sent %s to %s\n", p.ID, p.Tier, m, strings.Join(running, ", ")) + if step.first { + state.First = sent + state.FirstAt = &now + p.State = inventory.PlanRolling + p.Note = fmt.Sprintf("tier %d built; sent %s to %s first", p.Tier, m, strings.Join(sent, ", ")) + fmt.Printf("%s: tier %d built; sent %s to %s first, the rest once it reports it applied\n", + p.ID, p.Tier, m, strings.Join(sent, ", ")) + return true, nil + } + state.SentAt = &now + fmt.Printf("%s: tier %d built; sent %s to %s\n", p.ID, p.Tier, m, strings.Join(sent, ", ")) return true, nil } + if len(pending) > 0 { + note := "tier " + fmt.Sprint(p.Tier) + " built; waiting for " + strings.Join(pending, "; ") + + " to report it applied before the rest are sent" + changed := p.State != inventory.PlanRolling || p.Note != note + p.State = inventory.PlanRolling + p.Note = note + return changed, nil + } // And wait for what the next tier needs running. needed := gates(*p, edges, rollsOut) if len(needed) > 0 { @@ -540,6 +599,77 @@ func advanceOnce(ctx context.Context, open *stores, p *inventory.Plan, return true, nil } +// rolloutStep is what a plan does next with one built module's machines (novox/hq issue 249). +type rolloutStep struct { + // send is the machines to send now; first, whether they are the first machine's send. + send []string + first bool + // waiting names the first machines whose report the rest wait for. + waiting string + // failed says how a first machine did not take it; rest is what is then left alone. + failed string + rest []string +} + +// nextRollout is the next step of one module's rollout in a plan (novox/hq issue 249, ADR 0218). +// +// Together, every machine running it at once, as the policy says. Otherwise one machine first — +// the first by name, so the choice is the same on every controller and every resume — and the rest +// once each machine the first send reached has reported, after that send, that it applied the +// declaration it was last sent. A report after the send that failed or refused it is the rollout's +// end: the rest are not sent. A machine that has not reported since is waited for; how long it has +// been is the plan's to say. +func nextRollout(s inventory.PlanModule, running []string, together bool, reports []inventory.Reported) rolloutStep { + if len(running) == 0 { + return rolloutStep{} + } + if together { + return rolloutStep{send: running} + } + if s.FirstAt == nil { + sorted := append([]string{}, running...) + sort.Strings(sorted) + return rolloutStep{send: sorted[:1], first: true} + } + sentFirst := map[string]bool{} + for _, n := range s.First { + sentFirst[n] = true + } + var rest []string + for _, n := range running { + if !sentFirst[n] { + rest = append(rest, n) + } + } + byNode := map[string]inventory.Reported{} + for _, r := range reports { + byNode[r.Node] = r + } + var waiting, failed []string + for _, n := range s.First { + r, said := byNode[n] + // Only a report after the send, about what it was last sent, says anything about this build. + if !said || r.At == nil || r.At.Before(*s.FirstAt) || !r.Current { + waiting = append(waiting, n) + continue + } + switch r.Outcome { + case inventory.OutcomeApplied: + case inventory.OutcomeFailed, inventory.OutcomeRefused: + failed = append(failed, n+" "+r.Outcome+" what it was sent") + default: + waiting = append(waiting, n) + } + } + if len(failed) > 0 { + return rolloutStep{failed: strings.Join(failed, "; "), rest: rest} + } + if len(waiting) > 0 { + return rolloutStep{waiting: strings.Join(waiting, ", ")} + } + return rolloutStep{send: rest} +} + // planTicker advances open plans on a timer, for the steps outcomes alone cannot take. func planTicker(ctx context.Context, open *stores) { advancePlans(ctx, open) @@ -796,6 +926,11 @@ func planWhatIf(ctx context.Context, inv *inventory.Inventory, repository string if u, err := inv.UpgradeOf(ctx, name); err == nil && u.RollOut { running, _ := inv.Running(ctx, name) how = "built, then sent to " + orNone(strings.Join(running, ", ")) + // One machine first unless the policy says together (novox/hq issue 249). + if first := nextRollout(inventory.PlanModule{}, running, u.Together, nil); first.first && len(running) > 1 { + how = fmt.Sprintf("built, then sent to %s first and to the rest once it has applied it", + first.send[0]) + } rolls[name] = how } fmt.Printf(" %-22s %s\n", name, how) diff --git a/cmd/mesh-controller/rollout_first_test.go b/cmd/mesh-controller/rollout_first_test.go new file mode 100644 index 0000000..15587bb --- /dev/null +++ b/cmd/mesh-controller/rollout_first_test.go @@ -0,0 +1,84 @@ +package main + +import ( + "reflect" + "strings" + "testing" + "time" + + "github.com/novox/mesh-controller/internal/inventory" +) + +// novox/hq issue 249, ADR 0218: a plan rolls a module out to one machine first and the rest only +// once that machine has reported it applied; a module whose policy says together goes everywhere at +// once, as before. +func TestAPlanSendsOneMachineFirstAndTheRestAfterItsReport(t *testing.T) { + running := []string{"novox", "ace", "g14"} + + // Together: every machine at once. + if step := nextRollout(inventory.PlanModule{}, running, true, nil); !reflect.DeepEqual(step.send, running) || step.first { + t.Fatalf("a together policy did not send every machine at once: %+v", step) + } + + // Otherwise the first by name, alone. + step := nextRollout(inventory.PlanModule{}, running, false, nil) + if !step.first || !reflect.DeepEqual(step.send, []string{"ace"}) { + t.Fatalf("the first send was %+v, wanted ace alone", step) + } + + sentAt := time.Date(2026, 10, 5, 12, 0, 0, 0, time.UTC) + before, after := sentAt.Add(-time.Minute), sentAt.Add(time.Minute) + state := inventory.PlanModule{First: []string{"ace"}, FirstAt: &sentAt} + report := func(at time.Time, outcome string, current bool) []inventory.Reported { + return []inventory.Reported{{Node: "ace", At: &at, Outcome: outcome, Current: current}, + {Node: "g14", At: &after, Outcome: inventory.OutcomeApplied, Current: true}} + } + + // No report yet, a report from before the send, or one about an older declaration: wait. + for what, reports := range map[string][]inventory.Reported{ + "no report": nil, + "a report before the send": report(before, inventory.OutcomeApplied, true), + "a report about older": report(after, inventory.OutcomeApplied, false), + } { + step := nextRollout(state, running, false, reports) + if len(step.send) != 0 || step.waiting != "ace" || step.failed != "" { + t.Errorf("%s: %+v, wanted to wait for ace", what, step) + } + } + + // Applied after the send: the rest, and only the rest. + step = nextRollout(state, running, false, report(after, inventory.OutcomeApplied, true)) + if step.first || !reflect.DeepEqual(step.send, []string{"novox", "g14"}) { + t.Fatalf("after ace applied it the plan sent %+v, wanted novox and g14", step) + } + + // Failed or refused after the send: stop, the rest untouched. + for _, outcome := range []string{inventory.OutcomeFailed, inventory.OutcomeRefused} { + step := nextRollout(state, running, false, report(after, outcome, true)) + if len(step.send) != 0 || !strings.Contains(step.failed, "ace "+outcome) || + !reflect.DeepEqual(step.rest, []string{"novox", "g14"}) { + t.Errorf("a first machine that %s it: %+v", outcome, step) + } + } + + // The machine holding the bus went with the first send: the rest wait for it too, and it is + // not sent again. + both := inventory.PlanModule{First: []string{"novox", "ace"}, FirstAt: &sentAt} + half := report(after, inventory.OutcomeApplied, true) + if step := nextRollout(both, running, false, half); step.waiting != "novox" { + t.Fatalf("the plan did not wait for the bus's machine sent first: %+v", step) + } + all := append(half, inventory.Reported{Node: "novox", At: &after, Outcome: inventory.OutcomeApplied, Current: true}) + if step := nextRollout(both, running, false, all); !reflect.DeepEqual(step.send, []string{"g14"}) { + t.Fatalf("after both applied it the plan sent %+v, wanted g14 alone", step) + } + + // One machine, or none: nothing is waited for that cannot come. + if step := nextRollout(inventory.PlanModule{}, nil, false, nil); len(step.send) != 0 || step.first { + t.Fatalf("a module nothing runs was sent: %+v", step) + } + only := inventory.PlanModule{First: []string{"ace"}, FirstAt: &sentAt} + if step := nextRollout(only, []string{"ace"}, false, report(after, inventory.OutcomeApplied, true)); len(step.send) != 0 || step.waiting != "" { + t.Fatalf("a module on one machine waited for more: %+v", step) + } +} diff --git a/internal/inventory/plans.go b/internal/inventory/plans.go index c68bafe..170ef31 100644 --- a/internal/inventory/plans.go +++ b/internal/inventory/plans.go @@ -37,8 +37,15 @@ type PlanModule struct { // later tier is built by it (ADR 0163's gate): the reports that open the gate are the ones // after this. SentAt *time.Time `json:"sent_at,omitempty"` - Commit string `json:"commit,omitempty"` - Why string `json:"why,omitempty"` + // First is the machines the plan sent the new build to first, and FirstAt when (novox/hq issue + // 249, ADR 0218): unless the module's policy rolls it out together, one machine takes it before + // the rest, and the rest are sent once that one reports it applied. Kept so a controller + // replaced while the plan waits on that report resumes the wait rather than sending again. The + // machine holding the bus is among them when its user list had to go first. + First []string `json:"first,omitempty"` + FirstAt *time.Time `json:"first_at,omitempty"` + Commit string `json:"commit,omitempty"` + Why string `json:"why,omitempty"` } // The states a plan passes through.