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_MESH_BUILD_MACHINE") 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() // 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.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) } 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_MESH_BUILD_MACHINE") 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_MESH_BUILD_MACHINE") 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