diff --git a/cmd/mesh-controller/hold_test.go b/cmd/mesh-controller/hold_test.go new file mode 100644 index 0000000..9c2b480 --- /dev/null +++ b/cmd/mesh-controller/hold_test.go @@ -0,0 +1,45 @@ +package main + +import ( + "context" + "errors" + "testing" + "time" +) + +// A cascade round gives its hold back on every way out: a declaration that cannot be marshalled +// and a send that fails leave nobody waiting for the node. +func TestASendRoundGivesItsHoldBackOnEveryWayOut(t *testing.T) { + open := aMesh(t) + ctx := t.Context() + unmarshallable := func(context.Context, string) (sendable, error) { + return sendable{Resources: []map[string]any{{"id": "x", "bad": make(chan int)}}}, nil + } + plain := func(context.Context, string) (sendable, error) { + return sendable{Resources: []map[string]any{{"id": "x"}}}, nil + } + failing := func(readyNode, []byte) error { return errors.New("the broker went away") } + fine := func(readyNode, []byte) error { return nil } + + for name, round := range map[string]func() error{ + "a body that cannot be marshalled": func() error { + _, err := sendRound(ctx, open, []string{"anchor"}, unmarshallable, fine) + return err + }, + "a send that fails": func() error { + _, err := sendRound(ctx, open, []string{"anchor"}, plain, failing) + return err + }, + } { + if err := round(); err == nil { + t.Fatalf("%s was not an error", name) + } + waiting, cancel := context.WithTimeout(ctx, 2*time.Second) + release, err := open.inventory.HoldNodes(waiting, []string{"anchor"}) + cancel() + if err != nil { + t.Fatalf("after %s the node is still held: %v", name, err) + } + release() + } +} diff --git a/cmd/mesh-controller/push.go b/cmd/mesh-controller/push.go index 759cb34..2319213 100644 --- a/cmd/mesh-controller/push.go +++ b/cmd/mesh-controller/push.go @@ -383,40 +383,34 @@ func pushCommand(ctx context.Context, args []string) error { // push and skipped its --wait, the very intolerance the main path exists to avoid. // Held for this round only, and after the last round's were given back, so two pushes // cascading into each other's machines never each wait on the other. - held, release, err := holdNodes(ctx, open, also) + refused, err := sendRound(ctx, open, also, + func(held context.Context, node string) (sendable, error) { + plan, settings, err := planFor(held, open, node) + if err != nil { + return sendable{}, err + } + reportUnhostable(node, plan) + return declarationWith(held, open, node, plan, settings, gens, Allocating) + }, + func(s readyNode, body []byte) error { + if err := link.Declare(ctx, server.Channel(), ident, s.node, body, + 15*time.Second); err != nil { + return err + } + record, err := inv.NodeByName(ctx, s.node) + if err != nil { + return err + } + if err := inv.RecordSent(ctx, record.ID, digestOf(body)); err != nil { + return err + } + fmt.Printf("sent %s %d resource(s)\n", s.node, len(s.declared.Resources)) + return nil + }) + refusals = append(refusals, refused...) if err != nil { return err } - sending, refused := composeEach(also, func(node string) (sendable, error) { - plan, settings, err := planFor(held, open, node) - if err != nil { - return sendable{}, err - } - reportUnhostable(node, plan) - return declarationWith(held, open, node, plan, settings, gens, Allocating) - }) - refusals = append(refusals, refused...) - for _, s := range sending { - body, err := s.declared.Body() - if err != nil { - return err - } - if err := link.Declare(ctx, server.Channel(), ident, s.node, body, 15*time.Second); err != nil { - release() - return err - } - record, err := inv.NodeByName(ctx, s.node) - if err != nil { - release() - return err - } - if err := inv.RecordSent(ctx, record.ID, digestOf(body)); err != nil { - release() - return err - } - fmt.Printf("sent %s %d resource(s)\n", s.node, len(s.declared.Resources)) - } - release() // Every candidate this round is marked handled — the sent ones so they are not // re-listed, and the refused ones so a machine that cannot be composed does not make // the loop spin on it for ever. Its refusal is already in the report. @@ -516,6 +510,33 @@ func composeEach(names []string, return sending, refusals } +// sendRound holds the named nodes, composes each and sends each that composed, and gives the hold +// back on every way out — a body that cannot be marshalled and a send that fails included +// (novox/hq ADR 0100). A node that cannot be composed is a refusal, not an error: the others are +// still sent. +func sendRound(ctx context.Context, open *stores, names []string, + compose func(held context.Context, node string) (sendable, error), + send func(s readyNode, body []byte) error) ([]string, error) { + held, release, err := holdNodes(ctx, open, names) + if err != nil { + return nil, err + } + defer release() + sending, refused := composeEach(names, func(node string) (sendable, error) { + return compose(held, node) + }) + for _, s := range sending { + body, err := s.declared.Body() + if err != nil { + return refused, err + } + if err := send(s, body); err != nil { + return refused, err + } + } + return refused, nil +} + // couldNotBeResolved is what a push ends with when some machines could not be worked out. // // **After the rest have been sent, never instead of sending them.** It is still an error, because