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): } }