Give a cascade round's hold back on every return, a body that cannot be marshalled included (hq ADR 0100)
This commit is contained in:
@@ -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()
|
||||||
|
}
|
||||||
|
}
|
||||||
+38
-17
@@ -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.
|
// 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
|
// 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.
|
// 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,
|
||||||
if err != nil {
|
func(held context.Context, node string) (sendable, error) {
|
||||||
return err
|
|
||||||
}
|
|
||||||
sending, refused := composeEach(also, func(node string) (sendable, error) {
|
|
||||||
plan, settings, err := planFor(held, open, node)
|
plan, settings, err := planFor(held, open, node)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return sendable{}, err
|
return sendable{}, err
|
||||||
}
|
}
|
||||||
reportUnhostable(node, plan)
|
reportUnhostable(node, plan)
|
||||||
return declarationWith(held, open, node, plan, settings, gens, Allocating)
|
return declarationWith(held, open, node, plan, settings, gens, Allocating)
|
||||||
})
|
},
|
||||||
refusals = append(refusals, refused...)
|
func(s readyNode, body []byte) error {
|
||||||
for _, s := range sending {
|
if err := link.Declare(ctx, server.Channel(), ident, s.node, body,
|
||||||
body, err := s.declared.Body()
|
15*time.Second); err != nil {
|
||||||
if err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
if err := link.Declare(ctx, server.Channel(), ident, s.node, body, 15*time.Second); err != nil {
|
|
||||||
release()
|
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
record, err := inv.NodeByName(ctx, s.node)
|
record, err := inv.NodeByName(ctx, s.node)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
release()
|
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
if err := inv.RecordSent(ctx, record.ID, digestOf(body)); err != nil {
|
if err := inv.RecordSent(ctx, record.ID, digestOf(body)); err != nil {
|
||||||
release()
|
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
fmt.Printf("sent %s %d resource(s)\n", s.node, len(s.declared.Resources))
|
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
|
||||||
}
|
}
|
||||||
release()
|
|
||||||
// Every candidate this round is marked handled — the sent ones so they are not
|
// 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
|
// 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.
|
// 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
|
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.
|
// 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
|
// **After the rest have been sent, never instead of sending them.** It is still an error, because
|
||||||
|
|||||||
Reference in New Issue
Block a user