From 17f7cb0d9c5f5d54a70894663c07b68cc5b3f13c Mon Sep 17 00:00:00 2001 From: jochen Date: Thu, 1 Oct 2026 00:43:07 +0200 Subject: [PATCH] A build says what it does on the bus, as it happens (novox/hq ADR 0157) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The build-machine seat emits `started` and `log.` 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 ` 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. --- cmd/mesh-builder/main.go | 32 +++++++++---- cmd/mesh-controller/build.go | 66 ++++++++++++++++++++++++++- cmd/mesh-controller/seatverbs.go | 3 ++ cmd/mesh-controller/seatverbs_test.go | 12 +++++ internal/broker/nats_test.go | 12 +++++ internal/broker/streams.go | 5 +- internal/builder/builder.go | 28 ++++++++++-- internal/catalogue/seats.go | 5 +- internal/catalogue/verbs.go | 8 +++- internal/link/builds.go | 41 +++++++++++++++++ internal/link/builds_nats.go | 34 ++++++++++++++ internal/link/builds_nats_test.go | 32 +++++++++++++ 12 files changed, 259 insertions(+), 19 deletions(-) diff --git a/cmd/mesh-builder/main.go b/cmd/mesh-builder/main.go index c49ceca..76b5855 100644 --- a/cmd/mesh-builder/main.go +++ b/cmd/mesh-builder/main.go @@ -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)) } } diff --git a/cmd/mesh-controller/build.go b/cmd/mesh-controller/build.go index 72b6d99..3f5c74e 100644 --- a/cmd/mesh-controller/build.go +++ b/cmd/mesh-controller/build.go @@ -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 `. + 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 +} diff --git a/cmd/mesh-controller/seatverbs.go b/cmd/mesh-controller/seatverbs.go index 1182dc6..69b06d2 100644 --- a/cmd/mesh-controller/seatverbs.go +++ b/cmd/mesh-controller/seatverbs.go @@ -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 } diff --git a/cmd/mesh-controller/seatverbs_test.go b/cmd/mesh-controller/seatverbs_test.go index ef2dc10..8343040 100644 --- a/cmd/mesh-controller/seatverbs_test.go +++ b/cmd/mesh-controller/seatverbs_test.go @@ -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"`) { diff --git a/internal/broker/nats_test.go b/internal/broker/nats_test.go index 46703ca..9456370 100644 --- a/internal/broker/nats_test.go +++ b/internal/broker/nats_test.go @@ -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) { diff --git a/internal/broker/streams.go b/internal/broker/streams.go index 1be7620..8b588d0 100644 --- a/internal/broker/streams.go +++ b/internal/broker/streams.go @@ -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, diff --git a/internal/builder/builder.go b/internal/builder/builder.go index f247890..43ed91f 100644 --- a/internal/builder/builder.go +++ b/internal/builder/builder.go @@ -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 } diff --git a/internal/catalogue/seats.go b/internal/catalogue/seats.go index 72ee46d..39de060 100644 --- a/internal/catalogue/seats.go +++ b/internal/catalogue/seats.go @@ -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.` 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"}, diff --git a/internal/catalogue/verbs.go b/internal/catalogue/verbs.go index ef91cc2..725d9c5 100644 --- a/internal/catalogue/verbs.go +++ b/internal/catalogue/verbs.go @@ -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.", diff --git a/internal/link/builds.go b/internal/link/builds.go index 8963bc2..1396e3a 100644 --- a/internal/link/builds.go +++ b/internal/link/builds.go @@ -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 diff --git a/internal/link/builds_nats.go b/internal/link/builds_nats.go index 039d340..4e2b63d 100644 --- a/internal/link/builds_nats.go +++ b/internal/link/builds_nats.go @@ -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) } diff --git a/internal/link/builds_nats_test.go b/internal/link/builds_nats_test.go index 2672258..a87d6a7 100644 --- a/internal/link/builds_nats_test.go +++ b/internal/link/builds_nats_test.go @@ -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) } -- 2.54.0