Roll a plan's module out to one machine first (hq issue 249, ADR 0218)
A plan sent every machine running a module at once, ignoring the module's upgrade policy. Unless the policy says together, the first machine by name is sent, recorded in the plan, and the rest follow only once its report after the send says it applied; a failed first machine stops the plan.
This commit is contained in:
@@ -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)
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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.
|
||||
|
||||
Reference in New Issue
Block a user