On the control node the store-bound packages of the controller's suite ran five to seven times slower than on any other holder, and every controller check that landed there ran past the suite's thirty minutes: a throwaway store flushing to a disk the mesh's own store, bus and forge keep busy. A store that lives for one check needs no crash safety. A holder recreated mid-check left the check's store and bus running, and the redelivery went to another machine, so nothing removed them: eleven pairs across four machines. A starting holder has taken nothing, so every container labelled with an ask of the seat is an earlier holder's.
403 lines
14 KiB
Go
403 lines
14 KiB
Go
package main
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"os"
|
|
"os/exec"
|
|
"path/filepath"
|
|
"strings"
|
|
"sync"
|
|
"time"
|
|
|
|
"github.com/nats-io/nats.go/micro"
|
|
|
|
"github.com/novox/mesh-controller/internal/builder"
|
|
"github.com/novox/mesh-controller/internal/catalogue"
|
|
"github.com/novox/mesh-controller/internal/link"
|
|
)
|
|
|
|
// What a holder of the build seat answers for, on its own machine (novox/hq ADR 0219).
|
|
//
|
|
// **The queue is the controller's; the build running here is this machine's.** The controller can
|
|
// see and change what waits in the seat's queue, but an ask a holder already took is a process tree
|
|
// on this machine and containers in this machine's runtime, and only this machine can end them. So
|
|
// the holder serves four verbs on the seat's subjects for this machine: what it is building, kill
|
|
// it, pause, resume.
|
|
//
|
|
// **One build at a time** (ADR 0190), which is also what builder.Said assumes — a package global
|
|
// set per build — so "the build running here" is one or none, and kill names it by id so a call
|
|
// that arrives as one build ends and the next begins cannot end the wrong one.
|
|
|
|
// holder is this machine's state as a holder of the build seat.
|
|
type holder struct {
|
|
on, seat string
|
|
// workspace is where the paused flag is kept, so a holder restarted while paused stays paused
|
|
// rather than silently taking work again.
|
|
workspace string
|
|
// say publishes this machine's state: whether it takes work (link.HolderState).
|
|
say func(link.HolderState) error
|
|
// remove runs what a kill needs outside the build: the containers left behind.
|
|
remove builder.Runner
|
|
|
|
mu sync.Mutex
|
|
paused bool
|
|
running *running
|
|
}
|
|
|
|
// running is the build this machine is doing.
|
|
type running struct {
|
|
request link.BuildRequest
|
|
step string
|
|
started time.Time
|
|
cancel context.CancelFunc
|
|
killed bool
|
|
// returned is the build's own work having ended, before a kill or not; outcome is what was then
|
|
// announced, and announced whether it went out.
|
|
returned bool
|
|
outcome string
|
|
announced bool
|
|
done chan struct{}
|
|
}
|
|
|
|
// pausedFile is where the flag lives in the workspace.
|
|
func pausedFile(workspace string) string { return filepath.Join(workspace, ".mesh-builder-paused") }
|
|
|
|
// newHolder reads the paused flag the workspace keeps.
|
|
func newHolder(on, seat, workspace string, say func(link.HolderState) error) *holder {
|
|
h := &holder{on: on, seat: seat, workspace: workspace, say: say, remove: plainRun}
|
|
if _, err := os.Stat(pausedFile(workspace)); err == nil {
|
|
h.paused = true
|
|
}
|
|
return h
|
|
}
|
|
|
|
// Paused is asked by the taking loop before every fetch.
|
|
func (h *holder) Paused() bool {
|
|
h.mu.Lock()
|
|
defer h.mu.Unlock()
|
|
return h.paused
|
|
}
|
|
|
|
// setPaused records the flag in the workspace first and then in memory, so what this holder says
|
|
// it is and what it would be after a restart never differ.
|
|
func (h *holder) setPaused(paused bool) error {
|
|
path := pausedFile(h.workspace)
|
|
if paused {
|
|
if err := os.MkdirAll(h.workspace, 0o755); err != nil {
|
|
return err
|
|
}
|
|
if err := os.WriteFile(path, []byte(time.Now().UTC().Format(time.RFC3339)+"\n"), 0o644); err != nil {
|
|
return fmt.Errorf("cannot keep the paused flag in %s: %w", path, err)
|
|
}
|
|
} else if err := os.Remove(path); err != nil && !errors.Is(err, os.ErrNotExist) {
|
|
return fmt.Errorf("cannot remove the paused flag %s: %w", path, err)
|
|
}
|
|
h.mu.Lock()
|
|
h.paused = paused
|
|
h.mu.Unlock()
|
|
h.announce()
|
|
return nil
|
|
}
|
|
|
|
// announce says whether this machine takes work. Never fatal: the flag is kept either way, and the
|
|
// controller reading an older state is a plan read as late rather than a build lost.
|
|
func (h *holder) announce() {
|
|
if h.say == nil {
|
|
return
|
|
}
|
|
if err := h.say(link.HolderState{On: h.on, Paused: h.Paused(), At: time.Now().UTC().Format(time.RFC3339Nano)}); err != nil {
|
|
fmt.Fprintf(os.Stderr, "cannot say whether this machine takes builds: %v\n", err)
|
|
}
|
|
}
|
|
|
|
// stateWhilePaused says this machine's state at start, and again every while it stays paused, so
|
|
// a pause outlives the events stream's retention.
|
|
func (h *holder) stateWhilePaused(ctx context.Context, every time.Duration) {
|
|
h.announce()
|
|
tick := time.NewTicker(every)
|
|
defer tick.Stop()
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
return
|
|
case <-tick.C:
|
|
if h.Paused() {
|
|
h.announce()
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
// begin records the build this machine took, with the cancel that ends it.
|
|
func (h *holder) begin(request link.BuildRequest, cancel context.CancelFunc) *running {
|
|
r := &running{request: request, started: time.Now(), cancel: cancel, done: make(chan struct{})}
|
|
h.mu.Lock()
|
|
h.running = r
|
|
h.mu.Unlock()
|
|
return r
|
|
}
|
|
|
|
// end clears it, and lets a kill waiting on it know it is over.
|
|
func (h *holder) end(r *running) {
|
|
h.mu.Lock()
|
|
if h.running == r {
|
|
h.running = nil
|
|
}
|
|
h.mu.Unlock()
|
|
close(r.done)
|
|
}
|
|
|
|
// stepped records the step a build is at, for `current`.
|
|
func (h *holder) stepped(r *running, step string) {
|
|
h.mu.Lock()
|
|
r.step = step
|
|
h.mu.Unlock()
|
|
}
|
|
|
|
// currentBuild is what `current` answers.
|
|
type currentBuild struct {
|
|
On string `json:"on"`
|
|
Paused bool `json:"paused"`
|
|
Running *struct {
|
|
ID string `json:"id"`
|
|
Repository string `json:"repository"`
|
|
Path string `json:"path,omitempty"`
|
|
Ref string `json:"ref,omitempty"`
|
|
Step string `json:"step,omitempty"`
|
|
Started string `json:"started"`
|
|
Elapsed string `json:"elapsed"`
|
|
} `json:"running,omitempty"`
|
|
Said string `json:"said"`
|
|
}
|
|
|
|
func (h *holder) current() currentBuild {
|
|
h.mu.Lock()
|
|
defer h.mu.Unlock()
|
|
out := currentBuild{On: h.on, Paused: h.paused}
|
|
taking := "taking builds"
|
|
if h.paused {
|
|
taking = "paused, taking no new build"
|
|
}
|
|
if h.running == nil {
|
|
out.Said = fmt.Sprintf("%s is building nothing; %s", h.on, taking)
|
|
return out
|
|
}
|
|
r := h.running
|
|
elapsed := time.Since(r.started).Round(time.Second)
|
|
out.Running = &struct {
|
|
ID string `json:"id"`
|
|
Repository string `json:"repository"`
|
|
Path string `json:"path,omitempty"`
|
|
Ref string `json:"ref,omitempty"`
|
|
Step string `json:"step,omitempty"`
|
|
Started string `json:"started"`
|
|
Elapsed string `json:"elapsed"`
|
|
}{r.request.ID, r.request.Repository, r.request.Path, r.request.Ref, r.step,
|
|
r.started.UTC().Format(time.RFC3339), elapsed.String()}
|
|
out.Said = fmt.Sprintf("%s is building %s (%s) at %s for %s; %s",
|
|
h.on, r.request.Repository, r.request.ID, orNothing(r.step), elapsed, taking)
|
|
return out
|
|
}
|
|
|
|
// The bounds a kill keeps: each pass removing containers, and the wait for the build to end between
|
|
// them. The worst case — 15s, 20s, 15s — is inside what the controller waits for the answer
|
|
// (killAnswer in the controller's queue.go, 75s).
|
|
var (
|
|
killRemoves = 15 * time.Second
|
|
killWaits = 20 * time.Second
|
|
)
|
|
|
|
// kill ends the build with this id, if it is the one running here and still working: its context
|
|
// cancelled — which kills each command's process group — and the containers it started removed by
|
|
// their label, once at once and again after the build has ended, so one created while it was being
|
|
// killed is not left. The build's own goroutine announces it failed, killed by hand, and settles the
|
|
// ask so it is not redelivered; the answer says whether that happened.
|
|
func (h *holder) kill(id string) (string, error) {
|
|
h.mu.Lock()
|
|
r := h.running
|
|
if r == nil || r.request.ID != id {
|
|
doing := "nothing"
|
|
if r != nil {
|
|
doing = r.request.ID
|
|
}
|
|
h.mu.Unlock()
|
|
return "", fmt.Errorf("%s is not building %s; it is building %s. `queue` says where an ask is", h.on, id, doing)
|
|
}
|
|
if r.returned {
|
|
h.mu.Unlock()
|
|
return "", fmt.Errorf("%s on %s has already ended on its own and is saying how; `builds` shows it", id, h.on)
|
|
}
|
|
r.killed = true
|
|
cancel, done := r.cancel, r.done
|
|
h.mu.Unlock()
|
|
|
|
cancel()
|
|
removed, removeErr := h.removeContainers(id)
|
|
ended := false
|
|
select {
|
|
case <-done:
|
|
ended = true
|
|
case <-time.After(killWaits):
|
|
}
|
|
again, againErr := h.removeContainers(id)
|
|
removed += again
|
|
if removeErr == nil {
|
|
removeErr = againErr
|
|
}
|
|
containers := fmt.Sprintf("%d container(s) it started removed", removed)
|
|
if removeErr != nil {
|
|
containers = "its containers could not all be listed or removed: " + removeErr.Error()
|
|
}
|
|
if !ended {
|
|
return fmt.Sprintf("killed %s on %s: %s; the build has not finished ending yet — `builds` says when "+
|
|
"its outcome is in", id, h.on, containers), nil
|
|
}
|
|
h.mu.Lock()
|
|
outcome, announced := r.outcome, r.announced
|
|
h.mu.Unlock()
|
|
if !announced {
|
|
return fmt.Sprintf("killed %s (%s) on %s: its commands ended and %s, and its outcome could not be "+
|
|
"announced — the ask is not settled and will be handed out again", id, r.request.Repository, h.on,
|
|
containers), nil
|
|
}
|
|
return fmt.Sprintf("killed %s (%s) on %s: its commands ended, %s, and its outcome announced as failed, %s — "+
|
|
"settled, so it is not handed to another machine", id, r.request.Repository, h.on, containers, outcome), nil
|
|
}
|
|
|
|
// removeContainers is one pass of removing what the build left, bounded.
|
|
func (h *holder) removeContainers(id string) (int, error) {
|
|
cleanup, stop := context.WithTimeout(context.Background(), killRemoves)
|
|
defer stop()
|
|
return builder.RemoveContainersOf(cleanup, h.remove, id)
|
|
}
|
|
|
|
// removeLeftBehind removes, once, what earlier holders on this machine left — called before anything is
|
|
// taken, so none of it is this holder's — and says what it did. A runtime that cannot be asked is said
|
|
// and the holder takes work all the same: what is left is waste, not a reason to stop building.
|
|
func (h *holder) removeLeftBehind() int {
|
|
cleanup, stop := context.WithTimeout(context.Background(), killRemoves)
|
|
defer stop()
|
|
n, err := builder.RemoveLeftBehind(cleanup, h.remove)
|
|
switch {
|
|
case err != nil:
|
|
fmt.Fprintf(os.Stderr, "could not look for containers earlier builds left on %s: %v\n", h.on, err)
|
|
case n > 0:
|
|
fmt.Fprintf(os.Stderr, "removed %d container(s) earlier builds left on %s\n", n, h.on)
|
|
}
|
|
return n
|
|
}
|
|
|
|
// returned records that the build's work ended, and says whether a kill came first — only then is
|
|
// the build killed; an error it ended with on its own is its own outcome.
|
|
func (h *holder) returned(r *running) bool {
|
|
h.mu.Lock()
|
|
defer h.mu.Unlock()
|
|
r.returned = true
|
|
return r.killed
|
|
}
|
|
|
|
// said records the outcome announced, and whether it went out, for a kill to answer with.
|
|
func (h *holder) said(r *running, outcome string, announced bool) {
|
|
h.mu.Lock()
|
|
r.outcome, r.announced = outcome, announced
|
|
h.mu.Unlock()
|
|
}
|
|
|
|
// handlers are the seat's verbs, as this machine answers them.
|
|
func (h *holder) handlers() map[string]link.ToolHandler {
|
|
return map[string]link.ToolHandler{
|
|
"current": func(context.Context, json.RawMessage) (any, error) { return h.current(), nil },
|
|
"kill": func(_ context.Context, raw json.RawMessage) (any, error) {
|
|
var args struct {
|
|
ID string `json:"id"`
|
|
}
|
|
if err := json.Unmarshal(raw, &args); err != nil || strings.TrimSpace(args.ID) == "" {
|
|
return nil, errors.New("kill needs the build's id")
|
|
}
|
|
said, err := h.kill(strings.TrimSpace(args.ID))
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
return map[string]any{"said": said}, nil
|
|
},
|
|
"pause": func(context.Context, json.RawMessage) (any, error) {
|
|
if err := h.setPaused(true); err != nil {
|
|
return nil, err
|
|
}
|
|
said := h.on + " is paused: it takes no new build until resumed"
|
|
if c := h.current(); c.Running != nil {
|
|
said += "; " + c.Running.ID + " runs on and finishes"
|
|
}
|
|
return map[string]any{"said": said, "paused": true}, nil
|
|
},
|
|
"resume": func(context.Context, json.RawMessage) (any, error) {
|
|
if err := h.setPaused(false); err != nil {
|
|
return nil, err
|
|
}
|
|
return map[string]any{"said": h.on + " takes builds again", "paused": false}, nil
|
|
},
|
|
}
|
|
}
|
|
|
|
func orNothing(s string) string {
|
|
if s == "" {
|
|
return "its start"
|
|
}
|
|
return s
|
|
}
|
|
|
|
// plainRun runs a command outside any build: no line reaches a build's log, because the build whose
|
|
// log it would be is the one being ended.
|
|
func plainRun(ctx context.Context, dir, name string, args ...string) (string, error) {
|
|
cmd := exec.CommandContext(ctx, name, args...)
|
|
cmd.Dir = dir
|
|
out, err := cmd.CombinedOutput()
|
|
if err != nil {
|
|
return string(out), fmt.Errorf("%s %s: %w: %s", name, strings.Join(args, " "), err, strings.TrimSpace(string(out)))
|
|
}
|
|
return string(out), nil
|
|
}
|
|
|
|
// announcement is what this holder answers discovery with (novox/hq ADR 0195, ADR 0197): this
|
|
// machine's verbs of the seat, in the shape every tool runtime announces a seat's verb — kind seat,
|
|
// the module answering, the seat, scope node, the machine, its description and argument schema — so
|
|
// the console finds `<node>/node-build-agent.kill` by searching, as it finds any seat's verb.
|
|
//
|
|
// One service per machine, named for the seat and identified by the machine, so the answer is this
|
|
// machine's four verbs and nothing more. The verbs are the compiled seat row's: a build machine has no
|
|
// store, and what it serves is what this binary was built to serve.
|
|
func announcement(seat, module, node string) micro.Info {
|
|
s, _ := catalogue.SeatNamed(seat)
|
|
var endpoints []micro.EndpointInfo
|
|
for _, v := range s.Serves {
|
|
schema, _ := json.Marshal(v.Input)
|
|
endpoints = append(endpoints, micro.EndpointInfo{
|
|
Name: seat + "__" + v.Name,
|
|
Subject: link.NodeSeatToolSubject(seat, v.Name, node),
|
|
// The queue group the verbs are served in, as every runtime announces its own.
|
|
QueueGroup: "seat." + seat,
|
|
Metadata: map[string]string{
|
|
"kind": "seat", "module": module, "tool": v.Name, "seat": seat, "scope": "node",
|
|
"node": node, "interchangeable": "false", "description": v.Description, "schema": string(schema),
|
|
},
|
|
})
|
|
}
|
|
return micro.Info{
|
|
ServiceIdentity: micro.ServiceIdentity{Name: seat, ID: node, Version: "0.1.0",
|
|
Metadata: map[string]string{"seat": seat, "scope": "node", "node": node, "module": module}},
|
|
Description: "what the build running on " + node + " is, and ending, pausing and resuming it (novox/hq ADR 0219)",
|
|
Endpoints: endpoints,
|
|
}
|
|
}
|
|
|
|
// moduleOf is the module a credential was issued for: its user is `<node>.<module>`.
|
|
func moduleOf(user string) string {
|
|
if _, module, ok := strings.Cut(user, "."); ok && module != "" {
|
|
return module
|
|
}
|
|
return "build-agent"
|
|
}
|