A controller that asked node-build-agent from its first run would queue every build where nothing pulls, and the build that registers build-agent — the first holder — would be among them. So the role is chosen at ask time from the catalogue: the current role when any assigned module claims it, the retired one while only the builder does, the current one when neither. Outcomes are followed on both seats, the controller may publish to both, and a build's log is read under whichever role did it; a machine on the retired role is proven on the bus to take that role's asks. The switch order is written where the role is named, and the retired half is marked for removal with the seat row.
364 lines
12 KiB
Go
364 lines
12 KiB
Go
package link
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"io"
|
|
"log"
|
|
"os"
|
|
"sync"
|
|
"testing"
|
|
"time"
|
|
|
|
"github.com/nats-io/nats.go"
|
|
|
|
"github.com/novox/mesh-controller/internal/broker"
|
|
)
|
|
|
|
// A build, end to end, against a real server.
|
|
//
|
|
// The claim worth checking is the one ADR 0121 rests on: **one publish reaches three audiences**.
|
|
// Whoever asked matches the outcome by the id their request carried; the controller records it; the
|
|
// catalogue places it. On the old bus that fan-out came from a shared exchange, and it would be easy
|
|
// to write a version where only the asker hears it and nobody notices for weeks.
|
|
//
|
|
// docker run -d --rm --name t -p 14230:4222 nats:2.10-alpine -js
|
|
// MESH_TEST_NATS=nats://127.0.0.1:14230 go test ./internal/link/ -run TestNatsABuild
|
|
|
|
func aBusWithTheBuildRole(t *testing.T) *broker.JetStream {
|
|
t.Helper()
|
|
url := os.Getenv("MESH_TEST_NATS")
|
|
if url == "" {
|
|
t.Skip("MESH_TEST_NATS unset")
|
|
}
|
|
js, err := broker.Dial(url)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
t.Cleanup(js.Close)
|
|
|
|
seats := []broker.DeclaredSeat{{Name: TheBuildMachine, Accepts: []string{"build"},
|
|
Emits: []string{"built"}}}
|
|
if err := broker.AssertMeshStreams(js); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := broker.RaiseSeats(js, seats, map[string]broker.Holder{
|
|
TheBuildMachine: {Node: "anchor", Module: "builder"},
|
|
}); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
clean := func() {
|
|
_ = js.Context().DeleteStream("SEAT_NODE_BUILD_AGENT")
|
|
for _, s := range broker.MeshStreams() {
|
|
_ = js.Context().PurgeStream(s.Name)
|
|
}
|
|
}
|
|
t.Cleanup(clean)
|
|
return js
|
|
}
|
|
|
|
// The whole round trip: asked, taken, built, and the outcome heard by the asker and by a consumer of
|
|
// the role's event who never asked for anything.
|
|
func TestNatsABuildIsTakenAndItsOutcomeReachesEverybody(t *testing.T) {
|
|
js := aBusWithTheBuildRole(t)
|
|
ctx, stop := context.WithCancel(context.Background())
|
|
defer stop()
|
|
|
|
// A third party on the role's event — what the catalogue is. Subscribed first, so nothing is
|
|
// missed.
|
|
heard := make(chan BuildResult, 4)
|
|
watching, err := js.Conn().Subscribe(BuildOutcome(), func(msg *nats.Msg) {
|
|
var r BuildResult
|
|
if json.Unmarshal(msg.Data, &r) == nil {
|
|
heard <- r
|
|
}
|
|
})
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
defer func() { _ = watching.Unsubscribe() }()
|
|
_ = js.Conn().Flush()
|
|
|
|
// And a reader following this one build by its subject alone (novox/hq ADR 0157).
|
|
lines, err := js.Conn().SubscribeSync(BuildLog("b-1"))
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
defer func() { _ = lines.Unsubscribe() }()
|
|
_ = js.Conn().Flush()
|
|
|
|
// A build machine holding the role.
|
|
machine := MachineOverNATS(js, "anchor")
|
|
defer machine.Close()
|
|
failed := make(chan error, 1)
|
|
go func() {
|
|
failed <- machine.Take(ctx, func(ctx context.Context, work Build) {
|
|
r := work.Request()
|
|
_ = work.Began(ctx)
|
|
work.Say("clone", "cloning /r")
|
|
work.Say("image", "building shop")
|
|
_ = work.Announce(ctx, BuildResult{
|
|
ID: r.ID, Repository: r.Repository, On: "anchor", Commit: "abc1234",
|
|
Manifest: json.RawMessage(`{"module":"shop"}`),
|
|
})
|
|
_ = work.Done()
|
|
})
|
|
}()
|
|
|
|
ask := &natsBuilds{js: js}
|
|
result, err := ask.Submit(ctx, BuildRequest{ID: "b-1", Repository: "/r"}, 15*time.Second)
|
|
if err != nil {
|
|
select {
|
|
case why := <-failed:
|
|
t.Fatalf("the machine could not take work: %v", why)
|
|
default:
|
|
}
|
|
t.Fatalf("the asker never got an outcome: %v", err)
|
|
}
|
|
// The reader heard the build as it went, in order, under its id.
|
|
for want := 1; want <= 2; want++ {
|
|
msg, err := lines.NextMsg(5 * time.Second)
|
|
if err != nil {
|
|
t.Fatalf("line %d of the build never reached its subject: %v", want, err)
|
|
}
|
|
var line BuildLine
|
|
if err := json.Unmarshal(msg.Data, &line); err != nil || line.ID != "b-1" || line.Seq != want {
|
|
t.Fatalf("line %d came back as %s (%v)", want, msg.Data, err)
|
|
}
|
|
}
|
|
// And it is in the stream for a reader who comes later.
|
|
info, err := js.Context().StreamInfo(broker.EventsStream, &nats.StreamInfoRequest{SubjectsFilter: BuildLog("b-1")})
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if info.State.Subjects[BuildLog("b-1")] != 2 {
|
|
t.Fatalf("the stream holds %v under the build's subject, want 2", info.State.Subjects)
|
|
}
|
|
if err == nil {
|
|
}
|
|
if result.ID != "b-1" || result.Commit != "abc1234" {
|
|
t.Fatalf("the asker got %+v", result)
|
|
}
|
|
// Named in the outcome, because only the manifest says what was built and the catalogue reading
|
|
// this event needs to know.
|
|
if result.Module != "shop" {
|
|
t.Errorf("the outcome names module %q, so a catalogue reading it cannot place the build",
|
|
result.Module)
|
|
}
|
|
|
|
select {
|
|
case also := <-heard:
|
|
if also.ID != "b-1" {
|
|
t.Fatalf("a third party heard %+v", also)
|
|
}
|
|
case <-time.After(5 * time.Second):
|
|
t.Fatal("nobody but the asker heard the outcome, so the catalogue would never place the build")
|
|
}
|
|
|
|
// And the work left the queue: a request a machine took and settled must not be given to another.
|
|
deadline := time.Now().Add(5 * time.Second)
|
|
for time.Now().Before(deadline) {
|
|
info, err := js.Context().StreamInfo("SEAT_NODE_BUILD_AGENT")
|
|
if err == nil && info.State.Msgs == 0 {
|
|
return
|
|
}
|
|
time.Sleep(20 * time.Millisecond)
|
|
}
|
|
t.Fatal("the request is still queued after being settled, so another machine would build it again")
|
|
}
|
|
|
|
// Work submitted with no machine holding the role waits rather than failing, and is done when one
|
|
// arrives. **That is what a queue is for**, and the alternative — refusing because nobody is there
|
|
// yet — would make installing a build machine an ordering problem.
|
|
func TestNatsABuildWaitsForAMachineRatherThanFailing(t *testing.T) {
|
|
js := aBusWithTheBuildRole(t)
|
|
|
|
body, _ := json.Marshal(BuildRequest{ID: "b-2", Repository: "/r"})
|
|
if _, err := js.Context().Publish(BuildWork(), body); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
info, err := js.Context().StreamInfo("SEAT_NODE_BUILD_AGENT")
|
|
if err != nil || info.State.Msgs != 1 {
|
|
t.Fatalf("the work did not queue: %+v %v", info, err)
|
|
}
|
|
|
|
// Now a machine arrives and finds it waiting.
|
|
ctx, stop := context.WithCancel(context.Background())
|
|
defer stop()
|
|
took := make(chan string, 1)
|
|
machine := MachineOverNATS(js, "anchor")
|
|
defer machine.Close()
|
|
go func() {
|
|
_ = machine.Take(ctx, func(ctx context.Context, work Build) {
|
|
took <- work.Request().ID
|
|
_ = work.Announce(ctx, BuildResult{ID: work.Request().ID, On: "anchor", Failed: "no"})
|
|
_ = work.Done()
|
|
})
|
|
}()
|
|
select {
|
|
case id := <-took:
|
|
if id != "b-2" {
|
|
t.Fatalf("the machine took %q", id)
|
|
}
|
|
case <-time.After(10 * time.Second):
|
|
t.Fatal("a machine that arrived after the work did never got it, so the backlog was lost")
|
|
}
|
|
}
|
|
|
|
// A machine that dies before saying anything leaves the work for another.
|
|
func TestNatsWorkAMachineDidNotAnswerGoesBackToTheQueue(t *testing.T) {
|
|
js := aBusWithTheBuildRole(t)
|
|
ctx, stop := context.WithCancel(context.Background())
|
|
defer stop()
|
|
|
|
body, _ := json.Marshal(BuildRequest{ID: "b-3", Repository: "/r"})
|
|
if _, err := js.Context().Publish(BuildWork(), body); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
|
|
var once sync.Once
|
|
handed := make(chan struct{}, 2)
|
|
machine := MachineOverNATS(js, "anchor")
|
|
defer machine.Close()
|
|
go func() {
|
|
_ = machine.Take(ctx, func(ctx context.Context, work Build) {
|
|
handed <- struct{}{}
|
|
// The first time, hand it straight back — a machine that stopped mid-build.
|
|
var settled bool
|
|
once.Do(func() { _ = work.Hold(200 * time.Millisecond); settled = true })
|
|
if !settled {
|
|
_ = work.Announce(ctx, BuildResult{ID: work.Request().ID, On: "anchor"})
|
|
_ = work.Done()
|
|
}
|
|
})
|
|
}()
|
|
|
|
for i := 0; i < 2; i++ {
|
|
select {
|
|
case <-handed:
|
|
case <-time.After(10 * time.Second):
|
|
t.Fatalf("the work was handed over %d time(s); unanswered work must come back", i)
|
|
}
|
|
}
|
|
}
|
|
|
|
func quietLog() *log.Logger { return log.New(io.Discard, "", 0) }
|
|
|
|
var _ = quietLog
|
|
|
|
// Two machines holding the role share one queue (novox/hq ADR 0190): three asks, each machine takes
|
|
// one and the third waits until one of them is done; an ask is never handed to a machine that is
|
|
// busy; and a machine that stops mid-ask leaves its ask to the other.
|
|
func TestNatsTwoMachinesShareTheWorkAndNeitherIsHandedMoreThanItCanTake(t *testing.T) {
|
|
js := aBusWithTheBuildRole(t)
|
|
ctx, stop := context.WithCancel(context.Background())
|
|
defer stop()
|
|
|
|
for _, id := range []string{"w-1", "w-2", "w-3"} {
|
|
body, _ := json.Marshal(BuildRequest{ID: id, Repository: "/r"})
|
|
if _, err := js.Context().Publish(BuildWork(), body); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
}
|
|
|
|
type taken struct{ machine, id string }
|
|
took := make(chan taken, 8)
|
|
release := map[string]chan struct{}{"anchor": make(chan struct{}), "laptop": make(chan struct{})}
|
|
machines := map[string]BuildMachine{}
|
|
for _, name := range []string{"anchor", "laptop"} {
|
|
name := name
|
|
m := MachineOverNATS(js, name)
|
|
machines[name] = m
|
|
defer m.Close()
|
|
go func() {
|
|
_ = m.Take(ctx, func(ctx context.Context, work Build) {
|
|
took <- taken{name, work.Request().ID}
|
|
<-release[name]
|
|
_ = work.Announce(ctx, BuildResult{ID: work.Request().ID, On: name})
|
|
_ = work.Done()
|
|
})
|
|
}()
|
|
}
|
|
|
|
// Each machine took exactly one, and they are different asks.
|
|
first := map[string]string{}
|
|
for i := 0; i < 2; i++ {
|
|
select {
|
|
case got := <-took:
|
|
if _, twice := first[got.machine]; twice {
|
|
t.Fatalf("%s was handed a second ask while busy with its first", got.machine)
|
|
}
|
|
first[got.machine] = got.id
|
|
case <-time.After(10 * time.Second):
|
|
t.Fatalf("only %d machine(s) took work; two idle holders should both have", len(first))
|
|
}
|
|
}
|
|
if first["anchor"] == first["laptop"] {
|
|
t.Fatalf("both machines took %q: the queue is not shared, it is copied", first["anchor"])
|
|
}
|
|
// The third waits: nobody is free.
|
|
select {
|
|
case got := <-took:
|
|
t.Fatalf("%s was handed %s while both machines were busy", got.machine, got.id)
|
|
case <-time.After(2 * time.Second):
|
|
}
|
|
// One finishes, and only then is the third taken — by that machine, the one that is free.
|
|
close(release["anchor"])
|
|
release["anchor"] = make(chan struct{})
|
|
select {
|
|
case got := <-took:
|
|
if got.machine != "anchor" {
|
|
t.Fatalf("the third ask went to %s, which is still busy", got.machine)
|
|
}
|
|
case <-time.After(10 * time.Second):
|
|
t.Fatal("the third ask was never taken after a machine became free")
|
|
}
|
|
// A machine that stops mid-ask leaves its ask unacknowledged, and the ack wait brings it round
|
|
// to whoever is left — the path TestNatsWorkAMachineDidNotAnswerGoesBackToTheQueue proves with
|
|
// an explicit hand-back, because the real wait is a minute. Here: the laptop goes, anchor
|
|
// finishes, and with nothing queued nothing more is taken by the machine that is left.
|
|
machines["laptop"].Close()
|
|
close(release["anchor"])
|
|
select {
|
|
case got := <-took:
|
|
t.Fatalf("%s took %s; the queue should be empty", got.machine, got.id)
|
|
case <-time.After(2 * time.Second):
|
|
}
|
|
}
|
|
|
|
// During the handover (ADR 0190) two build roles exist. A machine whose credential claims the retired
|
|
// one takes an ask published to that seat and answers as that seat; the asker of that seat hears it.
|
|
func TestNatsAMachineOnTheRetiredBuildRoleTakesThatRolesAsks(t *testing.T) {
|
|
js := aBusWithTheBuildRole(t)
|
|
seats := []broker.DeclaredSeat{{Name: TheBuildMachineBefore, Accepts: []string{"build"}, Emits: []string{"built"}}}
|
|
if err := broker.RaiseSeats(js, seats, map[string]broker.Holder{TheBuildMachineBefore: {Node: "anchor", Module: "builder"}}); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
t.Cleanup(func() { _ = js.Context().DeleteStream("SEAT_MESH_BUILD_MACHINE") })
|
|
ctx, stop := context.WithCancel(context.Background())
|
|
defer stop()
|
|
|
|
machine := MachineOverNATSOn(js, "anchor", TheBuildMachineBefore)
|
|
defer machine.Close()
|
|
go func() {
|
|
_ = machine.Take(ctx, func(ctx context.Context, work Build) {
|
|
_ = work.Began(ctx)
|
|
_ = work.Announce(ctx, BuildResult{ID: work.Request().ID, Repository: work.Request().Repository, On: "anchor", Commit: "abc"})
|
|
_ = work.Done()
|
|
})
|
|
}()
|
|
|
|
asker, err := BuildsOverNATSOn(os.Getenv("MESH_TEST_NATS"), TheBuildMachineBefore)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
defer asker.Close()
|
|
result, err := asker.Submit(ctx, BuildRequest{ID: "build-old-seat", Repository: "r"}, 20*time.Second)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if result.On != "anchor" || result.ID != "build-old-seat" {
|
|
t.Errorf("the retired role's holder did not answer: %+v", result)
|
|
}
|
|
}
|