Hold a node from composing to sending, so a push composed before converge or adopt is never sent after it (hq ADR 0100)
This commit is contained in:
@@ -385,3 +385,49 @@ func TestTheFlipActsOnlyOnThePreviewTheOperatorSaw(t *testing.T) {
|
|||||||
t.Fatal("the flip did not converge the node")
|
t.Fatal("the flip did not converge the node")
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// The flip holds the node from its checks to its send, so a push composed meanwhile waits and is
|
||||||
|
// composed after it — never sent after it with the node still adopted. And the send inside the
|
||||||
|
// flip is not made to wait on the flip's own hold.
|
||||||
|
func TestTheFlipHoldsTheNodeWhileItSends(t *testing.T) {
|
||||||
|
open, _ := anAdoptedAnchor(t)
|
||||||
|
ctx := t.Context()
|
||||||
|
if _, err := take(ctx, open, "anchor", "hello-web"); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
reportsHolding(t, open, heldFile)
|
||||||
|
preview, err := converge(ctx, open, "anchor", false, "", "")
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
var heldElsewhere, heldHere error
|
||||||
|
sendNodes = func(inner context.Context, open *stores, names []string) error {
|
||||||
|
// Another caller cannot hold the node while the flip sends it.
|
||||||
|
waiting, cancel := context.WithTimeout(ctx, 300*time.Millisecond)
|
||||||
|
defer cancel()
|
||||||
|
if release, err := open.inventory.HoldNodes(waiting, names); err == nil {
|
||||||
|
release()
|
||||||
|
heldElsewhere = errors.New("another caller held the node while the flip sent it")
|
||||||
|
}
|
||||||
|
// The flip's own send holds it without waiting on itself.
|
||||||
|
_, release, err := holdNodes(inner, open, names)
|
||||||
|
if err != nil {
|
||||||
|
heldHere = err
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
release()
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
if _, err := converge(ctx, open, "anchor", true, digestIn(t, preview), ""); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if heldElsewhere != nil || heldHere != nil {
|
||||||
|
t.Fatalf("%v %v", heldElsewhere, heldHere)
|
||||||
|
}
|
||||||
|
// And given back once it is done.
|
||||||
|
release, err := open.inventory.HoldNodes(ctx, []string{"anchor"})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
release()
|
||||||
|
}
|
||||||
|
|||||||
@@ -152,6 +152,16 @@ func converge(ctx context.Context, open *stores, node string, yes bool, digest s
|
|||||||
if filter == "" {
|
if filter == "" {
|
||||||
filter = DefaultFilter
|
filter = DefaultFilter
|
||||||
}
|
}
|
||||||
|
if yes {
|
||||||
|
// Held from the checks to the send, so no push composed before the flip is sent after it
|
||||||
|
// and returns the node to adopted.
|
||||||
|
held, release, err := holdNodes(ctx, open, []string{node})
|
||||||
|
if err != nil {
|
||||||
|
return "", err
|
||||||
|
}
|
||||||
|
defer release()
|
||||||
|
ctx = held
|
||||||
|
}
|
||||||
record, err := inv.NodeByName(ctx, node)
|
record, err := inv.NodeByName(ctx, node)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return "", err
|
return "", err
|
||||||
@@ -435,6 +445,12 @@ func adopt(ctx context.Context, open *stores, node string) (string, error) {
|
|||||||
if record.Adopted {
|
if record.Adopted {
|
||||||
return "", fmt.Errorf("%s is adopted already", node)
|
return "", fmt.Errorf("%s is adopted already", node)
|
||||||
}
|
}
|
||||||
|
held, release, err := holdNodes(ctx, open, []string{node})
|
||||||
|
if err != nil {
|
||||||
|
return "", err
|
||||||
|
}
|
||||||
|
defer release()
|
||||||
|
ctx = held
|
||||||
if err := inv.SetAdopted(ctx, node, true); err != nil {
|
if err := inv.SetAdopted(ctx, node, true); err != nil {
|
||||||
return "", err
|
return "", err
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -0,0 +1,37 @@
|
|||||||
|
package main
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
)
|
||||||
|
|
||||||
|
// heldKey carries the nodes this call already holds, so an act that holds a node and then sends
|
||||||
|
// through sendTo does not wait on itself.
|
||||||
|
type heldKey struct{}
|
||||||
|
|
||||||
|
// holdNodes holds the named nodes for composing and sending their declarations (novox/hq ADR
|
||||||
|
// 0100), skipping any the context already holds, and returns a context that says it holds them.
|
||||||
|
// Release gives back only what this call took.
|
||||||
|
func holdNodes(ctx context.Context, open *stores, names []string) (context.Context, func(), error) {
|
||||||
|
already, _ := ctx.Value(heldKey{}).(map[string]bool)
|
||||||
|
var take []string
|
||||||
|
for _, n := range names {
|
||||||
|
if !already[n] {
|
||||||
|
take = append(take, n)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if len(take) == 0 {
|
||||||
|
return ctx, func() {}, nil
|
||||||
|
}
|
||||||
|
release, err := open.inventory.HoldNodes(ctx, take)
|
||||||
|
if err != nil {
|
||||||
|
return ctx, nil, err
|
||||||
|
}
|
||||||
|
held := map[string]bool{}
|
||||||
|
for n := range already {
|
||||||
|
held[n] = true
|
||||||
|
}
|
||||||
|
for _, n := range take {
|
||||||
|
held[n] = true
|
||||||
|
}
|
||||||
|
return context.WithValue(ctx, heldKey{}, held), release, nil
|
||||||
|
}
|
||||||
@@ -282,8 +282,14 @@ func pushCommand(ctx context.Context, args []string) error {
|
|||||||
asked = append(asked, n.Name)
|
asked = append(asked, n.Name)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Held from composing to sending, so a converge on one of them cannot send between the two
|
||||||
|
// and be overtaken by what was composed before it (novox/hq ADR 0100).
|
||||||
|
held, release, err := holdNodes(ctx, open, asked)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
sending, refusals := composeEach(asked, func(node string) (sendable, error) {
|
sending, refusals := composeEach(asked, func(node string) (sendable, error) {
|
||||||
plan, settings, err := planFor(ctx, open, node)
|
plan, settings, err := planFor(held, open, node)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return sendable{}, err
|
return sendable{}, err
|
||||||
}
|
}
|
||||||
@@ -294,10 +300,11 @@ func pushCommand(ctx context.Context, args []string) error {
|
|||||||
// The private network is in here with everything else. It used to be composed separately
|
// The private network is in here with everything else. It used to be composed separately
|
||||||
// and prepended, which meant every machine with an address was on it and no machine could
|
// and prepended, which meant every machine with an address was on it and no machine could
|
||||||
// be kept off. It is a module now, so it arrives the way a module does.
|
// be kept off. It is a module now, so it arrives the way a module does.
|
||||||
return declarationWith(ctx, open, node, plan, settings, gens, Allocating)
|
return declarationWith(held, open, node, plan, settings, gens, Allocating)
|
||||||
})
|
})
|
||||||
|
|
||||||
sentDigest := map[string]string{}
|
sentDigest := map[string]string{}
|
||||||
|
defer release()
|
||||||
for _, s := range sending {
|
for _, s := range sending {
|
||||||
body, err := s.declared.Body()
|
body, err := s.declared.Body()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -319,6 +326,7 @@ func pushCommand(ctx context.Context, args []string) error {
|
|||||||
sentDigest[s.node] = digest
|
sentDigest[s.node] = digest
|
||||||
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))
|
||||||
}
|
}
|
||||||
|
release()
|
||||||
fmt.Printf("\n%d node(s) told\n", len(sending))
|
fmt.Printf("\n%d node(s) told\n", len(sending))
|
||||||
|
|
||||||
// **A named push leaves the mesh consistent, not just the machine it named** (novox/hq
|
// **A named push leaves the mesh consistent, not just the machine it named** (novox/hq
|
||||||
@@ -373,13 +381,19 @@ func pushCommand(ctx context.Context, args []string) error {
|
|||||||
// (novox/hq ADR 0066). The earlier cut routed these through sendTo, which is
|
// (novox/hq ADR 0066). The earlier cut routed these through sendTo, which is
|
||||||
// all-or-nothing — so one swept machine's compose error failed the operator's named
|
// all-or-nothing — so one swept machine's compose error failed the operator's named
|
||||||
// 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
|
||||||
|
// cascading into each other's machines never each wait on the other.
|
||||||
|
held, release, err := holdNodes(ctx, open, also)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
sending, refused := composeEach(also, func(node string) (sendable, error) {
|
sending, refused := composeEach(also, func(node string) (sendable, error) {
|
||||||
plan, settings, err := planFor(ctx, 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(ctx, open, node, plan, settings, gens, Allocating)
|
return declarationWith(held, open, node, plan, settings, gens, Allocating)
|
||||||
})
|
})
|
||||||
refusals = append(refusals, refused...)
|
refusals = append(refusals, refused...)
|
||||||
for _, s := range sending {
|
for _, s := range sending {
|
||||||
@@ -388,17 +402,21 @@ func pushCommand(ctx context.Context, args []string) error {
|
|||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
if err := link.Declare(ctx, server.Channel(), ident, s.node, body, 15*time.Second); err != nil {
|
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))
|
||||||
}
|
}
|
||||||
|
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.
|
||||||
@@ -532,6 +550,14 @@ func sendTo(ctx context.Context, open *stores, names []string) error {
|
|||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Held from composing to sending (novox/hq ADR 0100); a caller that holds them already —
|
||||||
|
// converge, which flips the node and then sends it — is not made to wait on itself.
|
||||||
|
ctx, release, err := holdNodes(ctx, open, names)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
defer release()
|
||||||
|
|
||||||
var sending []readyNode
|
var sending []readyNode
|
||||||
var refusals []string
|
var refusals []string
|
||||||
for _, name := range names {
|
for _, name := range names {
|
||||||
|
|||||||
@@ -0,0 +1,47 @@
|
|||||||
|
package inventory
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"slices"
|
||||||
|
"sync"
|
||||||
|
)
|
||||||
|
|
||||||
|
// HoldNodes serialises composing and sending a declaration per node (novox/hq ADR 0100): while
|
||||||
|
// one caller holds a node, another asking for it waits. Without it a push that composed a node as
|
||||||
|
// adopted could send that declaration after `converge --yes` sent the converged one, and the node
|
||||||
|
// would return to adopted with nobody having asked.
|
||||||
|
//
|
||||||
|
// Session-level advisory locks on one connection, taken in name order so two callers holding
|
||||||
|
// overlapping sets cannot each wait on the other. Release gives every one back, and may be called
|
||||||
|
// more than once; a caller whose context ends while waiting holds nothing.
|
||||||
|
func (i *Inventory) HoldNodes(ctx context.Context, names []string) (func(), error) {
|
||||||
|
sorted := slices.Clone(names)
|
||||||
|
slices.Sort(sorted)
|
||||||
|
sorted = slices.Compact(sorted)
|
||||||
|
conn, err := i.store.Pool().Acquire(ctx)
|
||||||
|
if err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
var once sync.Once
|
||||||
|
release := func() {
|
||||||
|
once.Do(func() {
|
||||||
|
// Unlocking all of this session's advisory locks, then handing the connection back: a
|
||||||
|
// connection returned still holding one would hold it for whoever borrows it next.
|
||||||
|
_, err := conn.Exec(context.WithoutCancel(ctx), `select pg_advisory_unlock_all()`)
|
||||||
|
if err != nil {
|
||||||
|
// The session's locks die with the session: close it rather than return it.
|
||||||
|
_ = conn.Conn().Close(context.WithoutCancel(ctx))
|
||||||
|
}
|
||||||
|
conn.Release()
|
||||||
|
})
|
||||||
|
}
|
||||||
|
for _, name := range sorted {
|
||||||
|
if _, err := conn.Exec(ctx,
|
||||||
|
`select pg_advisory_lock(hashtext('mesh-node-declaration:' || $1)::bigint)`,
|
||||||
|
name); err != nil {
|
||||||
|
release()
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return release, nil
|
||||||
|
}
|
||||||
@@ -0,0 +1,44 @@
|
|||||||
|
package inventory
|
||||||
|
|
||||||
|
import (
|
||||||
|
"testing"
|
||||||
|
"time"
|
||||||
|
)
|
||||||
|
|
||||||
|
// Composing and sending one node's declaration is serialised: a second holder waits for the first.
|
||||||
|
func TestHoldingANodeMakesTheNextHolderWait(t *testing.T) {
|
||||||
|
inv := fresh(t)
|
||||||
|
ctx := t.Context()
|
||||||
|
release, err := inv.HoldNodes(ctx, []string{"anchor", "laptop"})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
got := make(chan func(), 1)
|
||||||
|
go func() {
|
||||||
|
second, err := inv.HoldNodes(ctx, []string{"laptop"})
|
||||||
|
if err != nil {
|
||||||
|
t.Error(err)
|
||||||
|
got <- func() {}
|
||||||
|
return
|
||||||
|
}
|
||||||
|
got <- second
|
||||||
|
}()
|
||||||
|
select {
|
||||||
|
case <-got:
|
||||||
|
t.Fatal("a node held by one caller was held by another at the same time")
|
||||||
|
case <-time.After(300 * time.Millisecond):
|
||||||
|
}
|
||||||
|
// Another node is not held up.
|
||||||
|
other, err := inv.HoldNodes(ctx, []string{"joiner"})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
other()
|
||||||
|
release()
|
||||||
|
select {
|
||||||
|
case second := <-got:
|
||||||
|
second()
|
||||||
|
case <-time.After(5 * time.Second):
|
||||||
|
t.Fatal("releasing the node did not let the next holder in")
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user