A build says what it does on the bus, as it happens (novox/hq ADR 0157)
The build-machine seat emits `started` and `log.<build id>` beside `built`. Every line the builder speaks — each step, each command with its duration, and on failure the command's own output — goes to stderr as before and onto the bus under the build's id, one subject per build, kept a week in EVENTS with every other event. `builds --log <id>` reads it back from the stream with a consumer that is gone when the reading is done, on the command line and as the controller's seat verb; `builds` lists each build's id and `build` says the id it asked with. Lines are core publishes with a sequence number, so a build is not slowed by an ack per line and a gap is visible; `started` and `built` are awaited into the stream. The seat protocol widens additively at the controller's next start; the holder's grant follows on the broker node's next composition.
This commit is contained in:
+22
-10
@@ -148,20 +148,34 @@ func answer(ctx context.Context, publisher builder.Publisher, on, workspace stri
|
||||
// it either finishes or fails is indistinguishable from one that never arrived — which cost a long
|
||||
// diagnosis against a running mesh, chasing "the handler never fired" when the truth was only that
|
||||
// the handler said nothing until the end.
|
||||
fmt.Fprintf(os.Stderr, "a build request arrived for %s\n", request.Repository)
|
||||
fmt.Fprintf(os.Stderr, "a build request arrived for %s (%s)\n", request.Repository, request.ID)
|
||||
|
||||
// **Everything a build says goes two ways**: to stderr, as always, and onto the bus as the
|
||||
// role's own events under the build's id (novox/hq ADR 0157) — so whoever asked, and anybody
|
||||
// watching, reads the same lines this container's log holds, live, and after the fact from the
|
||||
// stream. Said first, before anything runs, so a build that hangs is one that visibly started.
|
||||
say := func(step, message string) {
|
||||
fmt.Fprintf(os.Stderr, " [%s] %s\n", step, message)
|
||||
work.Say(step, message)
|
||||
}
|
||||
builder.Said = say
|
||||
defer func() { builder.Said = nil }()
|
||||
if err := work.Began(ctx); err != nil {
|
||||
fmt.Fprintf(os.Stderr, "cannot say a build started: %v\n", err)
|
||||
}
|
||||
|
||||
result := link.BuildResult{
|
||||
ID: request.ID, Repository: request.Repository, Path: request.Path,
|
||||
Ref: request.Ref, On: on,
|
||||
}
|
||||
fmt.Fprintf(os.Stderr, "building %s", request.Repository)
|
||||
what := "building " + request.Repository
|
||||
if request.Path != "" {
|
||||
fmt.Fprintf(os.Stderr, " at %s", request.Path)
|
||||
what += " at " + request.Path
|
||||
}
|
||||
if request.Ref != "" {
|
||||
fmt.Fprintf(os.Stderr, " at %s", request.Ref)
|
||||
what += " on " + request.Ref
|
||||
}
|
||||
fmt.Fprintln(os.Stderr)
|
||||
say("build", what)
|
||||
|
||||
npmrc, err := packagesFrom()
|
||||
var built builder.Result
|
||||
@@ -171,15 +185,13 @@ func answer(ctx context.Context, publisher builder.Publisher, on, workspace stri
|
||||
// after a clone that then fails at npm ci.
|
||||
built, err = builder.Build(ctx, builder.Command, publisher,
|
||||
request.Repository, request.Path, request.Ref, workspace, request.Held, npmrc,
|
||||
forgeFrom(), func(step, message string) {
|
||||
fmt.Fprintf(os.Stderr, " [%s] %s\n", step, message)
|
||||
}, request.Seats)
|
||||
forgeFrom(), say, request.Seats)
|
||||
}
|
||||
if err != nil {
|
||||
// A failure is a result. A build that fails and says nothing is indistinguishable from a
|
||||
// builder that is not running, and those want completely different responses.
|
||||
result.Failed = err.Error()
|
||||
fmt.Fprintf(os.Stderr, " failed: %v\n", err)
|
||||
say("failed", err.Error())
|
||||
} else {
|
||||
manifest, marshalErr := json.Marshal(built.Manifest)
|
||||
if marshalErr != nil {
|
||||
@@ -196,7 +208,7 @@ func answer(ctx context.Context, publisher builder.Publisher, on, workspace stri
|
||||
for _, r := range built.Read {
|
||||
result.Read = append(result.Read, link.ReadRepository{Repository: r.Repository, Ref: r.Ref})
|
||||
}
|
||||
fmt.Fprintf(os.Stderr, " built %s from %s\n", built.Manifest.Module, short(built.Commit))
|
||||
say("built", built.Manifest.Module+" from "+short(built.Commit))
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -10,6 +10,8 @@ import (
|
||||
"strings"
|
||||
"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"
|
||||
@@ -175,10 +177,14 @@ func buildFrom(result link.BuildResult) inventory.Build {
|
||||
func buildsCommand(ctx context.Context, args []string) error {
|
||||
set := flag.NewFlagSet("builds", flag.ContinueOnError)
|
||||
limit := set.Int("n", 20, "how many to show")
|
||||
logOf := set.String("log", "", "a build's id: print what the build machine said, line by line")
|
||||
positionals, err := parseAround(set, args)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if *logOf != "" {
|
||||
return buildLog(ctx, *logOf)
|
||||
}
|
||||
module := ""
|
||||
if len(positionals) == 1 {
|
||||
module = positionals[0]
|
||||
@@ -219,8 +225,8 @@ func buildsCommand(ctx context.Context, args []string) error {
|
||||
if !b.Worked() {
|
||||
outcome = "failed"
|
||||
}
|
||||
fmt.Printf("%-18s %-14s %-10s %s\n",
|
||||
what, outcome, b.On, b.At.Local().Format("2006-01-02 15:04"))
|
||||
fmt.Printf("%-18s %-14s %-10s %s %s\n",
|
||||
what, outcome, b.On, b.At.Local().Format("2006-01-02 15:04"), b.ID)
|
||||
fmt.Printf(" %s", b.Repository)
|
||||
if b.Ref != "" {
|
||||
fmt.Printf(" at %s", b.Ref)
|
||||
@@ -414,6 +420,8 @@ func buildOne(ctx context.Context, source buildSource, path, ref string, wait ti
|
||||
if source.Seat != "" {
|
||||
fmt.Printf(" (%s)", repository)
|
||||
}
|
||||
// The id is how a person follows this build while it runs: `builds --log <id>`.
|
||||
fmt.Printf(" as %s", request.ID)
|
||||
if path != "" {
|
||||
fmt.Printf(" at %s", path)
|
||||
}
|
||||
@@ -625,3 +633,57 @@ func askOver(_ *link.Server) (link.Builders, error) {
|
||||
}
|
||||
return link.BuildsOverNATS(address)
|
||||
}
|
||||
|
||||
// buildLog prints everything a build machine said about one build, read back from the bus.
|
||||
//
|
||||
// **From the stream, not from a record** (novox/hq ADR 0157). A build's lines are the role's own
|
||||
// events under the build's id, retained with every other event; the mesh keeps no second copy. Read
|
||||
// with a consumer of its own that is gone when this returns, so nothing accumulates in the server
|
||||
// for the reading, and filtered by subject, so one build's lines are all that travel.
|
||||
func buildLog(ctx context.Context, id string) error {
|
||||
address, err := broker.BusAddress()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
js, err := broker.Dial(address)
|
||||
if err != nil {
|
||||
return fmt.Errorf("cannot reach the bus to read a build's log: %w", err)
|
||||
}
|
||||
defer js.Close()
|
||||
|
||||
sub, err := js.Context().PullSubscribe(link.BuildLog(id), "",
|
||||
nats.BindStream(broker.EventsStream), nats.DeliverAll(), nats.AckNone())
|
||||
if err != nil {
|
||||
return fmt.Errorf("cannot read %s from the bus: %w", link.BuildLog(id), err)
|
||||
}
|
||||
defer func() { _ = sub.Unsubscribe() }()
|
||||
|
||||
printed := 0
|
||||
for {
|
||||
batch, err := sub.Fetch(200, nats.MaxWait(2*time.Second))
|
||||
if err != nil && !errors.Is(err, nats.ErrTimeout) && !errors.Is(err, context.DeadlineExceeded) {
|
||||
return fmt.Errorf("reading a build's log: %w", err)
|
||||
}
|
||||
for _, msg := range batch {
|
||||
var line link.BuildLine
|
||||
if err := json.Unmarshal(msg.Data, &line); err != nil {
|
||||
fmt.Printf(" ? %s\n", string(msg.Data))
|
||||
continue
|
||||
}
|
||||
at := line.At
|
||||
if t, err := time.Parse(time.RFC3339Nano, line.At); err == nil {
|
||||
at = t.Local().Format("15:04:05")
|
||||
}
|
||||
fmt.Printf("%s %4d [%s] %s\n", at, line.Seq, line.Step, line.Message)
|
||||
printed++
|
||||
}
|
||||
if len(batch) < 200 {
|
||||
break
|
||||
}
|
||||
}
|
||||
if printed == 0 {
|
||||
fmt.Printf("nothing on the bus for build %s: no build by that id in the last week, or a build "+
|
||||
"machine older than this that said nothing while building\n", id)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -62,6 +62,9 @@ func argvFor(verb string, args map[string]any) ([]string, error) {
|
||||
case "seats":
|
||||
return []string{"seats", "--json"}, nil
|
||||
case "builds":
|
||||
if id := str("log"); id != "" {
|
||||
return []string{"builds", "--log", id}, nil
|
||||
}
|
||||
if m := str("module"); m != "" {
|
||||
return []string{"builds", m}, nil
|
||||
}
|
||||
|
||||
@@ -30,6 +30,18 @@ func TestEveryDeclaredVerbHasACommandLine(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// `builds` given a build's id reads that build's log from the bus rather than listing builds
|
||||
// (novox/hq ADR 0157).
|
||||
func TestBuildsWithAnIdReadsThatBuildsLog(t *testing.T) {
|
||||
argv, err := argvFor("builds", map[string]any{"log": "build-17"})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if strings.Join(argv, " ") != "builds --log build-17" {
|
||||
t.Fatalf("builds with a log id became %q", strings.Join(argv, " "))
|
||||
}
|
||||
}
|
||||
|
||||
// A required argument missing is refused in the verb's own words, before anything runs.
|
||||
func TestAVerbMissingWhatItNeedsIsRefused(t *testing.T) {
|
||||
if _, err := argvFor("node", map[string]any{}); err == nil || !strings.Contains(err.Error(), `node needs "node"`) {
|
||||
|
||||
@@ -71,6 +71,18 @@ func TestHoldingASeatIsTheMirrorOfUsingIt(t *testing.T) {
|
||||
hasNot(t, perms.Publish, "mesh.seat.telegram-sender.accept.send")
|
||||
}
|
||||
|
||||
// A build machine may say everything about a build as it happens (novox/hq ADR 0157): that it
|
||||
// started, and every line under the build's own id — the seat's `log.*` becomes a publish over
|
||||
// one token, so a reader follows one build by subject and the holder can name no other subject.
|
||||
func TestTheBuildMachineMaySayWhatItDoesUnderTheBuildsId(t *testing.T) {
|
||||
seat := Seat{Name: "mesh-build-machine", Accepts: []string{"build"}, Emits: []string{"started", "built", "log.*"}}
|
||||
perms, _ := PermissionsFor(Principal{Kind: KindModule, Node: "anchor", Module: "builder",
|
||||
Holds: []Seat{seat}, PasswordHash: "x"})
|
||||
has(t, perms.Publish, "mesh.seat.mesh-build-machine.event.started")
|
||||
has(t, perms.Publish, "mesh.seat.mesh-build-machine.event.log.*")
|
||||
hasNot(t, perms.Publish, "mesh.seat.mesh-build-machine.event.>")
|
||||
}
|
||||
|
||||
// Without an ack permission a durable consumer never really consumes: every message it receives is
|
||||
// redelivered forever, refused by the permission list it already has (design 25 §4).
|
||||
func TestAModuleMayAckItsOwnDeliveriesAndNoOthers(t *testing.T) {
|
||||
|
||||
@@ -79,7 +79,7 @@ func MeshStreams() []Stream {
|
||||
"n-1 by construction (issue 107)",
|
||||
},
|
||||
{
|
||||
Name: "EVENTS",
|
||||
Name: EventsStream,
|
||||
// A seat's own events ride here too: they are 1:many like any event, and the
|
||||
// `event` token keeps them clear of both the seat's work queue (`accept`) and its
|
||||
// tools (`tool`), which must not be persisted.
|
||||
@@ -220,6 +220,9 @@ func seatEventSubject(seat, verb string) string {
|
||||
return "mesh.seat." + seat + ".event." + verb
|
||||
}
|
||||
|
||||
// EventsStream holds every module's and every role's events, a build's log among them.
|
||||
const EventsStream = "EVENTS"
|
||||
|
||||
// MeshConsumers is what the controller consumes, in the order a person reads it.
|
||||
//
|
||||
// **Unlimited redelivery on CONTROL, deliberately.** The store window's bound is the controller's,
|
||||
|
||||
@@ -719,6 +719,21 @@ func short(commit string) string {
|
||||
return commit
|
||||
}
|
||||
|
||||
// Said is where the lines Command speaks go, beside the build's own Log: what runs, how long it
|
||||
// took, and that it failed. Nil prints them to stderr, as a build machine with nobody listening
|
||||
// should. The machine sets it per build so every line reaches the bus too (novox/hq ADR 0157) —
|
||||
// the step is "run", and the message is the line as it has always been printed.
|
||||
var Said Log
|
||||
|
||||
func tell(step, format string, args ...any) {
|
||||
message := fmt.Sprintf(format, args...)
|
||||
if Said == nil {
|
||||
fmt.Fprintf(os.Stderr, " %s\n", message)
|
||||
return
|
||||
}
|
||||
Said(step, message)
|
||||
}
|
||||
|
||||
// Command is a Runner that actually runs things.
|
||||
func Command(ctx context.Context, dir, name string, args ...string) (string, error) {
|
||||
// **Every command is echoed before it runs**, with where. On a build that hangs, the last line
|
||||
@@ -726,16 +741,23 @@ func Command(ctx context.Context, dir, name string, args ...string) (string, err
|
||||
// nothing" and "git clone is waiting on a network that will not answer". Silent on success is
|
||||
// what made an empty workspace unreadable.
|
||||
started := timeNow()
|
||||
fmt.Fprintf(os.Stderr, " $ (%s) %s %s\n", short(filepath.Base(dir)), name, strings.Join(args, " "))
|
||||
tell("run", "$ (%s) %s %s", short(filepath.Base(dir)), name, strings.Join(args, " "))
|
||||
cmd := exec.CommandContext(ctx, name, args...)
|
||||
cmd.Dir = dir
|
||||
out, err := cmd.CombinedOutput()
|
||||
if err != nil {
|
||||
fmt.Fprintf(os.Stderr, " ! %s %s failed after %s\n", name, args[0], since(started))
|
||||
tell("run", "! %s %s failed after %s", name, args[0], since(started))
|
||||
// The command's own output is part of what a reader needs — the compiler's error, the
|
||||
// clone's refusal — and a line per output line keeps it readable on the bus.
|
||||
for _, line := range strings.Split(strings.TrimSpace(string(out)), "\n") {
|
||||
if line != "" {
|
||||
tell("output", "%s", line)
|
||||
}
|
||||
}
|
||||
return string(out), fmt.Errorf("%s %s: %w\n%s",
|
||||
name, strings.Join(args, " "), err, strings.TrimSpace(string(out)))
|
||||
}
|
||||
fmt.Fprintf(os.Stderr, " ✓ %s %s (%s)\n", name, firstArg(args), since(started))
|
||||
tell("run", "✓ %s %s (%s)", name, firstArg(args), since(started))
|
||||
return string(out), nil
|
||||
}
|
||||
|
||||
|
||||
@@ -80,8 +80,11 @@ var defaultSeats = []Seat{
|
||||
// A build is work submitted to this role and its outcome is the role's own event (ADR 0129).
|
||||
// One publish reaches whoever asked, the controller that records it, and the catalogue that
|
||||
// places it in the graph — what the old bus's shared exchange did for free.
|
||||
// A build says what it does as it does it (novox/hq ADR 0157): `started` when work is taken,
|
||||
// `log.<build id>` for every line, `built` for the outcome. The log's tail token is the build's
|
||||
// id, so a reader follows one build by subject alone.
|
||||
{Name: "mesh-build-machine", Scope: ScopeMesh,
|
||||
Accepts: []string{"build"}, Emits: []string{"built"}, Decision: "novox/hq ADR 0121"},
|
||||
Accepts: []string{"build"}, Emits: []string{"started", "built", "log.*"}, Decision: "novox/hq ADR 0121"},
|
||||
{Name: "node-dns-resolver", Scope: ScopeNode, Decision: "novox/hq ADR 0121"},
|
||||
{Name: "node-intrusion-prevention", Scope: ScopeNode, Decision: "novox/hq ADR 0121"},
|
||||
{Name: "node-packet-filter", Scope: ScopeNode, Decision: "novox/hq ADR 0121"},
|
||||
|
||||
@@ -85,8 +85,12 @@ var ControllerVerbs = []Verb{
|
||||
Input: schema(nil, nil)},
|
||||
{Name: "seats", Description: "Every seat the mesh defines, what it delivers, and who holds it.",
|
||||
Input: schema(nil, nil)},
|
||||
{Name: "builds", Description: "What has been built lately and what came of it, for every module or for one.",
|
||||
Input: schema(map[string]string{"module": "one module's name; every module when absent"}, nil)},
|
||||
{Name: "builds", Description: "What has been built lately and what came of it, for every module or for one; " +
|
||||
"or, given a build's id, everything the build machine said while building it, line by line, from the bus.",
|
||||
Input: schema(map[string]string{
|
||||
"module": "one module's name; every module when absent",
|
||||
"log": "a build's id (as `builds` lists it): print what the build machine said, line by line",
|
||||
}, nil)},
|
||||
{Name: "plan", Description: "What one machine would run, and why: the declaration the mesh would send it.",
|
||||
Input: schema(map[string]string{"node": "the machine's name"}, []string{"node"})},
|
||||
{Name: "assign", Description: "Put a module on a machine. Refused with the mesh's own words when it cannot resolve there.",
|
||||
|
||||
@@ -27,6 +27,40 @@ const TheBuildMachine = "mesh-build-machine"
|
||||
func BuildWork() string { return "mesh.seat." + TheBuildMachine + ".accept.build" }
|
||||
func BuildOutcome() string { return "mesh.seat." + TheBuildMachine + ".event.built" }
|
||||
|
||||
// BuildStarted is where a build machine says it has taken a build, and BuildLog is where it says
|
||||
// what it is doing, one line per message, under the build's own id (novox/hq ADR 0157).
|
||||
//
|
||||
// **The whole build is on the bus as it happens.** The outcome alone told a person that a build
|
||||
// failed and its first line why; everything between — which command, how long, where it hung —
|
||||
// lived in one container's stderr on one machine. Every line is now an event of the role, retained
|
||||
// with the rest of the mesh's events, so a reader follows a build live by subscribing its subject,
|
||||
// or reads it back afterwards from the stream, and a viewer is a subscriber and nothing more.
|
||||
func BuildStarted() string { return "mesh.seat." + TheBuildMachine + ".event.started" }
|
||||
func BuildLog(id string) string { return "mesh.seat." + TheBuildMachine + ".event.log." + id }
|
||||
|
||||
// BuildStart is what a build machine says the moment it takes a build.
|
||||
type BuildStart struct {
|
||||
ID string `json:"id"`
|
||||
Repository string `json:"repository"`
|
||||
Path string `json:"path,omitempty"`
|
||||
Ref string `json:"ref,omitempty"`
|
||||
On string `json:"on"`
|
||||
At string `json:"at"`
|
||||
}
|
||||
|
||||
// BuildLine is one thing a build said while building.
|
||||
type BuildLine struct {
|
||||
ID string `json:"id"`
|
||||
// Seq counts the lines of one build from 1, so a reader that joined late or read two copies
|
||||
// can order them and see a gap.
|
||||
Seq int `json:"seq"`
|
||||
At string `json:"at"`
|
||||
// Step is which part of the build spoke — clone, context, image, run, failed — and Message is
|
||||
// what it said, as the builder's own log prints it.
|
||||
Step string `json:"step"`
|
||||
Message string `json:"message"`
|
||||
}
|
||||
|
||||
// KeyRoleBuilt is the build outcome under the role's name, on the bus the mesh runs on today.
|
||||
//
|
||||
// The same event as KeyModuleBuilt and published beside it, because a catalogue installed before this
|
||||
@@ -57,6 +91,13 @@ type BuildMachine interface {
|
||||
|
||||
// Build is one request a machine has been handed.
|
||||
type Build interface {
|
||||
// Began says the build has been taken and is under way, before anything runs.
|
||||
Began(ctx context.Context) error
|
||||
|
||||
// Say publishes one line of what the build is doing. Never fails the build: a line the bus
|
||||
// did not take is a line lost, and the outcome still comes.
|
||||
Say(step, message string)
|
||||
|
||||
// Request is what to build.
|
||||
Request() BuildRequest
|
||||
|
||||
|
||||
@@ -169,6 +169,7 @@ type natsBuild struct {
|
||||
msg *nats.Msg
|
||||
on string
|
||||
js *broker.JetStream
|
||||
seq int
|
||||
}
|
||||
|
||||
func (b *natsBuild) Request() BuildRequest { return b.request }
|
||||
@@ -203,4 +204,37 @@ func (b *natsBuild) Announce(ctx context.Context, result BuildResult) error {
|
||||
|
||||
func (b *natsBuild) Done() error { return b.msg.Ack() }
|
||||
|
||||
// Began publishes that this machine has taken the build, into the stream like the outcome, so a
|
||||
// reader that asks afterwards sees when it started as well as how it ended.
|
||||
func (b *natsBuild) Began(ctx context.Context) error {
|
||||
body, err := json.Marshal(BuildStart{
|
||||
ID: b.request.ID, Repository: b.request.Repository, Path: b.request.Path,
|
||||
Ref: b.request.Ref, On: b.on, At: time.Now().UTC().Format(time.RFC3339Nano),
|
||||
})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
publish, cancel := context.WithTimeout(ctx, 30*time.Second)
|
||||
defer cancel()
|
||||
if _, err := b.js.Context().Publish(BuildStarted(), body, nats.Context(publish)); err != nil {
|
||||
return fmt.Errorf("cannot say a build started: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// Say publishes one line under the build's id. Core publish, unawaited: the stream that holds the
|
||||
// role's events captures it on its way through, and a build must not slow to the pace of an ack
|
||||
// per line. A line the bus did not take is counted anyway, so the gap is visible to a reader.
|
||||
func (b *natsBuild) Say(step, message string) {
|
||||
b.seq++
|
||||
body, err := json.Marshal(BuildLine{
|
||||
ID: b.request.ID, Seq: b.seq, At: time.Now().UTC().Format(time.RFC3339Nano),
|
||||
Step: step, Message: message,
|
||||
})
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
_ = b.js.Conn().Publish(BuildLog(b.request.ID), body)
|
||||
}
|
||||
|
||||
func (b *natsBuild) Hold(after time.Duration) error { return b.msg.NakWithDelay(after) }
|
||||
|
||||
@@ -79,6 +79,14 @@ func TestNatsABuildIsTakenAndItsOutcomeReachesEverybody(t *testing.T) {
|
||||
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()
|
||||
@@ -86,6 +94,9 @@ func TestNatsABuildIsTakenAndItsOutcomeReachesEverybody(t *testing.T) {
|
||||
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"}`),
|
||||
@@ -104,6 +115,27 @@ func TestNatsABuildIsTakenAndItsOutcomeReachesEverybody(t *testing.T) {
|
||||
}
|
||||
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)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user