Self-upgrade installer: SDK-by-version, mesh-controller rename, foundation adopted #26
@@ -168,6 +168,11 @@ func pushCommand(ctx context.Context, args []string) error {
|
||||
// eventually watches.
|
||||
behind := set.Bool("behind", false,
|
||||
"only machines whose last declaration was refused or partly failed")
|
||||
// For a named node, wait until it reports applying exactly what it was sent, so `push <node>`
|
||||
// means "this node is now what it was told" — a command right after does not race the apply
|
||||
// (novox/hq ADR 0010). 0 waits for nothing, which is the old fire-and-forget.
|
||||
wait := set.Duration("wait", 0,
|
||||
"for a named node, how long to wait for it to report applying what it was sent (0: do not wait)")
|
||||
positionals, err := parseAround(set, args)
|
||||
if err != nil {
|
||||
return err
|
||||
@@ -292,6 +297,7 @@ func pushCommand(ctx context.Context, args []string) error {
|
||||
return declarationWith(ctx, open, node, plan, settings, gens, Allocating)
|
||||
})
|
||||
|
||||
sentDigest := map[string]string{}
|
||||
for _, s := range sending {
|
||||
body, err := json.Marshal(map[string]any{"declaration": 1, "resources": s.resources})
|
||||
if err != nil {
|
||||
@@ -306,15 +312,62 @@ func pushCommand(ctx context.Context, args []string) error {
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if err := inv.RecordSent(ctx, record.ID, digestOf(body)); err != nil {
|
||||
digest := digestOf(body)
|
||||
if err := inv.RecordSent(ctx, record.ID, digest); err != nil {
|
||||
return err
|
||||
}
|
||||
sentDigest[s.node] = digest
|
||||
fmt.Printf("sent %s %d resource(s)\n", s.node, len(s.resources))
|
||||
}
|
||||
fmt.Printf("\n%d node(s) told\n", len(sending))
|
||||
|
||||
// A named node is a request to make THAT node current now, so it waits for the node to say it
|
||||
// applied exactly this. A whole-mesh or --behind push does not wait: it is a sweep, and blocking
|
||||
// on the slowest machine would hold back the report on all the others.
|
||||
if *wait > 0 && len(args) == 1 {
|
||||
if err := waitForApplied(ctx, inv, args[0], sentDigest[args[0]], *wait); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
return couldNotBeResolved(refusals, len(sending))
|
||||
}
|
||||
|
||||
// waitForApplied blocks until the node reports it applied exactly the declaration just sent, or the
|
||||
// wait runs out. A report of failure or refusal for that same declaration ends the wait at once —
|
||||
// there is nothing to wait for, and the reason is the node's own.
|
||||
func waitForApplied(ctx context.Context, inv *inventory.Inventory, node, digest string, wait time.Duration) error {
|
||||
if digest == "" {
|
||||
return nil // nothing was sent to this node
|
||||
}
|
||||
deadline := time.Now().Add(wait)
|
||||
for {
|
||||
doing, said, err := inv.DoingOf(ctx, node)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if said && doing.Declared == digest {
|
||||
switch doing.Outcome {
|
||||
case inventory.OutcomeApplied:
|
||||
fmt.Printf("%s applied it\n", node)
|
||||
return nil
|
||||
case inventory.OutcomeFailed:
|
||||
return fmt.Errorf("%s applied what it was sent but %d resource(s) failed", node, len(doing.Failed))
|
||||
case inventory.OutcomeRefused:
|
||||
return fmt.Errorf("%s refused what it was sent: %s", node, doing.Refused)
|
||||
}
|
||||
}
|
||||
if time.Now().After(deadline) {
|
||||
return fmt.Errorf("%s did not report applying what it was sent within %s "+
|
||||
"(it may still be converging; check `status`)", node, wait)
|
||||
}
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return ctx.Err()
|
||||
case <-time.After(500 * time.Millisecond):
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// readyNode is one machine and the declaration it would be sent.
|
||||
type readyNode struct {
|
||||
node string
|
||||
|
||||
Reference in New Issue
Block a user