Let a build agent be paused, have a build killed, and end an ask cancelled as it took it (hq ADR 0219)
A queued ask could only be waited out and a running build only ended by stopping the machine, which redelivered it elsewhere. The holder now serves current, kill, pause and resume on its own machine's subjects; a kill ends the build's process group and labelled containers and settles the ask as failed, killed by hand; pause is kept in the workspace across a restart and said on the bus. The controller writes cancelled ids to a cancelled set the holder reads on taking an ask, closing the race a delete alone leaves.
This commit is contained in:
@@ -0,0 +1,439 @@
|
||||
package link
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"regexp"
|
||||
"sort"
|
||||
"time"
|
||||
|
||||
"github.com/nats-io/nats.go"
|
||||
"github.com/nats-io/nats.go/jetstream"
|
||||
|
||||
"github.com/novox/mesh-controller/internal/broker"
|
||||
)
|
||||
|
||||
// The build queue, read and changed by hand (novox/hq ADR 0219).
|
||||
//
|
||||
// **A queue nobody can see is a queue nobody can trust.** Until this, the build seat's work queue
|
||||
// was visible only as its consequences: a plan LATE with nothing to say why, a build that never
|
||||
// started because five asks ahead of it were waiting on a laptop that was shut. ADR 0219 makes the
|
||||
// queue the controller's to show and change, and the build running on one machine that machine's
|
||||
// holder's to end.
|
||||
//
|
||||
// Three kinds of ask are in the seat's stream at any moment, told apart by the worker consumer
|
||||
// every holder pulls from:
|
||||
//
|
||||
// - **waiting** — after the last one the worker handed out (stream seq > delivered.stream_seq);
|
||||
// - **in flight** — handed out and not yet settled, which the holder that took it said with its
|
||||
// `started` event, and which the worker counts among its acks pending;
|
||||
// - **dead** — handed out as many times as the worker allows and never settled. A work queue
|
||||
// drops what it settles, so one still in the stream behind the worker is either in flight or
|
||||
// dead; the started events and the worker's pending count say which.
|
||||
//
|
||||
// Everything that drops an ask leaves a failed outcome for it, recorded the way a failed build's is
|
||||
// (`cancelled by hand`, `killed by hand`), so a plan waiting on it fails visibly rather than waiting
|
||||
// for ever on an ask nobody will answer.
|
||||
|
||||
// The words an outcome says when a person ended the ask.
|
||||
const (
|
||||
CancelledByHand = "cancelled by hand"
|
||||
KilledByHand = "killed by hand"
|
||||
)
|
||||
|
||||
// Ask states in a queue.
|
||||
const (
|
||||
AskWaiting = "waiting"
|
||||
AskInFlight = "in flight"
|
||||
AskDead = "dead"
|
||||
)
|
||||
|
||||
// QueuedAsk is one ask in the seat's work queue.
|
||||
type QueuedAsk struct {
|
||||
Seq uint64 `json:"seq"`
|
||||
Request BuildRequest `json:"-"`
|
||||
ID string `json:"id"`
|
||||
// Repository, Path and Ref are the ask's, as the request carried them; Held never travels out of
|
||||
// here — it is every artifact the mesh has built, and a queue listing is not where that belongs.
|
||||
Repository string `json:"repository"`
|
||||
Path string `json:"path,omitempty"`
|
||||
Ref string `json:"ref,omitempty"`
|
||||
AskedAt time.Time `json:"asked-at,omitempty"`
|
||||
State string `json:"state"`
|
||||
// On and Started are the holder that took an in-flight ask and when, from its started event.
|
||||
On string `json:"on,omitempty"`
|
||||
Started time.Time `json:"started,omitempty"`
|
||||
}
|
||||
|
||||
// Queue is the seat's work queue as it stands.
|
||||
type Queue struct {
|
||||
Seat string `json:"seat"`
|
||||
// Delivered is the last stream sequence the worker handed out; AckPending how many it handed out
|
||||
// and has not had settled; MaxDeliver how many times it hands one ask out.
|
||||
Delivered uint64 `json:"delivered"`
|
||||
AckPending int `json:"ack-pending"`
|
||||
MaxDeliver int `json:"max-deliver"`
|
||||
Asks []QueuedAsk `json:"asks"`
|
||||
}
|
||||
|
||||
// Of is the asks in one state, in queue order.
|
||||
func (q Queue) Of(state string) []QueuedAsk {
|
||||
var out []QueuedAsk
|
||||
for _, a := range q.Asks {
|
||||
if a.State == state {
|
||||
out = append(out, a)
|
||||
}
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
// Find is the ask with this id.
|
||||
func (q Queue) Find(id string) (QueuedAsk, bool) {
|
||||
for _, a := range q.Asks {
|
||||
if a.ID == id {
|
||||
return a, true
|
||||
}
|
||||
}
|
||||
return QueuedAsk{}, false
|
||||
}
|
||||
|
||||
// BuildHeard is what the seat's events say about builds: the latest start of each id, and every id
|
||||
// whose outcome was heard.
|
||||
type BuildHeard struct {
|
||||
Started map[string]BuildStart
|
||||
Outcomes map[string]bool
|
||||
}
|
||||
|
||||
// ClassifyQueue says which state each ask is in. Pure: the stream's asks in sequence order, what the
|
||||
// worker says, and what the seat's events said.
|
||||
//
|
||||
// **In flight is what a machine is building now.** A holder builds one ask at a time (ADR 0190), so
|
||||
// what a machine is building is its latest start, if no outcome has been heard for it. An ask
|
||||
// behind the worker that is some machine's current build is in flight; at most AckPending of them
|
||||
// are, newest start first — a machine that died mid-build has a latest start forever, and the
|
||||
// worker's own count is what bounds that. An ask behind the worker that no machine has said it
|
||||
// started is in flight only while the worker counts more pending than starts explain (taken, and
|
||||
// not yet said); otherwise it is dead.
|
||||
//
|
||||
// **The worker's count is the server's, and it lags in one case**: an ask handed out its last time
|
||||
// and handed back is still counted pending until some holder pulls again, when the server finds it
|
||||
// past its deliveries and lets it go. Holders pull whenever they are idle, so this is a moment —
|
||||
// except while every holder is paused, when such an ask reads as in flight with no machine named.
|
||||
// Cancel takes an ask in flight that no machine said it started, for exactly that reason.
|
||||
func ClassifyQueue(asks []QueuedAsk, delivered uint64, ackPending int, heard BuildHeard) []QueuedAsk {
|
||||
// Each machine's latest start.
|
||||
latest := map[string]BuildStart{}
|
||||
for _, s := range heard.Started {
|
||||
at := startedAt(s)
|
||||
if was, ok := latest[s.On]; !ok || at.After(startedAt(was)) {
|
||||
latest[s.On] = s
|
||||
}
|
||||
}
|
||||
current := map[string]BuildStart{}
|
||||
for _, s := range latest {
|
||||
if !heard.Outcomes[s.ID] {
|
||||
current[s.ID] = s
|
||||
}
|
||||
}
|
||||
|
||||
out := make([]QueuedAsk, len(asks))
|
||||
copy(out, asks)
|
||||
var running, unsaid []int
|
||||
for i := range out {
|
||||
a := &out[i]
|
||||
if a.Seq > delivered {
|
||||
a.State = AskWaiting
|
||||
continue
|
||||
}
|
||||
a.State = AskDead
|
||||
if s, ok := current[a.ID]; ok {
|
||||
a.On, a.Started = s.On, startedAt(s)
|
||||
running = append(running, i)
|
||||
continue
|
||||
}
|
||||
if _, ever := heard.Started[a.ID]; !ever {
|
||||
unsaid = append(unsaid, i)
|
||||
}
|
||||
}
|
||||
sort.SliceStable(running, func(x, y int) bool { return out[running[x]].Started.After(out[running[y]].Started) })
|
||||
slots := ackPending
|
||||
for _, i := range running {
|
||||
if slots == 0 {
|
||||
// A machine's latest start the worker no longer counts: it died with the ask, and the
|
||||
// ask was handed out as often as it may be.
|
||||
out[i].On, out[i].Started = "", time.Time{}
|
||||
continue
|
||||
}
|
||||
out[i].State = AskInFlight
|
||||
slots--
|
||||
}
|
||||
for _, i := range unsaid {
|
||||
if slots == 0 {
|
||||
break
|
||||
}
|
||||
out[i].State = AskInFlight
|
||||
slots--
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
func startedAt(s BuildStart) time.Time {
|
||||
t, _ := time.Parse(time.RFC3339Nano, s.At)
|
||||
return t
|
||||
}
|
||||
|
||||
// ReadQueue reads a build seat's work queue: the stream's asks, the worker's position, and what the
|
||||
// seat's events said about the builds behind it. Through the JetStream API alone (STREAM.INFO,
|
||||
// CONSUMER.INFO, STREAM.MSG.GET), which only the controller's account reaches.
|
||||
func ReadQueue(ctx context.Context, js *broker.JetStream, seat string) (Queue, error) {
|
||||
worker, _ := broker.HolderConsumerFor("", "", broker.DeclaredSeat{Name: seat, Accepts: []string{"build"}})
|
||||
q := Queue{Seat: seat}
|
||||
api, err := jetstream.New(js.Conn())
|
||||
if err != nil {
|
||||
return q, err
|
||||
}
|
||||
stream, err := api.Stream(ctx, worker.Stream)
|
||||
if err != nil {
|
||||
return q, fmt.Errorf("the %s work queue (%s) cannot be read: %w", seat, worker.Stream, err)
|
||||
}
|
||||
consumer, err := stream.Consumer(ctx, worker.Name)
|
||||
switch {
|
||||
case errors.Is(err, jetstream.ErrConsumerNotFound):
|
||||
// No holder has ever been given a worker: everything in the stream is waiting.
|
||||
case err != nil:
|
||||
return q, fmt.Errorf("the worker of %s cannot be read: %w", seat, err)
|
||||
default:
|
||||
info, err := consumer.Info(ctx)
|
||||
if err != nil {
|
||||
return q, fmt.Errorf("the worker of %s cannot be read: %w", seat, err)
|
||||
}
|
||||
q.Delivered = info.Delivered.Stream
|
||||
q.AckPending = info.NumAckPending
|
||||
q.MaxDeliver = info.Config.MaxDeliver
|
||||
}
|
||||
|
||||
// Every ask, by next-by-subject from the first: the stream holds only what is unsettled, so this
|
||||
// is a handful of reads, never the history.
|
||||
var asks []QueuedAsk
|
||||
var seq uint64 = 1
|
||||
for n := 0; n < 10000; n++ {
|
||||
msg, err := stream.GetMsg(ctx, seq, jetstream.WithGetMsgSubject(BuildWorkOf(seat)))
|
||||
if errors.Is(err, jetstream.ErrMsgNotFound) {
|
||||
break
|
||||
}
|
||||
if err != nil {
|
||||
return q, fmt.Errorf("reading the %s work queue: %w", seat, err)
|
||||
}
|
||||
ask := QueuedAsk{Seq: msg.Sequence}
|
||||
if json.Unmarshal(msg.Data, &ask.Request) == nil {
|
||||
r := ask.Request
|
||||
ask.ID, ask.Repository, ask.Path, ask.Ref = r.ID, r.Repository, r.Path, r.Ref
|
||||
if at, ok := BuildAskedAt(r.ID); ok {
|
||||
ask.AskedAt = at
|
||||
}
|
||||
}
|
||||
asks = append(asks, ask)
|
||||
seq = msg.Sequence + 1
|
||||
}
|
||||
|
||||
// What the events say, from shortly before the oldest ask behind the worker: its start, if any,
|
||||
// came after it was asked.
|
||||
var since time.Time
|
||||
for _, a := range asks {
|
||||
if a.Seq <= q.Delivered && (since.IsZero() || (!a.AskedAt.IsZero() && a.AskedAt.Before(since))) {
|
||||
since = a.AskedAt
|
||||
}
|
||||
}
|
||||
heard := BuildHeard{Started: map[string]BuildStart{}, Outcomes: map[string]bool{}}
|
||||
if len(asks) > 0 && q.Delivered > 0 {
|
||||
if since.IsZero() {
|
||||
since = time.Now().Add(-7 * 24 * time.Hour)
|
||||
}
|
||||
if heard, err = ReadBuildEvents(ctx, js, seat, since.Add(-time.Minute)); err != nil {
|
||||
return q, err
|
||||
}
|
||||
}
|
||||
q.Asks = ClassifyQueue(asks, q.Delivered, q.AckPending, heard)
|
||||
return q, nil
|
||||
}
|
||||
|
||||
// ReadBuildEvents is every start and outcome the seat announced since a moment, from the events
|
||||
// stream, with a consumer of its own that is gone when this returns.
|
||||
func ReadBuildEvents(ctx context.Context, js *broker.JetStream, seat string, since time.Time) (BuildHeard, error) {
|
||||
heard := BuildHeard{Started: map[string]BuildStart{}, Outcomes: map[string]bool{}}
|
||||
// One token after `event`: started and built, never a build's log lines, which are two.
|
||||
subject := "mesh.seat." + seat + ".event.*"
|
||||
sub, err := js.Context().PullSubscribe(subject, "",
|
||||
nats.BindStream(broker.EventsStream), nats.StartTime(since), nats.AckNone())
|
||||
if err != nil {
|
||||
return heard, fmt.Errorf("cannot read %s from the bus: %w", subject, err)
|
||||
}
|
||||
defer func() { _ = sub.Unsubscribe() }()
|
||||
for {
|
||||
batch, err := sub.Fetch(500, nats.MaxWait(2*time.Second))
|
||||
if err != nil && !errors.Is(err, nats.ErrTimeout) && !errors.Is(err, context.DeadlineExceeded) {
|
||||
return heard, fmt.Errorf("reading %s: %w", subject, err)
|
||||
}
|
||||
for _, msg := range batch {
|
||||
switch msg.Subject {
|
||||
case BuildStartedOf(seat):
|
||||
var s BuildStart
|
||||
if json.Unmarshal(msg.Data, &s) == nil && s.ID != "" {
|
||||
if was, ok := heard.Started[s.ID]; !ok || startedAt(s).After(startedAt(was)) {
|
||||
heard.Started[s.ID] = s
|
||||
}
|
||||
}
|
||||
case BuildOutcomeOf(seat):
|
||||
var r BuildResult
|
||||
if json.Unmarshal(msg.Data, &r) == nil && r.ID != "" {
|
||||
heard.Outcomes[r.ID] = true
|
||||
}
|
||||
}
|
||||
}
|
||||
if len(batch) < 500 {
|
||||
break
|
||||
}
|
||||
if ctx.Err() != nil {
|
||||
return heard, ctx.Err()
|
||||
}
|
||||
}
|
||||
return heard, nil
|
||||
}
|
||||
|
||||
// --- the cancelled set ----------------------------------------------------------------------
|
||||
|
||||
// safeKey is an id that can be a key, and a subject token, as it is: a build id, never a pattern.
|
||||
var safeKey = regexp.MustCompile(`^[A-Za-z0-9_-]+$`)
|
||||
|
||||
// Cancelled is what the controller writes for an ask it cancelled.
|
||||
type Cancelled struct {
|
||||
At string `json:"at"`
|
||||
Why string `json:"why"`
|
||||
}
|
||||
|
||||
// MarkCancelled puts an ask's id into the seat's cancelled set (novox/hq ADR 0219). The controller's
|
||||
// act, before it deletes the ask from the queue.
|
||||
func MarkCancelled(ctx context.Context, js *broker.JetStream, seat, id string) error {
|
||||
if !safeKey.MatchString(id) {
|
||||
return fmt.Errorf("%q is not a build id", id)
|
||||
}
|
||||
api, err := jetstream.New(js.Conn())
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
kv, err := api.KeyValue(ctx, broker.CancelledSetName(seat))
|
||||
if err != nil {
|
||||
return fmt.Errorf("the cancelled set of %s is not on the bus — the controller asserts it at its "+
|
||||
"start, so one older than this has not: %w", seat, err)
|
||||
}
|
||||
body, err := json.Marshal(Cancelled{At: time.Now().UTC().Format(time.RFC3339Nano), Why: CancelledByHand})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
_, err = kv.Put(ctx, id, body)
|
||||
return err
|
||||
}
|
||||
|
||||
// IsCancelled asks the seat's cancelled set whether an ask was cancelled: one direct read of one key,
|
||||
// which is all a holder is granted of it.
|
||||
func IsCancelled(conn *nats.Conn, seat, id string) (bool, error) {
|
||||
if !safeKey.MatchString(id) {
|
||||
// Not a key the controller could have written.
|
||||
return false, nil
|
||||
}
|
||||
set := broker.CancelledSetName(seat)
|
||||
reply, err := conn.Request("$JS.API.DIRECT.GET.KV_"+set+".$KV."+set+"."+id, nil, 3*time.Second)
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
return cancelledFrom(reply)
|
||||
}
|
||||
|
||||
// cancelledFrom reads a direct get's answer: a value is a cancel; not found, or a deletion marker,
|
||||
// is not.
|
||||
func cancelledFrom(reply *nats.Msg) (bool, error) {
|
||||
if reply.Header != nil {
|
||||
switch reply.Header.Get("Status") {
|
||||
case "":
|
||||
case "404":
|
||||
return false, nil
|
||||
default:
|
||||
return false, fmt.Errorf("the cancelled set answered %s %s",
|
||||
reply.Header.Get("Status"), reply.Header.Get("Description"))
|
||||
}
|
||||
if op := reply.Header.Get("KV-Operation"); op == "DEL" || op == "PURGE" {
|
||||
return false, nil
|
||||
}
|
||||
}
|
||||
return len(reply.Data) > 0, nil
|
||||
}
|
||||
|
||||
// --- a holder's verbs -----------------------------------------------------------------------
|
||||
|
||||
// NodeSeatToolSubject is where a node-scoped seat's verb is asked of one machine's holder (design 33
|
||||
// §4): the seat's verb subject with the machine as its last token.
|
||||
func NodeSeatToolSubject(seat, verb, node string) string {
|
||||
return SeatToolSubject(seat, verb) + "." + node
|
||||
}
|
||||
|
||||
// AskSeatTool asks one machine's holder of a node seat one of its verbs and reads its answer.
|
||||
func AskSeatTool(ctx context.Context, conn *nats.Conn, seat, verb, node string, args any,
|
||||
timeout time.Duration) (Answer, error) {
|
||||
body, err := json.Marshal(args)
|
||||
if err != nil {
|
||||
return Answer{}, err
|
||||
}
|
||||
asking, cancel := context.WithTimeout(ctx, timeout)
|
||||
defer cancel()
|
||||
reply, err := conn.RequestWithContext(asking, NodeSeatToolSubject(seat, verb, node), body)
|
||||
switch {
|
||||
case errors.Is(err, nats.ErrNoResponders):
|
||||
return Answer{}, fmt.Errorf("nothing on %s answers %s.%s: its holder is not running, or is "+
|
||||
"older than the verb", node, seat, verb)
|
||||
case errors.Is(err, context.DeadlineExceeded), errors.Is(err, nats.ErrTimeout):
|
||||
return Answer{}, fmt.Errorf("%s did not answer %s.%s within %s", node, seat, verb, timeout)
|
||||
case err != nil:
|
||||
return Answer{}, err
|
||||
}
|
||||
var answer Answer
|
||||
if err := json.Unmarshal(reply.Data, &answer); err != nil {
|
||||
return Answer{}, fmt.Errorf("%s answered %s.%s with something unreadable: %w", node, seat, verb, err)
|
||||
}
|
||||
return answer, nil
|
||||
}
|
||||
|
||||
// --- a holder saying whether it takes work ----------------------------------------------------
|
||||
|
||||
// HolderState is what a holder says about itself whenever it is paused or resumed, at its start,
|
||||
// and while it stays paused: whether it takes work (novox/hq ADR 0219). Retained on the events
|
||||
// stream under the machine's own subject, so the last one is the machine's state.
|
||||
type HolderState struct {
|
||||
On string `json:"on"`
|
||||
Paused bool `json:"paused"`
|
||||
At string `json:"at"`
|
||||
}
|
||||
|
||||
// BuildPausedOf is where one machine's holder of a build seat says whether it takes work.
|
||||
func BuildPausedOf(seat, node string) string { return "mesh.seat." + seat + ".event.paused." + node }
|
||||
|
||||
// PausedSaid reads what each machine last said about taking work. A machine that never said is not
|
||||
// in the answer — a holder older than ADR 0219 takes work and never pauses.
|
||||
func PausedSaid(js *broker.JetStream, seat string, nodes []string) (map[string]HolderState, error) {
|
||||
out := map[string]HolderState{}
|
||||
for _, node := range nodes {
|
||||
msg, err := js.Context().GetLastMsg(broker.EventsStream, BuildPausedOf(seat, node))
|
||||
if errors.Is(err, nats.ErrMsgNotFound) {
|
||||
continue
|
||||
}
|
||||
if err != nil {
|
||||
return out, err
|
||||
}
|
||||
var s HolderState
|
||||
if json.Unmarshal(msg.Data, &s) == nil {
|
||||
out[node] = s
|
||||
}
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
Reference in New Issue
Block a user