diff --git a/cmd/mesh-control/push.go b/cmd/mesh-control/push.go index 8eb13f2..1def023 100644 --- a/cmd/mesh-control/push.go +++ b/cmd/mesh-control/push.go @@ -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 ` + // 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