A plan sends what it waits for
A tier whose module a later tier is built by waited for the machines running it to report after the build — and relied on the catalogue's moved event to send them. A rebuild from the same source commit is not a move the catalogue announces: the build machine rebuilt for a controller change kept its commit, nothing sent it, and the plan waited on a report that would never come (2026-10-01, 19:45Z). The plan now sends the machines running a gated module once, records when, and waits for the reports after that.
This commit is contained in:
@@ -256,21 +256,21 @@ func (b builds) Built(ctx context.Context, result link.BuildResult) error {
|
||||
case err != nil && result.Failed != "":
|
||||
fmt.Printf("%s: %v\n", result.ID, err)
|
||||
if result.Module != "" {
|
||||
planBuilt(ctx, b.inv, result.Module, result.Commit, result.Failed)
|
||||
planBuilt(ctx, b.open, result.Module, result.Commit, result.Failed)
|
||||
} else {
|
||||
planFailedBuild(ctx, b.inv, result)
|
||||
planFailedBuild(ctx, b.open, result)
|
||||
}
|
||||
return nil
|
||||
case err != nil:
|
||||
fmt.Printf("%s: heard and recorded, and not registered: %v\n", result.ID, err)
|
||||
if manifest.Module != "" {
|
||||
planBuilt(ctx, b.inv, manifest.Module, result.Commit, err.Error())
|
||||
planBuilt(ctx, b.open, manifest.Module, result.Commit, err.Error())
|
||||
}
|
||||
return nil
|
||||
}
|
||||
fmt.Printf("%s: %s %s registered, built on %s from %s\n",
|
||||
result.ID, manifest.Module, manifest.Version, result.On, short(result.Commit))
|
||||
saysWhenThePolicyActs(ctx, b.inv, manifest.Module)
|
||||
planBuilt(ctx, b.inv, manifest.Module, result.Commit, "")
|
||||
planBuilt(ctx, b.open, manifest.Module, result.Commit, "")
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -113,10 +113,10 @@ func serve(ctx context.Context) error {
|
||||
// while everything else about it looks correct.
|
||||
// And build results nobody was waiting for. A build triggered any other way than `build`
|
||||
// would otherwise be reported into the void, which is the same as not reporting it.
|
||||
server.Records(builds{inv})
|
||||
server.Records(builds{inv, open})
|
||||
// Open plans move on a timer as well as on outcomes (novox/hq ADR 0162): a tier waiting for
|
||||
// machines to report moves when they have, and a plan left by a replaced controller resumes.
|
||||
go planTicker(ctx, inv)
|
||||
go planTicker(ctx, open)
|
||||
// And what the catalogue decided a build meant. The builder's own result is already handled
|
||||
// above; this is the other half — the control plane is the only one of the three that knows
|
||||
// which machines run the thing, so it is the one that acts (novox/hq ADR 0072).
|
||||
|
||||
@@ -272,7 +272,8 @@ func askTier(ctx context.Context, inv *inventory.Inventory, p *inventory.Plan) e
|
||||
|
||||
// planBuilt marks a module built (or failed) in every open plan whose current tier holds it, and
|
||||
// advances what that completes. Called from the daemon's take-in of every outcome.
|
||||
func planBuilt(ctx context.Context, inv *inventory.Inventory, module, commit, failed string) {
|
||||
func planBuilt(ctx context.Context, open *stores, module, commit, failed string) {
|
||||
inv := open.inventory
|
||||
plans, err := inv.OpenPlans(ctx)
|
||||
if err != nil {
|
||||
fmt.Printf("plans: cannot read them: %v\n", err)
|
||||
@@ -316,13 +317,14 @@ func planBuilt(ctx context.Context, inv *inventory.Inventory, module, commit, fa
|
||||
fmt.Printf("%s: %s; the tiers after it are not asked\n", p.ID, p.Note)
|
||||
}
|
||||
}
|
||||
advancePlans(ctx, inv)
|
||||
advancePlans(ctx, open)
|
||||
}
|
||||
|
||||
// advancePlans moves every open plan as far as the facts allow: a tier whose modules are all built
|
||||
// and whose gates are applied gives way to the next; the last tier done is the plan done. Called
|
||||
// after every outcome and on a timer, so a plan waiting on a machine's report moves when it comes.
|
||||
func advancePlans(ctx context.Context, inv *inventory.Inventory) {
|
||||
func advancePlans(ctx context.Context, open *stores) {
|
||||
inv := open.inventory
|
||||
plans, err := inv.OpenPlans(ctx)
|
||||
if err != nil {
|
||||
fmt.Printf("plans: cannot read them: %v\n", err)
|
||||
@@ -343,7 +345,7 @@ func advancePlans(ctx context.Context, inv *inventory.Inventory) {
|
||||
for i := range plans {
|
||||
p := &plans[i]
|
||||
for p.Open() {
|
||||
moved, err := advanceOnce(ctx, inv, p, edges, rollsOut)
|
||||
moved, err := advanceOnce(ctx, open, p, edges, rollsOut)
|
||||
if err != nil {
|
||||
fmt.Printf("%s: %v\n", p.ID, err)
|
||||
break
|
||||
@@ -360,8 +362,9 @@ func advancePlans(ctx context.Context, inv *inventory.Inventory) {
|
||||
}
|
||||
|
||||
// advanceOnce takes one step of one plan and says whether anything changed.
|
||||
func advanceOnce(ctx context.Context, inv *inventory.Inventory, p *inventory.Plan,
|
||||
func advanceOnce(ctx context.Context, open *stores, p *inventory.Plan,
|
||||
edges []inventory.Edge, rollsOut func(string) bool) (bool, error) {
|
||||
inv := open.inventory
|
||||
if p.Tier >= len(p.Tiers) {
|
||||
p.State = inventory.PlanDone
|
||||
fmt.Printf("%s: done — %s at %s, %d tier(s)\n", p.ID, p.Repository, short(p.Commit), len(p.Tiers))
|
||||
@@ -405,11 +408,34 @@ func advanceOnce(ctx context.Context, inv *inventory.Inventory, p *inventory.Pla
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
builtAt := latest
|
||||
if s := p.Modules[m]; s != nil && s.BuiltAt != nil {
|
||||
builtAt = *s.BuiltAt
|
||||
state := p.Modules[m]
|
||||
if state == nil {
|
||||
state = &inventory.PlanModule{}
|
||||
p.Modules[m] = state
|
||||
}
|
||||
if ok, on := applied(m, builtAt, running, reports); !ok {
|
||||
// The plan sends what it waits for. A rebuild from the same source commit is not a
|
||||
// move the catalogue announces — the build machine rebuilt for a controller change
|
||||
// is one — so the roll-out that opens this gate is the plan's to make, once, and
|
||||
// the reports that open it are the ones after the send.
|
||||
if state.SentAt == nil && len(running) > 0 {
|
||||
if err := sendTo(ctx, open, running); err != nil {
|
||||
return false, fmt.Errorf("sending %s to %s so tier %d can be built by it: %w",
|
||||
m, strings.Join(running, ", "), p.Tier+1, err)
|
||||
}
|
||||
now := time.Now().UTC()
|
||||
state.SentAt = &now
|
||||
fmt.Printf("%s: tier %d built; sent %s to %s, and tier %d waits until it is applied\n",
|
||||
p.ID, p.Tier, m, strings.Join(running, ", "), p.Tier+1)
|
||||
return true, nil
|
||||
}
|
||||
since := latest
|
||||
if state.BuiltAt != nil {
|
||||
since = *state.BuiltAt
|
||||
}
|
||||
if state.SentAt != nil && state.SentAt.After(since) {
|
||||
since = *state.SentAt
|
||||
}
|
||||
if ok, on := applied(m, since, running, reports); !ok {
|
||||
waiting = append(waiting, fmt.Sprintf("%s on %s", m, strings.Join(on, ", ")))
|
||||
}
|
||||
}
|
||||
@@ -431,8 +457,8 @@ func advanceOnce(ctx context.Context, inv *inventory.Inventory, p *inventory.Pla
|
||||
}
|
||||
|
||||
// planTicker advances open plans on a timer, for the steps outcomes alone cannot take.
|
||||
func planTicker(ctx context.Context, inv *inventory.Inventory) {
|
||||
advancePlans(ctx, inv)
|
||||
func planTicker(ctx context.Context, open *stores) {
|
||||
advancePlans(ctx, open)
|
||||
tick := time.NewTicker(30 * time.Second)
|
||||
defer tick.Stop()
|
||||
for {
|
||||
@@ -440,7 +466,7 @@ func planTicker(ctx context.Context, inv *inventory.Inventory) {
|
||||
case <-ctx.Done():
|
||||
return
|
||||
case <-tick.C:
|
||||
advancePlans(ctx, inv)
|
||||
advancePlans(ctx, open)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -468,14 +494,14 @@ func planLine(p inventory.Plan, now time.Time) string {
|
||||
|
||||
// planFailedBuild marks the module a failed build was for when the result names no module: by the
|
||||
// repository and path the plan's modules were asked at.
|
||||
func planFailedBuild(ctx context.Context, inv *inventory.Inventory, result link.BuildResult) {
|
||||
entries, err := inv.Catalogued(ctx)
|
||||
func planFailedBuild(ctx context.Context, open *stores, result link.BuildResult) {
|
||||
entries, err := open.inventory.Catalogued(ctx)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
for _, e := range entries {
|
||||
if repositoryMatches(e.Source.Repository, result.Repository) && e.Source.Path == result.Path {
|
||||
planBuilt(ctx, inv, e.Manifest.Module, result.Commit, result.Failed)
|
||||
planBuilt(ctx, open, e.Manifest.Module, result.Commit, result.Failed)
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
@@ -301,7 +301,10 @@ func firstLine(s string) string {
|
||||
// A type of its own rather than a method on the enrolment, because they are unrelated things
|
||||
// arriving on one queue and an implementation of one should not have to say anything about the
|
||||
// other.
|
||||
type builds struct{ inv *inventory.Inventory }
|
||||
type builds struct {
|
||||
inv *inventory.Inventory
|
||||
open *stores
|
||||
}
|
||||
|
||||
// theThreeQuestions reads what anything answering "is the mesh alright" needs.
|
||||
//
|
||||
|
||||
@@ -33,8 +33,12 @@ type PlanModule struct {
|
||||
State string `json:"state,omitempty"`
|
||||
AskedAt *time.Time `json:"asked_at,omitempty"`
|
||||
BuiltAt *time.Time `json:"built_at,omitempty"`
|
||||
Commit string `json:"commit,omitempty"`
|
||||
Why string `json:"why,omitempty"`
|
||||
// SentAt is when the plan sent the machines running this module its new build, because a
|
||||
// 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"`
|
||||
}
|
||||
|
||||
// The states a plan passes through.
|
||||
|
||||
Reference in New Issue
Block a user