Adoption mode: a node in use is adopted before it is converged (hq ADR 0100–0103) #20
+13
-1
@@ -799,10 +799,22 @@ func applyDeclared(ctx context.Context, opts options, raw []byte, sched *apply.S
|
|||||||
return applyAndKeep(ctx, opts, raw, nil, sched)
|
return applyAndKeep(ctx, opts, raw, nil, sched)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// applying serialises applies on this node.
|
||||||
|
//
|
||||||
|
// **Two things apply here: the link and the reconcile loop**, and each reads the node's state,
|
||||||
|
// acts on the machine, and writes the state back. Run at the same time they interleave, and the
|
||||||
|
// one that saves last writes a state read before the other acted — losing what the first recorded:
|
||||||
|
// a hold, the firewall found here, a resource just applied. The machine would then be one thing
|
||||||
|
// and its record another, which is the fault every read-back in this package exists to prevent.
|
||||||
|
var applying sync.Mutex
|
||||||
|
|
||||||
// applyAndKeep applies a declaration and, when it came from the mesh, keeps it so this node can
|
// applyAndKeep applies a declaration and, when it came from the mesh, keeps it so this node can
|
||||||
// go on obeying it while disconnected.
|
// go on obeying it while disconnected. One at a time, whoever asks.
|
||||||
func applyAndKeep(ctx context.Context, opts options, raw []byte, signed *store.Declared,
|
func applyAndKeep(ctx context.Context, opts options, raw []byte, signed *store.Declared,
|
||||||
sched *apply.Scheduler) link.Report {
|
sched *apply.Scheduler) link.Report {
|
||||||
|
applying.Lock()
|
||||||
|
defer applying.Unlock()
|
||||||
|
|
||||||
declared, err := declaration.Parse(raw)
|
declared, err := declaration.Parse(raw)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return link.Report{Refused: err.Error()}
|
return link.Report{Refused: err.Error()}
|
||||||
|
|||||||
@@ -2,10 +2,16 @@ package main
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
"github.com/novox/mesh-host/internal/link"
|
"errors"
|
||||||
"github.com/novox/mesh-host/internal/store"
|
"os"
|
||||||
|
"path/filepath"
|
||||||
"testing"
|
"testing"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
|
"github.com/novox/mesh-host/internal/apply"
|
||||||
|
"github.com/novox/mesh-host/internal/link"
|
||||||
|
"github.com/novox/mesh-host/internal/store"
|
||||||
|
"github.com/novox/mesh-host/internal/system"
|
||||||
)
|
)
|
||||||
|
|
||||||
// Argument handling gets tests because it already failed silently once: `mesh-host inventory
|
// Argument handling gets tests because it already failed silently once: `mesh-host inventory
|
||||||
@@ -190,3 +196,51 @@ func TestAReconcileSpeaksWhenWhatIsReachableChanged(t *testing.T) {
|
|||||||
t.Error("a newly published port was not said")
|
t.Error("a newly published port was not said")
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Defends the node's own record: the link and the reconcile loop both apply, and each reads the
|
||||||
|
// state, acts, and writes it back — so they must not run at the same time, or the last save loses
|
||||||
|
// what the other recorded.
|
||||||
|
func TestOnlyOneApplyRunsAtATime(t *testing.T) {
|
||||||
|
// A host is built for one system at link time, and a test binary has no link time: this asks
|
||||||
|
// the machine it runs on, and stands aside where the answer is no.
|
||||||
|
built, err := system.For("arch")
|
||||||
|
if err != nil || built.Confirm(context.Background(), apply.ExecRunner) != nil {
|
||||||
|
t.Skip("this machine is not one these tests can apply on")
|
||||||
|
}
|
||||||
|
was := builtFor
|
||||||
|
builtFor = "arch"
|
||||||
|
t.Cleanup(func() { builtFor = was })
|
||||||
|
dir := t.TempDir()
|
||||||
|
opts := options{state: filepath.Join(dir, "state.json")}
|
||||||
|
raw := []byte(`{"declaration":1,"resources":[{"id":"a","type":"file","path":"` +
|
||||||
|
filepath.Join(dir, "a.conf") + `","content":"x\n"}]}`)
|
||||||
|
|
||||||
|
// Whatever else is applying — the link, while this is the reconcile — this waits for it.
|
||||||
|
applying.Lock()
|
||||||
|
done := make(chan link.Report, 1)
|
||||||
|
go func() { done <- applyAndKeep(context.Background(), opts, raw, nil, nil) }()
|
||||||
|
select {
|
||||||
|
case report := <-done:
|
||||||
|
applying.Unlock()
|
||||||
|
t.Fatalf("an apply ran while another held the node: %+v", report)
|
||||||
|
case <-time.After(50 * time.Millisecond):
|
||||||
|
}
|
||||||
|
if _, err := os.Stat(filepath.Join(dir, "a.conf")); !errors.Is(err, os.ErrNotExist) {
|
||||||
|
applying.Unlock()
|
||||||
|
t.Fatal("the waiting apply had already touched the machine")
|
||||||
|
}
|
||||||
|
applying.Unlock()
|
||||||
|
|
||||||
|
select {
|
||||||
|
case report := <-done:
|
||||||
|
if report.Refused != "" {
|
||||||
|
t.Fatalf("refused: %s", report.Refused)
|
||||||
|
}
|
||||||
|
case <-time.After(10 * time.Second):
|
||||||
|
t.Fatal("the apply never ran once the node was free")
|
||||||
|
}
|
||||||
|
known, loadErr := store.Load(opts.state)
|
||||||
|
if loadErr != nil || len(known.Resources) != 1 {
|
||||||
|
t.Errorf("the apply recorded %d resource(s): %v", len(known.Resources), loadErr)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user