plans said a build queued for minutes had been building for a few seconds: the line counted from the last save, and every advance saved the plan whether or not it moved, bumping its revision and saying plan-moved on the bus. The line, status and LATE now count from when the plan entered its tier with the stalled condition's bound, and an advance that changes nothing writes nothing.
773 lines
33 KiB
Go
773 lines
33 KiB
Go
package main
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"fmt"
|
|
"reflect"
|
|
"strings"
|
|
"testing"
|
|
"time"
|
|
|
|
"github.com/nats-io/nats.go"
|
|
|
|
"github.com/novox/mesh-controller/internal/broker"
|
|
"github.com/novox/mesh-controller/internal/catalogue"
|
|
"github.com/novox/mesh-controller/internal/inventory"
|
|
"github.com/novox/mesh-controller/internal/link"
|
|
|
|
"github.com/novox/mesh-controller/internal/testbus"
|
|
)
|
|
|
|
// The build queue, controlled by hand (novox/hq ADR 0219), and the plans that follow it.
|
|
//
|
|
// make postgres PG_PORT=55566 PG_CONTAINER=bq-pg
|
|
// MESH_TEST_POSTGRES='postgres://postgres:check@127.0.0.1:55566/postgres?sslmode=disable' \
|
|
// go test ./cmd/mesh-controller/ -run 'Queue|Cancel|Clear|Retry|Rebuild|Replay|Paused' (each test on a bus of its own)
|
|
|
|
// asksRecorded makes every ask a plan makes return the next id in a row, and says which were asked.
|
|
func asksRecorded(t *testing.T) *[]string {
|
|
t.Helper()
|
|
var asked []string
|
|
was := askABuild
|
|
askABuild = func(_ context.Context, source buildSource, path, ref string) (string, error) {
|
|
id := link.NewBuildID(time.Now().Add(time.Duration(len(asked)) * time.Millisecond))
|
|
asked = append(asked, source.Repository+"#"+id)
|
|
return id, nil
|
|
}
|
|
t.Cleanup(func() { askABuild = was })
|
|
return &asked
|
|
}
|
|
|
|
// twoTiers registers two modules, b standing on a, each with a source a plan asks.
|
|
func twoTiers(t *testing.T, open *stores) {
|
|
t.Helper()
|
|
for _, name := range []string{"a", "b"} {
|
|
if err := open.inventory.RegisterModule(t.Context(), catalogue.Manifest{Module: name, Version: "1"},
|
|
inventory.Source{Repository: "novox/" + name, Seat: "git", Ref: "main", BuiltFrom: "c0ffee", Head: "c0ffee"}); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
}
|
|
}
|
|
|
|
// A failed plan is retried: its failed module asked again under a new id, the plan building at that
|
|
// tier, and when that build comes in the plan goes on and asks its next tier.
|
|
func TestAFailedPlanIsRetriedAndGoesOnThroughItsLaterTiers(t *testing.T) {
|
|
open := aMesh(t)
|
|
ctx := t.Context()
|
|
asked := asksRecorded(t)
|
|
twoTiers(t, open)
|
|
before := time.Now().UTC().Add(-time.Hour)
|
|
failed := inventory.Plan{ID: "plan-retry", Repository: "novox/a", Branch: "main", Commit: "c0ffee",
|
|
Created: before, State: inventory.PlanFailed, Tier: 0, Tiers: [][]string{{"a"}, {"b"}},
|
|
Note: "a failed to build in tier 0",
|
|
Modules: map[string]*inventory.PlanModule{"a": {State: "failed", AskedAt: &before, Build: "build-1",
|
|
Why: link.KilledByHand}}}
|
|
if err := open.inventory.SavePlan(ctx, &failed); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
|
|
said, err := retryPlan(ctx, open, failed.ID)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
p, err := open.inventory.PlanByID(ctx, failed.ID)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
a := p.Modules["a"]
|
|
if p.State != inventory.PlanBuilding || a.State != "asked" || a.Build == "" || a.Build == "build-1" || a.Why != "" {
|
|
t.Fatalf("after retry the plan is %s and a is %+v", p.State, a)
|
|
}
|
|
if !strings.Contains(said, a.Build) {
|
|
t.Errorf("retry does not say the new id: %q", said)
|
|
}
|
|
if len(*asked) != 1 {
|
|
t.Fatalf("asked %v", *asked)
|
|
}
|
|
|
|
// The killed build's own late outcome is not this ask's; the new one's is, and tier 1 follows.
|
|
planBuilt(ctx, open, "a", "c0ffee", link.KilledByHand, before, "build-1")
|
|
if p, _ = open.inventory.PlanByID(ctx, failed.ID); p.State != inventory.PlanBuilding {
|
|
t.Fatalf("the old ask's outcome failed the retried plan: %s %q", p.State, p.Note)
|
|
}
|
|
asking, _ := link.BuildAskedAt(a.Build)
|
|
planBuilt(ctx, open, "a", "c0ffee", "", asking, a.Build)
|
|
p, err = open.inventory.PlanByID(ctx, failed.ID)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if p.Tier != 1 || p.Modules["b"] == nil || p.Modules["b"].State != "asked" || p.Modules["b"].Build == "" {
|
|
t.Fatalf("the retried plan did not go on to tier 1: tier %d, %s, b %+v", p.Tier, p.State, p.Modules["b"])
|
|
}
|
|
if len(*asked) != 2 {
|
|
t.Fatalf("asked %v", *asked)
|
|
}
|
|
}
|
|
|
|
// What retry refuses, and says why.
|
|
func TestRetryRefusesWhatItCannotResume(t *testing.T) {
|
|
at := time.Date(2026, 10, 5, 12, 0, 0, 0, time.UTC)
|
|
plan := func(id, state string, created time.Time) inventory.Plan {
|
|
return inventory.Plan{ID: id, Repository: "novox/mesh-catalog", Branch: "main", Commit: id + "c0ffee",
|
|
Created: created, State: state, Tiers: [][]string{{"a"}},
|
|
Modules: map[string]*inventory.PlanModule{"a": {State: "failed"}}}
|
|
}
|
|
failed := plan("plan-1", inventory.PlanFailed, at)
|
|
for _, c := range []struct {
|
|
p inventory.Plan
|
|
others []inventory.Plan
|
|
says string
|
|
}{
|
|
{plan("plan-d", inventory.PlanDone, at), nil, "is done"},
|
|
{plan("plan-s", inventory.PlanSuperseded, at), nil, "was"},
|
|
{plan("plan-o", inventory.PlanBuilding, at), nil, "still building"},
|
|
{failed, []inventory.Plan{plan("plan-2", inventory.PlanRolling, at.Add(time.Hour))}, "plan-2 supersedes it"},
|
|
} {
|
|
err := retryRefusal(c.p, c.others)
|
|
if err == nil || !strings.Contains(err.Error(), c.says) {
|
|
t.Errorf("%s: %v, wanted it to say %q", c.p.ID, err, c.says)
|
|
}
|
|
}
|
|
stopped := failed
|
|
stopped.Modules = map[string]*inventory.PlanModule{"a": {State: "built"}}
|
|
stopped.Note = "a stopped at its first machine"
|
|
if err := retryRefusal(stopped, nil); err == nil || !strings.Contains(err.Error(), "nothing in tier 0") {
|
|
t.Errorf("a plan with nothing failed to build was retried: %v", err)
|
|
}
|
|
// Another branch's newer plan, an older one, and a failed one do not supersede it.
|
|
other := plan("plan-3", inventory.PlanBuilding, at.Add(time.Hour))
|
|
other.Branch = "release"
|
|
if err := retryRefusal(failed, []inventory.Plan{other, plan("plan-0", inventory.PlanBuilding, at.Add(-time.Hour)),
|
|
plan("plan-4", inventory.PlanFailed, at.Add(time.Hour))}); err != nil {
|
|
t.Errorf("refused for a plan that does not supersede it: %v", err)
|
|
}
|
|
}
|
|
|
|
// `rebuild` of a module a failed plan holds joins that plan; one held by nothing runs alone.
|
|
func TestARebuildJoinsThePlanHoldingTheModule(t *testing.T) {
|
|
open := aMesh(t)
|
|
ctx := t.Context()
|
|
asked := asksRecorded(t)
|
|
twoTiers(t, open)
|
|
before := time.Now().UTC().Add(-time.Hour)
|
|
failed := inventory.Plan{ID: "plan-join", Repository: "novox/a", Branch: "main", Commit: "c0ffee",
|
|
Created: before, State: inventory.PlanFailed, Tiers: [][]string{{"a"}, {"b"}},
|
|
Note: "a failed to build in tier 0",
|
|
Modules: map[string]*inventory.PlanModule{"a": {State: "failed", AskedAt: &before, Build: "build-1", Why: link.CancelledByHand}}}
|
|
if err := open.inventory.SavePlan(ctx, &failed); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := rebuildCommand(ctx, []string{"a"}); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
p, err := open.inventory.PlanByID(ctx, failed.ID)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if p.State != inventory.PlanBuilding || p.Modules["a"].State != "asked" || p.Modules["a"].Build == "build-1" {
|
|
t.Fatalf("the rebuild did not join the failed plan: %s %+v", p.State, p.Modules["a"])
|
|
}
|
|
if len(*asked) != 1 {
|
|
t.Fatalf("asked %v", *asked)
|
|
}
|
|
// b is in a tier not yet reached: nothing holds it in its current tier, so it is asked alone.
|
|
if err := rebuildCommand(ctx, []string{"b"}); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if p, _ = open.inventory.PlanByID(ctx, failed.ID); p.Modules["b"] != nil {
|
|
t.Fatalf("a module the plan has not reached joined it: %+v", p.Modules["b"])
|
|
}
|
|
if len(*asked) != 2 {
|
|
t.Fatalf("asked %v", *asked)
|
|
}
|
|
}
|
|
|
|
// A failed plan an open plan of its repository supersedes is not joined: the open one holds the module.
|
|
func TestARebuildJoinsTheOpenPlanBeforeASupersededFailedOne(t *testing.T) {
|
|
at := time.Date(2026, 10, 5, 12, 0, 0, 0, time.UTC)
|
|
failed := inventory.Plan{ID: "plan-old", Repository: "novox/a", Branch: "main", Created: at, State: inventory.PlanFailed,
|
|
Tiers: [][]string{{"a"}}, Modules: map[string]*inventory.PlanModule{"a": {State: "failed"}}}
|
|
newer := inventory.Plan{ID: "plan-new", Repository: "novox/a", Branch: "main", Created: at.Add(time.Hour),
|
|
State: inventory.PlanBuilding, Tiers: [][]string{{"x"}, {"a"}}, Modules: map[string]*inventory.PlanModule{}}
|
|
if _, found := planHolding("a", []inventory.Plan{newer}, []inventory.Plan{newer, failed}); found {
|
|
t.Fatal("joined a failed plan a newer open one supersedes")
|
|
}
|
|
if p, found := planHolding("a", nil, []inventory.Plan{failed}); !found || p.ID != "plan-old" {
|
|
t.Fatalf("did not join the failed plan holding it: %v %s", found, p.ID)
|
|
}
|
|
built := failed
|
|
built.Modules = map[string]*inventory.PlanModule{"a": {State: "built"}}
|
|
if _, found := planHolding("a", nil, []inventory.Plan{built}); found {
|
|
t.Fatal("joined a plan that has the module built")
|
|
}
|
|
}
|
|
|
|
// replay is a dry run unless registered, and registering an older commit than one registered since
|
|
// is refused unless --older says it is meant (novox/hq issue 207).
|
|
func TestReplayRefusesToRollAnOlderCommitOutUnlessToldTo(t *testing.T) {
|
|
asked := time.Date(2026, 10, 5, 12, 0, 0, 0, time.UTC)
|
|
old := inventory.Build{ID: "build-old", Module: "a", Commit: "0ldc0mm1t", Asked: asked}
|
|
newer := inventory.Build{ID: "build-new", Module: "a", Commit: "n3wc0mm1t", Asked: asked.Add(time.Hour)}
|
|
history := []inventory.Build{newer, old}
|
|
|
|
if err := replayRefusal(old, "a", history, false, false, nil); err != nil {
|
|
t.Errorf("a dry run was refused: %v", err)
|
|
}
|
|
err := replayRefusal(old, "a", history, true, false, nil)
|
|
if err == nil || !strings.Contains(err.Error(), "build-new") || !strings.Contains(err.Error(), "--older") {
|
|
t.Errorf("registering an older commit than the one registered was not refused: %v", err)
|
|
}
|
|
if err := replayRefusal(old, "a", history, true, true, nil); err != nil {
|
|
t.Errorf("--register --older was refused: %v", err)
|
|
}
|
|
if err := replayRefusal(newer, "a", history, true, false, nil); err != nil {
|
|
t.Errorf("registering the newest build again was refused: %v", err)
|
|
}
|
|
// A newer build of the same commit, or one that failed, is not a newer version to roll back from.
|
|
same := inventory.Build{ID: "build-same", Module: "a", Commit: old.Commit, Asked: asked.Add(2 * time.Hour)}
|
|
broken := inventory.Build{ID: "build-broken", Module: "a", Failed: "no", Asked: asked.Add(3 * time.Hour)}
|
|
if err := replayRefusal(old, "a", []inventory.Build{broken, same, old}, true, false, nil); err != nil {
|
|
t.Errorf("refused for a newer build of the same commit or a failed one: %v", err)
|
|
}
|
|
if err := replayRefusal(inventory.Build{ID: "build-x", Failed: "clone"}, "", nil, false, false, nil); err == nil {
|
|
t.Error("a build that recorded no commit was replayed")
|
|
}
|
|
if err := replayRefusal(old, "a", history, false, true, nil); err == nil {
|
|
t.Error("--older without --register was taken")
|
|
}
|
|
// **Nothing newer outstanding**, --older or not: an ask of the module in the queue, or an open plan
|
|
// holding it unbuilt, would be replaced by a replay asked now (issue 219).
|
|
q := link.Queue{Asks: []link.QueuedAsk{
|
|
{ID: "build-queued", Repository: "https://forge.example/novox/a.git", State: link.AskWaiting},
|
|
{ID: "build-elsewhere", Repository: "https://forge.example/novox/b.git", State: link.AskWaiting},
|
|
{ID: "build-other-path", Repository: "https://forge.example/novox/a.git", Path: "sub", State: link.AskWaiting},
|
|
}}
|
|
plans := []inventory.Plan{
|
|
{ID: "plan-holds", State: inventory.PlanBuilding, Tiers: [][]string{{"x"}, {"a"}}, Modules: map[string]*inventory.PlanModule{}},
|
|
{ID: "plan-built", State: inventory.PlanRolling, Tiers: [][]string{{"a"}}, Modules: map[string]*inventory.PlanModule{"a": {State: "built"}}},
|
|
{ID: "plan-failed", State: inventory.PlanFailed, Tiers: [][]string{{"a"}}, Modules: map[string]*inventory.PlanModule{}},
|
|
}
|
|
outstanding := outstandingFor("a", "novox/a", "", q, plans)
|
|
if len(outstanding) != 2 || !strings.Contains(outstanding[0], "build-queued") || !strings.Contains(outstanding[1], "plan-holds") {
|
|
t.Fatalf("outstanding: %v", outstanding)
|
|
}
|
|
if err := replayRefusal(old, "a", history, true, true, outstanding); err == nil || !strings.Contains(err.Error(), "build-queued") {
|
|
t.Errorf("registered over an outstanding ask: %v", err)
|
|
}
|
|
if err := replayRefusal(old, "a", history, false, false, outstanding); err != nil {
|
|
t.Errorf("a dry run was refused for what is outstanding: %v", err)
|
|
}
|
|
if said := replaySaid(old, "a", "build-1", false); !strings.Contains(said, "dry run") || !strings.Contains(said, "--register") {
|
|
t.Errorf("a dry replay does not say what it is: %q", said)
|
|
}
|
|
if said := replaySaid(old, "a", "build-1", true); !strings.Contains(said, "registered") || !strings.Contains(said, "rolled out") {
|
|
t.Errorf("a registered replay does not say what it does: %q", said)
|
|
}
|
|
}
|
|
|
|
// A plan waiting on builds of a seat whose every holder is paused says so and is not late; a seat
|
|
// paused on some holders only is not a reason the plan is waiting.
|
|
func TestAPlanWaitingOnAPausedSeatSaysSoAndIsNotLate(t *testing.T) {
|
|
now := time.Date(2026, 10, 5, 12, 0, 0, 0, time.UTC)
|
|
asked := now.Add(-12 * time.Minute)
|
|
p := inventory.Plan{ID: "plan-p", Repository: "novox/a", Commit: "c0ffee", State: inventory.PlanBuilding,
|
|
Updated: now.Add(-2 * time.Hour), Tiers: [][]string{{"a"}},
|
|
Modules: map[string]*inventory.PlanModule{"a": {State: "asked", AskedAt: &asked}}}
|
|
|
|
all := pauseOf([]string{"g14", "ace"}, map[string]link.HolderState{"ace": {Paused: true}, "g14": {Paused: true}})
|
|
line := planLineWith(p, now, all, tierAtLeast)
|
|
if !strings.Contains(line, "waiting: the build seat is paused on ace, g14 (asked 12m0s ago)") || strings.Contains(line, "LATE") {
|
|
t.Errorf("a plan on a paused seat reads %q", line)
|
|
}
|
|
if st := planStatuses([]inventory.Plan{p}, now, all, nil)[0]; st.Late || !strings.Contains(st.Waiting, "paused on ace, g14") {
|
|
t.Errorf("status --json says %+v", st)
|
|
}
|
|
if _, late := openPlans([]inventory.Plan{p}, now, all, nil); late != 0 {
|
|
t.Errorf("a plan waiting on a paused seat counted late")
|
|
}
|
|
|
|
longAgo := now.Add(-25 * time.Hour)
|
|
forgotten := p
|
|
forgotten.Modules = map[string]*inventory.PlanModule{"a": {State: "asked", AskedAt: &longAgo}}
|
|
if line := planLineWith(forgotten, now, all, tierAtLeast); !strings.Contains(line, "PAUSED OVER A DAY") || strings.Contains(line, "LATE") {
|
|
t.Errorf("a plan paused over a day reads %q", line)
|
|
}
|
|
|
|
some := pauseOf([]string{"g14", "ace"}, map[string]link.HolderState{"ace": {Paused: true}})
|
|
if some.All || !reflect.DeepEqual(some.Nodes, []string{"ace"}) {
|
|
t.Fatalf("%+v", some)
|
|
}
|
|
if line := planLineWith(p, now, some, tierAtLeast); !strings.Contains(line, "LATE") {
|
|
t.Errorf("a seat paused on one holder of two made the plan not late: %q", line)
|
|
}
|
|
if st := planStatuses([]inventory.Plan{p}, now, some, nil)[0]; !st.Late {
|
|
t.Errorf("status --json: %+v", st)
|
|
}
|
|
// A plan rolling out, or with nothing asked, is not waiting on the seat.
|
|
rolling := p
|
|
rolling.State = inventory.PlanRolling
|
|
if _, paused := pausedWaiting(rolling, all, now); paused {
|
|
t.Error("a rolling plan reads as waiting on the build seat")
|
|
}
|
|
}
|
|
|
|
// The queue's verbs are the controller seat's, each to the command it names.
|
|
func TestTheQueueVerbsRunTheirCommands(t *testing.T) {
|
|
for _, c := range []struct {
|
|
verb string
|
|
args map[string]any
|
|
want []string
|
|
}{
|
|
{"queue", nil, []string{"queue"}},
|
|
{"cancel", map[string]any{"id": "build-1"}, []string{"cancel", "build-1"}},
|
|
{"clear", nil, []string{"clear"}},
|
|
{"clear", map[string]any{"dead": "true"}, []string{"clear", "--dead"}},
|
|
{"rebuild", map[string]any{"what": "gitea"}, []string{"rebuild", "gitea"}},
|
|
{"replay", map[string]any{"id": "build-1"}, []string{"replay", "build-1"}},
|
|
{"replay", map[string]any{"id": "build-1", "register": "true", "older": "true"}, []string{"replay", "build-1", "--register", "--older"}},
|
|
{"kill", map[string]any{"id": "build-1"}, []string{"kill", "build-1"}},
|
|
{"pause", nil, []string{"pause"}},
|
|
{"resume", map[string]any{"node": "ace"}, []string{"resume", "ace"}},
|
|
{"plans", map[string]any{"retry": "plan-1"}, []string{"plans", "retry", "plan-1"}},
|
|
} {
|
|
got, err := argvFor(c.verb, c.args)
|
|
if err != nil || !reflect.DeepEqual(got, c.want) {
|
|
t.Errorf("%s %v: %v %v, want %v", c.verb, c.args, got, err, c.want)
|
|
}
|
|
}
|
|
if _, err := argvFor("cancel", nil); err == nil {
|
|
t.Error("cancel without an id was taken")
|
|
}
|
|
declared := map[string]bool{}
|
|
for _, v := range catalogue.ControllerVerbs {
|
|
declared[v.Name] = true
|
|
}
|
|
for _, v := range []string{"queue", "cancel", "clear", "rebuild", "replay", "kill", "pause", "resume"} {
|
|
if !declared[v] {
|
|
t.Errorf("%s is not a verb of the controller seat", v)
|
|
}
|
|
}
|
|
}
|
|
|
|
// --- against a real bus -----------------------------------------------------------------------
|
|
|
|
// aBuildQueue is the build seat's queue, worker and cancelled set on a real server, the controller
|
|
// pointed at it, and a function that asks the seat one build.
|
|
func aBuildQueue(t *testing.T) (*broker.JetStream, func(id, repository string) link.BuildRequest) {
|
|
t.Helper()
|
|
url := testbus.URL(t)
|
|
t.Setenv(broker.NATSVar, url)
|
|
js, err := broker.Dial(url)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
t.Cleanup(js.Close)
|
|
seat := broker.DeclaredSeat{Name: link.TheBuildMachine, Accepts: []string{"build"},
|
|
Emits: []string{"started", "built", "log.*", "paused.*"}}
|
|
if err := broker.AssertMeshStreams(js); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
_ = js.Context().DeleteStream("SEAT_NODE_BUILD_AGENT")
|
|
if err := broker.RaiseSeats(js, []broker.DeclaredSeat{seat}, map[string]broker.Holder{
|
|
link.TheBuildMachine: {Node: "anchor", Module: "build-agent"}}); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := broker.RaiseCancelledSets(js, []broker.DeclaredSeat{seat}); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
t.Cleanup(func() {
|
|
_ = js.Context().DeleteStream("SEAT_NODE_BUILD_AGENT")
|
|
_ = js.Context().DeleteKeyValue(broker.CancelledSetName(link.TheBuildMachine))
|
|
_ = js.Context().PurgeStream(broker.EventsStream)
|
|
})
|
|
ask := func(id, repository string) link.BuildRequest {
|
|
r := link.BuildRequest{ID: id, Repository: repository, Held: map[string]string{"x/y": "secret-ish"}}
|
|
body, _ := json.Marshal(r)
|
|
if _, err := js.Context().Publish(link.BuildWork(), body); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
return r
|
|
}
|
|
return js, ask
|
|
}
|
|
|
|
// deadOne takes the oldest ask from the worker and hands it back as often as the worker allows.
|
|
func deadOne(t *testing.T, js *broker.JetStream) {
|
|
t.Helper()
|
|
worker, _ := broker.HolderConsumerFor("", "", broker.DeclaredSeat{Name: link.TheBuildMachine, Accepts: []string{"build"}})
|
|
sub, err := js.Context().PullSubscribe(worker.Filters[0], worker.Name, nats.Bind(worker.Stream, worker.Name), nats.ManualAck())
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
defer func() { _ = sub.Unsubscribe() }()
|
|
for i := 0; i < worker.MaxDeliver; i++ {
|
|
msgs, err := sub.Fetch(1, nats.MaxWait(3*time.Second))
|
|
if err != nil {
|
|
t.Fatalf("delivery %d: %v", i+1, err)
|
|
}
|
|
_ = msgs[0].Nak()
|
|
}
|
|
// One more pull, as a holder always has one waiting: the server finds the ask past its deliveries
|
|
// then, and stops counting it pending.
|
|
if msgs, _ := sub.Fetch(1, nats.MaxWait(time.Second)); len(msgs) > 0 {
|
|
t.Fatalf("an ask past its deliveries was delivered again")
|
|
}
|
|
}
|
|
|
|
// cancel drops a waiting ask from the bus and records it failed, cancelled by hand — and the plan
|
|
// that asked for it, matched by the id, fails with it; an ask the plan's records name no module for
|
|
// is still found.
|
|
func TestCancelDeletesTheAskAndFailsThePlanThatAskedIt(t *testing.T) {
|
|
js, ask := aBuildQueue(t)
|
|
open := aMesh(t)
|
|
ctx := t.Context()
|
|
twoTiers(t, open)
|
|
id := link.NewBuildID(time.Now())
|
|
ask(id, "https://forge.example/novox/a.git")
|
|
asked, _ := link.BuildAskedAt(id)
|
|
plan := inventory.Plan{ID: "plan-cancel", Repository: "novox/a", Commit: "c0ffee", Created: asked,
|
|
State: inventory.PlanBuilding, Tiers: [][]string{{"a"}, {"b"}},
|
|
Modules: map[string]*inventory.PlanModule{"a": {State: "asked", AskedAt: &asked, Build: id}}}
|
|
if err := open.inventory.SavePlan(ctx, &plan); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
|
|
q, err := link.ReadQueue(ctx, js, link.TheBuildMachine)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if len(q.Asks) != 1 || q.Asks[0].State != link.AskWaiting || q.Asks[0].ID != id {
|
|
t.Fatalf("the queue reads %+v", q)
|
|
}
|
|
if text := queueText(q, time.Now()); !strings.Contains(text, "1 waiting, 0 in flight, 0 dead") ||
|
|
strings.Contains(text, "secret-ish") {
|
|
t.Errorf("the queue says:\n%s", text)
|
|
}
|
|
|
|
if err := cancelCommand(ctx, []string{id}); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if q, _ = link.ReadQueue(ctx, js, link.TheBuildMachine); len(q.Asks) != 0 {
|
|
t.Fatalf("the ask is still queued: %+v", q.Asks)
|
|
}
|
|
if cancelled, err := link.IsCancelled(js.Conn(), link.TheBuildMachine, id); err != nil || !cancelled {
|
|
t.Errorf("the cancelled set does not hold it: %v %v", cancelled, err)
|
|
}
|
|
b, found, err := open.inventory.BuildByID(ctx, id)
|
|
if err != nil || !found || b.Failed != link.CancelledByHand {
|
|
t.Fatalf("the cancel is recorded as %+v (%v %v)", b, found, err)
|
|
}
|
|
p, err := open.inventory.PlanByID(ctx, plan.ID)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if p.State != inventory.PlanFailed || p.Modules["a"].State != "failed" || p.Modules["a"].Why != link.CancelledByHand {
|
|
t.Fatalf("the plan that asked is %s, a %+v", p.State, p.Modules["a"])
|
|
}
|
|
if err := cancelCommand(ctx, []string{id}); err == nil {
|
|
t.Error("an ask cancelled already was cancelled again")
|
|
}
|
|
}
|
|
|
|
// clear cancels every waiting ask and leaves the dead ones unless told, and never one in flight.
|
|
func TestClearCancelsTheWaitingAndTheDeadOnlyWhenTold(t *testing.T) {
|
|
js, ask := aBuildQueue(t)
|
|
open := aMesh(t)
|
|
ctx := t.Context()
|
|
start := time.Now()
|
|
dead := link.NewBuildID(start)
|
|
ask(dead, "https://forge.example/novox/dead.git")
|
|
deadOne(t, js)
|
|
var waiting []string
|
|
for i := 1; i <= 2; i++ {
|
|
id := link.NewBuildID(start.Add(time.Duration(i) * time.Millisecond))
|
|
waiting = append(waiting, id)
|
|
ask(id, fmt.Sprintf("https://forge.example/novox/w%d.git", i))
|
|
}
|
|
|
|
q, err := link.ReadQueue(ctx, js, link.TheBuildMachine)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if len(q.Of(link.AskDead)) != 1 || len(q.Of(link.AskWaiting)) != 2 || q.MaxDeliver != 5 {
|
|
t.Fatalf("the queue reads %+v", q)
|
|
}
|
|
said, err := clearQueue(ctx, js, open, link.TheBuildMachine, false)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if !strings.Contains(said, "2 waiting ask(s) cancelled") || !strings.Contains(said, "1 dead, left") {
|
|
t.Errorf("clear said:\n%s", said)
|
|
}
|
|
q, _ = link.ReadQueue(ctx, js, link.TheBuildMachine)
|
|
if len(q.Asks) != 1 || q.Asks[0].ID != dead || q.Asks[0].State != link.AskDead {
|
|
t.Fatalf("after clear the queue is %+v", q.Asks)
|
|
}
|
|
if _, err := clearQueue(ctx, js, open, link.TheBuildMachine, true); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if q, _ = link.ReadQueue(ctx, js, link.TheBuildMachine); len(q.Asks) != 0 {
|
|
t.Fatalf("clear --dead left %+v", q.Asks)
|
|
}
|
|
for _, id := range append(waiting, dead) {
|
|
if b, found, _ := open.inventory.BuildByID(ctx, id); !found || b.Failed != link.CancelledByHand {
|
|
t.Errorf("%s is recorded as %+v", id, b)
|
|
}
|
|
}
|
|
}
|
|
|
|
// An ask in flight is not cancelled: kill ends it where it runs.
|
|
func TestCancelRefusesAnAskInFlight(t *testing.T) {
|
|
js, ask := aBuildQueue(t)
|
|
open := aMesh(t)
|
|
ctx := t.Context()
|
|
id := link.NewBuildID(time.Now())
|
|
ask(id, "https://forge.example/novox/a.git")
|
|
worker, _ := broker.HolderConsumerFor("", "", broker.DeclaredSeat{Name: link.TheBuildMachine, Accepts: []string{"build"}})
|
|
sub, err := js.Context().PullSubscribe(worker.Filters[0], worker.Name, nats.Bind(worker.Stream, worker.Name), nats.ManualAck())
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
defer func() { _ = sub.Unsubscribe() }()
|
|
if _, err := sub.Fetch(1, nats.MaxWait(3*time.Second)); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
started, _ := json.Marshal(link.BuildStart{ID: id, On: "ace", At: time.Now().UTC().Format(time.RFC3339Nano)})
|
|
if _, err := js.Context().Publish(link.BuildStarted(), started); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
q, err := link.ReadQueue(ctx, js, link.TheBuildMachine)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
a, _ := q.Find(id)
|
|
if a.State != link.AskInFlight || a.On != "ace" {
|
|
t.Fatalf("the taken ask reads %+v", a)
|
|
}
|
|
_, err = cancelAsk(ctx, js, open, link.TheBuildMachine, a)
|
|
if err == nil || !strings.Contains(err.Error(), "kill "+id) {
|
|
t.Fatalf("an ask in flight was cancelled: %v", err)
|
|
}
|
|
if _, found, _ := open.inventory.BuildByID(ctx, id); found {
|
|
t.Error("a refused cancel recorded an outcome")
|
|
}
|
|
}
|
|
|
|
// kill finds the machine from the build's start and asks that machine's holder; pause asks the
|
|
// machine named. Each prints what the holder answered, and a refusal is the command's failure.
|
|
func TestKillAndPauseAskTheHolderOnTheMachine(t *testing.T) {
|
|
js, _ := aBuildQueue(t)
|
|
ctx := t.Context()
|
|
asked := map[string]string{}
|
|
handlers := map[string]link.ToolHandler{
|
|
"kill": func(_ context.Context, raw json.RawMessage) (any, error) {
|
|
var args struct{ ID string }
|
|
_ = json.Unmarshal(raw, &args)
|
|
asked["kill"] = args.ID
|
|
if args.ID != "build-running" {
|
|
return nil, fmt.Errorf("ace is not building %s", args.ID)
|
|
}
|
|
return map[string]any{"said": "killed " + args.ID}, nil
|
|
},
|
|
"pause": func(context.Context, json.RawMessage) (any, error) {
|
|
asked["pause"] = "ace"
|
|
return map[string]any{"said": "ace is paused"}, nil
|
|
},
|
|
}
|
|
stop, err := link.OverNATS{Conn: js.Conn()}.ServeNodeSeatTools(link.TheBuildMachine, "ace", handlers, nil)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
defer stop()
|
|
for _, id := range []string{"build-running", "build-other"} {
|
|
started, _ := json.Marshal(link.BuildStart{ID: id, On: "ace", At: time.Now().UTC().Format(time.RFC3339Nano)})
|
|
if _, err := js.Context().Publish(link.BuildStarted(), started); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
}
|
|
if err := killCommand(ctx, []string{"build-running"}); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if asked["kill"] != "build-running" {
|
|
t.Fatalf("the holder was asked %v", asked)
|
|
}
|
|
if err := killCommand(ctx, []string{"build-other"}); err == nil {
|
|
t.Error("the holder's refusal was not the command's")
|
|
}
|
|
if err := killCommand(ctx, []string{"build-never"}); err == nil || !strings.Contains(err.Error(), "cancel build-never") {
|
|
t.Errorf("a build nobody started: %v", err)
|
|
}
|
|
if err := pauseCommand(ctx, "pause", []string{"ace"}); err != nil || asked["pause"] != "ace" {
|
|
t.Fatalf("pause: %v %v", err, asked)
|
|
}
|
|
if err := pauseCommand(ctx, "resume", []string{"g14"}); err == nil {
|
|
t.Error("a machine nothing answers on was resumed")
|
|
}
|
|
}
|
|
|
|
// A plan that stopped at its first machine is retried: the module sent to that machine again, the
|
|
// send recorded as the first anew, and the plan goes on — unless a newer plan holds the module.
|
|
func TestAPlanStoppedAtItsFirstMachineIsRetried(t *testing.T) {
|
|
open := aMesh(t)
|
|
ctx := t.Context()
|
|
asked := asksRecorded(t)
|
|
twoTiers(t, open)
|
|
var sentTo [][]string
|
|
was := sendRollout
|
|
sendRollout = func(_ context.Context, _ *stores, names []string) ([]string, error) {
|
|
sentTo = append(sentTo, names)
|
|
return names, nil
|
|
}
|
|
t.Cleanup(func() { sendRollout = was })
|
|
|
|
// a records: a person's choice, since the default rolls out (novox/hq ADR 0236).
|
|
if err := open.inventory.SetUpgradeOf(ctx, "a", inventory.Upgrade{Why: "test"}); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
long := time.Now().UTC().Add(-2 * time.Hour)
|
|
stopped := inventory.Plan{ID: "plan-rollout", Repository: "novox/a", Branch: "main", Commit: "c0ffee",
|
|
Created: long, State: inventory.PlanFailed, Tiers: [][]string{{"a"}, {"b"}},
|
|
Note: "a stopped at its first machine in tier 0: laptop refused what it was sent",
|
|
Modules: map[string]*inventory.PlanModule{"a": {State: "built", BuiltAt: &long, Commit: "c0ffee",
|
|
First: []string{"laptop"}, FirstAt: &long, Why: "laptop refused what it was sent"}}}
|
|
if err := open.inventory.SavePlan(ctx, &stopped); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
|
|
// A newer plan holding a refuses it: sending the older build would put it back.
|
|
newer := inventory.Plan{ID: "plan-newer", Repository: "novox/other", Commit: "d00d", Created: long.Add(time.Hour),
|
|
State: inventory.PlanDone, Tiers: [][]string{{"a"}}, Modules: map[string]*inventory.PlanModule{}}
|
|
if err := open.inventory.SavePlan(ctx, &newer); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if _, err := retryPlan(ctx, open, stopped.ID); err == nil || !strings.Contains(err.Error(), "plan-newer") {
|
|
t.Fatalf("retried under a newer plan: %v", err)
|
|
}
|
|
newer.State = inventory.PlanSuperseded
|
|
if err := open.inventory.SavePlan(ctx, &newer); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
|
|
said, err := retryPlan(ctx, open, stopped.ID)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if len(sentTo) != 1 || !reflect.DeepEqual(sentTo[0], []string{"laptop"}) || !strings.Contains(said, "laptop") {
|
|
t.Fatalf("sent %v; said %q", sentTo, said)
|
|
}
|
|
p, err := open.inventory.PlanByID(ctx, stopped.ID)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
a := p.Modules["a"]
|
|
if p.State != inventory.PlanRolling || a.FirstAt == nil || !a.FirstAt.After(long) || a.Why != "" {
|
|
t.Fatalf("after retry the plan is %s, a %+v", p.State, a)
|
|
}
|
|
// And it goes on: a records (its policy sends nothing more), so the next tier is asked.
|
|
advancePlans(ctx, open)
|
|
if p, _ = open.inventory.PlanByID(ctx, stopped.ID); p.Tier != 1 || p.Modules["b"] == nil || p.Modules["b"].State != "asked" {
|
|
t.Fatalf("the retried plan did not go on: tier %d %s %+v", p.Tier, p.State, p.Modules["b"])
|
|
}
|
|
if len(*asked) != 1 {
|
|
t.Fatalf("asked %v", *asked)
|
|
}
|
|
}
|
|
|
|
// A plan module asked under an id is settled by that id's outcome alone: a replay or a rebuild beside
|
|
// the plan, asked later, never answers it (novox/hq ADR 0219).
|
|
func TestAPlanIsAnsweredOnlyByTheBuildItAskedFor(t *testing.T) {
|
|
open := aMesh(t)
|
|
ctx := t.Context()
|
|
asked := time.Now().UTC().Add(-time.Minute)
|
|
plan := inventory.Plan{ID: "plan-own", Repository: "novox/a", Commit: "c0ffee", Created: asked,
|
|
State: inventory.PlanBuilding, Tiers: [][]string{{"a"}},
|
|
Modules: map[string]*inventory.PlanModule{"a": {State: "asked", AskedAt: &asked, Build: "build-own"}}}
|
|
if err := open.inventory.SavePlan(ctx, &plan); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
planBuilt(ctx, open, "a", "0ldc0mm1t", "", time.Now().UTC(), "build-replay")
|
|
p, _ := open.inventory.PlanByID(ctx, plan.ID)
|
|
if p.Modules["a"].State != "asked" {
|
|
t.Fatalf("a replay asked after the plan settled it: %+v", p.Modules["a"])
|
|
}
|
|
// From the records too: a later build of the module recorded is not the plan's.
|
|
recorded := map[string][]inventory.Build{"a": {{ID: "build-replay", Commit: "0ldc0mm1t", Asked: time.Now(), At: time.Now()}}}
|
|
if settleFromRecords(&p, p.Tiers[0], recorded, nil) {
|
|
t.Fatalf("the records settled it with another build: %+v", p.Modules["a"])
|
|
}
|
|
planBuilt(ctx, open, "a", "c0ffee", "", asked, "build-own")
|
|
if p, _ = open.inventory.PlanByID(ctx, plan.ID); p.Modules["a"].State != "built" || p.Modules["a"].Commit != "c0ffee" {
|
|
t.Fatalf("its own build did not settle it: %+v", p.Modules["a"])
|
|
}
|
|
}
|
|
|
|
// rebuild of a build made at a commit asks what the module follows now, never the commit.
|
|
func TestARebuildOfACommitAsksWhatTheModuleFollows(t *testing.T) {
|
|
open := aMesh(t)
|
|
ctx := t.Context()
|
|
asked := asksRecorded(t)
|
|
twoTiers(t, open)
|
|
var refs []string
|
|
was := askABuild
|
|
askABuild = func(c context.Context, source buildSource, path, ref string) (string, error) {
|
|
refs = append(refs, ref)
|
|
return was(c, source, path, ref)
|
|
}
|
|
if err := open.inventory.RecordBuild(ctx, inventory.Build{ID: "build-at-commit", Repository: "novox/a",
|
|
Ref: "0123456789abcdef0123456789abcdef01234567", Module: "a", Commit: "0123456789abcdef0123456789abcdef01234567"}); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err := rebuildCommand(ctx, []string{"build-at-commit"}); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if len(refs) != 1 || refs[0] != "main" || len(*asked) != 1 {
|
|
t.Fatalf("asked %v at %v", *asked, refs)
|
|
}
|
|
}
|
|
|
|
// An ask the worker counts out that no machine said it started is cancelled only if no start comes
|
|
// while a holder looks: one that does withdraws the cancel, and kill is what ends it.
|
|
func TestCancelOfAnAskNobodySaidIsWithdrawnWhenItStarts(t *testing.T) {
|
|
js, ask := aBuildQueue(t)
|
|
open := aMesh(t)
|
|
ctx := t.Context()
|
|
was := holderLooks
|
|
holderLooks = 300 * time.Millisecond
|
|
t.Cleanup(func() { holderLooks = was })
|
|
id := link.NewBuildID(time.Now())
|
|
ask(id, "https://forge.example/novox/a.git")
|
|
worker, _ := broker.HolderConsumerFor("", "", broker.DeclaredSeat{Name: link.TheBuildMachine, Accepts: []string{"build"}})
|
|
sub, err := js.Context().PullSubscribe(worker.Filters[0], worker.Name, nats.Bind(worker.Stream, worker.Name), nats.ManualAck())
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
defer func() { _ = sub.Unsubscribe() }()
|
|
if _, err := sub.Fetch(1, nats.MaxWait(3*time.Second)); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
q, err := link.ReadQueue(ctx, js, link.TheBuildMachine)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
a, _ := q.Find(id)
|
|
if a.State != link.AskInFlight || a.On != "" {
|
|
t.Fatalf("the taken ask reads %+v", a)
|
|
}
|
|
// The holder that took it says it started, while the cancel waits.
|
|
go func() {
|
|
time.Sleep(100 * time.Millisecond)
|
|
started, _ := json.Marshal(link.BuildStart{ID: id, On: "ace", At: time.Now().UTC().Format(time.RFC3339Nano)})
|
|
_, _ = js.Context().Publish(link.BuildStarted(), started)
|
|
}()
|
|
if _, err := cancelAsk(ctx, js, open, link.TheBuildMachine, a); err == nil || !strings.Contains(err.Error(), "started on ace") {
|
|
t.Fatalf("a build that started was cancelled: %v", err)
|
|
}
|
|
if cancelled, _ := link.IsCancelled(js.Conn(), link.TheBuildMachine, id); cancelled {
|
|
t.Error("the withdrawn cancel is still marked")
|
|
}
|
|
if _, found, _ := open.inventory.BuildByID(ctx, id); found {
|
|
t.Error("a withdrawn cancel recorded an outcome")
|
|
}
|
|
}
|