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.
440 lines
16 KiB
Go
440 lines
16 KiB
Go
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
|
|
}
|